From 22d640701c1468285571edee50f8ea6775944e53 Mon Sep 17 00:00:00 2001 From: wei Date: Wed, 22 Jul 2026 16:49:34 +0800 Subject: [PATCH 1/5] fix(pvc): reject terminating handoff dependencies Treat terminating populated PVs, helper PVCs, and selected Nodes as topology-not-ready before external-populator handoff commits. Cover each dependency with focused no-write regressions. --- .../persistentvolumeclaims/syncer.go | 49 ++++- .../persistentvolumeclaims/syncer_test.go | 175 ++++++++++++++++++ 2 files changed, 220 insertions(+), 4 deletions(-) diff --git a/pkg/controllers/resources/persistentvolumeclaims/syncer.go b/pkg/controllers/resources/persistentvolumeclaims/syncer.go index e054fb815d..0b59d3de6a 100644 --- a/pkg/controllers/resources/persistentvolumeclaims/syncer.go +++ b/pkg/controllers/resources/persistentvolumeclaims/syncer.go @@ -452,6 +452,9 @@ func (s *persistentVolumeClaimSyncer) ensureExternalPopulatorHostMaterialization } return false, err } + if !s.externalPopulatorHostPVReady(vObj, hostPV) { + return false, nil + } currentTarget, targetReady, err := s.externalPopulatorCurrentHostTarget(ctx, pObj, vObj, hostPVName) if err != nil || !targetReady { return false, err @@ -498,6 +501,9 @@ func (s *persistentVolumeClaimSyncer) ensureExternalPopulatorHostMaterialization } return false, err } + if !s.externalPopulatorHostPVReady(vObj, currentHostPV) { + return false, nil + } if !claimRefMatchesPersistentVolumeClaim(currentHostPV.Spec.ClaimRef, pObj) { s.recordExternalPopulatorTopologyEvent( vObj, @@ -604,6 +610,17 @@ func (s *persistentVolumeClaimSyncer) externalPopulatorHandoffTopologyReady( ) return false, nil } + if hostHelperPVC.DeletionTimestamp != nil { + s.recordExternalPopulatorTopologyEvent( + vObj, + externalPopulatorTopologyNotReadyEventReason, + "External-populator handoff is waiting because host helper PVC %s/%s is terminating since %s", + hostHelperPVC.Namespace, + hostHelperPVC.Name, + hostHelperPVC.DeletionTimestamp.Time.UTC().Format(time.RFC3339Nano), + ) + return false, nil + } helperSelectedNode := hostHelperPVC.Annotations[selectedNodeAnnotation] if helperSelectedNode == "" { @@ -631,10 +648,6 @@ func (s *persistentVolumeClaimSyncer) externalPopulatorHandoffTopologyReady( return false, nil } - if hostPV.Spec.NodeAffinity == nil || hostPV.Spec.NodeAffinity.Required == nil { - return true, nil - } - hostNode := &corev1.Node{} err = ctx.HostClient.Get(ctx.Context, types.NamespacedName{Name: targetSelectedNode}, hostNode) if err != nil { @@ -649,6 +662,19 @@ func (s *persistentVolumeClaimSyncer) externalPopulatorHandoffTopologyReady( } return false, fmt.Errorf("get host target Node %q for external-populator handoff: %w", targetSelectedNode, err) } + if hostNode.DeletionTimestamp != nil { + s.recordExternalPopulatorTopologyEvent( + vObj, + externalPopulatorTopologyNotReadyEventReason, + "External-populator handoff is waiting because selected host Node %q is terminating since %s", + targetSelectedNode, + hostNode.DeletionTimestamp.Time.UTC().Format(time.RFC3339Nano), + ) + return false, nil + } + if hostPV.Spec.NodeAffinity == nil || hostPV.Spec.NodeAffinity.Required == nil { + return true, nil + } matches, err := nodeaffinity.NewLazyErrorNodeSelector(hostPV.Spec.NodeAffinity.Required).Match(hostNode) if err != nil { @@ -760,6 +786,21 @@ func (s *persistentVolumeClaimSyncer) recordExternalPopulatorTopologyEvent(vObj ) } +func (s *persistentVolumeClaimSyncer) externalPopulatorHostPVReady(vObj *corev1.PersistentVolumeClaim, hostPV *corev1.PersistentVolume) bool { + if hostPV.DeletionTimestamp == nil { + return true + } + + s.recordExternalPopulatorTopologyEvent( + vObj, + externalPopulatorTopologyNotReadyEventReason, + "External-populator handoff is waiting because populated host PV %q is terminating since %s", + hostPV.Name, + hostPV.DeletionTimestamp.Time.UTC().Format(time.RFC3339Nano), + ) + return false +} + func (s *persistentVolumeClaimSyncer) externalPopulatorCurrentHostTarget( ctx *synccontext.SyncContext, pObj, vObj *corev1.PersistentVolumeClaim, diff --git a/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go b/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go index 48427c8ec6..cb8dd5c98f 100644 --- a/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go +++ b/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go @@ -581,6 +581,16 @@ func TestSync(t *testing.T) { dataProtectionWFFCHostPVNode4BoundToHelper := withRequiredNodeAffinity(dataProtectionHostPVBoundToHelper, "node4") dataProtectionWFFCHostPVNode4BoundToTarget := dataProtectionWFFCHostPVNode4BoundToHelper.DeepCopy() dataProtectionWFFCHostPVNode4BoundToTarget.Spec.ClaimRef = dataProtectionHostPVBoundToTarget.Spec.ClaimRef.DeepCopy() + terminatingDependencyTimestamp := metav1.NewTime(time.Unix(1_700_000_000, 0).UTC()) + dataProtectionWFFCTerminatingHostPVNode4BoundToHelper := dataProtectionWFFCHostPVNode4BoundToHelper.DeepCopy() + dataProtectionWFFCTerminatingHostPVNode4BoundToHelper.DeletionTimestamp = &terminatingDependencyTimestamp + dataProtectionWFFCTerminatingHostPVNode4BoundToHelper.Finalizers = []string{"kubernetes.io/pv-protection"} + dataProtectionHostWFFCTerminatingHelperNode4Pvc := dataProtectionHostWFFCHelperNode4Pvc.DeepCopy() + dataProtectionHostWFFCTerminatingHelperNode4Pvc.DeletionTimestamp = &terminatingDependencyTimestamp + dataProtectionHostWFFCTerminatingHelperNode4Pvc.Finalizers = []string{"kubernetes.io/pvc-protection"} + terminatingNode4 := node4.DeepCopy() + terminatingNode4.DeletionTimestamp = &terminatingDependencyTimestamp + terminatingNode4.Finalizers = []string{"test.vcluster.loft.sh/review-hold"} dataProtectionWFFCTargetWithoutSelectedNodePvc := withStorageClassAndSelectedNode(dataProtectionBackupPendingPvcWithVolumeName, waitForFirstConsumerStorageClassName, "", false) dataProtectionHostWFFCTargetWithoutSelectedNodePvc := withStorageClassAndSelectedNode(dataProtectionHostPendingPvcWithUID, waitForFirstConsumerStorageClassName, "", true) @@ -1583,6 +1593,171 @@ func TestSync(t *testing.T) { ) }, }, + { + Name: "Keep WFFC external-populator handoff pending when host PV is terminating", + InitialVirtualState: []runtime.Object{ + waitForFirstConsumerStorageClass.DeepCopy(), + dataProtectionWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionWFFCVirtualPVNode4BoundToHelper.DeepCopy(), + dataProtectionWFFCHelperNode4Pvc.DeepCopy(), + }, + InitialPhysicalState: []runtime.Object{ + dataProtectionHostWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionHostWFFCHelperNode4Pvc.DeepCopy(), + dataProtectionWFFCTerminatingHostPVNode4BoundToHelper.DeepCopy(), + node4.DeepCopy(), + }, + ExpectedVirtualState: map[schema.GroupVersionKind][]runtime.Object{ + corev1.SchemeGroupVersion.WithKind("PersistentVolumeClaim"): { + dataProtectionWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionWFFCHelperNode4Pvc.DeepCopy(), + }, + corev1.SchemeGroupVersion.WithKind("PersistentVolume"): {dataProtectionWFFCVirtualPVNode4BoundToHelper.DeepCopy()}, + storagev1.SchemeGroupVersion.WithKind("StorageClass"): {waitForFirstConsumerStorageClass.DeepCopy()}, + }, + ExpectedPhysicalState: map[schema.GroupVersionKind][]runtime.Object{ + corev1.SchemeGroupVersion.WithKind("PersistentVolumeClaim"): { + dataProtectionHostWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionHostWFFCHelperNode4Pvc.DeepCopy(), + }, + corev1.SchemeGroupVersion.WithKind("PersistentVolume"): {dataProtectionWFFCTerminatingHostPVNode4BoundToHelper.DeepCopy()}, + corev1.SchemeGroupVersion.WithKind("Node"): {node4.DeepCopy()}, + }, + Sync: func(ctx *synccontext.RegisterContext) { + syncCtx, objectSyncer := syncertesting.FakeStartSyncer(t, ctx, New) + pvcSyncer := objectSyncer.(*persistentVolumeClaimSyncer) + pvcSyncer.useFakePersistentVolumes = true + recorder := installPVCEventRecorder(pvcSyncer) + + result, err := pvcSyncer.Sync(syncCtx, synccontext.NewSyncEventWithOld( + dataProtectionHostWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionHostWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionWFFCTargetNode4Pvc.DeepCopy(), + )) + assert.NilError(t, err) + assert.Check(t, result.RequeueAfter == 2*time.Second, "expected a 2s topology requeue, got %s", result.RequeueAfter) + assertSinglePVCEvent(t, recorder, "Warning ExternalPopulatorTopologyNotReady", dataProtectionWFFCTerminatingHostPVNode4BoundToHelper.Name, "terminating") + checkExternalPopulatorHandoffPending( + t, + syncCtx, + types.NamespacedName{Namespace: dataProtectionHostWFFCTargetNode4Pvc.Namespace, Name: dataProtectionHostWFFCTargetNode4Pvc.Name}, + types.NamespacedName{Namespace: dataProtectionWFFCTargetNode4Pvc.Namespace, Name: dataProtectionWFFCTargetNode4Pvc.Name}, + types.NamespacedName{Namespace: dataProtectionHostWFFCHelperNode4Pvc.Namespace, Name: dataProtectionHostWFFCHelperNode4Pvc.Name}, + dataProtectionWFFCTerminatingHostPVNode4BoundToHelper.Name, + ) + }, + }, + { + Name: "Keep WFFC external-populator handoff pending when host helper PVC is terminating", + InitialVirtualState: []runtime.Object{ + waitForFirstConsumerStorageClass.DeepCopy(), + dataProtectionWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionWFFCVirtualPVNode4BoundToHelper.DeepCopy(), + dataProtectionWFFCHelperNode4Pvc.DeepCopy(), + }, + InitialPhysicalState: []runtime.Object{ + dataProtectionHostWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionHostWFFCTerminatingHelperNode4Pvc.DeepCopy(), + dataProtectionWFFCHostPVNode4BoundToHelper.DeepCopy(), + node4.DeepCopy(), + }, + ExpectedVirtualState: map[schema.GroupVersionKind][]runtime.Object{ + corev1.SchemeGroupVersion.WithKind("PersistentVolumeClaim"): { + dataProtectionWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionWFFCHelperNode4Pvc.DeepCopy(), + }, + corev1.SchemeGroupVersion.WithKind("PersistentVolume"): {dataProtectionWFFCVirtualPVNode4BoundToHelper.DeepCopy()}, + storagev1.SchemeGroupVersion.WithKind("StorageClass"): {waitForFirstConsumerStorageClass.DeepCopy()}, + }, + ExpectedPhysicalState: map[schema.GroupVersionKind][]runtime.Object{ + corev1.SchemeGroupVersion.WithKind("PersistentVolumeClaim"): { + dataProtectionHostWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionHostWFFCTerminatingHelperNode4Pvc.DeepCopy(), + }, + corev1.SchemeGroupVersion.WithKind("PersistentVolume"): {dataProtectionWFFCHostPVNode4BoundToHelper.DeepCopy()}, + corev1.SchemeGroupVersion.WithKind("Node"): {node4.DeepCopy()}, + }, + Sync: func(ctx *synccontext.RegisterContext) { + syncCtx, objectSyncer := syncertesting.FakeStartSyncer(t, ctx, New) + pvcSyncer := objectSyncer.(*persistentVolumeClaimSyncer) + pvcSyncer.useFakePersistentVolumes = true + recorder := installPVCEventRecorder(pvcSyncer) + + result, err := pvcSyncer.Sync(syncCtx, synccontext.NewSyncEventWithOld( + dataProtectionHostWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionHostWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionWFFCTargetNode4Pvc.DeepCopy(), + )) + assert.NilError(t, err) + assert.Check(t, result.RequeueAfter == 2*time.Second, "expected a 2s topology requeue, got %s", result.RequeueAfter) + assertSinglePVCEvent(t, recorder, "Warning ExternalPopulatorTopologyNotReady", dataProtectionHostWFFCTerminatingHelperNode4Pvc.Name, "terminating") + checkExternalPopulatorHandoffPending( + t, + syncCtx, + types.NamespacedName{Namespace: dataProtectionHostWFFCTargetNode4Pvc.Namespace, Name: dataProtectionHostWFFCTargetNode4Pvc.Name}, + types.NamespacedName{Namespace: dataProtectionWFFCTargetNode4Pvc.Namespace, Name: dataProtectionWFFCTargetNode4Pvc.Name}, + types.NamespacedName{Namespace: dataProtectionHostWFFCTerminatingHelperNode4Pvc.Namespace, Name: dataProtectionHostWFFCTerminatingHelperNode4Pvc.Name}, + dataProtectionWFFCHostPVNode4BoundToHelper.Name, + ) + }, + }, + { + Name: "Keep WFFC external-populator handoff pending when selected host Node is terminating", + InitialVirtualState: []runtime.Object{ + waitForFirstConsumerStorageClass.DeepCopy(), + dataProtectionWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionWFFCVirtualPVNode4BoundToHelper.DeepCopy(), + dataProtectionWFFCHelperNode4Pvc.DeepCopy(), + }, + InitialPhysicalState: []runtime.Object{ + dataProtectionHostWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionHostWFFCHelperNode4Pvc.DeepCopy(), + dataProtectionWFFCHostPVNode4BoundToHelper.DeepCopy(), + terminatingNode4.DeepCopy(), + }, + ExpectedVirtualState: map[schema.GroupVersionKind][]runtime.Object{ + corev1.SchemeGroupVersion.WithKind("PersistentVolumeClaim"): { + dataProtectionWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionWFFCHelperNode4Pvc.DeepCopy(), + }, + corev1.SchemeGroupVersion.WithKind("PersistentVolume"): {dataProtectionWFFCVirtualPVNode4BoundToHelper.DeepCopy()}, + storagev1.SchemeGroupVersion.WithKind("StorageClass"): {waitForFirstConsumerStorageClass.DeepCopy()}, + }, + ExpectedPhysicalState: map[schema.GroupVersionKind][]runtime.Object{ + corev1.SchemeGroupVersion.WithKind("PersistentVolumeClaim"): { + dataProtectionHostWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionHostWFFCHelperNode4Pvc.DeepCopy(), + }, + corev1.SchemeGroupVersion.WithKind("PersistentVolume"): {dataProtectionWFFCHostPVNode4BoundToHelper.DeepCopy()}, + corev1.SchemeGroupVersion.WithKind("Node"): {terminatingNode4.DeepCopy()}, + }, + Sync: func(ctx *synccontext.RegisterContext) { + syncCtx, objectSyncer := syncertesting.FakeStartSyncer(t, ctx, New) + pvcSyncer := objectSyncer.(*persistentVolumeClaimSyncer) + pvcSyncer.useFakePersistentVolumes = true + recorder := installPVCEventRecorder(pvcSyncer) + + result, err := pvcSyncer.Sync(syncCtx, synccontext.NewSyncEventWithOld( + dataProtectionHostWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionHostWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionWFFCTargetNode4Pvc.DeepCopy(), + dataProtectionWFFCTargetNode4Pvc.DeepCopy(), + )) + assert.NilError(t, err) + assert.Check(t, result.RequeueAfter == 2*time.Second, "expected a 2s topology requeue, got %s", result.RequeueAfter) + assertSinglePVCEvent(t, recorder, "Warning ExternalPopulatorTopologyNotReady", terminatingNode4.Name, "terminating") + checkExternalPopulatorHandoffPending( + t, + syncCtx, + types.NamespacedName{Namespace: dataProtectionHostWFFCTargetNode4Pvc.Namespace, Name: dataProtectionHostWFFCTargetNode4Pvc.Name}, + types.NamespacedName{Namespace: dataProtectionWFFCTargetNode4Pvc.Namespace, Name: dataProtectionWFFCTargetNode4Pvc.Name}, + types.NamespacedName{Namespace: dataProtectionHostWFFCHelperNode4Pvc.Namespace, Name: dataProtectionHostWFFCHelperNode4Pvc.Name}, + dataProtectionWFFCHostPVNode4BoundToHelper.Name, + ) + }, + }, { Name: "Bridge WFFC external-populator handoff when target helper and PV all select node4", InitialVirtualState: []runtime.Object{ From 82b48a2ee57b2edc19b4f77769e3a6ac96e2e679 Mon Sep 17 00:00:00 2001 From: Gina Date: Fri, 24 Jul 2026 14:04:12 +0800 Subject: [PATCH 2/5] fix(storageclasses): project host defaults for synced PVCs --- chart/templates/_rbac.tpl | 19 ++- chart/templates/clusterrole.yaml | 9 +- chart/tests/clusterrole_test.yaml | 114 +++++++++++++++--- chart/values.schema.json | 4 +- chart/values.yaml | 2 +- config/config.go | 2 +- pkg/config/validation.go | 11 +- pkg/config/validation_test.go | 56 +++++++++ .../persistentvolumeclaims/syncer_test.go | 36 ++++++ .../storageclasses/host_syncer_test.go | 16 ++- 10 files changed, 233 insertions(+), 36 deletions(-) diff --git a/chart/templates/_rbac.tpl b/chart/templates/_rbac.tpl index 86716379ed..ddcf036e6b 100644 --- a/chart/templates/_rbac.tpl +++ b/chart/templates/_rbac.tpl @@ -6,6 +6,22 @@ {{- printf "vc-mn-%s-v-%s" .Release.Name .Release.Namespace | trunc 63 | trimSuffix "-" -}} {{- end -}} +{{/* + Whether storage classes are synced from the host after resolving auto. +*/}} +{{- define "vcluster.syncFromHostStorageClasses" -}} +{{- if or + (eq (toString .Values.sync.fromHost.storageClasses.enabled) "true") + (and + (eq (toString .Values.sync.fromHost.storageClasses.enabled) "auto") + .Values.sync.toHost.persistentVolumeClaims.enabled + (not .Values.sync.toHost.storageClasses.enabled) + ) + -}} +{{- true -}} +{{- end -}} +{{- end -}} + {{/* Whether to create a cluster role or not */}} @@ -33,7 +49,7 @@ .Values.sync.toHost.pods.hybridScheduling.enabled .Values.sync.fromHost.ingressClasses.enabled .Values.sync.fromHost.runtimeClasses.enabled - (eq (toString .Values.sync.fromHost.storageClasses.enabled) "true") + (include "vcluster.syncFromHostStorageClasses" .) (eq (toString .Values.sync.fromHost.csiNodes.enabled) "true") (eq (toString .Values.sync.fromHost.csiDrivers.enabled) "true") (eq (toString .Values.sync.fromHost.csiStorageCapacities.enabled) "true") @@ -255,4 +271,3 @@ {{- end }} {{- end }} {{- end }} - diff --git a/chart/templates/clusterrole.yaml b/chart/templates/clusterrole.yaml index 8e10372ac7..c0c1923d01 100644 --- a/chart/templates/clusterrole.yaml +++ b/chart/templates/clusterrole.yaml @@ -48,7 +48,14 @@ rules: resources: ["storageclasses", "csinodes", "csidrivers", "csistoragecapacities"] verbs: ["get", "watch", "list"] {{- end }} - {{- if eq (toString .Values.sync.fromHost.storageClasses.enabled) "true" }} + {{- if and + (include "vcluster.syncFromHostStorageClasses" .) + (not (or + .Values.controlPlane.distro.k8s.scheduler.enabled + .Values.controlPlane.advanced.virtualScheduler.enabled + .Values.sync.toHost.pods.hybridScheduling.enabled + )) + }} - apiGroups: ["storage.k8s.io"] resources: ["storageclasses"] verbs: ["get", "watch", "list"] diff --git a/chart/tests/clusterrole_test.yaml b/chart/tests/clusterrole_test.yaml index 4dab1a1867..3c6a611d3c 100644 --- a/chart/tests/clusterrole_test.yaml +++ b/chart/tests/clusterrole_test.yaml @@ -58,7 +58,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 2 + count: 3 - contains: path: rules content: @@ -72,6 +72,80 @@ tests: resources: [ "nodes" ] verbs: [ "get" ] + - it: auto projects host storage classes for translated PVCs + set: + rbac: + enableVolumeSnapshotRules: + enabled: false + sync: + toHost: + persistentVolumeClaims: + enabled: true + persistentVolumes: + enabled: false + asserts: + - hasDocuments: + count: 1 + - contains: + path: rules + content: + apiGroups: [ "storage.k8s.io" ] + resources: [ "storageclasses" ] + verbs: [ "get", "watch", "list" ] + + - it: explicit false does not project host storage classes + set: + rbac: + enableVolumeSnapshotRules: + enabled: false + sync: + toHost: + persistentVolumeClaims: + enabled: true + persistentVolumes: + enabled: false + fromHost: + storageClasses: + enabled: false + asserts: + - hasDocuments: + count: 1 + - notContains: + path: rules + content: + apiGroups: [ "storage.k8s.io" ] + resources: [ "storageclasses" ] + verbs: [ "get", "watch", "list" ] + + - it: auto does not project host storage classes when guest classes sync to host + set: + rbac: + enableVolumeSnapshotRules: + enabled: false + sync: + toHost: + persistentVolumeClaims: + enabled: true + persistentVolumes: + enabled: false + storageClasses: + enabled: true + asserts: + - hasDocuments: + count: 1 + - contains: + path: rules + content: + apiGroups: [ "storage.k8s.io" ] + resources: [ "storageclasses" ] + verbs: [ "create", "delete", "patch", "update", "get", "watch", "list" ] + - notContains: + path: rules + content: + apiGroups: [ "storage.k8s.io" ] + resources: [ "storageclasses" ] + verbs: [ "get", "watch", "list" ] + - it: force enable set: rbac: @@ -196,7 +270,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 6 + count: 7 - contains: path: rules content: @@ -246,7 +320,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 7 + count: 8 - contains: path: rules content: @@ -277,7 +351,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 7 + count: 8 - contains: path: rules content: @@ -303,7 +377,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 6 + count: 7 - contains: path: rules content: @@ -322,7 +396,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 6 + count: 7 - contains: path: rules content: @@ -352,7 +426,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 5 + count: 6 - contains: path: rules content: @@ -393,7 +467,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 7 + count: 8 - contains: path: rules content: @@ -453,7 +527,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 6 + count: 7 - contains: path: rules content: @@ -473,7 +547,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 7 + count: 8 - contains: path: rules content: @@ -499,7 +573,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 7 + count: 8 - contains: path: rules content: @@ -521,7 +595,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 6 + count: 7 - contains: path: rules content: @@ -548,7 +622,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 8 + count: 9 - contains: path: rules content: @@ -622,7 +696,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 8 + count: 9 - contains: path: rules content: @@ -658,7 +732,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 8 + count: 9 - contains: path: rules content: @@ -694,7 +768,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 8 + count: 9 - contains: path: rules content: @@ -741,7 +815,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 8 + count: 9 - contains: path: rules content: @@ -777,7 +851,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 8 + count: 9 - contains: path: rules content: @@ -813,7 +887,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 8 + count: 9 - contains: path: rules content: @@ -847,7 +921,7 @@ tests: count: 1 - lengthEqual: path: rules - count: 6 + count: 7 - contains: path: rules content: diff --git a/chart/values.schema.json b/chart/values.schema.json index 6d60463ae6..e55dc4652f 100755 --- a/chart/values.schema.json +++ b/chart/values.schema.json @@ -5091,7 +5091,7 @@ }, "storageClasses": { "$ref": "#/$defs/EnableAutoSwitchWithPatchesAndSelector", - "description": "StorageClasses defines if storage classes should get synced from the host cluster to the virtual cluster, but not back. If auto, is automatically enabled when the virtual scheduler is enabled." + "description": "StorageClasses defines if storage classes should get synced from the host cluster to the virtual cluster, but not back. If auto, is automatically enabled when persistent volume claims are synced to the host and storage classes are not synced to the host." }, "csiNodes": { "$ref": "#/$defs/EnableAutoSwitchWithPatches", @@ -5903,4 +5903,4 @@ "additionalProperties": false, "type": "object", "description": "Config is the vCluster config." -} \ No newline at end of file +} diff --git a/chart/values.yaml b/chart/values.yaml index e03596b1eb..9cc706e852 100644 --- a/chart/values.yaml +++ b/chart/values.yaml @@ -196,7 +196,7 @@ sync: csiStorageCapacities: # Enabled defines if this option should be enabled. enabled: auto - # StorageClasses defines if storage classes should get synced from the host cluster to the virtual cluster, but not back. If auto, is automatically enabled when the virtual scheduler is enabled. + # StorageClasses defines if storage classes should get synced from the host cluster to the virtual cluster, but not back. If auto, is automatically enabled when persistent volume claims are synced to the host and storage classes are not synced to the host. storageClasses: # Enabled defines if this option should be enabled. enabled: auto diff --git a/config/config.go b/config/config.go index c2cfbbbc22..0f3c643c9f 100644 --- a/config/config.go +++ b/config/config.go @@ -1369,7 +1369,7 @@ type SyncFromHost struct { // PriorityClasses defines if priority classes classes should get synced from the host cluster to the virtual cluster, but not back. PriorityClasses EnableSwitchWithPatchesAndSelector `json:"priorityClasses,omitempty"` - // StorageClasses defines if storage classes should get synced from the host cluster to the virtual cluster, but not back. If auto, is automatically enabled when the virtual scheduler is enabled. + // StorageClasses defines if storage classes should get synced from the host cluster to the virtual cluster, but not back. If auto, is automatically enabled when persistent volume claims are synced to the host and storage classes are not synced to the host. StorageClasses EnableAutoSwitchWithPatchesAndSelector `json:"storageClasses,omitempty"` // CSINodes defines if csi nodes should get synced from the host cluster to the virtual cluster, but not back. If auto, is automatically enabled when the virtual scheduler is enabled. diff --git a/pkg/config/validation.go b/pkg/config/validation.go index 87936e7d01..1a51bb685e 100644 --- a/pkg/config/validation.go +++ b/pkg/config/validation.go @@ -60,6 +60,14 @@ func ValidateConfigAndSetDefaults(vConfig *VirtualClusterConfig) error { vConfig.Sync.FromHost.Nodes.Selector.All = true } + // Keep the guest's effective storage classes aligned with the host whenever + // guest PVCs are translated to the host. + if vConfig.Sync.ToHost.PersistentVolumeClaims.Enabled && + vConfig.Sync.FromHost.StorageClasses.Enabled == "auto" && + !vConfig.Sync.ToHost.StorageClasses.Enabled { + vConfig.Sync.FromHost.StorageClasses.Enabled = "true" + } + // enable additional controllers required for scheduling with storage if vConfig.SchedulingInVirtualClusterEnabled() && vConfig.Sync.ToHost.PersistentVolumeClaims.Enabled { if vConfig.Sync.FromHost.CSINodes.Enabled == "auto" { @@ -71,9 +79,6 @@ func ValidateConfigAndSetDefaults(vConfig *VirtualClusterConfig) error { if vConfig.Sync.FromHost.CSIDrivers.Enabled == "auto" { vConfig.Sync.FromHost.CSIDrivers.Enabled = "true" } - if vConfig.Sync.FromHost.StorageClasses.Enabled == "auto" && !vConfig.Sync.ToHost.StorageClasses.Enabled { - vConfig.Sync.FromHost.StorageClasses.Enabled = "true" - } } // check if embedded database and multiple replicas diff --git a/pkg/config/validation_test.go b/pkg/config/validation_test.go index 20dfc1f010..b6463b98e4 100644 --- a/pkg/config/validation_test.go +++ b/pkg/config/validation_test.go @@ -1982,6 +1982,62 @@ func TestValidateCustomResourceSyncProxyConflicts(t *testing.T) { } } +func TestStorageClassAutoSyncFollowsPVCTranslation(t *testing.T) { + tests := []struct { + name string + persistentVolumeClaims bool + toHostStorageClasses bool + fromHostStorageClasses string + expectedFromHostStorageClass string + }{ + { + name: "auto enables effective host classes for translated PVCs", + persistentVolumeClaims: true, + fromHostStorageClasses: "auto", + expectedFromHostStorageClass: "true", + }, + { + name: "auto remains disabled without translated PVCs", + fromHostStorageClasses: "auto", + expectedFromHostStorageClass: "auto", + }, + { + name: "to-host storage classes take precedence", + persistentVolumeClaims: true, + toHostStorageClasses: true, + fromHostStorageClasses: "auto", + expectedFromHostStorageClass: "auto", + }, + { + name: "explicit false remains disabled", + persistentVolumeClaims: true, + fromHostStorageClasses: "false", + expectedFromHostStorageClass: "false", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + defaultConfig, err := config.NewDefaultConfig() + if err != nil { + t.Fatalf("create default config: %v", err) + } + + defaultConfig.Sync.ToHost.PersistentVolumeClaims.Enabled = tt.persistentVolumeClaims + defaultConfig.Sync.ToHost.StorageClasses.Enabled = tt.toHostStorageClasses + defaultConfig.Sync.FromHost.StorageClasses.Enabled = config.StrBool(tt.fromHostStorageClasses) + + vConfig := &VirtualClusterConfig{Config: *defaultConfig} + if err := ValidateConfigAndSetDefaults(vConfig); err != nil { + t.Fatalf("validate config: %v", err) + } + if got := string(vConfig.Sync.FromHost.StorageClasses.Enabled); got != tt.expectedFromHostStorageClass { + t.Fatalf("expected sync.fromHost.storageClasses.enabled=%q, got %q", tt.expectedFromHostStorageClass, got) + } + }) + } +} + func TestValidateExperimentalProxyCustomResourcesConfig(t *testing.T) { cases := []struct { name string diff --git a/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go b/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go index cb8dd5c98f..540db604d2 100644 --- a/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go +++ b/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go @@ -72,6 +72,42 @@ func assertNoPVCEvent(t *testing.T, recorder *events.FakeRecorder) { } } +func TestTranslateSelectorPreservesNilAndExplicitEmptyStorageClass(t *testing.T) { + empty := "" + tests := []struct { + name string + storageClassName *string + }{ + { + name: "nil remains nil", + storageClassName: nil, + }, + { + name: "explicit empty remains empty", + storageClassName: &empty, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + pvc := &corev1.PersistentVolumeClaim{ + Spec: corev1.PersistentVolumeClaimSpec{ + StorageClassName: tt.storageClassName, + }, + } + + (&persistentVolumeClaimSyncer{}).translateSelector(nil, pvc) + + if tt.storageClassName == nil { + assert.Assert(t, pvc.Spec.StorageClassName == nil) + return + } + assert.Assert(t, pvc.Spec.StorageClassName != nil) + assert.Equal(t, *pvc.Spec.StorageClassName, "") + }) + } +} + func checkExternalPopulatorHandoffPending( t *testing.T, ctx *synccontext.SyncContext, diff --git a/pkg/controllers/resources/storageclasses/host_syncer_test.go b/pkg/controllers/resources/storageclasses/host_syncer_test.go index 418e4368bf..a83445d8c9 100644 --- a/pkg/controllers/resources/storageclasses/host_syncer_test.go +++ b/pkg/controllers/resources/storageclasses/host_syncer_test.go @@ -27,8 +27,9 @@ func TestFromHostSync(t *testing.T) { "example.com/label-b": "test-2", }, Annotations: map[string]string{ - "example.com/annotation-a": "test-1", - "example.com/annotation-b": "test-2", + DefaultStorageClassAnnotation: "true", + "example.com/annotation-a": "test-1", + "example.com/annotation-b": "test-2", }, }, Provisioner: "my-provisioner", @@ -42,8 +43,9 @@ func TestFromHostSync(t *testing.T) { "example.com/label-b": "test-2", }, Annotations: map[string]string{ - "example.com/annotation-a": "test-1", - "example.com/annotation-b": "test-2", + DefaultStorageClassAnnotation: "true", + "example.com/annotation-a": "test-1", + "example.com/annotation-b": "test-2", }, }, Provisioner: "my-provisioner", @@ -60,6 +62,8 @@ func TestFromHostSync(t *testing.T) { vObjectUpdated.Parameters = map[string]string{ "test": "value", } + vObjectWithoutDefault := vObject.DeepCopy() + delete(vObjectWithoutDefault.Annotations, DefaultStorageClassAnnotation) syncertesting.RunTests(t, []*syncertesting.SyncTest{ { @@ -80,7 +84,7 @@ func TestFromHostSync(t *testing.T) { { Name: "Sync host changes to virtual", InitialPhysicalState: []runtime.Object{pObjectUpdated.DeepCopy()}, // host resource has been updated - InitialVirtualState: []runtime.Object{vObject.DeepCopy()}, // virtual resource has old values + InitialVirtualState: []runtime.Object{vObjectWithoutDefault}, // virtual resource is missing the host default annotation ExpectedPhysicalState: map[schema.GroupVersionKind][]runtime.Object{ storagev1.SchemeGroupVersion.WithKind("StorageClass"): {pObjectUpdated}, // host resource did not change }, @@ -89,7 +93,7 @@ func TestFromHostSync(t *testing.T) { }, Sync: func(ctx *synccontext.RegisterContext) { syncerCtx, syncer := newFakeSyncer(t, ctx) - _, err := syncer.Sync(syncerCtx, synccontext.NewSyncEvent(pObjectUpdated, vObject.DeepCopy())) + _, err := syncer.Sync(syncerCtx, synccontext.NewSyncEvent(pObjectUpdated, vObjectWithoutDefault.DeepCopy())) assert.NilError(t, err) }, }, From 42a462faadada201fec1bd539d83f419e4822f81 Mon Sep 17 00:00:00 2001 From: wei Date: Fri, 24 Jul 2026 20:53:42 +0800 Subject: [PATCH 3/5] test(pvc): expose missing external populator dependency wakeups Capture guest PV creation, populate-helper creation, and PV Bound visibility transitions after the target PVC pre-gate returns (nil, false, nil). Require an authoritative mapped target request while preserving zero timer retries. --- .../persistentvolumeclaims/syncer_test.go | 329 ++++++++++++++++++ 1 file changed, 329 insertions(+) diff --git a/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go b/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go index 540db604d2..dd42a52429 100644 --- a/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go +++ b/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go @@ -1,6 +1,7 @@ package persistentvolumeclaims import ( + "context" "strings" "testing" "time" @@ -18,10 +19,13 @@ import ( corev1 "k8s.io/api/core/v1" storagev1 "k8s.io/api/storage/v1" + apiequality "k8s.io/apimachinery/pkg/api/equality" "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/runtime/schema" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" ) type eventRecordingTranslator struct { @@ -72,6 +76,331 @@ func assertNoPVCEvent(t *testing.T, recorder *events.FakeRecorder) { } } +func TestSyncExternalPopulatorPreGateTransientSchedulesTarget(t *testing.T) { + type dependencyRequestMapper interface { + externalPopulatorDependencyRequests(context.Context, client.Object) []ctrl.Request + } + type fixture struct { + virtualTarget *corev1.PersistentVolumeClaim + virtualPV *corev1.PersistentVolume + virtualHelper *corev1.PersistentVolumeClaim + hostTarget *corev1.PersistentVolumeClaim + hostPV *corev1.PersistentVolume + } + + newFixture := func() *fixture { + const ( + virtualNamespace = "testns" + virtualTargetUID = types.UID("target-pvc-uid") + virtualPVName = "restore-populated-pv" + ) + + dataProtectionGroup := "dataprotection.kubeblocks.io" + virtualTarget := &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{ + Name: "testpvc", + Namespace: virtualNamespace, + UID: virtualTargetUID, + }, + Spec: corev1.PersistentVolumeClaimSpec{ + VolumeName: virtualPVName, + DataSourceRef: &corev1.TypedObjectReference{ + APIGroup: &dataProtectionGroup, + Kind: "Backup", + Name: "backup-1", + }, + }, + Status: corev1.PersistentVolumeClaimStatus{ + Phase: corev1.ClaimPending, + }, + } + virtualPV := &corev1.PersistentVolume{ + ObjectMeta: metav1.ObjectMeta{Name: virtualPVName}, + Spec: corev1.PersistentVolumeSpec{ + AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteOnce}, + Capacity: corev1.ResourceList{ + corev1.ResourceStorage: resource.MustParse("1Gi"), + }, + ClaimRef: &corev1.ObjectReference{ + Namespace: virtualNamespace, + Name: virtualTarget.Name, + UID: virtualTargetUID, + }, + }, + Status: corev1.PersistentVolumeStatus{ + Phase: corev1.VolumeBound, + }, + } + virtualHelper := &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{ + Name: externalPopulatorPopulateHelperPrefix + string(virtualTargetUID), + Namespace: virtualNamespace, + UID: types.UID("populate-helper-pvc-uid"), + }, + Spec: corev1.PersistentVolumeClaimSpec{ + VolumeName: virtualPVName, + }, + } + + hostTranslator := translate.NewSingleNamespaceTranslator(testingutil.DefaultTestTargetNamespace) + hostTargetName := hostTranslator.HostName( + nil, + virtualTarget.Name, + virtualTarget.Namespace, + ) + hostTarget := &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{ + Name: hostTargetName.Name, + Namespace: testingutil.DefaultTestTargetNamespace, + UID: types.UID("host-target-pvc-uid"), + Annotations: map[string]string{ + translate.NameAnnotation: virtualTarget.Name, + translate.NamespaceAnnotation: virtualTarget.Namespace, + translate.UIDAnnotation: string(virtualTarget.UID), + translate.KindAnnotation: corev1.SchemeGroupVersion.WithKind("PersistentVolumeClaim").String(), + translate.HostNamespaceAnnotation: testingutil.DefaultTestTargetNamespace, + translate.HostNameAnnotation: hostTargetName.Name, + }, + Labels: map[string]string{ + translate.MarkerLabel: translate.VClusterName, + translate.NamespaceLabel: virtualTarget.Namespace, + }, + }, + Status: corev1.PersistentVolumeClaimStatus{ + Phase: corev1.ClaimPending, + Capacity: corev1.ResourceList{}, + }, + } + hostPV := &corev1.PersistentVolume{ + ObjectMeta: metav1.ObjectMeta{Name: virtualPVName}, + Spec: corev1.PersistentVolumeSpec{ + ClaimRef: &corev1.ObjectReference{ + APIVersion: corev1.SchemeGroupVersion.Version, + Kind: "PersistentVolumeClaim", + Namespace: hostTarget.Namespace, + Name: hostTarget.Name, + UID: hostTarget.UID, + }, + }, + Status: corev1.PersistentVolumeStatus{ + Phase: corev1.VolumeBound, + }, + } + + return &fixture{ + virtualTarget: virtualTarget, + virtualPV: virtualPV, + virtualHelper: virtualHelper, + hostTarget: hostTarget, + hostPV: hostPV, + } + } + + tests := []struct { + name string + trigger string + initialVirtual func(*fixture) []runtime.Object + converge func(*testing.T, *synccontext.SyncContext, *fixture) client.Object + }{ + { + name: "guest fake pv cache not found", + trigger: "guest PV create", + initialVirtual: func(f *fixture) []runtime.Object { + return []runtime.Object{f.virtualTarget.DeepCopy()} + }, + converge: func(t *testing.T, ctx *synccontext.SyncContext, f *fixture) client.Object { + t.Helper() + converged := f.virtualPV.DeepCopy() + assert.NilError(t, ctx.VirtualClient.Create(ctx.Context, converged)) + return converged + }, + }, + { + name: "guest helper identity missing", + trigger: "guest helper PVC create", + initialVirtual: func(f *fixture) []runtime.Object { + pvBoundToHelper := f.virtualPV.DeepCopy() + pvBoundToHelper.Spec.ClaimRef = &corev1.ObjectReference{ + Namespace: f.virtualHelper.Namespace, + Name: f.virtualHelper.Name, + UID: f.virtualHelper.UID, + } + return []runtime.Object{f.virtualTarget.DeepCopy(), pvBoundToHelper} + }, + converge: func(t *testing.T, ctx *synccontext.SyncContext, f *fixture) client.Object { + t.Helper() + converged := f.virtualHelper.DeepCopy() + assert.NilError(t, ctx.VirtualClient.Create(ctx.Context, converged)) + return converged + }, + }, + { + name: "pvc pv bound visibility mismatch", + trigger: "guest PV phase update", + initialVirtual: func(f *fixture) []runtime.Object { + pendingPV := f.virtualPV.DeepCopy() + pendingPV.Status.Phase = corev1.VolumePending + return []runtime.Object{f.virtualTarget.DeepCopy(), pendingPV} + }, + converge: func(t *testing.T, ctx *synccontext.SyncContext, f *fixture) client.Object { + t.Helper() + converged := &corev1.PersistentVolume{} + assert.NilError(t, ctx.VirtualClient.Get( + ctx.Context, + types.NamespacedName{Name: f.virtualPV.Name}, + converged, + )) + converged.Status.Phase = corev1.VolumeBound + assert.NilError(t, ctx.VirtualClient.Status().Update(ctx.Context, converged)) + return converged + }, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + f := newFixture() + test := &syncertesting.SyncTest{ + Name: tt.name, + InitialVirtualState: tt.initialVirtual(f), + InitialPhysicalState: []runtime.Object{f.hostTarget.DeepCopy(), f.hostPV.DeepCopy()}, + Sync: func(registerCtx *synccontext.RegisterContext) { + syncCtx, objectSyncer := syncertesting.FakeStartSyncer(t, registerCtx, New) + pvcSyncer := objectSyncer.(*persistentVolumeClaimSyncer) + pvcSyncer.useFakePersistentVolumes = true + + targetKey := types.NamespacedName{ + Namespace: f.virtualTarget.Namespace, + Name: f.virtualTarget.Name, + } + hostTargetKey := types.NamespacedName{ + Namespace: f.hostTarget.Namespace, + Name: f.hostTarget.Name, + } + targetReconciles := 0 + timerRequeues := 0 + mappedTargetRequests := 0 + + reconcileTarget := func() ctrl.Result { + t.Helper() + currentVirtual := &corev1.PersistentVolumeClaim{} + assert.NilError(t, syncCtx.VirtualClient.Get(syncCtx.Context, targetKey, currentVirtual)) + currentHost := &corev1.PersistentVolumeClaim{} + assert.NilError(t, syncCtx.HostClient.Get(syncCtx.Context, hostTargetKey, currentHost)) + result, err := pvcSyncer.Sync(syncCtx, synccontext.NewSyncEventWithOld( + currentHost.DeepCopy(), + currentHost.DeepCopy(), + currentVirtual.DeepCopy(), + currentVirtual.DeepCopy(), + )) + assert.NilError(t, err) + targetReconciles++ + if result.Requeue || result.RequeueAfter > 0 { + timerRequeues++ + } + return result + } + + // Pin the exact pre-gate return before exercising the full Sync + // boundary: every fixture must be a transient (nil, false, nil), + // rather than the genuine non-applicable entry gate. + currentVirtual := &corev1.PersistentVolumeClaim{} + assert.NilError(t, syncCtx.VirtualClient.Get(syncCtx.Context, targetKey, currentVirtual)) + currentHost := &corev1.PersistentVolumeClaim{} + assert.NilError(t, syncCtx.HostClient.Get(syncCtx.Context, hostTargetKey, currentHost)) + preGatePV, preGateReady, err := pvcSyncer.externalPopulatorPersistentVolume( + syncCtx, + currentHost, + currentVirtual, + ) + assert.NilError(t, err) + if preGatePV != nil || preGateReady { + t.Fatalf( + "%s did not enter a transient (nil, false, nil) pre-gate: pv=%#v ready=%t", + tt.trigger, + preGatePV, + preGateReady, + ) + } + + // The first target reconcile observes that transient and + // intentionally schedules no polling timer. + firstResult := reconcileTarget() + assert.Check(t, firstResult.IsZero(), "%s first result was %#v", tt.trigger, firstResult) + + virtualBeforeDependency := &corev1.PersistentVolumeClaim{} + assert.NilError(t, syncCtx.VirtualClient.Get(syncCtx.Context, targetKey, virtualBeforeDependency)) + convergedDependency := tt.converge(t, syncCtx, f) + virtualAfterDependency := &corev1.PersistentVolumeClaim{} + assert.NilError(t, syncCtx.VirtualClient.Get(syncCtx.Context, targetKey, virtualAfterDependency)) + assert.Assert( + t, + apiequality.Semantic.DeepEqual(virtualAfterDependency, virtualBeforeDependency), + "dependency convergence changed the target PVC before its mapped reconcile", + ) + + // Model only requests produced by the dependency event. A PVC's + // primary watch would enqueue the helper key, not this target key. + mapper, hasDependencyMapper := any(pvcSyncer).(dependencyRequestMapper) + _, hasControllerModifier := any(pvcSyncer).(syncertypes.ControllerModifier) + var requests []ctrl.Request + if hasDependencyMapper { + requests = mapper.externalPopulatorDependencyRequests(syncCtx.Context, convergedDependency) + } + for _, request := range requests { + if request.NamespacedName != targetKey { + t.Fatalf( + "%s mapped unexpected request %s/%s; want target %s/%s", + tt.trigger, + request.Namespace, + request.Name, + targetKey.Namespace, + targetKey.Name, + ) + } + mappedTargetRequests++ + reconcileTarget() + } + + actualHostTarget := &corev1.PersistentVolumeClaim{} + assert.NilError(t, syncCtx.HostClient.Get(syncCtx.Context, hostTargetKey, actualHostTarget)) + actualHostPV := &corev1.PersistentVolume{} + assert.NilError(t, syncCtx.HostClient.Get( + syncCtx.Context, + types.NamespacedName{Name: f.hostPV.Name}, + actualHostPV, + )) + actualVirtualTarget := &corev1.PersistentVolumeClaim{} + assert.NilError(t, syncCtx.VirtualClient.Get(syncCtx.Context, targetKey, actualVirtualTarget)) + + if timerRequeues != 0 || + !hasControllerModifier || + mappedTargetRequests != 1 || + targetReconciles != 2 || + actualHostTarget.Spec.VolumeName != f.virtualPV.Name || + !claimRefMatchesPersistentVolumeClaim(actualHostPV.Spec.ClaimRef, actualHostTarget) || + actualVirtualTarget.Status.Phase != corev1.ClaimBound { + t.Fatalf( + "%s convergence contract: timer_requeues=%d controller_modifier=%t mapped_target_requests=%d target_reconciles=%d host_volume=%q host_claim_ref=%#v virtual_phase=%q; want 0/true/1/2/%q/target/%q", + tt.trigger, + timerRequeues, + hasControllerModifier, + mappedTargetRequests, + targetReconciles, + actualHostTarget.Spec.VolumeName, + actualHostPV.Spec.ClaimRef, + actualVirtualTarget.Status.Phase, + f.virtualPV.Name, + corev1.ClaimBound, + ) + } + }, + } + test.Run(t, syncertesting.NewFakeRegisterContext) + }) + } +} + func TestTranslateSelectorPreservesNilAndExplicitEmptyStorageClass(t *testing.T) { empty := "" tests := []struct { From f2206ed0319f9bf519dc4eeeebe85724f5968a78 Mon Sep 17 00:00:00 2001 From: Gina Date: Fri, 24 Jul 2026 22:06:27 +0800 Subject: [PATCH 4/5] fix(pvc): wake external populator targets on dependency changes --- .../persistentvolumeclaims/syncer.go | 281 ++++++- .../persistentvolumeclaims/syncer_test.go | 775 +++++++++++++++++- pkg/patcher/patcher.go | 15 + pkg/patcher/patcher_test.go | 73 ++ 4 files changed, 1141 insertions(+), 3 deletions(-) create mode 100644 pkg/patcher/patcher_test.go diff --git a/pkg/controllers/resources/persistentvolumeclaims/syncer.go b/pkg/controllers/resources/persistentvolumeclaims/syncer.go index 0b59d3de6a..b41e764a42 100644 --- a/pkg/controllers/resources/persistentvolumeclaims/syncer.go +++ b/pkg/controllers/resources/persistentvolumeclaims/syncer.go @@ -1,6 +1,7 @@ package persistentvolumeclaims import ( + "context" "fmt" "strings" "time" @@ -27,9 +28,13 @@ import ( kerrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/types" + "k8s.io/client-go/util/workqueue" "k8s.io/component-helpers/scheduling/corev1/nodeaffinity" ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/builder" "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/event" + "sigs.k8s.io/controller-runtime/pkg/handler" "github.com/loft-sh/vcluster/pkg/util/loghelper" ) @@ -50,6 +55,7 @@ const ( dataProtectionBackupKind = "Backup" externalPopulatorPopulateHelperPrefix = "kb-populate-" + externalPopulatorPVCByUIDIndex = "externalPopulatorPVCByUID" externalPopulatorRestoreConditionType = corev1.PersistentVolumeClaimConditionType("Restore") externalPopulatorPopulateConditionType = corev1.PersistentVolumeClaimConditionType("Populating") @@ -78,6 +84,7 @@ func New(ctx *synccontext.RegisterContext) (syncertypes.Object, error) { storageClassesEnabled: ctx.Config.Sync.ToHost.StorageClasses.Enabled, schedulerEnabled: ctx.Config.SchedulingInVirtualClusterEnabled(), useFakePersistentVolumes: !ctx.Config.Sync.ToHost.PersistentVolumes.Enabled, + virtualClient: ctx.VirtualManager.GetClient(), }, nil } @@ -90,6 +97,68 @@ type persistentVolumeClaimSyncer struct { storageClassesEnabled bool schedulerEnabled bool useFakePersistentVolumes bool + virtualClient client.Client +} + +var _ syncertypes.IndicesRegisterer = &persistentVolumeClaimSyncer{} + +func (s *persistentVolumeClaimSyncer) RegisterIndices(ctx *synccontext.RegisterContext) error { + return ctx.VirtualManager.GetFieldIndexer().IndexField(ctx, &corev1.PersistentVolumeClaim{}, externalPopulatorPVCByUIDIndex, func(rawObj client.Object) []string { + pvc, ok := rawObj.(*corev1.PersistentVolumeClaim) + if !ok || pvc.UID == "" { + return nil + } + + return []string{string(pvc.UID)} + }) +} + +var _ syncertypes.ControllerModifier = &persistentVolumeClaimSyncer{} + +func (s *persistentVolumeClaimSyncer) ModifyController(_ *synccontext.RegisterContext, controllerBuilder *builder.Builder) (*builder.Builder, error) { + enqueueDependency := func( + ctx context.Context, + object client.Object, + allowDeletingDependency bool, + queue workqueue.TypedRateLimitingInterface[ctrl.Request], + ) { + requests, err := s.externalPopulatorDependencyRequestsWithOptions(ctx, object, allowDeletingDependency) + if err != nil { + ctrl.LoggerFrom(ctx).Error( + err, + "failed to map external populator dependency", + "kind", + fmt.Sprintf("%T", object), + "namespace", + object.GetNamespace(), + "name", + object.GetName(), + ) + return + } + + for _, request := range requests { + queue.Add(request) + } + } + dependencyHandler := &handler.Funcs{ + CreateFunc: func(ctx context.Context, createEvent event.CreateEvent, queue workqueue.TypedRateLimitingInterface[ctrl.Request]) { + enqueueDependency(ctx, createEvent.Object, false, queue) + }, + UpdateFunc: func(ctx context.Context, updateEvent event.UpdateEvent, queue workqueue.TypedRateLimitingInterface[ctrl.Request]) { + enqueueDependency(ctx, updateEvent.ObjectNew, false, queue) + }, + DeleteFunc: func(ctx context.Context, deleteEvent event.DeleteEvent, queue workqueue.TypedRateLimitingInterface[ctrl.Request]) { + enqueueDependency(ctx, deleteEvent.Object, true, queue) + }, + GenericFunc: func(ctx context.Context, genericEvent event.GenericEvent, queue workqueue.TypedRateLimitingInterface[ctrl.Request]) { + enqueueDependency(ctx, genericEvent.Object, false, queue) + }, + } + + return controllerBuilder. + Watches(&corev1.PersistentVolume{}, dependencyHandler). + Watches(&corev1.PersistentVolumeClaim{}, dependencyHandler), nil } var _ syncertypes.OptionsProvider = &persistentVolumeClaimSyncer{} @@ -266,7 +335,7 @@ func (s *persistentVolumeClaimSyncer) Sync(ctx *synccontext.SyncContext, event * retErr = utilerrors.NewAggregate([]error{retErr, err}) } - if kerrors.IsConflict(retErr) { + if containsConflictError(retErr) { result = ctrl.Result{RequeueAfter: time.Second} retErr = nil return @@ -294,10 +363,20 @@ func (s *persistentVolumeClaimSyncer) Sync(ctx *synccontext.SyncContext, event * return ctrl.Result{}, err } if preserveVirtualStatus { + hostResourceVersionBeforeMaterialization := event.Host.ResourceVersion hostConverged, err := s.ensureExternalPopulatorHostMaterialization(ctx, event.Host, event.Virtual, vPV) if err != nil { return ctrl.Result{}, err } + // A successful optimistic-lock materialization patch returns a new + // resourceVersion in event.Host. Rebase only when that direct write or a + // fresher target read changed the snapshot captured by NewSyncerPatcher. + if event.Host.ResourceVersion != hostResourceVersionBeforeMaterialization { + err = patch.RebaseHost(event.Host) + if err != nil { + return ctrl.Result{}, err + } + } if hostConverged { ensureExternalPopulatorVirtualPopulateStatus(event.Virtual, vPV) } else { @@ -326,6 +405,24 @@ func (s *persistentVolumeClaimSyncer) Sync(ctx *synccontext.SyncContext, event * return result, nil } +func containsConflictError(err error) bool { + if err == nil { + return false + } + if kerrors.IsConflict(err) { + return true + } + if aggregate, ok := err.(utilerrors.Aggregate); ok { + for _, aggregateErr := range aggregate.Errors() { + if containsConflictError(aggregateErr) { + return true + } + } + } + + return false +} + func (s *persistentVolumeClaimSyncer) SyncToVirtual(ctx *synccontext.SyncContext, event *synccontext.SyncToVirtualEvent[*corev1.PersistentVolumeClaim]) (_ ctrl.Result, retErr error) { if event.VirtualOld != nil || translate.ShouldDeleteHostObject(event.Host) { // the virtual populate helper PVC can be fully gone (finalizer removed) while @@ -1016,6 +1113,188 @@ func claimRefReferencesPersistentVolumeClaim(ref *corev1.ObjectReference, pvc *c return ref.Namespace == pvc.Namespace && ref.Name == pvc.Name } +func (s *persistentVolumeClaimSyncer) externalPopulatorDependencyRequests(ctx context.Context, object client.Object) ([]ctrl.Request, error) { + return s.externalPopulatorDependencyRequestsWithOptions(ctx, object, false) +} + +func (s *persistentVolumeClaimSyncer) externalPopulatorDependencyRequestsWithOptions( + ctx context.Context, + object client.Object, + allowDeletingDependency bool, +) ([]ctrl.Request, error) { + switch dependency := object.(type) { + case *corev1.PersistentVolume: + return s.externalPopulatorPersistentVolumeRequests(ctx, dependency, allowDeletingDependency) + case *corev1.PersistentVolumeClaim: + return s.externalPopulatorHelperPVCRequests(ctx, dependency, allowDeletingDependency) + default: + return nil, nil + } +} + +func (s *persistentVolumeClaimSyncer) externalPopulatorPersistentVolumeRequests( + ctx context.Context, + vPV *corev1.PersistentVolume, + allowDeletingDependency bool, +) ([]ctrl.Request, error) { + if vPV == nil || + vPV.Spec.ClaimRef == nil || + vPV.Spec.ClaimRef.Namespace == "" || + vPV.Spec.ClaimRef.Name == "" || + vPV.Spec.ClaimRef.UID == "" { + return nil, nil + } + + ref := vPV.Spec.ClaimRef + referencedPVC := &corev1.PersistentVolumeClaim{} + err := s.virtualClient.Get(ctx, types.NamespacedName{Namespace: ref.Namespace, Name: ref.Name}, referencedPVC) + if err != nil { + if kerrors.IsNotFound(err) { + return nil, nil + } + return nil, err + } + + if isExternalPopulatorWakeTarget(referencedPVC) && + referencedPVC.Spec.VolumeName == vPV.Name && + isExternalPopulatorPersistentVolumeForPVC(vPV, referencedPVC, false) && + vPV.Spec.ClaimRef.UID == referencedPVC.UID { + return externalPopulatorTargetRequest(referencedPVC), nil + } + + if referencedPVC.Spec.VolumeName != vPV.Name || + referencedPVC.UID == "" || + !isExternalPopulatorPersistentVolumeForPVC(vPV, referencedPVC, false) || + vPV.Spec.ClaimRef.UID != referencedPVC.UID { + return nil, nil + } + + target, found, err := s.externalPopulatorTargetForHelper(ctx, referencedPVC, allowDeletingDependency) + if err != nil || !found { + return nil, err + } + + return externalPopulatorTargetRequest(target), nil +} + +func (s *persistentVolumeClaimSyncer) externalPopulatorHelperPVCRequests( + ctx context.Context, + helperPVC *corev1.PersistentVolumeClaim, + allowDeletingDependency bool, +) ([]ctrl.Request, error) { + if helperPVC == nil || + (!allowDeletingDependency && helperPVC.DeletionTimestamp != nil) || + helperPVC.Spec.VolumeName == "" || + !strings.HasPrefix(helperPVC.Name, externalPopulatorPopulateHelperPrefix) { + return nil, nil + } + + vPV := &corev1.PersistentVolume{} + err := s.virtualClient.Get(ctx, types.NamespacedName{Name: helperPVC.Spec.VolumeName}, vPV) + if err != nil { + if kerrors.IsNotFound(err) { + return nil, nil + } + return nil, err + } + if helperPVC.UID == "" || + !isExternalPopulatorPersistentVolumeForPVC(vPV, helperPVC, false) || + vPV.Spec.ClaimRef.UID != helperPVC.UID { + return nil, nil + } + + target, found, err := s.externalPopulatorTargetForHelper(ctx, helperPVC, allowDeletingDependency) + if err != nil || !found { + return nil, err + } + + return externalPopulatorTargetRequest(target), nil +} + +func (s *persistentVolumeClaimSyncer) externalPopulatorTargetForHelper( + ctx context.Context, + helperPVC *corev1.PersistentVolumeClaim, + allowDeletingDependency bool, +) (*corev1.PersistentVolumeClaim, bool, error) { + if helperPVC == nil || + (!allowDeletingDependency && helperPVC.DeletionTimestamp != nil) || + helperPVC.Spec.VolumeName == "" || + !strings.HasPrefix(helperPVC.Name, externalPopulatorPopulateHelperPrefix) { + return nil, false, nil + } + + targetUID := strings.TrimPrefix(helperPVC.Name, externalPopulatorPopulateHelperPrefix) + if targetUID == "" { + return nil, false, nil + } + + targets := &corev1.PersistentVolumeClaimList{} + err := s.virtualClient.List( + ctx, + targets, + client.InNamespace(helperPVC.Namespace), + client.MatchingFields{externalPopulatorPVCByUIDIndex: targetUID}, + ) + if err != nil { + return nil, false, err + } + + var match *corev1.PersistentVolumeClaim + for i := range targets.Items { + target := &targets.Items[i] + if !isExternalPopulatorWakeTarget(target) || + target.Spec.VolumeName != helperPVC.Spec.VolumeName || + !isExternalPopulatorWakeHelperForTarget(helperPVC, target, allowDeletingDependency) { + continue + } + if match != nil { + return nil, false, fmt.Errorf( + "multiple external populator target persistent volume claims match helper pvc %s/%s", + helperPVC.Namespace, + helperPVC.Name, + ) + } + match = target.DeepCopy() + } + if match == nil { + return nil, false, nil + } + + return match, true, nil +} + +func isExternalPopulatorWakeHelperForTarget( + helperPVC, targetPVC *corev1.PersistentVolumeClaim, + allowDeletingDependency bool, +) bool { + if !allowDeletingDependency || helperPVC == nil || helperPVC.DeletionTimestamp == nil { + return isExternalPopulatorPopulateHelperPVCForTarget(helperPVC, targetPVC) + } + if targetPVC == nil || targetPVC.UID == "" { + return false + } + + return helperPVC.Namespace == targetPVC.Namespace && + helperPVC.Name == externalPopulatorPopulateHelperPrefix+string(targetPVC.UID) +} + +func isExternalPopulatorWakeTarget(pvc *corev1.PersistentVolumeClaim) bool { + return pvc != nil && + pvc.DeletionTimestamp == nil && + pvc.UID != "" && + pvc.Spec.VolumeName != "" && + isExternalPopulatorPVC(pvc) +} + +func externalPopulatorTargetRequest(pvc *corev1.PersistentVolumeClaim) []ctrl.Request { + return []ctrl.Request{{ + NamespacedName: types.NamespacedName{ + Namespace: pvc.Namespace, + Name: pvc.Name, + }, + }} +} + func (s *persistentVolumeClaimSyncer) externalPopulatorPersistentVolume(ctx *synccontext.SyncContext, pObj, vObj *corev1.PersistentVolumeClaim) (*corev1.PersistentVolume, bool, error) { if !isExternalPopulatorPVC(vObj) || vObj.Spec.VolumeName == "" || !isHostPVCWaitingForVolume(pObj) { return nil, false, nil diff --git a/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go b/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go index dd42a52429..dba5e6bc0b 100644 --- a/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go +++ b/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go @@ -2,11 +2,15 @@ package persistentvolumeclaims import ( "context" + "errors" + "reflect" "strings" + "sync" "testing" "time" "github.com/loft-sh/vcluster/pkg/config" + "github.com/loft-sh/vcluster/pkg/scheme" "github.com/loft-sh/vcluster/pkg/syncer/synccontext" syncertesting "github.com/loft-sh/vcluster/pkg/syncer/testing" syncertypes "github.com/loft-sh/vcluster/pkg/syncer/types" @@ -24,8 +28,11 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/runtime/schema" + toolscache "k8s.io/client-go/tools/cache" ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/cache" "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/reconcile" ) type eventRecordingTranslator struct { @@ -76,9 +83,221 @@ func assertNoPVCEvent(t *testing.T, recorder *events.FakeRecorder) { } } +type externalPopulatorDependencyFixture struct { + target *corev1.PersistentVolumeClaim + helper *corev1.PersistentVolumeClaim + pv *corev1.PersistentVolume +} + +func newExternalPopulatorDependencyFixture() *externalPopulatorDependencyFixture { + const ( + namespace = "testns" + pvName = "restore-populated-pv" + ) + targetUID := types.UID("target-pvc-uid") + helperUID := types.UID("populate-helper-pvc-uid") + dataProtectionGroup := dataProtectionAPIGroup + target := &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{ + Name: "target", + Namespace: namespace, + UID: targetUID, + }, + Spec: corev1.PersistentVolumeClaimSpec{ + VolumeName: pvName, + DataSourceRef: &corev1.TypedObjectReference{ + APIGroup: &dataProtectionGroup, + Kind: dataProtectionBackupKind, + Name: "backup", + }, + }, + } + helper := &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{ + Name: externalPopulatorPopulateHelperPrefix + string(targetUID), + Namespace: namespace, + UID: helperUID, + }, + Spec: corev1.PersistentVolumeClaimSpec{ + VolumeName: pvName, + }, + } + pv := &corev1.PersistentVolume{ + ObjectMeta: metav1.ObjectMeta{Name: pvName}, + Spec: corev1.PersistentVolumeSpec{ + ClaimRef: &corev1.ObjectReference{ + Namespace: namespace, + Name: helper.Name, + UID: helperUID, + }, + }, + Status: corev1.PersistentVolumeStatus{Phase: corev1.VolumeBound}, + } + + return &externalPopulatorDependencyFixture{ + target: target, + helper: helper, + pv: pv, + } +} + +func newExternalPopulatorMapperTestSyncer( + t *testing.T, + virtualObjects ...runtime.Object, +) (*synccontext.SyncContext, *persistentVolumeClaimSyncer) { + t.Helper() + + hostClient := testingutil.NewFakeClient(scheme.Scheme) + virtualClient := testingutil.NewFakeClient(scheme.Scheme, virtualObjects...) + registerCtx := syncertesting.NewFakeRegisterContext(testingutil.NewFakeConfig(), hostClient, virtualClient) + syncCtx, objectSyncer := syncertesting.FakeStartSyncer(t, registerCtx, New) + return syncCtx, objectSyncer.(*persistentVolumeClaimSyncer) +} + +type externalPopulatorMapperErrorClient struct { + client.Client + getErr error + listErr error +} + +func (c *externalPopulatorMapperErrorClient) Get(ctx context.Context, key client.ObjectKey, obj client.Object, opts ...client.GetOption) error { + if c.getErr != nil { + return c.getErr + } + return c.Client.Get(ctx, key, obj, opts...) +} + +func (c *externalPopulatorMapperErrorClient) List(ctx context.Context, list client.ObjectList, opts ...client.ListOption) error { + if c.listErr != nil { + return c.listErr + } + return c.Client.List(ctx, list, opts...) +} + +type externalPopulatorConcurrentWriterClient struct { + client.Client + target types.NamespacedName + + once sync.Once + directRV string + concurrentRV string + concurrentErr error +} + +func (c *externalPopulatorConcurrentWriterClient) Patch( + ctx context.Context, + obj client.Object, + patch client.Patch, + opts ...client.PatchOption, +) error { + if err := c.Client.Patch(ctx, obj, patch, opts...); err != nil { + return err + } + + pvc, ok := obj.(*corev1.PersistentVolumeClaim) + if !ok || + pvc.Namespace != c.target.Namespace || + pvc.Name != c.target.Name || + pvc.Spec.VolumeName == "" { + return nil + } + + c.once.Do(func() { + c.directRV = pvc.ResourceVersion + current := &corev1.PersistentVolumeClaim{} + c.concurrentErr = c.Client.Get(ctx, c.target, current) + if c.concurrentErr != nil { + return + } + if current.Annotations == nil { + current.Annotations = map[string]string{} + } + current.Annotations["example.test/concurrent-writer"] = "preserve" + current.Spec.Resources.Requests = corev1.ResourceList{ + corev1.ResourceStorage: resource.MustParse("2Gi"), + } + c.concurrentErr = c.Client.Update(ctx, current) + if c.concurrentErr == nil { + c.concurrentRV = current.ResourceVersion + } + }) + + return nil +} + +type externalPopulatorRecordingManager struct { + ctrl.Manager + recordingCache cache.Cache +} + +func (m *externalPopulatorRecordingManager) GetCache() cache.Cache { + return m.recordingCache +} + +type externalPopulatorRecordingCache struct { + cache.Cache + + mu sync.Mutex + handlers map[reflect.Type][]toolscache.ResourceEventHandler +} + +func newExternalPopulatorRecordingCache(delegate cache.Cache) *externalPopulatorRecordingCache { + return &externalPopulatorRecordingCache{ + Cache: delegate, + handlers: map[reflect.Type][]toolscache.ResourceEventHandler{}, + } +} + +func (c *externalPopulatorRecordingCache) GetInformer( + ctx context.Context, + object client.Object, + opts ...cache.InformerGetOption, +) (cache.Informer, error) { + informer, err := c.Cache.GetInformer(ctx, object, opts...) + if err != nil { + return nil, err + } + + objectType := reflect.TypeOf(object) + return &externalPopulatorRecordingInformer{ + Informer: informer, + record: func(handler toolscache.ResourceEventHandler) { + c.mu.Lock() + defer c.mu.Unlock() + c.handlers[objectType] = append(c.handlers[objectType], handler) + }, + }, nil +} + +func (c *externalPopulatorRecordingCache) handlersFor(object client.Object) []toolscache.ResourceEventHandler { + c.mu.Lock() + defer c.mu.Unlock() + return append([]toolscache.ResourceEventHandler(nil), c.handlers[reflect.TypeOf(object)]...) +} + +type externalPopulatorRecordingInformer struct { + cache.Informer + record func(toolscache.ResourceEventHandler) +} + +func (i *externalPopulatorRecordingInformer) AddEventHandler(handler toolscache.ResourceEventHandler) (toolscache.ResourceEventHandlerRegistration, error) { + i.record(handler) + return i.Informer.AddEventHandler(handler) +} + +func (i *externalPopulatorRecordingInformer) AddEventHandlerWithResyncPeriod(handler toolscache.ResourceEventHandler, resyncPeriod time.Duration) (toolscache.ResourceEventHandlerRegistration, error) { + i.record(handler) + return i.Informer.AddEventHandlerWithResyncPeriod(handler, resyncPeriod) +} + +func (i *externalPopulatorRecordingInformer) AddEventHandlerWithOptions(handler toolscache.ResourceEventHandler, options toolscache.HandlerOptions) (toolscache.ResourceEventHandlerRegistration, error) { + i.record(handler) + return i.Informer.AddEventHandlerWithOptions(handler, options) +} + func TestSyncExternalPopulatorPreGateTransientSchedulesTarget(t *testing.T) { type dependencyRequestMapper interface { - externalPopulatorDependencyRequests(context.Context, client.Object) []ctrl.Request + externalPopulatorDependencyRequests(context.Context, client.Object) ([]ctrl.Request, error) } type fixture struct { virtualTarget *corev1.PersistentVolumeClaim @@ -345,7 +564,9 @@ func TestSyncExternalPopulatorPreGateTransientSchedulesTarget(t *testing.T) { _, hasControllerModifier := any(pvcSyncer).(syncertypes.ControllerModifier) var requests []ctrl.Request if hasDependencyMapper { - requests = mapper.externalPopulatorDependencyRequests(syncCtx.Context, convergedDependency) + var err error + requests, err = mapper.externalPopulatorDependencyRequests(syncCtx.Context, convergedDependency) + assert.NilError(t, err) } for _, request := range requests { if request.NamespacedName != targetKey { @@ -401,6 +622,556 @@ func TestSyncExternalPopulatorPreGateTransientSchedulesTarget(t *testing.T) { } } +func TestExternalPopulatorDependencyMapperRejectsInvalidDependencies(t *testing.T) { + tests := []struct { + name string + prepare func(*externalPopulatorDependencyFixture) client.Object + }{ + { + name: "non-applicable target", + prepare: func(f *externalPopulatorDependencyFixture) client.Object { + f.target.Spec.DataSourceRef.Kind = "PersistentVolumeClaim" + f.pv.Spec.ClaimRef = &corev1.ObjectReference{ + Namespace: f.target.Namespace, + Name: f.target.Name, + UID: f.target.UID, + } + return f.pv + }, + }, + { + name: "terminating helper", + prepare: func(f *externalPopulatorDependencyFixture) client.Object { + now := metav1.Now() + f.helper.DeletionTimestamp = &now + f.helper.Finalizers = []string{"test"} + return f.helper + }, + }, + { + name: "invalid helper identity", + prepare: func(f *externalPopulatorDependencyFixture) client.Object { + f.helper.Name = externalPopulatorPopulateHelperPrefix + "other-target" + f.pv.Spec.ClaimRef.Name = f.helper.Name + return f.helper + }, + }, + { + name: "helper claimRef UID mismatch", + prepare: func(f *externalPopulatorDependencyFixture) client.Object { + f.pv.Spec.ClaimRef.UID = types.UID("stale-helper-uid") + return f.helper + }, + }, + { + name: "helper claimRef UID absent", + prepare: func(f *externalPopulatorDependencyFixture) client.Object { + f.pv.Spec.ClaimRef.UID = "" + return f.helper + }, + }, + { + name: "target claimRef UID mismatch", + prepare: func(f *externalPopulatorDependencyFixture) client.Object { + f.pv.Spec.ClaimRef = &corev1.ObjectReference{ + Namespace: f.target.Namespace, + Name: f.target.Name, + UID: types.UID("stale-target-uid"), + } + return f.pv + }, + }, + { + name: "target claimRef UID absent", + prepare: func(f *externalPopulatorDependencyFixture) client.Object { + f.pv.Spec.ClaimRef = &corev1.ObjectReference{ + Namespace: f.target.Namespace, + Name: f.target.Name, + } + return f.pv + }, + }, + { + name: "helper absent", + prepare: func(f *externalPopulatorDependencyFixture) client.Object { + return f.pv + }, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + f := newExternalPopulatorDependencyFixture() + event := tt.prepare(f) + virtualObjects := []runtime.Object{f.target.DeepCopy(), f.pv.DeepCopy()} + if tt.name != "helper absent" { + virtualObjects = append(virtualObjects, f.helper.DeepCopy()) + } + syncCtx, pvcSyncer := newExternalPopulatorMapperTestSyncer(t, virtualObjects...) + + requests, err := pvcSyncer.externalPopulatorDependencyRequests(syncCtx.Context, event) + assert.NilError(t, err) + assert.Equal(t, len(requests), 0) + }) + } +} + +func TestExternalPopulatorDependencyMapperPropagatesLookupErrors(t *testing.T) { + tests := []struct { + name string + getErr error + listErr error + }{ + { + name: "persistent volume lookup", + getErr: errors.New("injected persistent volume lookup failure"), + }, + { + name: "target index lookup", + listErr: errors.New("injected target index lookup failure"), + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + f := newExternalPopulatorDependencyFixture() + syncCtx, pvcSyncer := newExternalPopulatorMapperTestSyncer( + t, + f.target.DeepCopy(), + f.helper.DeepCopy(), + f.pv.DeepCopy(), + ) + pvcSyncer.virtualClient = &externalPopulatorMapperErrorClient{ + Client: syncCtx.VirtualClient, + getErr: tt.getErr, + listErr: tt.listErr, + } + + requests, err := pvcSyncer.externalPopulatorDependencyRequests(syncCtx.Context, f.helper) + assert.Equal(t, len(requests), 0) + expectedErr := tt.getErr + if expectedErr == nil { + expectedErr = tt.listErr + } + if !errors.Is(err, expectedErr) { + t.Fatalf("mapper error = %v, want wrapped %v", err, expectedErr) + } + }) + } +} + +func TestExternalPopulatorDependencyMapperIgnoresUnrelatedOrIncompleteEventsBeforeLookup(t *testing.T) { + f := newExternalPopulatorDependencyFixture() + syncCtx, pvcSyncer := newExternalPopulatorMapperTestSyncer( + t, + f.target.DeepCopy(), + f.helper.DeepCopy(), + f.pv.DeepCopy(), + ) + pvcSyncer.virtualClient = &externalPopulatorMapperErrorClient{ + Client: syncCtx.VirtualClient, + getErr: errors.New("ordinary pvc must not trigger a dependency lookup"), + } + ordinaryPVC := f.target.DeepCopy() + ordinaryPVC.Name = "ordinary" + ordinaryPVC.Spec.DataSourceRef = nil + + incompletePV := f.pv.DeepCopy() + incompletePV.Spec.ClaimRef.UID = "" + + for _, tt := range []struct { + name string + object client.Object + }{ + {name: "ordinary pvc", object: ordinaryPVC}, + {name: "persistent volume missing claim UID", object: incompletePV}, + } { + t.Run(tt.name, func(t *testing.T) { + requests, err := pvcSyncer.externalPopulatorDependencyRequests(syncCtx.Context, tt.object) + assert.NilError(t, err) + assert.Equal(t, len(requests), 0) + }) + } +} + +func TestExternalPopulatorDependencyModifyControllerWiring(t *testing.T) { + f := newExternalPopulatorDependencyFixture() + hostClient := testingutil.NewFakeClient(scheme.Scheme) + virtualClient := testingutil.NewFakeClient( + scheme.Scheme, + f.target.DeepCopy(), + f.helper.DeepCopy(), + f.pv.DeepCopy(), + ) + registerCtx := syncertesting.NewFakeRegisterContext(testingutil.NewFakeConfig(), hostClient, virtualClient) + baseManager := testingutil.NewFakeManager(virtualClient) + recordingCache := newExternalPopulatorRecordingCache(baseManager.GetCache()) + recordingManager := &externalPopulatorRecordingManager{ + Manager: baseManager, + recordingCache: recordingCache, + } + registerCtx.VirtualManager = recordingManager + + objectSyncer, err := New(registerCtx) + assert.NilError(t, err) + pvcSyncer := objectSyncer.(*persistentVolumeClaimSyncer) + assert.NilError(t, pvcSyncer.RegisterIndices(registerCtx)) + + requests := make(chan ctrl.Request, 32) + controllerBuilder := ctrl.NewControllerManagedBy(recordingManager). + Named("external-populator-dependency-wiring-test"). + For(&corev1.PersistentVolumeClaim{}) + controllerBuilder, err = pvcSyncer.ModifyController(registerCtx, controllerBuilder) + assert.NilError(t, err) + controller, err := controllerBuilder.Build(reconcile.Func(func(_ context.Context, request ctrl.Request) (ctrl.Result, error) { + requests <- request + return ctrl.Result{}, nil + })) + assert.NilError(t, err) + + runCtx, cancel := context.WithCancel(context.Background()) + defer cancel() + started := make(chan error, 1) + go func() { + started <- controller.Start(runCtx) + }() + + waitDeadline := time.NewTimer(2 * time.Second) + defer waitDeadline.Stop() + waitTicker := time.NewTicker(5 * time.Millisecond) + defer waitTicker.Stop() + for { + if len(recordingCache.handlersFor(&corev1.PersistentVolume{})) == 1 && + len(recordingCache.handlersFor(&corev1.PersistentVolumeClaim{})) == 2 { + break + } + select { + case err := <-started: + t.Fatalf("controller stopped before watch registration: %v", err) + case <-waitDeadline.C: + t.Fatalf( + "watch registration timed out: pv handlers=%d pvc handlers=%d", + len(recordingCache.handlersFor(&corev1.PersistentVolume{})), + len(recordingCache.handlersFor(&corev1.PersistentVolumeClaim{})), + ) + case <-waitTicker.C: + } + } + + targetKey := types.NamespacedName{Namespace: f.target.Namespace, Name: f.target.Name} + drainRequests := func() { + for { + select { + case <-requests: + default: + return + } + } + } + expectTarget := func(label string) { + t.Helper() + deadline := time.NewTimer(2 * time.Second) + defer deadline.Stop() + for { + select { + case request := <-requests: + if request.NamespacedName == targetKey { + return + } + case <-deadline.C: + t.Fatalf("%s did not enqueue exact target %s", label, targetKey) + } + } + } + expectNoTarget := func(label string) { + t.Helper() + deadline := time.NewTimer(150 * time.Millisecond) + defer deadline.Stop() + for { + select { + case request := <-requests: + if request.NamespacedName == targetKey { + t.Fatalf("%s unexpectedly enqueued exact target %s", label, targetKey) + } + case <-deadline.C: + return + } + } + } + fireAdd := func(object client.Object) { + t.Helper() + handlers := recordingCache.handlersFor(object) + if len(handlers) == 0 { + t.Fatalf("no registered handler for %T", object) + } + for _, eventHandler := range handlers { + eventHandler.OnAdd(object, false) + } + } + fireUpdate := func(oldObject, newObject client.Object) { + t.Helper() + handlers := recordingCache.handlersFor(newObject) + if len(handlers) == 0 { + t.Fatalf("no registered handler for %T", newObject) + } + for _, eventHandler := range handlers { + eventHandler.OnUpdate(oldObject, newObject) + } + } + fireDelete := func(object client.Object) { + t.Helper() + handlers := recordingCache.handlersFor(object) + if len(handlers) == 0 { + t.Fatalf("no registered handler for %T", object) + } + for _, eventHandler := range handlers { + eventHandler.OnDelete(object) + } + } + + directPV := f.pv.DeepCopy() + directPV.Spec.ClaimRef = &corev1.ObjectReference{ + Namespace: f.target.Namespace, + Name: f.target.Name, + UID: f.target.UID, + } + fireAdd(directPV) + expectTarget("PV create") + drainRequests() + pendingDirectPV := directPV.DeepCopy() + pendingDirectPV.Status.Phase = corev1.VolumePending + fireUpdate(pendingDirectPV, directPV) + expectTarget("PV update") + drainRequests() + fireDelete(directPV) + expectTarget("PV delete") + drainRequests() + + fireAdd(f.helper.DeepCopy()) + expectTarget("helper create") + drainRequests() + updatedHelper := f.helper.DeepCopy() + updatedHelper.Labels = map[string]string{"generation": "2"} + fireUpdate(f.helper.DeepCopy(), updatedHelper) + expectTarget("helper update") + drainRequests() + + assert.NilError(t, virtualClient.Delete(runCtx, f.helper.DeepCopy())) + deletedHelper := f.helper.DeepCopy() + now := metav1.Now() + deletedHelper.DeletionTimestamp = &now + fireDelete(deletedHelper) + expectTarget("helper delete") + drainRequests() + fireUpdate(f.pv.DeepCopy(), f.pv.DeepCopy()) + expectNoTarget("PV update while helper absent") + drainRequests() + + recreatedHelper := f.helper.DeepCopy() + recreatedHelper.ResourceVersion = "" + assert.NilError(t, virtualClient.Create(runCtx, recreatedHelper)) + fireAdd(recreatedHelper.DeepCopy()) + expectTarget("helper recreate") + drainRequests() + + terminatingHelper := recreatedHelper.DeepCopy() + terminatingHelper.DeletionTimestamp = &now + fireAdd(terminatingHelper) + expectNoTarget("terminating helper") + drainRequests() + + assert.NilError(t, virtualClient.Delete(runCtx, f.pv.DeepCopy())) + fireUpdate(recreatedHelper.DeepCopy(), recreatedHelper.DeepCopy()) + expectNoTarget("helper update while PV absent") + drainRequests() + recreatedPV := f.pv.DeepCopy() + recreatedPV.ResourceVersion = "" + assert.NilError(t, virtualClient.Create(runCtx, recreatedPV)) + fireAdd(recreatedPV.DeepCopy()) + expectTarget("PV recreate") + + cancel() + select { + case err := <-started: + if err != nil && !errors.Is(err, context.Canceled) { + t.Fatalf("controller stop: %v", err) + } + case <-time.After(2 * time.Second): + t.Fatal("controller did not stop within 2s") + } +} + +func TestSyncExternalPopulatorDirectMaterializationConcurrentWriterRetriesFromFreshState(t *testing.T) { + const ( + virtualNamespace = "testns" + virtualPVName = "restore-populated-pv" + ) + + dataProtectionGroup := dataProtectionAPIGroup + virtualTarget := &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{ + Name: "target", + Namespace: virtualNamespace, + UID: types.UID("target-pvc-uid"), + Annotations: map[string]string{ + "example.test/virtual-writer": "preserve", + }, + }, + Spec: corev1.PersistentVolumeClaimSpec{ + VolumeName: virtualPVName, + DataSourceRef: &corev1.TypedObjectReference{ + APIGroup: &dataProtectionGroup, + Kind: dataProtectionBackupKind, + Name: "backup", + }, + Resources: corev1.VolumeResourceRequirements{ + Requests: corev1.ResourceList{ + corev1.ResourceStorage: resource.MustParse("1Gi"), + }, + }, + }, + Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimPending}, + } + virtualPV := &corev1.PersistentVolume{ + ObjectMeta: metav1.ObjectMeta{Name: virtualPVName}, + Spec: corev1.PersistentVolumeSpec{ + AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteOnce}, + Capacity: corev1.ResourceList{ + corev1.ResourceStorage: resource.MustParse("1Gi"), + }, + ClaimRef: &corev1.ObjectReference{ + Namespace: virtualTarget.Namespace, + Name: virtualTarget.Name, + UID: virtualTarget.UID, + }, + }, + Status: corev1.PersistentVolumeStatus{Phase: corev1.VolumeBound}, + } + + hostTranslator := translate.NewSingleNamespaceTranslator(testingutil.DefaultTestTargetNamespace) + hostTargetName := hostTranslator.HostName(nil, virtualTarget.Name, virtualTarget.Namespace) + hostTarget := &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{ + Name: hostTargetName.Name, + Namespace: testingutil.DefaultTestTargetNamespace, + UID: types.UID("host-target-pvc-uid"), + Annotations: map[string]string{ + translate.NameAnnotation: virtualTarget.Name, + translate.NamespaceAnnotation: virtualTarget.Namespace, + translate.UIDAnnotation: string(virtualTarget.UID), + translate.KindAnnotation: corev1.SchemeGroupVersion.WithKind("PersistentVolumeClaim").String(), + translate.HostNamespaceAnnotation: testingutil.DefaultTestTargetNamespace, + translate.HostNameAnnotation: hostTargetName.Name, + }, + Labels: map[string]string{ + translate.MarkerLabel: translate.VClusterName, + translate.NamespaceLabel: virtualTarget.Namespace, + }, + }, + Spec: corev1.PersistentVolumeClaimSpec{ + Resources: corev1.VolumeResourceRequirements{ + Requests: corev1.ResourceList{ + corev1.ResourceStorage: resource.MustParse("512Mi"), + }, + }, + }, + Status: corev1.PersistentVolumeClaimStatus{ + Phase: corev1.ClaimPending, + Capacity: corev1.ResourceList{}, + }, + } + hostPV := &corev1.PersistentVolume{ + ObjectMeta: metav1.ObjectMeta{Name: virtualPVName}, + Spec: corev1.PersistentVolumeSpec{ + ClaimRef: &corev1.ObjectReference{ + APIVersion: corev1.SchemeGroupVersion.Version, + Kind: "PersistentVolumeClaim", + Namespace: hostTarget.Namespace, + Name: hostTarget.Name, + UID: hostTarget.UID, + }, + }, + Status: corev1.PersistentVolumeStatus{Phase: corev1.VolumeBound}, + } + + test := &syncertesting.SyncTest{ + Name: "direct materialization rebases before deferred patch", + InitialVirtualState: []runtime.Object{virtualTarget.DeepCopy(), virtualPV.DeepCopy()}, + InitialPhysicalState: []runtime.Object{hostTarget.DeepCopy(), hostPV.DeepCopy()}, + Sync: func(registerCtx *synccontext.RegisterContext) { + syncCtx, objectSyncer := syncertesting.FakeStartSyncer(t, registerCtx, New) + pvcSyncer := objectSyncer.(*persistentVolumeClaimSyncer) + pvcSyncer.useFakePersistentVolumes = true + + targetKey := types.NamespacedName{ + Namespace: virtualTarget.Namespace, + Name: virtualTarget.Name, + } + hostTargetKey := types.NamespacedName{ + Namespace: hostTarget.Namespace, + Name: hostTarget.Name, + } + concurrentClient := &externalPopulatorConcurrentWriterClient{ + Client: syncCtx.HostClient, + target: hostTargetKey, + } + syncCtx.HostClient = concurrentClient + + reconcileTarget := func() ctrl.Result { + t.Helper() + currentVirtual := &corev1.PersistentVolumeClaim{} + assert.NilError(t, syncCtx.VirtualClient.Get(syncCtx.Context, targetKey, currentVirtual)) + currentHost := &corev1.PersistentVolumeClaim{} + assert.NilError(t, syncCtx.HostClient.Get(syncCtx.Context, hostTargetKey, currentHost)) + result, err := pvcSyncer.Sync(syncCtx, synccontext.NewSyncEventWithOld( + currentHost.DeepCopy(), + currentHost.DeepCopy(), + currentVirtual.DeepCopy(), + currentVirtual.DeepCopy(), + )) + assert.NilError(t, err) + return result + } + + firstResult := reconcileTarget() + assert.Equal(t, firstResult.RequeueAfter, time.Second) + assert.NilError(t, concurrentClient.concurrentErr) + assert.Assert(t, concurrentClient.directRV != "") + assert.Assert(t, concurrentClient.concurrentRV != "") + assert.Assert(t, concurrentClient.directRV != concurrentClient.concurrentRV) + + hostAfterConflict := &corev1.PersistentVolumeClaim{} + assert.NilError(t, syncCtx.HostClient.Get(syncCtx.Context, hostTargetKey, hostAfterConflict)) + assert.Equal(t, hostAfterConflict.Spec.VolumeName, virtualPVName) + assert.Equal(t, hostAfterConflict.Annotations["example.test/concurrent-writer"], "preserve") + hostStorageAfterConflict := hostAfterConflict.Spec.Resources.Requests[corev1.ResourceStorage] + assert.Equal(t, hostStorageAfterConflict.String(), "2Gi") + + virtualAfterConflict := &corev1.PersistentVolumeClaim{} + assert.NilError(t, syncCtx.VirtualClient.Get(syncCtx.Context, targetKey, virtualAfterConflict)) + assert.Equal(t, virtualAfterConflict.Status.Phase, corev1.ClaimBound) + virtualCapacityAfterConflict := virtualAfterConflict.Status.Capacity[corev1.ResourceStorage] + assert.Equal(t, virtualCapacityAfterConflict.String(), "1Gi") + assert.Equal(t, virtualAfterConflict.Annotations["example.test/virtual-writer"], "preserve") + + secondResult := reconcileTarget() + assert.Check(t, secondResult.IsZero()) + + hostAfterRetry := &corev1.PersistentVolumeClaim{} + assert.NilError(t, syncCtx.HostClient.Get(syncCtx.Context, hostTargetKey, hostAfterRetry)) + assert.Equal(t, hostAfterRetry.Annotations["example.test/concurrent-writer"], "preserve") + hostStorageAfterRetry := hostAfterRetry.Spec.Resources.Requests[corev1.ResourceStorage] + assert.Equal(t, hostStorageAfterRetry.String(), "1Gi") + + virtualAfterRetry := &corev1.PersistentVolumeClaim{} + assert.NilError(t, syncCtx.VirtualClient.Get(syncCtx.Context, targetKey, virtualAfterRetry)) + assert.Equal(t, virtualAfterRetry.Status.Phase, corev1.ClaimBound) + assert.Equal(t, virtualAfterRetry.Annotations["example.test/virtual-writer"], "preserve") + }, + } + test.Run(t, syncertesting.NewFakeRegisterContext) +} + func TestTranslateSelectorPreservesNilAndExplicitEmptyStorageClass(t *testing.T) { empty := "" tests := []struct { diff --git a/pkg/patcher/patcher.go b/pkg/patcher/patcher.go index 84f9a0f409..278310c1ac 100644 --- a/pkg/patcher/patcher.go +++ b/pkg/patcher/patcher.go @@ -86,6 +86,21 @@ func (h *SyncerPatcher) Patch(ctx *synccontext.SyncContext, pObj, vObj client.Ob return nil } +// RebaseHost records host changes that were committed directly after this +// SyncerPatcher was created. Later deferred patches then start from the +// committed resourceVersion instead of replaying against a stale snapshot. +func (h *SyncerPatcher) RebaseHost(obj client.Object) error { + if clienthelper.IsNilObject(obj) { + return fmt.Errorf("rebase host patcher: expected non-nil object") + } + if obj.GetResourceVersion() == "" { + return fmt.Errorf("rebase host patcher for %s/%s: committed resourceVersion is empty", obj.GetNamespace(), obj.GetName()) + } + + h.pPatcher.beforeObject = obj.DeepCopyObject().(client.Object) + return nil +} + // Patcher is a utility for ensuring the proper patching of objects. type Patcher struct { client client.Client diff --git a/pkg/patcher/patcher_test.go b/pkg/patcher/patcher_test.go new file mode 100644 index 0000000000..cfdd9ae0db --- /dev/null +++ b/pkg/patcher/patcher_test.go @@ -0,0 +1,73 @@ +package patcher + +import ( + "context" + "strings" + "testing" + + "github.com/loft-sh/vcluster/pkg/scheme" + "github.com/loft-sh/vcluster/pkg/syncer/synccontext" + testingutil "github.com/loft-sh/vcluster/pkg/util/testing" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +func TestSyncerPatcherRebaseHostFailClosed(t *testing.T) { + host := &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{ + Name: "host", + Namespace: "test", + ResourceVersion: "1", + }, + } + virtual := &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{ + Name: "virtual", + Namespace: "test", + ResourceVersion: "1", + }, + } + ctx := &synccontext.SyncContext{ + Context: context.Background(), + HostClient: testingutil.NewFakeClient(scheme.Scheme), + VirtualClient: testingutil.NewFakeClient(scheme.Scheme), + } + syncerPatcher, err := NewSyncerPatcher(ctx, host, virtual) + if err != nil { + t.Fatalf("new syncer patcher: %v", err) + } + + tests := []struct { + name string + obj *corev1.PersistentVolumeClaim + want string + }{ + { + name: "nil committed object", + obj: nil, + want: "expected non-nil object", + }, + { + name: "missing committed resource version", + obj: &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{Name: "host", Namespace: "test"}, + }, + want: "committed resourceVersion is empty", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + err := syncerPatcher.RebaseHost(tt.obj) + if err == nil || !strings.Contains(err.Error(), tt.want) { + t.Fatalf("RebaseHost() error = %v, want substring %q", err, tt.want) + } + }) + } + + committed := host.DeepCopy() + committed.ResourceVersion = "2" + if err := syncerPatcher.RebaseHost(committed); err != nil { + t.Fatalf("RebaseHost() committed object: %v", err) + } +} From 35909724128e2c57c9f3d3aadd70cafa02649722 Mon Sep 17 00:00:00 2001 From: Gina Date: Fri, 24 Jul 2026 23:10:53 +0800 Subject: [PATCH 5/5] fix(pvc): harden external populator recovery --- .../persistentvolumeclaims/syncer.go | 150 ++++- .../persistentvolumeclaims/syncer_test.go | 603 ++++++++++++++++++ pkg/patcher/patcher.go | 10 +- pkg/patcher/patcher_test.go | 65 ++ 4 files changed, 805 insertions(+), 23 deletions(-) diff --git a/pkg/controllers/resources/persistentvolumeclaims/syncer.go b/pkg/controllers/resources/persistentvolumeclaims/syncer.go index b41e764a42..f992d4793a 100644 --- a/pkg/controllers/resources/persistentvolumeclaims/syncer.go +++ b/pkg/controllers/resources/persistentvolumeclaims/syncer.go @@ -35,6 +35,7 @@ import ( "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/event" "sigs.k8s.io/controller-runtime/pkg/handler" + "sigs.k8s.io/controller-runtime/pkg/manager" "github.com/loft-sh/vcluster/pkg/util/loghelper" ) @@ -67,6 +68,10 @@ const ( externalPopulatorTopologyMismatchEventReason = "ExternalPopulatorTopologyMismatch" externalPopulatorTopologyNotReadyEventReason = "ExternalPopulatorTopologyNotReady" externalPopulatorTopologyEventAction = "SyncPersistentVolumeClaim" + + externalPopulatorDependencyRetryBaseDelay = 250 * time.Millisecond + externalPopulatorDependencyRetryMaxDelay = 4 * time.Second + externalPopulatorDependencyMaxRetries = 5 ) func New(ctx *synccontext.RegisterContext) (syncertypes.Object, error) { @@ -115,7 +120,30 @@ func (s *persistentVolumeClaimSyncer) RegisterIndices(ctx *synccontext.RegisterC var _ syncertypes.ControllerModifier = &persistentVolumeClaimSyncer{} -func (s *persistentVolumeClaimSyncer) ModifyController(_ *synccontext.RegisterContext, controllerBuilder *builder.Builder) (*builder.Builder, error) { +type externalPopulatorDependencyRetry struct { + object client.Object + allowDeletingDependency bool + targetQueue workqueue.TypedRateLimitingInterface[ctrl.Request] +} + +func (s *persistentVolumeClaimSyncer) ModifyController(registerCtx *synccontext.RegisterContext, controllerBuilder *builder.Builder) (*builder.Builder, error) { + dependencyRetryQueue := workqueue.NewTypedRateLimitingQueueWithConfig( + workqueue.NewTypedItemExponentialFailureRateLimiter[*externalPopulatorDependencyRetry]( + externalPopulatorDependencyRetryBaseDelay, + externalPopulatorDependencyRetryMaxDelay, + ), + workqueue.TypedRateLimitingQueueConfig[*externalPopulatorDependencyRetry]{ + Name: "external-populator-dependency-mapper", + }, + ) + err := registerCtx.VirtualManager.Add(manager.RunnableFunc(func(ctx context.Context) error { + return s.runExternalPopulatorDependencyRetryWorker(ctx, dependencyRetryQueue) + })) + if err != nil { + dependencyRetryQueue.ShutDown() + return nil, fmt.Errorf("register external populator dependency retry worker: %w", err) + } + enqueueDependency := func( ctx context.Context, object client.Object, @@ -134,6 +162,11 @@ func (s *persistentVolumeClaimSyncer) ModifyController(_ *synccontext.RegisterCo "name", object.GetName(), ) + dependencyRetryQueue.AddRateLimited(&externalPopulatorDependencyRetry{ + object: object.DeepCopyObject().(client.Object), + allowDeletingDependency: allowDeletingDependency, + targetQueue: queue, + }) return } @@ -161,6 +194,59 @@ func (s *persistentVolumeClaimSyncer) ModifyController(_ *synccontext.RegisterCo Watches(&corev1.PersistentVolumeClaim{}, dependencyHandler), nil } +func (s *persistentVolumeClaimSyncer) runExternalPopulatorDependencyRetryWorker( + ctx context.Context, + queue workqueue.TypedRateLimitingInterface[*externalPopulatorDependencyRetry], +) error { + go func() { + <-ctx.Done() + queue.ShutDown() + }() + + for { + retry, shutdown := queue.Get() + if shutdown { + return nil + } + + func() { + defer queue.Done(retry) + + requests, err := s.externalPopulatorDependencyRequestsWithOptions( + ctx, + retry.object, + retry.allowDeletingDependency, + ) + if err != nil { + if queue.NumRequeues(retry) < externalPopulatorDependencyMaxRetries { + queue.AddRateLimited(retry) + return + } + + queue.Forget(retry) + ctrl.LoggerFrom(ctx).Error( + err, + "external populator dependency mapper retries exhausted", + "kind", + fmt.Sprintf("%T", retry.object), + "namespace", + retry.object.GetNamespace(), + "name", + retry.object.GetName(), + "retries", + externalPopulatorDependencyMaxRetries, + ) + return + } + + queue.Forget(retry) + for _, request := range requests { + retry.targetQueue.Add(request) + } + }() + } +} + var _ syncertypes.OptionsProvider = &persistentVolumeClaimSyncer{} func (s *persistentVolumeClaimSyncer) Options() *syncertypes.Options { @@ -335,7 +421,7 @@ func (s *persistentVolumeClaimSyncer) Sync(ctx *synccontext.SyncContext, event * retErr = utilerrors.NewAggregate([]error{retErr, err}) } - if containsConflictError(retErr) { + if containsOnlyConflictErrors(retErr) { result = ctrl.Result{RequeueAfter: time.Second} retErr = nil return @@ -364,19 +450,23 @@ func (s *persistentVolumeClaimSyncer) Sync(ctx *synccontext.SyncContext, event * } if preserveVirtualStatus { hostResourceVersionBeforeMaterialization := event.Host.ResourceVersion - hostConverged, err := s.ensureExternalPopulatorHostMaterialization(ctx, event.Host, event.Virtual, vPV) - if err != nil { - return ctrl.Result{}, err - } - // A successful optimistic-lock materialization patch returns a new - // resourceVersion in event.Host. Rebase only when that direct write or a - // fresher target read changed the snapshot captured by NewSyncerPatcher. + hostConverged, materializationErr := s.ensureExternalPopulatorHostMaterialization(ctx, event.Host, event.Virtual, vPV) + // A direct materialization write or fresher target read can change + // event.Host before a later topology lookup fails. Rebase before + // processing that error so the deferred host patch never replays the + // fresh snapshot as a controller-owned diff against an older baseline. if event.Host.ResourceVersion != hostResourceVersionBeforeMaterialization { - err = patch.RebaseHost(event.Host) - if err != nil { - return ctrl.Result{}, err + rebaseErr := patch.RebaseHost(event.Host) + if rebaseErr != nil { + if materializationErr != nil { + return ctrl.Result{}, fmt.Errorf("%w (host patch disabled after rebase failure: %v)", materializationErr, rebaseErr) + } + return ctrl.Result{}, rebaseErr } } + if materializationErr != nil { + return ctrl.Result{}, materializationErr + } if hostConverged { ensureExternalPopulatorVirtualPopulateStatus(event.Virtual, vPV) } else { @@ -405,22 +495,42 @@ func (s *persistentVolumeClaimSyncer) Sync(ctx *synccontext.SyncContext, event * return result, nil } -func containsConflictError(err error) bool { +func containsOnlyConflictErrors(err error) bool { + hasError, onlyConflicts := classifyConflictErrors(err) + return hasError && onlyConflicts +} + +func classifyConflictErrors(err error) (bool, bool) { if err == nil { - return false - } - if kerrors.IsConflict(err) { - return true + return false, true } if aggregate, ok := err.(utilerrors.Aggregate); ok { + hasError := false for _, aggregateErr := range aggregate.Errors() { - if containsConflictError(aggregateErr) { - return true + nestedHasError, nestedOnlyConflicts := classifyConflictErrors(aggregateErr) + hasError = hasError || nestedHasError + if !nestedOnlyConflicts { + return hasError, false } } + return hasError, true + } + if multiError, ok := err.(interface{ Unwrap() []error }); ok { + hasError := false + for _, nestedErr := range multiError.Unwrap() { + nestedHasError, nestedOnlyConflicts := classifyConflictErrors(nestedErr) + hasError = hasError || nestedHasError + if !nestedOnlyConflicts { + return hasError, false + } + } + return hasError, true + } + if unwrapped := errors.Unwrap(err); unwrapped != nil { + return classifyConflictErrors(unwrapped) } - return false + return true, kerrors.IsConflict(err) } func (s *persistentVolumeClaimSyncer) SyncToVirtual(ctx *synccontext.SyncContext, event *synccontext.SyncToVirtualEvent[*corev1.PersistentVolumeClaim]) (_ ctrl.Result, retErr error) { diff --git a/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go b/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go index dba5e6bc0b..78ac40252a 100644 --- a/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go +++ b/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go @@ -3,9 +3,11 @@ package persistentvolumeclaims import ( "context" "errors" + "fmt" "reflect" "strings" "sync" + "sync/atomic" "testing" "time" @@ -24,17 +26,23 @@ import ( corev1 "k8s.io/api/core/v1" storagev1 "k8s.io/api/storage/v1" apiequality "k8s.io/apimachinery/pkg/api/equality" + kerrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/runtime/schema" + utilerrors "k8s.io/apimachinery/pkg/util/errors" toolscache "k8s.io/client-go/tools/cache" + "k8s.io/client-go/util/workqueue" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/cache" "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/manager" "sigs.k8s.io/controller-runtime/pkg/reconcile" ) +var externalPopulatorControllerTestID atomic.Uint64 + type eventRecordingTranslator struct { syncertypes.GenericTranslator recorder events.EventRecorder @@ -141,6 +149,56 @@ func newExternalPopulatorDependencyFixture() *externalPopulatorDependencyFixture } } +func newExternalPopulatorHostHandoffObjects( + f *externalPopulatorDependencyFixture, +) (*corev1.PersistentVolumeClaim, *corev1.PersistentVolumeClaim, *corev1.PersistentVolume) { + hostTranslator := translate.NewSingleNamespaceTranslator(testingutil.DefaultTestTargetNamespace) + hostTargetName := hostTranslator.HostName(nil, f.target.Name, f.target.Namespace) + hostTarget := &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{ + Name: hostTargetName.Name, + Namespace: testingutil.DefaultTestTargetNamespace, + UID: types.UID("host-target-pvc-uid"), + Annotations: map[string]string{ + translate.NameAnnotation: f.target.Name, + translate.NamespaceAnnotation: f.target.Namespace, + translate.UIDAnnotation: string(f.target.UID), + translate.KindAnnotation: corev1.SchemeGroupVersion.WithKind("PersistentVolumeClaim").String(), + translate.HostNamespaceAnnotation: testingutil.DefaultTestTargetNamespace, + translate.HostNameAnnotation: hostTargetName.Name, + }, + Labels: map[string]string{ + translate.MarkerLabel: translate.VClusterName, + translate.NamespaceLabel: f.target.Namespace, + }, + }, + Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimPending}, + } + hostHelperName := hostTranslator.HostName(nil, f.helper.Name, f.helper.Namespace) + hostHelper := &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{ + Name: hostHelperName.Name, + Namespace: testingutil.DefaultTestTargetNamespace, + UID: types.UID("host-helper-pvc-uid"), + }, + Spec: corev1.PersistentVolumeClaimSpec{VolumeName: f.pv.Name}, + } + hostPV := &corev1.PersistentVolume{ + ObjectMeta: metav1.ObjectMeta{Name: f.pv.Name}, + Spec: corev1.PersistentVolumeSpec{ + ClaimRef: &corev1.ObjectReference{ + APIVersion: corev1.SchemeGroupVersion.Version, + Kind: "PersistentVolumeClaim", + Namespace: hostHelper.Namespace, + Name: hostHelper.Name, + UID: hostHelper.UID, + }, + }, + Status: corev1.PersistentVolumeStatus{Phase: corev1.VolumeBound}, + } + return hostTarget, hostHelper, hostPV +} + func newExternalPopulatorMapperTestSyncer( t *testing.T, virtualObjects ...runtime.Object, @@ -174,6 +232,109 @@ func (c *externalPopulatorMapperErrorClient) List(ctx context.Context, list clie return c.Client.List(ctx, list, opts...) } +type externalPopulatorFlakyMapperClient struct { + client.Client + + mu sync.Mutex + remainingGetErrors int + injectedGetErrors int + getErr error +} + +func (c *externalPopulatorFlakyMapperClient) Get( + ctx context.Context, + key client.ObjectKey, + obj client.Object, + opts ...client.GetOption, +) error { + c.mu.Lock() + if c.remainingGetErrors > 0 { + c.remainingGetErrors-- + c.injectedGetErrors++ + err := c.getErr + c.mu.Unlock() + return err + } + c.mu.Unlock() + + return c.Client.Get(ctx, key, obj, opts...) +} + +func (c *externalPopulatorFlakyMapperClient) injectedErrors() int { + c.mu.Lock() + defer c.mu.Unlock() + return c.injectedGetErrors +} + +type externalPopulatorTopologyListErrorAfterWriterClient struct { + client.Client + hostClient client.Client + hostTarget types.NamespacedName + topologyErr error + updateErr error + + mu sync.Mutex + pvcListCalls int + injectedUpdateErrors int + writerErr error +} + +func (c *externalPopulatorTopologyListErrorAfterWriterClient) List( + ctx context.Context, + list client.ObjectList, + opts ...client.ListOption, +) error { + if _, ok := list.(*corev1.PersistentVolumeClaimList); !ok { + return c.Client.List(ctx, list, opts...) + } + + c.mu.Lock() + c.pvcListCalls++ + trigger := c.pvcListCalls == 2 + c.mu.Unlock() + if !trigger { + return c.Client.List(ctx, list, opts...) + } + + current := &corev1.PersistentVolumeClaim{} + err := c.hostClient.Get(ctx, c.hostTarget, current) + if err == nil { + if current.Annotations == nil { + current.Annotations = map[string]string{} + } + current.Annotations["example.test/external-writer"] = "after-fresh-read" + current.Spec.Resources.Requests = corev1.ResourceList{ + corev1.ResourceStorage: resource.MustParse("3Gi"), + } + err = c.hostClient.Update(ctx, current) + } + + c.mu.Lock() + c.writerErr = err + c.mu.Unlock() + return c.topologyErr +} + +func (c *externalPopulatorTopologyListErrorAfterWriterClient) Update( + ctx context.Context, + obj client.Object, + opts ...client.UpdateOption, +) error { + if c.updateErr != nil { + c.mu.Lock() + c.injectedUpdateErrors++ + c.mu.Unlock() + return c.updateErr + } + return c.Client.Update(ctx, obj, opts...) +} + +func (c *externalPopulatorTopologyListErrorAfterWriterClient) state() (int, int, error) { + c.mu.Lock() + defer c.mu.Unlock() + return c.pvcListCalls, c.injectedUpdateErrors, c.writerErr +} + type externalPopulatorConcurrentWriterClient struct { client.Client target types.NamespacedName @@ -228,12 +389,28 @@ func (c *externalPopulatorConcurrentWriterClient) Patch( type externalPopulatorRecordingManager struct { ctrl.Manager recordingCache cache.Cache + + mu sync.Mutex + runnables []manager.Runnable } func (m *externalPopulatorRecordingManager) GetCache() cache.Cache { return m.recordingCache } +func (m *externalPopulatorRecordingManager) Add(runnable manager.Runnable) error { + m.mu.Lock() + m.runnables = append(m.runnables, runnable) + m.mu.Unlock() + return m.Manager.Add(runnable) +} + +func (m *externalPopulatorRecordingManager) registeredRunnables() []manager.Runnable { + m.mu.Lock() + defer m.mu.Unlock() + return append([]manager.Runnable(nil), m.runnables...) +} + type externalPopulatorRecordingCache struct { cache.Cache @@ -1001,6 +1178,432 @@ func TestExternalPopulatorDependencyModifyControllerWiring(t *testing.T) { } } +func TestExternalPopulatorDependencyModifyControllerRetriesTransientMapperError(t *testing.T) { + f := newExternalPopulatorDependencyFixture() + hostTarget, hostHelper, hostPV := newExternalPopulatorHostHandoffObjects(f) + + hostClient := testingutil.NewFakeClient( + scheme.Scheme, + hostTarget, + hostHelper, + hostPV, + ) + virtualClient := testingutil.NewFakeClient( + scheme.Scheme, + f.target.DeepCopy(), + f.pv.DeepCopy(), + ) + registerCtx := syncertesting.NewFakeRegisterContext(testingutil.NewFakeConfig(), hostClient, virtualClient) + baseManager := testingutil.NewFakeManager(virtualClient) + recordingCache := newExternalPopulatorRecordingCache(baseManager.GetCache()) + recordingManager := &externalPopulatorRecordingManager{ + Manager: baseManager, + recordingCache: recordingCache, + } + registerCtx.VirtualManager = recordingManager + + syncCtx, objectSyncer := syncertesting.FakeStartSyncer(t, registerCtx, New) + pvcSyncer := objectSyncer.(*persistentVolumeClaimSyncer) + pvcSyncer.useFakePersistentVolumes = true + + targetKey := types.NamespacedName{ + Namespace: f.target.Namespace, + Name: f.target.Name, + } + hostTargetKey := types.NamespacedName{ + Namespace: hostTarget.Namespace, + Name: hostTarget.Name, + } + reconcileTarget := func() (ctrl.Result, error) { + currentVirtual := &corev1.PersistentVolumeClaim{} + if err := syncCtx.VirtualClient.Get(syncCtx.Context, targetKey, currentVirtual); err != nil { + return ctrl.Result{}, err + } + currentHost := &corev1.PersistentVolumeClaim{} + if err := syncCtx.HostClient.Get(syncCtx.Context, hostTargetKey, currentHost); err != nil { + return ctrl.Result{}, err + } + return pvcSyncer.Sync(syncCtx, synccontext.NewSyncEventWithOld( + currentHost.DeepCopy(), + currentHost.DeepCopy(), + currentVirtual.DeepCopy(), + currentVirtual.DeepCopy(), + )) + } + + firstResult, err := reconcileTarget() + assert.NilError(t, err) + assert.Check(t, firstResult.IsZero(), "first transient result was %#v", firstResult) + preRecoveryHostTarget := &corev1.PersistentVolumeClaim{} + assert.NilError(t, hostClient.Get(syncCtx.Context, hostTargetKey, preRecoveryHostTarget)) + assert.Equal(t, preRecoveryHostTarget.Spec.VolumeName, "") + preRecoveryHostPV := &corev1.PersistentVolume{} + assert.NilError(t, hostClient.Get( + syncCtx.Context, + types.NamespacedName{Name: f.pv.Name}, + preRecoveryHostPV, + )) + assert.Check( + t, + claimRefMatchesPersistentVolumeClaim(preRecoveryHostPV.Spec.ClaimRef, hostHelper), + "pre-recovery host PV claimRef = %#v, want helper %s/%s", + preRecoveryHostPV.Spec.ClaimRef, + hostHelper.Namespace, + hostHelper.Name, + ) + + assert.NilError(t, virtualClient.Create(syncCtx.Context, f.helper.DeepCopy())) + injectedErr := errors.New("injected first dependency mapper lookup failure") + flakyMapperClient := &externalPopulatorFlakyMapperClient{ + Client: virtualClient, + remainingGetErrors: 1, + getErr: injectedErr, + } + pvcSyncer.virtualClient = flakyMapperClient + + type reconcileOutcome struct { + request ctrl.Request + result ctrl.Result + err error + } + outcomes := make(chan reconcileOutcome, 4) + controllerBuilder := ctrl.NewControllerManagedBy(recordingManager). + Named(fmt.Sprintf( + "external-populator-dependency-retry-test-%d", + externalPopulatorControllerTestID.Add(1), + )). + For(&corev1.PersistentVolumeClaim{}) + controllerBuilder, err = pvcSyncer.ModifyController(registerCtx, controllerBuilder) + assert.NilError(t, err) + + retryRunnables := recordingManager.registeredRunnables() + assert.Equal(t, len(retryRunnables), 1) + retryRunnable := retryRunnables[0] + controller, err := controllerBuilder.Build(reconcile.Func(func(_ context.Context, request ctrl.Request) (ctrl.Result, error) { + if request.NamespacedName != targetKey { + return ctrl.Result{}, nil + } + result, reconcileErr := reconcileTarget() + outcomes <- reconcileOutcome{ + request: request, + result: result, + err: reconcileErr, + } + return result, reconcileErr + })) + assert.NilError(t, err) + + runCtx, cancel := context.WithCancel(context.Background()) + defer cancel() + retryStopped := make(chan error, 1) + go func() { + retryStopped <- retryRunnable.Start(runCtx) + }() + controllerStopped := make(chan error, 1) + go func() { + controllerStopped <- controller.Start(runCtx) + }() + + waitDeadline := time.NewTimer(2 * time.Second) + waitTicker := time.NewTicker(5 * time.Millisecond) + for { + if len(recordingCache.handlersFor(&corev1.PersistentVolumeClaim{})) == 2 { + break + } + select { + case err := <-controllerStopped: + t.Fatalf("controller stopped before watch registration: %v", err) + case err := <-retryStopped: + t.Fatalf("retry worker stopped before dependency event: %v", err) + case <-waitDeadline.C: + t.Fatalf( + "watch registration timed out: pvc handlers=%d", + len(recordingCache.handlersFor(&corev1.PersistentVolumeClaim{})), + ) + case <-waitTicker.C: + } + } + waitDeadline.Stop() + waitTicker.Stop() + + // Fire the helper creation exactly once. The fast mapper consumes the + // injected Get error; no second dependency event is sent after recovery. + for _, eventHandler := range recordingCache.handlersFor(&corev1.PersistentVolumeClaim{}) { + eventHandler.OnAdd(f.helper.DeepCopy(), false) + } + + select { + case outcome := <-outcomes: + if outcome.request.NamespacedName != targetKey { + t.Fatalf("recovered mapper enqueued %s, want %s", outcome.request.NamespacedName, targetKey) + } + assert.NilError(t, outcome.err) + assert.Check(t, outcome.result.IsZero(), "recovered target result was %#v", outcome.result) + case <-time.After(3 * time.Second): + t.Fatalf("recovered mapper did not enqueue exact target %s within 3s", targetKey) + } + + assert.Equal(t, flakyMapperClient.injectedErrors(), 1) + actualHostTarget := &corev1.PersistentVolumeClaim{} + assert.NilError(t, hostClient.Get(syncCtx.Context, hostTargetKey, actualHostTarget)) + actualHostPV := &corev1.PersistentVolume{} + assert.NilError(t, hostClient.Get( + syncCtx.Context, + types.NamespacedName{Name: f.pv.Name}, + actualHostPV, + )) + actualVirtualTarget := &corev1.PersistentVolumeClaim{} + assert.NilError(t, virtualClient.Get(syncCtx.Context, targetKey, actualVirtualTarget)) + if actualHostTarget.Spec.VolumeName != f.pv.Name || + !claimRefMatchesPersistentVolumeClaim(actualHostPV.Spec.ClaimRef, actualHostTarget) || + actualVirtualTarget.Status.Phase != corev1.ClaimBound { + t.Fatalf( + "recovered handoff: host_volume=%q host_claim_ref=%#v virtual_phase=%q; want %q/target/%q", + actualHostTarget.Spec.VolumeName, + actualHostPV.Spec.ClaimRef, + actualVirtualTarget.Status.Phase, + f.pv.Name, + corev1.ClaimBound, + ) + } + + cancel() + for label, stopped := range map[string]<-chan error{ + "controller": controllerStopped, + "retry worker": retryStopped, + } { + select { + case err := <-stopped: + if err != nil && !errors.Is(err, context.Canceled) { + t.Fatalf("%s stop: %v", label, err) + } + case <-time.After(2 * time.Second): + t.Fatalf("%s did not stop within 2s", label) + } + } +} + +func TestExternalPopulatorDependencyRetryExhaustsFiniteBudget(t *testing.T) { + f := newExternalPopulatorDependencyFixture() + syncCtx, pvcSyncer := newExternalPopulatorMapperTestSyncer( + t, + f.target.DeepCopy(), + f.pv.DeepCopy(), + ) + flakyMapperClient := &externalPopulatorFlakyMapperClient{ + Client: syncCtx.VirtualClient, + remainingGetErrors: 100, + getErr: errors.New("persistent dependency mapper failure"), + } + pvcSyncer.virtualClient = flakyMapperClient + + retryQueue := workqueue.NewTypedRateLimitingQueue( + workqueue.NewTypedItemExponentialFailureRateLimiter[*externalPopulatorDependencyRetry]( + time.Millisecond, + time.Millisecond, + ), + ) + targetQueue := workqueue.NewTypedRateLimitingQueue( + workqueue.NewTypedItemExponentialFailureRateLimiter[ctrl.Request]( + time.Millisecond, + time.Millisecond, + ), + ) + defer targetQueue.ShutDown() + + retry := &externalPopulatorDependencyRetry{ + object: f.helper.DeepCopy(), + targetQueue: targetQueue, + } + runCtx, cancel := context.WithCancel(context.Background()) + stopped := make(chan error, 1) + go func() { + stopped <- pvcSyncer.runExternalPopulatorDependencyRetryWorker(runCtx, retryQueue) + }() + retryQueue.AddRateLimited(retry) + + deadline := time.NewTimer(time.Second) + ticker := time.NewTicker(time.Millisecond) + defer deadline.Stop() + defer ticker.Stop() + for { + if flakyMapperClient.injectedErrors() == externalPopulatorDependencyMaxRetries && + retryQueue.NumRequeues(retry) == 0 && + retryQueue.Len() == 0 { + break + } + + select { + case err := <-stopped: + t.Fatalf("retry worker stopped before exhausting the retry budget: %v", err) + case <-deadline.C: + t.Fatalf( + "retry budget was not exhausted: injected_errors=%d requeues=%d queue_len=%d", + flakyMapperClient.injectedErrors(), + retryQueue.NumRequeues(retry), + retryQueue.Len(), + ) + case <-ticker.C: + } + } + + assert.Equal(t, targetQueue.Len(), 0) + cancel() + select { + case err := <-stopped: + assert.NilError(t, err) + case <-time.After(time.Second): + t.Fatal("retry worker did not stop within 1s") + } +} + +func TestSyncExternalPopulatorTopologyErrorDoesNotReplayFreshHostSnapshot(t *testing.T) { + f := newExternalPopulatorDependencyFixture() + f.target.Spec.Resources.Requests = corev1.ResourceList{ + corev1.ResourceStorage: resource.MustParse("1Gi"), + } + hostTarget, hostHelper, hostPV := newExternalPopulatorHostHandoffObjects(f) + hostTarget.Annotations["example.test/external-writer"] = "initial" + hostTarget.Annotations[bindCompletedAnnotation] = "yes" + hostTarget.Spec.Resources.Requests = corev1.ResourceList{ + corev1.ResourceStorage: resource.MustParse("1Gi"), + } + + hostClient := testingutil.NewFakeClient( + scheme.Scheme, + hostTarget, + hostHelper, + hostPV, + ) + virtualClient := testingutil.NewFakeClient( + scheme.Scheme, + f.target.DeepCopy(), + f.helper.DeepCopy(), + f.pv.DeepCopy(), + ) + registerCtx := syncertesting.NewFakeRegisterContext(testingutil.NewFakeConfig(), hostClient, virtualClient) + syncCtx, objectSyncer := syncertesting.FakeStartSyncer(t, registerCtx, New) + pvcSyncer := objectSyncer.(*persistentVolumeClaimSyncer) + pvcSyncer.useFakePersistentVolumes = true + + targetKey := types.NamespacedName{Namespace: f.target.Namespace, Name: f.target.Name} + hostTargetKey := types.NamespacedName{Namespace: hostTarget.Namespace, Name: hostTarget.Name} + eventHost := &corev1.PersistentVolumeClaim{} + assert.NilError(t, hostClient.Get(syncCtx.Context, hostTargetKey, eventHost)) + eventVirtual := &corev1.PersistentVolumeClaim{} + assert.NilError(t, virtualClient.Get(syncCtx.Context, targetKey, eventVirtual)) + initialHostResourceVersion := eventHost.ResourceVersion + + firstWriter := &corev1.PersistentVolumeClaim{} + assert.NilError(t, hostClient.Get(syncCtx.Context, hostTargetKey, firstWriter)) + firstWriter.Annotations["example.test/external-writer"] = "fresh-read" + firstWriter.Spec.Resources.Requests = corev1.ResourceList{ + corev1.ResourceStorage: resource.MustParse("2Gi"), + } + assert.NilError(t, hostClient.Update(syncCtx.Context, firstWriter)) + + topologyErr := errors.New("injected helper topology list failure after fresh host read") + deferredConflict := kerrors.NewConflict( + schema.GroupResource{Resource: "persistentvolumeclaims"}, + f.target.Name, + errors.New("injected deferred virtual patch conflict"), + ) + errorClient := &externalPopulatorTopologyListErrorAfterWriterClient{ + Client: virtualClient, + hostClient: hostClient, + hostTarget: hostTargetKey, + topologyErr: topologyErr, + updateErr: deferredConflict, + } + syncCtx.VirtualClient = errorClient + + result, err := pvcSyncer.Sync(syncCtx, synccontext.NewSyncEventWithOld( + eventHost.DeepCopy(), + eventHost, + eventVirtual.DeepCopy(), + eventVirtual, + )) + if !errors.Is(err, topologyErr) { + t.Fatalf("Sync() error = %v, want wrapped topology error %v", err, topologyErr) + } + assert.Check(t, result.IsZero(), "topology error result was %#v", result) + if eventHost.ResourceVersion == initialHostResourceVersion || + eventHost.Annotations["example.test/external-writer"] != "fresh-read" { + t.Fatalf( + "ensure did not install the fresh host snapshot: initial_rv=%q event_rv=%q writer=%q", + initialHostResourceVersion, + eventHost.ResourceVersion, + eventHost.Annotations["example.test/external-writer"], + ) + } + + listCalls, injectedUpdateErrors, writerErr := errorClient.state() + assert.Equal(t, listCalls, 2) + assert.Equal(t, injectedUpdateErrors, 1) + assert.NilError(t, writerErr) + + actualHostTarget := &corev1.PersistentVolumeClaim{} + assert.NilError(t, hostClient.Get(syncCtx.Context, hostTargetKey, actualHostTarget)) + if actualHostTarget.Annotations["example.test/external-writer"] != "after-fresh-read" || + actualHostTarget.Spec.Resources.Requests.Storage().Cmp(resource.MustParse("3Gi")) != 0 || + actualHostTarget.Spec.VolumeName != "" { + t.Fatalf( + "deferred host patch replayed a stale snapshot: writer=%q storage=%s volume=%q; want after-fresh-read/3Gi/empty", + actualHostTarget.Annotations["example.test/external-writer"], + actualHostTarget.Spec.Resources.Requests.Storage().String(), + actualHostTarget.Spec.VolumeName, + ) + } + actualHostPV := &corev1.PersistentVolume{} + assert.NilError(t, hostClient.Get( + syncCtx.Context, + types.NamespacedName{Name: f.pv.Name}, + actualHostPV, + )) + assert.Check( + t, + claimRefMatchesPersistentVolumeClaim(actualHostPV.Spec.ClaimRef, hostHelper), + "topology error changed host PV claimRef to %#v", + actualHostPV.Spec.ClaimRef, + ) +} + +func TestContainsOnlyConflictErrors(t *testing.T) { + sentinel := errors.New("sentinel non-conflict") + conflict := kerrors.NewConflict( + schema.GroupResource{Resource: "persistentvolumeclaims"}, + "target", + errors.New("conflict"), + ) + pureAggregate := utilerrors.NewAggregate([]error{ + conflict, + fmt.Errorf("wrapped: %w", conflict), + }) + mixedAggregate := utilerrors.NewAggregate([]error{ + sentinel, + conflict, + }) + + for _, tt := range []struct { + name string + err error + want bool + }{ + {name: "nil", err: nil, want: false}, + {name: "single conflict", err: conflict, want: true}, + {name: "wrapped conflict", err: fmt.Errorf("wrapped: %w", conflict), want: true}, + {name: "pure aggregate", err: pureAggregate, want: true}, + {name: "wrapped pure aggregate", err: fmt.Errorf("wrapped: %w", pureAggregate), want: true}, + {name: "single non-conflict", err: sentinel, want: false}, + {name: "mixed aggregate", err: mixedAggregate, want: false}, + {name: "wrapped mixed aggregate", err: fmt.Errorf("wrapped: %w", mixedAggregate), want: false}, + } { + t.Run(tt.name, func(t *testing.T) { + assert.Equal(t, containsOnlyConflictErrors(tt.err), tt.want) + }) + } +} + func TestSyncExternalPopulatorDirectMaterializationConcurrentWriterRetriesFromFreshState(t *testing.T) { const ( virtualNamespace = "testns" diff --git a/pkg/patcher/patcher.go b/pkg/patcher/patcher.go index 278310c1ac..fc232de1c1 100644 --- a/pkg/patcher/patcher.go +++ b/pkg/patcher/patcher.go @@ -86,14 +86,18 @@ func (h *SyncerPatcher) Patch(ctx *synccontext.SyncContext, pObj, vObj client.Ob return nil } -// RebaseHost records host changes that were committed directly after this -// SyncerPatcher was created. Later deferred patches then start from the -// committed resourceVersion instead of replaying against a stale snapshot. +// RebaseHost records host changes that were committed directly or read from a +// fresher snapshot after this SyncerPatcher was created. Later deferred patches +// then start from that resourceVersion instead of replaying against a stale +// snapshot. A failed rebase disables only the deferred host patch; the virtual +// patch remains enabled. func (h *SyncerPatcher) RebaseHost(obj client.Object) error { if clienthelper.IsNilObject(obj) { + h.pPatcher.SkipHostPatch = true return fmt.Errorf("rebase host patcher: expected non-nil object") } if obj.GetResourceVersion() == "" { + h.pPatcher.SkipHostPatch = true return fmt.Errorf("rebase host patcher for %s/%s: committed resourceVersion is empty", obj.GetNamespace(), obj.GetName()) } diff --git a/pkg/patcher/patcher_test.go b/pkg/patcher/patcher_test.go index cfdd9ae0db..ff5f05ff68 100644 --- a/pkg/patcher/patcher_test.go +++ b/pkg/patcher/patcher_test.go @@ -10,6 +10,7 @@ import ( testingutil "github.com/loft-sh/vcluster/pkg/util/testing" corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "sigs.k8s.io/controller-runtime/pkg/client" ) func TestSyncerPatcherRebaseHostFailClosed(t *testing.T) { @@ -71,3 +72,67 @@ func TestSyncerPatcherRebaseHostFailClosed(t *testing.T) { t.Fatalf("RebaseHost() committed object: %v", err) } } + +func TestSyncerPatcherFailedHostRebaseDisablesOnlyHostPatch(t *testing.T) { + hostClient := testingutil.NewFakeClient(scheme.Scheme, &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{ + Name: "host", + Namespace: "test", + Annotations: map[string]string{ + "example.test/external-writer": "preserve", + }, + }, + }) + virtualClient := testingutil.NewFakeClient(scheme.Scheme, &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{ + Name: "virtual", + Namespace: "test", + }, + }) + ctx := &synccontext.SyncContext{ + Context: context.Background(), + HostClient: hostClient, + VirtualClient: virtualClient, + } + host := &corev1.PersistentVolumeClaim{} + if err := hostClient.Get(ctx, client.ObjectKey{Namespace: "test", Name: "host"}, host); err != nil { + t.Fatalf("get host: %v", err) + } + virtual := &corev1.PersistentVolumeClaim{} + if err := virtualClient.Get(ctx, client.ObjectKey{Namespace: "test", Name: "virtual"}, virtual); err != nil { + t.Fatalf("get virtual: %v", err) + } + syncerPatcher, err := NewSyncerPatcher(ctx, host, virtual) + if err != nil { + t.Fatalf("new syncer patcher: %v", err) + } + + host.Annotations["example.test/stale-controller-diff"] = "must-not-apply" + virtual.Annotations = map[string]string{ + "example.test/virtual-diff": "apply", + } + freshWithoutResourceVersion := host.DeepCopy() + freshWithoutResourceVersion.ResourceVersion = "" + if err := syncerPatcher.RebaseHost(freshWithoutResourceVersion); err == nil { + t.Fatal("RebaseHost() succeeded with empty resourceVersion") + } + if err := syncerPatcher.Patch(ctx, host, virtual); err != nil { + t.Fatalf("Patch() after failed host rebase: %v", err) + } + + actualHost := &corev1.PersistentVolumeClaim{} + if err := hostClient.Get(ctx, client.ObjectKey{Namespace: "test", Name: "host"}, actualHost); err != nil { + t.Fatalf("get actual host: %v", err) + } + if actualHost.Annotations["example.test/external-writer"] != "preserve" || + actualHost.Annotations["example.test/stale-controller-diff"] != "" { + t.Fatalf("host patch was not disabled after failed rebase: annotations=%v", actualHost.Annotations) + } + actualVirtual := &corev1.PersistentVolumeClaim{} + if err := virtualClient.Get(ctx, client.ObjectKey{Namespace: "test", Name: "virtual"}, actualVirtual); err != nil { + t.Fatalf("get actual virtual: %v", err) + } + if actualVirtual.Annotations["example.test/virtual-diff"] != "apply" { + t.Fatalf("virtual patch was disabled with host patch: annotations=%v", actualVirtual.Annotations) + } +}