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.go b/pkg/controllers/resources/persistentvolumeclaims/syncer.go index e054fb815d..f992d4793a 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,14 @@ 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" + "sigs.k8s.io/controller-runtime/pkg/manager" "github.com/loft-sh/vcluster/pkg/util/loghelper" ) @@ -50,6 +56,7 @@ const ( dataProtectionBackupKind = "Backup" externalPopulatorPopulateHelperPrefix = "kb-populate-" + externalPopulatorPVCByUIDIndex = "externalPopulatorPVCByUID" externalPopulatorRestoreConditionType = corev1.PersistentVolumeClaimConditionType("Restore") externalPopulatorPopulateConditionType = corev1.PersistentVolumeClaimConditionType("Populating") @@ -61,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) { @@ -78,6 +89,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 +102,149 @@ 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{} + +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, + 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(), + ) + dependencyRetryQueue.AddRateLimited(&externalPopulatorDependencyRetry{ + object: object.DeepCopyObject().(client.Object), + allowDeletingDependency: allowDeletingDependency, + targetQueue: queue, + }) + 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 +} + +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{} @@ -266,7 +421,7 @@ func (s *persistentVolumeClaimSyncer) Sync(ctx *synccontext.SyncContext, event * retErr = utilerrors.NewAggregate([]error{retErr, err}) } - if kerrors.IsConflict(retErr) { + if containsOnlyConflictErrors(retErr) { result = ctrl.Result{RequeueAfter: time.Second} retErr = nil return @@ -294,9 +449,23 @@ func (s *persistentVolumeClaimSyncer) Sync(ctx *synccontext.SyncContext, event * return ctrl.Result{}, err } if preserveVirtualStatus { - hostConverged, err := s.ensureExternalPopulatorHostMaterialization(ctx, event.Host, event.Virtual, vPV) - if err != nil { - return ctrl.Result{}, err + hostResourceVersionBeforeMaterialization := event.Host.ResourceVersion + 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 { + 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) @@ -326,6 +495,44 @@ func (s *persistentVolumeClaimSyncer) Sync(ctx *synccontext.SyncContext, event * return result, nil } +func containsOnlyConflictErrors(err error) bool { + hasError, onlyConflicts := classifyConflictErrors(err) + return hasError && onlyConflicts +} + +func classifyConflictErrors(err error) (bool, bool) { + if err == nil { + return false, true + } + if aggregate, ok := err.(utilerrors.Aggregate); ok { + hasError := false + for _, aggregateErr := range aggregate.Errors() { + 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 true, kerrors.IsConflict(err) +} + 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 @@ -452,6 +659,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 +708,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 +817,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 +855,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 +869,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 +993,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, @@ -975,6 +1223,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 48427c8ec6..78ac40252a 100644 --- a/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go +++ b/pkg/controllers/resources/persistentvolumeclaims/syncer_test.go @@ -1,11 +1,18 @@ package persistentvolumeclaims import ( + "context" + "errors" + "fmt" + "reflect" "strings" + "sync" + "sync/atomic" "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" @@ -18,12 +25,24 @@ 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 @@ -33,42 +52,1762 @@ func (t *eventRecordingTranslator) EventRecorder() events.EventRecorder { return t.recorder } -func installPVCEventRecorder(syncer *persistentVolumeClaimSyncer) *events.FakeRecorder { - recorder := events.NewFakeRecorder(4) - syncer.GenericTranslator = &eventRecordingTranslator{ - GenericTranslator: syncer.GenericTranslator, - recorder: recorder, +func installPVCEventRecorder(syncer *persistentVolumeClaimSyncer) *events.FakeRecorder { + recorder := events.NewFakeRecorder(4) + syncer.GenericTranslator = &eventRecordingTranslator{ + GenericTranslator: syncer.GenericTranslator, + recorder: recorder, + } + return recorder +} + +func assertSinglePVCEvent(t *testing.T, recorder *events.FakeRecorder, fragments ...string) { + t.Helper() + + var event string + select { + case event = <-recorder.Events: + default: + t.Error("expected one PVC event, got none") + return + } + for _, fragment := range fragments { + assert.Check(t, strings.Contains(event, fragment), "event %q does not contain %q", event, fragment) + } + select { + case extra := <-recorder.Events: + t.Errorf("expected one PVC event, got extra event %q", extra) + default: + } +} + +func assertNoPVCEvent(t *testing.T, recorder *events.FakeRecorder) { + t.Helper() + + select { + case event := <-recorder.Events: + t.Errorf("expected no PVC event, got %q", event) + default: + } +} + +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 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, +) (*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 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 + + 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 + + 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 + + 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, error) + } + 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 { + var err error + requests, err = mapper.externalPopulatorDependencyRequests(syncCtx.Context, convergedDependency) + assert.NilError(t, err) + } + 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 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 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) + }) } - return recorder } -func assertSinglePVCEvent(t *testing.T, recorder *events.FakeRecorder, fragments ...string) { - t.Helper() +func TestSyncExternalPopulatorDirectMaterializationConcurrentWriterRetriesFromFreshState(t *testing.T) { + const ( + virtualNamespace = "testns" + virtualPVName = "restore-populated-pv" + ) - var event string - select { - case event = <-recorder.Events: - default: - t.Error("expected one PVC event, got none") - return + 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}, } - for _, fragment := range fragments { - assert.Check(t, strings.Contains(event, fragment), "event %q does not contain %q", event, fragment) + 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}, } - select { - case extra := <-recorder.Events: - t.Errorf("expected one PVC event, got extra event %q", extra) - default: + + 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 assertNoPVCEvent(t *testing.T, recorder *events.FakeRecorder) { - t.Helper() +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, + }, + } - select { - case event := <-recorder.Events: - t.Errorf("expected no PVC event, got %q", event) - default: + 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, "") + }) } } @@ -581,6 +2320,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 +3332,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{ 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) }, }, diff --git a/pkg/patcher/patcher.go b/pkg/patcher/patcher.go index 84f9a0f409..fc232de1c1 100644 --- a/pkg/patcher/patcher.go +++ b/pkg/patcher/patcher.go @@ -86,6 +86,25 @@ func (h *SyncerPatcher) Patch(ctx *synccontext.SyncContext, pObj, vObj client.Ob return nil } +// 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()) + } + + 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..ff5f05ff68 --- /dev/null +++ b/pkg/patcher/patcher_test.go @@ -0,0 +1,138 @@ +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" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +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) + } +} + +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) + } +}