diff --git a/packages/cluster-operator/cmd/manager/main.go b/packages/cluster-operator/cmd/manager/main.go index 0484a42a..f858b769 100644 --- a/packages/cluster-operator/cmd/manager/main.go +++ b/packages/cluster-operator/cmd/manager/main.go @@ -70,7 +70,7 @@ func main() { ctrl.Log.Error(err, "unable to register T4ClusterHost controller") os.Exit(1) } - if err := (&controllers.WorkspaceReconciler{Client: manager.GetClient(), Scheme: manager.GetScheme()}).SetupWithManager(manager); err != nil { + if err := (&controllers.WorkspaceReconciler{Client: manager.GetClient(), APIReader: manager.GetAPIReader(), Scheme: manager.GetScheme()}).SetupWithManager(manager); err != nil { ctrl.Log.Error(err, "unable to register T4Workspace controller") os.Exit(1) } diff --git a/packages/cluster-operator/controllers/helpers.go b/packages/cluster-operator/controllers/helpers.go index ec5681af..8fef978c 100644 --- a/packages/cluster-operator/controllers/helpers.go +++ b/packages/cluster-operator/controllers/helpers.go @@ -18,6 +18,7 @@ import ( const ( ReasonStorageClassNotFound = "StorageClassNotFound" ReasonStorageClassNotRWX = "StorageClassNotRWX" + ReasonStorageClassMismatch = "StorageClassMismatch" ReasonStorageReady = "StorageClassSupportsRWX" ) @@ -82,6 +83,13 @@ func pvcHasRWX(pvc *corev1.PersistentVolumeClaim) bool { return false } +func pvcStorageClassName(pvc *corev1.PersistentVolumeClaim) string { + if pvc.Spec.StorageClassName == nil { + return "" + } + return *pvc.Spec.StorageClassName +} + func hasString(values []string, wanted string) bool { for _, value := range values { if value == wanted { diff --git a/packages/cluster-operator/controllers/reconciler_test.go b/packages/cluster-operator/controllers/reconciler_test.go index 90ed6600..eccef0a5 100644 --- a/packages/cluster-operator/controllers/reconciler_test.go +++ b/packages/cluster-operator/controllers/reconciler_test.go @@ -2,8 +2,11 @@ package controllers_test import ( "context" + "errors" + "reflect" "strings" "testing" + "time" corev1 "k8s.io/api/core/v1" storagev1 "k8s.io/api/storage/v1" @@ -11,6 +14,7 @@ import ( apiresource "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" "k8s.io/apimachinery/pkg/types" utilvalidation "k8s.io/apimachinery/pkg/util/validation" ctrl "sigs.k8s.io/controller-runtime" @@ -69,6 +73,172 @@ func TestWorkspaceReconcileIsIdempotentAcrossDuplicateEvents(t *testing.T) { } } +func TestWorkspaceCreateAlreadyExistsRefetchesForeignPVC(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + workspace := testWorkspace(clusterv1alpha1.RetentionPolicyDelete) + workspace.UID = "workspace-uid" + base := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&clusterv1alpha1.T4Workspace{}, &corev1.PersistentVolumeClaim{}). + WithObjects(testHost(), rwxStorageClass(), workspace).Build() + c := &createAlreadyExistsClient{Client: base, raceKind: "PVC", hideWinnerFromCache: true} + r := &controllers.WorkspaceReconciler{Client: c, APIReader: base, Scheme: scheme} + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(workspace)}); err != nil { + t.Fatal(err) + } + if c.winner == nil { + t.Fatal("PVC create race was not exercised") + } + var pvc corev1.PersistentVolumeClaim + if err := base.Get(ctx, client.ObjectKeyFromObject(c.winner), &pvc); err != nil { + t.Fatalf("foreign PVC winner was deleted: %v", err) + } + want := c.winner.(*corev1.PersistentVolumeClaim) + if !reflect.DeepEqual(pvc.ObjectMeta, want.ObjectMeta) || !reflect.DeepEqual(pvc.Spec, want.Spec) { + t.Fatalf("foreign PVC winner was mutated: %#v", pvc) + } + var failed clusterv1alpha1.T4Workspace + if err := base.Get(ctx, client.ObjectKeyFromObject(workspace), &failed); err != nil { + t.Fatal(err) + } + storageReady := findCondition(failed.Status.Conditions, "StorageReady") + if failed.Status.PVCName != "" || storageReady == nil || storageReady.Status != metav1.ConditionFalse || storageReady.Reason != "PVCOwnershipConflict" || storageReady.ObservedGeneration != failed.Generation { + t.Fatalf("foreign PVC winner was published as authoritative: status=%#v StorageReady=%#v", failed.Status, storageReady) + } +} + +func TestWorkspaceReadinessRejectsAuthoritativeForeignPVCReplacement(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + workspace := testWorkspace(clusterv1alpha1.RetentionPolicyDelete) + workspace.UID = "workspace-uid" + cachedPVC := ownedWorkspacePVC(workspace) + cachedPVC.UID = "cached-pvc-uid" + authoritativePVC := cachedPVC.DeepCopy() + authoritativePVC.UID = "replacement-pvc-uid" + authoritativePVC.Annotations = nil + authoritativePVC.OwnerReferences = nil + cacheClient := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&clusterv1alpha1.T4Workspace{}, &corev1.PersistentVolumeClaim{}). + WithObjects(testHost(), rwxStorageClass(), workspace, cachedPVC).Build() + r := &controllers.WorkspaceReconciler{ + Client: cacheClient, + APIReader: &pvcOverrideReader{Reader: cacheClient, pvc: authoritativePVC}, + Scheme: scheme, + } + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(workspace)}); err != nil { + t.Fatal(err) + } + var got clusterv1alpha1.T4Workspace + if err := cacheClient.Get(ctx, client.ObjectKeyFromObject(workspace), &got); err != nil { + t.Fatal(err) + } + storageReady := findCondition(got.Status.Conditions, "StorageReady") + ready := findCondition(got.Status.Conditions, "Ready") + if got.Status.PVCName != "" || got.Status.PVCPhase != "" || !got.Status.Capacity.IsZero() || + storageReady == nil || storageReady.Status != metav1.ConditionFalse || storageReady.Reason != "PVCOwnershipConflict" || storageReady.ObservedGeneration != got.Generation || + ready == nil || ready.Status != metav1.ConditionFalse || ready.ObservedGeneration != got.Generation { + t.Fatalf("stale cached PVC published workspace authority: status=%#v StorageReady=%#v Ready=%#v", got.Status, storageReady, ready) + } +} + + +func TestWorkspacePendingPVCPolicyFailsBeforeAuthority(t *testing.T) { + for _, test := range []struct { + name string + mutate func(*corev1.PersistentVolumeClaim) + reason string + }{ + {name: "wrong class", reason: controllers.ReasonStorageClassMismatch, mutate: func(pvc *corev1.PersistentVolumeClaim) { pvc.Spec.StorageClassName = ptr("other-rwx") }}, + {name: "wrong access", reason: "PVCNotRWX", mutate: func(pvc *corev1.PersistentVolumeClaim) { pvc.Spec.AccessModes = []corev1.PersistentVolumeAccessMode{corev1.ReadWriteOnce} }}, + } { + t.Run(test.name, func(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + workspace := testWorkspace(clusterv1alpha1.RetentionPolicyDelete) + workspace.UID = "workspace-uid" + pvc := ownedWorkspacePVC(workspace) + pvc.Status.Phase = corev1.ClaimPending + test.mutate(pvc) + c := fake.NewClientBuilder().WithScheme(scheme).WithStatusSubresource(&clusterv1alpha1.T4Workspace{}, &corev1.PersistentVolumeClaim{}).WithObjects(testHost(), rwxStorageClass(), workspace, pvc).Build() + r := &controllers.WorkspaceReconciler{Client: c, APIReader: c, Scheme: scheme} + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(workspace)}); err != nil { + t.Fatal(err) + } + var failed clusterv1alpha1.T4Workspace + if err := c.Get(ctx, client.ObjectKeyFromObject(workspace), &failed); err != nil { + t.Fatal(err) + } + condition := findCondition(failed.Status.Conditions, "StorageReady") + if failed.Status.PVCName != "" || condition == nil || condition.Status != metav1.ConditionFalse || condition.Reason != test.reason { + t.Fatalf("incompatible Pending PVC published authority: status=%#v condition=%#v", failed.Status, condition) + } + }) + } +} + +func TestWorkspaceDeletionUsesAuthoritativePVCReader(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + workspace := testWorkspace(clusterv1alpha1.RetentionPolicyDelete) + workspace.UID = "workspace-uid" + workspace.Finalizers = []string{clusterv1alpha1.WorkspaceFinalizer} + pvc := ownedWorkspacePVC(workspace) + pvc.OwnerReferences = nil + base := fake.NewClientBuilder().WithScheme(scheme).WithStatusSubresource(&clusterv1alpha1.T4Workspace{}).WithObjects(workspace, pvc).Build() + if err := base.Delete(ctx, workspace); err != nil { + t.Fatal(err) + } + c := &createAlreadyExistsClient{Client: base, raceKind: "PVC", winner: pvc, hideWinnerFromCache: true} + r := &controllers.WorkspaceReconciler{Client: c, APIReader: base, Scheme: scheme} + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(workspace)}); err != nil { + t.Fatal(err) + } + var waiting clusterv1alpha1.T4Workspace + if err := base.Get(ctx, client.ObjectKeyFromObject(workspace), &waiting); err != nil { + t.Fatalf("workspace finalizer ignored authoritative PVC conflict: %v", err) + } + condition := findCondition(waiting.Status.Conditions, "Ready") + if condition == nil || condition.Reason != "CleanupOwnershipConflict" || !contains(waiting.Finalizers, clusterv1alpha1.WorkspaceFinalizer) { + t.Fatalf("authoritative PVC conflict not retained: %#v", waiting) + } +} + +func TestWorkspaceDeletionUsesAuthoritativeSessionReader(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + workspace := testWorkspace(clusterv1alpha1.RetentionPolicyDelete) + workspace.UID = "workspace-uid" + workspace.Finalizers = []string{clusterv1alpha1.WorkspaceFinalizer} + pvc := ownedWorkspacePVC(workspace) + session := testSession() + session.Spec.WorkspaceRef = workspace.Name + cacheClient := fake.NewClientBuilder().WithScheme(scheme).WithStatusSubresource(&clusterv1alpha1.T4Workspace{}).WithObjects(workspace, pvc).Build() + apiReader := fake.NewClientBuilder().WithScheme(scheme).WithObjects(pvc, session).Build() + if err := cacheClient.Delete(ctx, workspace); err != nil { + t.Fatal(err) + } + r := &controllers.WorkspaceReconciler{Client: cacheClient, APIReader: apiReader, Scheme: scheme} + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(workspace)}); err != nil { + t.Fatal(err) + } + var waiting clusterv1alpha1.T4Workspace + if err := cacheClient.Get(ctx, client.ObjectKeyFromObject(workspace), &waiting); err != nil { + t.Fatalf("workspace finalizer ignored authoritative session: %v", err) + } + ready := findCondition(waiting.Status.Conditions, "Ready") + if ready == nil || ready.Status != metav1.ConditionFalse || ready.Reason != "SessionsRemain" || ready.ObservedGeneration != waiting.Generation { + t.Fatalf("Ready = %#v, want current-generation False/SessionsRemain", ready) + } + if !contains(waiting.Finalizers, clusterv1alpha1.WorkspaceFinalizer) { + t.Fatal("workspace finalizer was removed while authoritative session remains") + } + var remainingPVC corev1.PersistentVolumeClaim + if err := cacheClient.Get(ctx, client.ObjectKeyFromObject(pvc), &remainingPVC); err != nil { + t.Fatalf("workspace PVC was deleted while authoritative session remains: %v", err) + } +} + func TestRetainWorkspaceCreatesPVCWithoutGarbageCollectableOwner(t *testing.T) { scheme := testScheme(t) workspace := testWorkspace(clusterv1alpha1.RetentionPolicyRetain) @@ -151,7 +321,7 @@ func TestRetainDeletionOrphansPVCBeforeRemovingFinalizer(t *testing.T) { APIVersion: clusterv1alpha1.GroupVersion.String(), Kind: "T4Workspace", Name: workspace.Name, UID: workspace.UID, Controller: ptr(true), }}, }, - Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, } c := fake.NewClientBuilder().WithScheme(scheme).WithStatusSubresource(&clusterv1alpha1.T4Workspace{}).WithObjects(testHost(), rwxStorageClass(), workspace, pvc).Build() if err := c.Delete(context.Background(), workspace); err != nil { @@ -185,7 +355,7 @@ func TestWorkspaceDeletionWaitsForSessionResources(t *testing.T) { Name: controllers.WorkspacePVCName(workspace), Namespace: workspace.Namespace, OwnerReferences: []metav1.OwnerReference{{APIVersion: clusterv1alpha1.GroupVersion.String(), Kind: "T4Workspace", Name: workspace.Name, UID: workspace.UID, Controller: ptr(true)}}, }, - Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, } session := testSession() session.Spec.WorkspaceRef = workspace.Name @@ -214,6 +384,107 @@ func TestWorkspaceDeletionWaitsForSessionResources(t *testing.T) { } } +func TestSessionPodCreateRevalidatesAuthoritativePVC(t *testing.T) { + tests := []struct { + name string + mutate func(*corev1.PersistentVolumeClaim) + }{ + {name: "replacement UID", mutate: func(pvc *corev1.PersistentVolumeClaim) { pvc.UID = "replacement-pvc-uid" }}, + {name: "foreign owner", mutate: func(pvc *corev1.PersistentVolumeClaim) { + pvc.OwnerReferences = []metav1.OwnerReference{{APIVersion: "v1", Kind: "Secret", Name: "foreign", UID: "foreign-uid", Controller: ptr(true)}} + }}, + {name: "storage class drift", mutate: func(pvc *corev1.PersistentVolumeClaim) { pvc.Spec.StorageClassName = ptr("other-rwx") }}, + {name: "access mode drift", mutate: func(pvc *corev1.PersistentVolumeClaim) { pvc.Spec.AccessModes = []corev1.PersistentVolumeAccessMode{corev1.ReadWriteOnce} }}, + {name: "unbound replacement", mutate: func(pvc *corev1.PersistentVolumeClaim) { pvc.Status.Phase = corev1.ClaimPending }}, + } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + workspace := testWorkspace(clusterv1alpha1.RetentionPolicyDelete) + workspace.UID = "workspace-uid" + workspace.Status.PVCName = controllers.WorkspacePVCName(workspace) + cachedPVC := ownedWorkspacePVC(workspace) + cachedPVC.UID = "cached-pvc-uid" + authoritativePVC := cachedPVC.DeepCopy() + test.mutate(authoritativePVC) + session := testSession() + session.UID = "session-uid" + cacheClient := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&clusterv1alpha1.T4Session{}, &corev1.PersistentVolumeClaim{}, &corev1.Pod{}). + WithObjects(testHost(), workspace, cachedPVC, session).Build() + r := configuredSessionReconciler(cacheClient, scheme) + r.APIReader = &pvcOverrideReader{Reader: cacheClient, pvc: authoritativePVC} + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(session)}); err != nil { + t.Fatal(err) + } + var pods corev1.PodList + if err := cacheClient.List(ctx, &pods, client.InNamespace(session.Namespace)); err != nil { + t.Fatal(err) + } + if len(pods.Items) != 0 { + t.Fatalf("created %d Pods after authoritative PVC %s", len(pods.Items), test.name) + } + var got clusterv1alpha1.T4Session + if err := cacheClient.Get(ctx, client.ObjectKeyFromObject(session), &got); err != nil { + t.Fatal(err) + } + workspaceReady := findCondition(got.Status.Conditions, "WorkspaceReady") + if got.Status.PodName != "" || workspaceReady == nil || workspaceReady.Status != metav1.ConditionFalse || workspaceReady.Reason != "PVCAuthorityChanged" || workspaceReady.ObservedGeneration != got.Generation { + t.Fatalf("authoritative PVC %s published session authority: status=%#v WorkspaceReady=%#v", test.name, got.Status, workspaceReady) + } + }) + } +} + +func TestSessionExistingPodRepairRejectsAuthoritativeForeignPVCReplacement(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + workspace := testWorkspace(clusterv1alpha1.RetentionPolicyDelete) + workspace.UID = "workspace-uid" + workspace.Status.PVCName = controllers.WorkspacePVCName(workspace) + cachedPVC := ownedWorkspacePVC(workspace) + cachedPVC.UID = "cached-pvc-uid" + session := testSession() + session.UID = "session-uid" + cacheClient := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&clusterv1alpha1.T4Session{}, &corev1.PersistentVolumeClaim{}, &corev1.Pod{}). + WithObjects(testHost(), workspace, cachedPVC, session).Build() + r := configuredSessionReconciler(cacheClient, scheme) + reconcileMany(t, 2, func() error { + _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(session)}) + return err + }) + var pod corev1.Pod + if err := cacheClient.Get(ctx, types.NamespacedName{Namespace: session.Namespace, Name: controllers.SessionPodName(session)}, &pod); err != nil { + t.Fatal(err) + } + pod.Labels = nil + if err := cacheClient.Update(ctx, &pod); err != nil { + t.Fatal(err) + } + authoritativePVC := cachedPVC.DeepCopy() + authoritativePVC.UID = "replacement-pvc-uid" + authoritativePVC.Annotations = nil + authoritativePVC.OwnerReferences = nil + r.APIReader = &pvcOverrideReader{Reader: cacheClient, pvc: authoritativePVC} + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(session)}); err != nil { + t.Fatal(err) + } + if err := cacheClient.Get(ctx, client.ObjectKeyFromObject(&pod), &pod); !apierrors.IsNotFound(err) { + t.Fatalf("existing Pod survived authoritative PVC replacement: %v", err) + } + var got clusterv1alpha1.T4Session + if err := cacheClient.Get(ctx, client.ObjectKeyFromObject(session), &got); err != nil { + t.Fatal(err) + } + workspaceReady := findCondition(got.Status.Conditions, "WorkspaceReady") + if got.Status.PodName != "" || workspaceReady == nil || workspaceReady.Status != metav1.ConditionFalse || workspaceReady.Reason != "PVCAuthorityChanged" || workspaceReady.ObservedGeneration != got.Generation { + t.Fatalf("existing Pod repair retained stale PVC authority: status=%#v WorkspaceReady=%#v", got.Status, workspaceReady) + } +} + + func TestSessionFailsClosedWhenAnyOMPReferenceIsMissing(t *testing.T) { for _, test := range []struct { name string @@ -230,7 +501,7 @@ func TestSessionFailsClosedWhenAnyOMPReferenceIsMissing(t *testing.T) { workspace.Status.PVCName = "workspace-a-data" pvc := &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: "team"}, - Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}, } session := testSession() @@ -307,7 +578,7 @@ func TestSessionRejectsCredentialBearingModelsConfiguration(t *testing.T) { workspace.Status.PVCName = "workspace-a-data" pvc := &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: "team"}, - Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}, } session := testSession() @@ -358,7 +629,7 @@ func TestSessionRejectsCredentialBearingSettingsConfiguration(t *testing.T) { workspace.Status.PVCName = "workspace-a-data" pvc := &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: "team"}, - Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}, } session := testSession() @@ -407,7 +678,7 @@ func TestSessionRuntimeImageMustBeImmutableDigest(t *testing.T) { workspace.Status.PVCName = "workspace-a-data" pvc := &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: "team"}, - Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}, } session := testSession() @@ -478,7 +749,7 @@ func TestSessionAuthorityRevocationDeletesOwnedPodAndService(t *testing.T) { workspace.Status.PVCName = "workspace-a-data" pvc := &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: "team"}, - Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}, } session := testSession() @@ -525,7 +796,7 @@ func TestSessionNamesProduceSafeRuntimeIdentities(t *testing.T) { workspace.Status.PVCName = "workspace-a-data" pvc := &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: "team"}, - Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}, } session := testSession() @@ -570,7 +841,7 @@ func TestSessionWaitsForBoundRWXThenCreatesExactlyOnePodAndService(t *testing.T) workspace.Status.Phase = clusterv1alpha1.InfrastructurePending pvc := &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: "team"}, - Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimPending}, } session := testSession() @@ -701,7 +972,7 @@ func TestSessionOMPModeOmitsCredentialReferences(t *testing.T) { workspace.Status.PVCName = "workspace-a-data" pvc := &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: "team"}, - Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}, } session := testSession() @@ -733,7 +1004,7 @@ func TestSessionRejectsUnownedDeterministicResources(t *testing.T) { workspace.Status.PVCName = "workspace-a-data" pvc := &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: "team"}, - Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}, } session := testSession() @@ -768,7 +1039,7 @@ func TestSessionRecreatesPodWhenImmutableDesiredStateChanges(t *testing.T) { workspace.Status.PVCName = "workspace-a-data" pvc := &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: "team"}, - Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}, } session := testSession() @@ -823,7 +1094,7 @@ func TestSessionPodHashIncludesEveryOMPReference(t *testing.T) { workspace.Status.PVCName = "workspace-a-data" pvc := &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: "team"}, - Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}, } session := testSession() @@ -869,7 +1140,7 @@ func TestSessionRecreatesPodWhenOMPResourceVersionChanges(t *testing.T) { workspace.Status.PVCName = "workspace-a-data" pvc := &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: "team"}, - Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}, } session := testSession() @@ -912,7 +1183,7 @@ func TestSessionRuntimeReferenceRevocationStopsAuthority(t *testing.T) { workspace.Status.PVCName = "workspace-a-data" pvc := &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: "team"}, - Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}, } session := testSession() @@ -966,7 +1237,7 @@ func TestSessionFailsClosedWhenOMPConfigMapIsMissing(t *testing.T) { scheme := testScheme(t) workspace := testWorkspace(clusterv1alpha1.RetentionPolicyDelete) workspace.Status.PVCName = "workspace-a-data" - pvc := &corev1.PersistentVolumeClaim{ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: "team"}, Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}} + pvc := &corev1.PersistentVolumeClaim{ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: "team"}, Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}} session := testSession() objects := []client.Object{testHost(), rwxStorageClass(), workspace, pvc, session} if test.configMap != nil { @@ -996,7 +1267,7 @@ func TestSessionRecreatesExternallyExposedOwnedService(t *testing.T) { workspace.Status.PVCName = "workspace-a-data" pvc := &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: "team"}, - Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}, } session := testSession() @@ -1044,7 +1315,7 @@ func TestSessionRestoresRequiredPodSelectorLabelsBeforeAvailability(t *testing.T workspace.Status.PVCName = "workspace-a-data" pvc := &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: "team"}, - Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}, } session := testSession() @@ -1086,7 +1357,7 @@ func TestSessionRestoresRequiredPodSelectorLabelsBeforeAvailability(t *testing.T func TestWorkspaceDeletionRefusesForeignDeterministicPVC(t *testing.T) { for _, policy := range []clusterv1alpha1.RetentionPolicy{clusterv1alpha1.RetentionPolicyRetain, clusterv1alpha1.RetentionPolicyDelete} { - for _, mismatch := range []string{"uid-annotation", "controller-owner"} { + for _, mismatch := range []string{"uid-annotation", "controller-owner", "foreign-non-controller-owner"} { t.Run(string(policy)+"/"+mismatch, func(t *testing.T) { scheme := testScheme(t) workspace := testWorkspace(policy) @@ -1098,10 +1369,14 @@ func TestWorkspaceDeletionRefusesForeignDeterministicPVC(t *testing.T) { }} if mismatch == "uid-annotation" { pvc.Annotations[clusterv1alpha1.WorkspaceUIDAnnotation] = "foreign-workspace-uid" - } else { + } else if mismatch == "controller-owner" { pvc.OwnerReferences = []metav1.OwnerReference{{ APIVersion: clusterv1alpha1.GroupVersion.String(), Kind: "T4Workspace", Name: "foreign", UID: "foreign-workspace-uid", Controller: ptr(true), }} + } else { + pvc.OwnerReferences = []metav1.OwnerReference{{ + APIVersion: "example.test/v1", Kind: "Foreign", Name: "foreign", UID: "foreign-uid", + }} } expectedOwnerCount := len(pvc.OwnerReferences) c := fake.NewClientBuilder().WithScheme(scheme).WithStatusSubresource(&clusterv1alpha1.T4Workspace{}).WithObjects(workspace, pvc).Build() @@ -1151,6 +1426,14 @@ func TestSessionDeletionRefusesForeignDeterministicResources(t *testing.T) { service.OwnerReferences = nil } c := fake.NewClientBuilder().WithScheme(scheme).WithStatusSubresource(&clusterv1alpha1.T4Session{}).WithObjects(session, pod, service).Build() + beforePod := &corev1.Pod{} + if err := c.Get(context.Background(), client.ObjectKeyFromObject(pod), beforePod); err != nil { + t.Fatal(err) + } + beforeService := &corev1.Service{} + if err := c.Get(context.Background(), client.ObjectKeyFromObject(service), beforeService); err != nil { + t.Fatal(err) + } if err := c.Delete(context.Background(), session); err != nil { t.Fatal(err) } @@ -1158,7 +1441,28 @@ func TestSessionDeletionRefusesForeignDeterministicResources(t *testing.T) { if _, err := r.Reconcile(context.Background(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(session)}); err != nil { t.Fatal(err) } - assertObjectCounts(t, c, 1, 1) + if foreignKind == "Pod" { + assertObjectCounts(t, c, 1, 0) + } else { + assertObjectCounts(t, c, 0, 1) + } + if foreignKind == "Pod" { + var got corev1.Pod + if err := c.Get(context.Background(), client.ObjectKeyFromObject(pod), &got); err != nil { + t.Fatalf("foreign Pod was deleted: %v", err) + } + if !reflect.DeepEqual(got.ObjectMeta, beforePod.ObjectMeta) || !reflect.DeepEqual(got.Spec, beforePod.Spec) { + t.Fatalf("foreign Pod was mutated during finalizer cleanup: %#v", got) + } + } else { + var got corev1.Service + if err := c.Get(context.Background(), client.ObjectKeyFromObject(service), &got); err != nil { + t.Fatalf("foreign Service was deleted: %v", err) + } + if !reflect.DeepEqual(got.ObjectMeta, beforeService.ObjectMeta) || !reflect.DeepEqual(got.Spec, beforeService.Spec) { + t.Fatalf("foreign Service was mutated during finalizer cleanup: %#v", got) + } + } var waiting clusterv1alpha1.T4Session if err := c.Get(context.Background(), client.ObjectKeyFromObject(session), &waiting); err != nil { t.Fatalf("session finalizer was removed on cleanup conflict: %v", err) @@ -1171,6 +1475,100 @@ func TestSessionDeletionRefusesForeignDeterministicResources(t *testing.T) { } } +func TestSessionDeletionUsesAuthoritativeChildReader(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + session := testSession() + session.UID = "session-uid" + session.Finalizers = []string{clusterv1alpha1.SessionFinalizer} + pod, service := ownedSessionResources(session) + pod.OwnerReferences = nil + base := fake.NewClientBuilder().WithScheme(scheme).WithStatusSubresource(&clusterv1alpha1.T4Session{}).WithObjects(session, pod, service).Build() + if err := base.Delete(ctx, session); err != nil { + t.Fatal(err) + } + c := &createAlreadyExistsClient{Client: base, raceKind: "Pod", winner: pod, hideWinnerFromCache: true} + r := configuredSessionReconciler(c, scheme) + r.APIReader = base + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(session)}); err != nil { + t.Fatal(err) + } + assertObjectCounts(t, base, 1, 0) + var waiting clusterv1alpha1.T4Session + if err := base.Get(ctx, client.ObjectKeyFromObject(session), &waiting); err != nil { + t.Fatalf("session finalizer ignored authoritative Pod conflict: %v", err) + } + condition := findCondition(waiting.Status.Conditions, "Available") + if condition == nil || condition.Reason != "CleanupOwnershipConflict" || !contains(waiting.Finalizers, clusterv1alpha1.SessionFinalizer) { + t.Fatalf("authoritative child conflict not retained: %#v", waiting) + } +} + +func TestSessionDeletionPreconditionsProtectSameNameReplacement(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + session := testSession() + session.UID = "session-uid" + session.Finalizers = []string{clusterv1alpha1.SessionFinalizer} + pod, service := ownedSessionResources(session) + pod.UID = "pod-uid-a" + pod.ResourceVersion = "7" + service.UID = "service-uid-a" + service.ResourceVersion = "8" + base := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&clusterv1alpha1.T4Session{}). + WithObjects(session, pod, service).Build() + var observedPod corev1.Pod + if err := base.Get(ctx, client.ObjectKeyFromObject(pod), &observedPod); err != nil { + t.Fatal(err) + } + if observedPod.UID == "" || observedPod.ResourceVersion == "" { + t.Fatalf("fake client discarded delete precondition identity: %#v", observedPod.ObjectMeta) + } + if err := base.Delete(ctx, session); err != nil { + t.Fatal(err) + } + racingClient := &replaceBeforeDeleteClient{Client: base, raceKind: "Pod", replacementUID: "pod-uid-b"} + r := configuredSessionReconciler(racingClient, scheme) + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(session)}); !apierrors.IsConflict(err) { + t.Fatalf("same-name replacement did not conflict with the stale delete: %v", err) + } + if !racingClient.raced || racingClient.observedUID != observedPod.UID || racingClient.observedResourceVersion != observedPod.ResourceVersion { + t.Fatalf("delete did not carry both observed preconditions: raced=%t uid=%q resourceVersion=%q", racingClient.raced, racingClient.observedUID, racingClient.observedResourceVersion) + } + var replacement corev1.Pod + if err := base.Get(ctx, client.ObjectKeyFromObject(pod), &replacement); err != nil { + t.Fatalf("same-name replacement Pod did not survive stale delete: %v", err) + } + if replacement.UID != racingClient.replacementUID { + t.Fatalf("surviving Pod UID = %q, want replacement %q", replacement.UID, racingClient.replacementUID) + } + var waiting clusterv1alpha1.T4Session + if err := base.Get(ctx, client.ObjectKeyFromObject(session), &waiting); err != nil { + t.Fatalf("session finalizer advanced after stale child delete: %v", err) + } + available := findCondition(waiting.Status.Conditions, "Available") + if !contains(waiting.Finalizers, clusterv1alpha1.SessionFinalizer) || waiting.Status.Phase != clusterv1alpha1.InfrastructureTerminating || available == nil || available.Reason != "Terminating" { + t.Fatalf("session advanced after stale child delete: %#v", waiting) + } +} + +func TestSessionDependencyCleanupUsesAuthoritativeChildReader(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + session := testSession() + session.UID = "session-uid" + pod, service := ownedSessionResources(session) + base := fake.NewClientBuilder().WithScheme(scheme).WithStatusSubresource(&clusterv1alpha1.T4Session{}, &corev1.Pod{}).WithObjects(session, pod, service).Build() + c := &createAlreadyExistsClient{Client: base, raceKind: "Pod", winner: pod, hideWinnerFromCache: true} + r := configuredSessionReconciler(c, scheme) + r.APIReader = base + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(session)}); err != nil { + t.Fatal(err) + } + assertObjectCounts(t, base, 0, 0) +} + func TestSessionDeletionCleansResourcesBeforeFinalizer(t *testing.T) { scheme := testScheme(t) session := testSession() @@ -1198,51 +1596,753 @@ func TestSessionDeletionCleansResourcesBeforeFinalizer(t *testing.T) { } } -func configuredSessionReconciler(c client.Client, scheme *runtime.Scheme) *controllers.SessionReconciler { - for _, object := range []client.Object{ - &corev1.ConfigMap{ObjectMeta: metav1.ObjectMeta{Name: "omp-runtime-config", Namespace: "team"}, Data: map[string]string{ - "provider-models": testOMPModels, "agent-settings": testOMPSettings, "other-models": otherTestOMPModels, "other-settings": otherTestOMPSettings, +func TestSessionDependencyRevocationCleansOwnedResourcesAndConvergesAfterRestart(t *testing.T) { + for _, test := range []struct { + name string + conditionType string + wantReason string + revoke func(context.Context, client.Client) error + }{ + {name: "missing Host", conditionType: "HostReady", wantReason: "HostNotFound", revoke: func(ctx context.Context, c client.Client) error { + return c.Delete(ctx, &clusterv1alpha1.T4ClusterHost{ObjectMeta: metav1.ObjectMeta{Name: "host-a", Namespace: "team"}}) }}, - &corev1.ConfigMap{ObjectMeta: metav1.ObjectMeta{Name: "other-omp-config", Namespace: "team"}, Data: map[string]string{ - "provider-models": testOMPModels, "agent-settings": testOMPSettings, "other-models": otherTestOMPModels, "other-settings": otherTestOMPSettings, + {name: "invalid Host runtime profile", conditionType: "RuntimeConfigured", wantReason: "RuntimeProfileNotAllowed", revoke: func(ctx context.Context, c client.Client) error { + var host clusterv1alpha1.T4ClusterHost + if err := c.Get(ctx, types.NamespacedName{Namespace: "team", Name: "host-a"}, &host); err != nil { + return err + } + host.Spec.RuntimeProfiles = nil + return c.Update(ctx, &host) + }}, + {name: "missing Workspace", conditionType: "WorkspaceReady", wantReason: "WorkspaceNotFound", revoke: func(ctx context.Context, c client.Client) error { + return c.Delete(ctx, &clusterv1alpha1.T4Workspace{ObjectMeta: metav1.ObjectMeta{Name: "workspace-a", Namespace: "team"}}) + }}, + {name: "mismatched Workspace Host", conditionType: "WorkspaceReady", wantReason: "HostMismatch", revoke: func(ctx context.Context, c client.Client) error { + var workspace clusterv1alpha1.T4Workspace + if err := c.Get(ctx, types.NamespacedName{Namespace: "team", Name: "workspace-a"}, &workspace); err != nil { + return err + } + workspace.Spec.HostRef = "host-b" + return c.Update(ctx, &workspace) + }}, + {name: "missing OMP ConfigMap", conditionType: "RuntimeConfigured", wantReason: "OMPConfigMapNotFound", revoke: func(ctx context.Context, c client.Client) error { + return c.Delete(ctx, &corev1.ConfigMap{ObjectMeta: metav1.ObjectMeta{Name: "omp-runtime-config", Namespace: "team"}}) + }}, + {name: "mismatched Host storage class", conditionType: "WorkspaceReady", wantReason: "StorageClassMismatch", revoke: func(ctx context.Context, c client.Client) error { + otherClass := &storagev1.StorageClass{ObjectMeta: metav1.ObjectMeta{Name: "other-rwx", Annotations: map[string]string{clusterv1alpha1.RWXStorageClassAnnotation: string(corev1.ReadWriteMany)}}, Provisioner: "example.invalid/csi"} + if err := c.Create(ctx, otherClass); err != nil { + return err + } + var host clusterv1alpha1.T4ClusterHost + if err := c.Get(ctx, types.NamespacedName{Namespace: "team", Name: "host-a"}, &host); err != nil { + return err + } + host.Spec.StorageClassName = otherClass.Name + return c.Update(ctx, &host) + }}, + {name: "classless Workspace PVC", conditionType: "WorkspaceReady", wantReason: "StorageClassMismatch", revoke: func(ctx context.Context, c client.Client) error { + var pvc corev1.PersistentVolumeClaim + if err := c.Get(ctx, types.NamespacedName{Namespace: "team", Name: "workspace-a-data"}, &pvc); err != nil { + return err + } + pvc.Spec.StorageClassName = nil + return c.Update(ctx, &pvc) }}, - rwxStorageClass(), } { - if err := c.Create(context.Background(), object); err != nil && !apierrors.IsAlreadyExists(err) { - panic(err) - } - } - return &controllers.SessionReconciler{ - Client: c, - APIReader: c, - Scheme: scheme, - RuntimeImage: testRuntimeImage, - OMPConfig: controllers.SessionOMPConfig{ - ConfigMapName: "omp-runtime-config", - ModelsKey: "provider-models", - SettingsKey: "agent-settings", - }, - } -} - -func testScheme(t *testing.T) *runtime.Scheme { - t.Helper() - scheme := runtime.NewScheme() - for _, add := range []func(*runtime.Scheme) error{corev1.AddToScheme, storagev1.AddToScheme, clusterv1alpha1.AddToScheme} { - if err := add(scheme); err != nil { - t.Fatal(err) - } - } - return scheme -} - -func testHost() *clusterv1alpha1.T4ClusterHost { - return &clusterv1alpha1.T4ClusterHost{ - ObjectMeta: metav1.ObjectMeta{Name: "host-a", Namespace: "team", UID: "host-uid"}, - Spec: clusterv1alpha1.T4ClusterHostSpec{StorageClassName: "portable-rwx", RuntimeProfiles: []string{"default"}}, - } -} - + t.Run(test.name, func(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + workspace := testWorkspace(clusterv1alpha1.RetentionPolicyDelete) + workspace.Status.PVCName = "workspace-a-data" + pvc := &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: "team"}, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}, + } + session := testSession() + session.UID = "session-uid" + foreignPod := &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: "foreign-pod", Namespace: "team"}} + foreignService := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: "foreign-service", Namespace: "team"}} + c := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&clusterv1alpha1.T4Session{}, &corev1.PersistentVolumeClaim{}, &corev1.Pod{}). + WithObjects(testHost(), workspace, pvc, session, foreignPod, foreignService).Build() + r := configuredSessionReconciler(c, scheme) + reconcileMany(t, 2, func() error { + _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(session)}) + return err + }) + if err := test.revoke(ctx, c); err != nil { + t.Fatal(err) + } + + result, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(session)}) + if err != nil { + t.Fatal(err) + } + if result.RequeueAfter <= 0 || result.RequeueAfter > 30*time.Second { + t.Fatalf("revoked dependency requeue = %s, want bounded positive retry", result.RequeueAfter) + } + for _, object := range []client.Object{ + &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: controllers.SessionPodName(session), Namespace: session.Namespace}}, + &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: controllers.SessionServiceName(session), Namespace: session.Namespace}}, + } { + if err := c.Get(ctx, client.ObjectKeyFromObject(object), object); !apierrors.IsNotFound(err) { + t.Fatalf("owned stale %T remained after dependency revocation: %v", object, err) + } + } + for _, object := range []client.Object{foreignPod.DeepCopy(), foreignService.DeepCopy()} { + if err := c.Get(ctx, client.ObjectKeyFromObject(object), object); err != nil { + t.Fatalf("unowned %T was removed: %v", object, err) + } + } + + var failed clusterv1alpha1.T4Session + if err := c.Get(ctx, client.ObjectKeyFromObject(session), &failed); err != nil { + t.Fatal(err) + } + condition := findCondition(failed.Status.Conditions, test.conditionType) + available := findCondition(failed.Status.Conditions, "Available") + if failed.Status.ObservedGeneration != failed.Generation || failed.Status.PodName != "" || failed.Status.ServiceName != "" || failed.Status.Phase != clusterv1alpha1.InfrastructureFailed || + condition == nil || condition.Status != metav1.ConditionFalse || condition.Reason != test.wantReason || condition.ObservedGeneration != failed.Generation || + available == nil || available.Status != metav1.ConditionFalse || available.Reason != test.wantReason || available.ObservedGeneration != failed.Generation { + t.Fatalf("revoked session did not converge: status=%#v condition=%#v available=%#v", failed.Status, condition, available) + } + stableStatus := failed.Status + restarted := *r + reconcileMany(t, 2, func() error { + _, err := restarted.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(session)}) + return err + }) + if err := c.Get(ctx, client.ObjectKeyFromObject(session), &failed); err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(failed.Status, stableStatus) { + t.Fatalf("duplicate/restart reconciliation changed converged status: got %#v, want %#v", failed.Status, stableStatus) + } + }) + } +} + +func TestWorkspaceHostStorageClassDriftFailsClosedWithoutRecreatingPVC(t *testing.T) { + oldClass := "portable-rwx" + for _, test := range []struct { + name string + claimClass *string + }{ + {name: "different class", claimClass: &oldClass}, + {name: "class omitted"}, + } { + t.Run(test.name, func(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + workspace := testWorkspace(clusterv1alpha1.RetentionPolicyDelete) + workspace.UID = "workspace-uid" + workspace.Status.ObservedGeneration = workspace.Generation + workspace.Status.PVCName = controllers.WorkspacePVCName(workspace) + workspace.Status.Phase = clusterv1alpha1.InfrastructureReady + workspace.Status.Conditions = []metav1.Condition{ + {Type: "StorageReady", Status: metav1.ConditionTrue, Reason: controllers.ReasonStorageReady, ObservedGeneration: workspace.Generation}, + {Type: "Ready", Status: metav1.ConditionTrue, Reason: "PVCBound", ObservedGeneration: workspace.Generation}, + } + pvc := &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{ + Name: workspace.Status.PVCName, Namespace: workspace.Namespace, + Annotations: map[string]string{clusterv1alpha1.WorkspaceUIDAnnotation: string(workspace.UID)}, + OwnerReferences: []metav1.OwnerReference{{APIVersion: clusterv1alpha1.GroupVersion.String(), Kind: "T4Workspace", Name: workspace.Name, UID: workspace.UID, Controller: ptr(true)}}, + }, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: test.claimClass, AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}, + } + host := testHost() + host.Spec.StorageClassName = "other-rwx" + otherClass := &storagev1.StorageClass{ObjectMeta: metav1.ObjectMeta{Name: "other-rwx", Annotations: map[string]string{clusterv1alpha1.RWXStorageClassAnnotation: string(corev1.ReadWriteMany)}}, Provisioner: "example.invalid/csi"} + c := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&clusterv1alpha1.T4Workspace{}, &corev1.PersistentVolumeClaim{}). + WithObjects(host, otherClass, workspace, pvc).Build() + r := &controllers.WorkspaceReconciler{Client: c, Scheme: scheme} + + result, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(workspace)}) + if err != nil { + t.Fatal(err) + } + if result.RequeueAfter <= 0 || result.RequeueAfter > 30*time.Second { + t.Fatalf("storage drift requeue = %s, want bounded positive retry", result.RequeueAfter) + } + var got clusterv1alpha1.T4Workspace + if err := c.Get(ctx, client.ObjectKeyFromObject(workspace), &got); err != nil { + t.Fatal(err) + } + storageReady := findCondition(got.Status.Conditions, "StorageReady") + ready := findCondition(got.Status.Conditions, "Ready") + if got.Status.Phase != clusterv1alpha1.InfrastructureFailed || storageReady == nil || storageReady.Status != metav1.ConditionFalse || storageReady.Reason != "StorageClassMismatch" || ready == nil || ready.Status != metav1.ConditionFalse || ready.Reason != "StorageClassMismatch" { + t.Fatalf("storage drift remained Ready: status=%#v StorageReady=%#v Ready=%#v", got.Status, storageReady, ready) + } + var retained corev1.PersistentVolumeClaim + if err := c.Get(ctx, client.ObjectKeyFromObject(pvc), &retained); err != nil { + t.Fatalf("storage drift removed the data PVC: %v", err) + } + if !reflect.DeepEqual(retained.Spec.StorageClassName, test.claimClass) { + t.Fatalf("storage drift recreated or mutated PVC class: got %#v, want %#v", retained.Spec.StorageClassName, test.claimClass) + } + }) + } +} + +func TestHostDependencyRecoveryReplacesFalseConditions(t *testing.T) { + t.Run("Workspace", func(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + workspace := testWorkspace(clusterv1alpha1.RetentionPolicyDelete) + workspace.UID = "workspace-uid" + c := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&clusterv1alpha1.T4Workspace{}, &corev1.PersistentVolumeClaim{}). + WithObjects(workspace).Build() + r := &controllers.WorkspaceReconciler{Client: c, Scheme: scheme} + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(workspace)}); err != nil { + t.Fatal(err) + } + if err := c.Create(ctx, testHost()); err != nil { + t.Fatal(err) + } + if err := c.Create(ctx, rwxStorageClass()); err != nil { + t.Fatal(err) + } + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(workspace)}); err != nil { + t.Fatal(err) + } + var pvc corev1.PersistentVolumeClaim + if err := c.Get(ctx, types.NamespacedName{Namespace: workspace.Namespace, Name: controllers.WorkspacePVCName(workspace)}, &pvc); err != nil { + t.Fatal(err) + } + pvc.Status.Phase = corev1.ClaimBound + if err := c.Status().Update(ctx, &pvc); err != nil { + t.Fatal(err) + } + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(workspace)}); err != nil { + t.Fatal(err) + } + var recovered clusterv1alpha1.T4Workspace + if err := c.Get(ctx, client.ObjectKeyFromObject(workspace), &recovered); err != nil { + t.Fatal(err) + } + hostReady := findCondition(recovered.Status.Conditions, "HostReady") + if hostReady == nil || hostReady.Status != metav1.ConditionTrue || hostReady.ObservedGeneration != recovered.Generation { + t.Fatalf("recovered Workspace retained stale HostReady: %#v", hostReady) + } + }) + + t.Run("Session", func(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + workspace := testWorkspace(clusterv1alpha1.RetentionPolicyDelete) + workspace.Status.PVCName = "workspace-a-data" + pvc := &corev1.PersistentVolumeClaim{ObjectMeta: metav1.ObjectMeta{Name: workspace.Status.PVCName, Namespace: workspace.Namespace}, Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}} + session := testSession() + session.UID = "session-uid" + c := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&clusterv1alpha1.T4Session{}, &corev1.PersistentVolumeClaim{}, &corev1.Pod{}). + WithObjects(workspace, pvc, session).Build() + r := configuredSessionReconciler(c, scheme) + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(session)}); err != nil { + t.Fatal(err) + } + if err := c.Create(ctx, testHost()); err != nil { + t.Fatal(err) + } + reconcileMany(t, 2, func() error { + _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(session)}) + return err + }) + var recovered clusterv1alpha1.T4Session + if err := c.Get(ctx, client.ObjectKeyFromObject(session), &recovered); err != nil { + t.Fatal(err) + } + hostReady := findCondition(recovered.Status.Conditions, "HostReady") + if hostReady == nil || hostReady.Status != metav1.ConditionTrue || hostReady.ObservedGeneration != recovered.Generation { + t.Fatalf("recovered Session retained stale HostReady: %#v", hostReady) + } + }) + + t.Run("static runtime failure", func(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + session := testSession() + session.Status.ObservedGeneration = session.Generation + session.Status.Phase = clusterv1alpha1.InfrastructureFailed + session.Status.Conditions = []metav1.Condition{{Type: "HostReady", Status: metav1.ConditionFalse, Reason: "HostNotFound", ObservedGeneration: session.Generation}} + c := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&clusterv1alpha1.T4Session{}). + WithObjects(testHost(), session).Build() + r := configuredSessionReconciler(c, scheme) + r.RuntimeImage = "registry.example/session:latest" + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(session)}); err != nil { + t.Fatal(err) + } + var failed clusterv1alpha1.T4Session + if err := c.Get(ctx, client.ObjectKeyFromObject(session), &failed); err != nil { + t.Fatal(err) + } + hostReady := findCondition(failed.Status.Conditions, "HostReady") + if hostReady == nil || hostReady.Status != metav1.ConditionUnknown || hostReady.Reason != "NotEvaluated" || hostReady.ObservedGeneration != failed.Generation { + t.Fatalf("static runtime failure retained stale HostReady: %#v", hostReady) + } + }) +} + +func TestWorkspaceFailureRevokesPublishedPVCAuthorityAndRefreshesConditions(t *testing.T) { + for _, test := range []struct { + name string + objects func(*clusterv1alpha1.T4Workspace) []client.Object + wantHostStatus metav1.ConditionStatus + wantStorageReason string + }{ + { + name: "missing Host", + objects: func(workspace *clusterv1alpha1.T4Workspace) []client.Object { + return []client.Object{workspace} + }, + wantHostStatus: metav1.ConditionFalse, + wantStorageReason: "NotEvaluated", + }, + { + name: "missing StorageClass", + objects: func(workspace *clusterv1alpha1.T4Workspace) []client.Object { + return []client.Object{testHost(), workspace} + }, + wantHostStatus: metav1.ConditionTrue, + wantStorageReason: controllers.ReasonStorageClassNotFound, + }, + { + name: "non-RWX StorageClass", + objects: func(workspace *clusterv1alpha1.T4Workspace) []client.Object { + class := rwxStorageClass() + class.Annotations = nil + return []client.Object{testHost(), class, workspace} + }, + wantHostStatus: metav1.ConditionTrue, + wantStorageReason: controllers.ReasonStorageClassNotRWX, + }, + { + name: "PVC class drift", + objects: func(workspace *clusterv1alpha1.T4Workspace) []client.Object { + pvc := ownedWorkspacePVC(workspace) + pvc.Spec.StorageClassName = ptr("old-rwx") + return []client.Object{testHost(), rwxStorageClass(), workspace, pvc} + }, + wantHostStatus: metav1.ConditionTrue, + wantStorageReason: controllers.ReasonStorageClassMismatch, + }, + { + name: "PVC ownership conflict", + objects: func(workspace *clusterv1alpha1.T4Workspace) []client.Object { + pvc := ownedWorkspacePVC(workspace) + pvc.OwnerReferences = []metav1.OwnerReference{{APIVersion: "example.test/v1", Kind: "Foreign", Name: "foreign", UID: "foreign-uid", Controller: ptr(true)}} + return []client.Object{testHost(), rwxStorageClass(), workspace, pvc} + }, + wantHostStatus: metav1.ConditionTrue, + wantStorageReason: "PVCOwnershipConflict", + }, + } { + t.Run(test.name, func(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + workspace := testWorkspace(clusterv1alpha1.RetentionPolicyDelete) + workspace.UID = "workspace-uid" + workspace.Generation = 7 + workspace.Status.ObservedGeneration = 6 + workspace.Status.PVCName = controllers.WorkspacePVCName(workspace) + workspace.Status.PVCPhase = corev1.ClaimBound + workspace.Status.Capacity = apiresource.MustParse("10Gi") + workspace.Status.Phase = clusterv1alpha1.InfrastructureReady + workspace.Status.Conditions = []metav1.Condition{ + {Type: "HostReady", Status: metav1.ConditionTrue, Reason: "Stale", ObservedGeneration: 6}, + {Type: "StorageReady", Status: metav1.ConditionTrue, Reason: "Stale", ObservedGeneration: 6}, + {Type: "Ready", Status: metav1.ConditionTrue, Reason: "Stale", ObservedGeneration: 6}, + } + c := fake.NewClientBuilder().WithScheme(scheme).WithStatusSubresource(&clusterv1alpha1.T4Workspace{}, &corev1.PersistentVolumeClaim{}).WithObjects(test.objects(workspace)...).Build() + r := &controllers.WorkspaceReconciler{Client: c, Scheme: scheme} + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(workspace)}); err != nil { + t.Fatal(err) + } + var failed clusterv1alpha1.T4Workspace + if err := c.Get(ctx, client.ObjectKeyFromObject(workspace), &failed); err != nil { + t.Fatal(err) + } + if failed.Status.PVCName != "" || failed.Status.PVCPhase != "" || !failed.Status.Capacity.IsZero() { + t.Fatalf("failed Workspace retained PVC authority: %#v", failed.Status) + } + hostReady := findCondition(failed.Status.Conditions, "HostReady") + storageReady := findCondition(failed.Status.Conditions, "StorageReady") + ready := findCondition(failed.Status.Conditions, "Ready") + if hostReady == nil || hostReady.Status != test.wantHostStatus || hostReady.ObservedGeneration != failed.Generation || + storageReady == nil || storageReady.Status == metav1.ConditionTrue || storageReady.Reason != test.wantStorageReason || storageReady.ObservedGeneration != failed.Generation || + ready == nil || ready.Status != metav1.ConditionFalse || ready.ObservedGeneration != failed.Generation { + t.Fatalf("failure conditions are stale: HostReady=%#v StorageReady=%#v Ready=%#v", hostReady, storageReady, ready) + } + }) + } +} + +func TestSessionRejectsWorkspacePVCWithoutExactIdentityAndOwnership(t *testing.T) { + for _, test := range []struct { + name string + mutatePVC func(*clusterv1alpha1.T4Workspace, *corev1.PersistentVolumeClaim) + }{ + {name: "foreign deterministic PVC", mutatePVC: func(_ *clusterv1alpha1.T4Workspace, pvc *corev1.PersistentVolumeClaim) { + pvc.OwnerReferences = []metav1.OwnerReference{{APIVersion: "example.test/v1", Kind: "Foreign", Name: "foreign", UID: "foreign-uid", Controller: ptr(true)}} + }}, + {name: "tampered Workspace status name", mutatePVC: func(workspace *clusterv1alpha1.T4Workspace, pvc *corev1.PersistentVolumeClaim) { + workspace.Status.PVCName = "foreign-data" + pvc.Name = workspace.Status.PVCName + }}, + } { + t.Run(test.name, func(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + workspace := testWorkspace(clusterv1alpha1.RetentionPolicyDelete) + workspace.UID = "workspace-uid" + workspace.Status.PVCName = controllers.WorkspacePVCName(workspace) + pvc := ownedWorkspacePVC(workspace) + test.mutatePVC(workspace, pvc) + session := testSession() + session.UID = "session-uid" + pod, service := ownedSessionResources(session) + c := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&clusterv1alpha1.T4Session{}, &corev1.PersistentVolumeClaim{}, &corev1.Pod{}). + WithObjects(testHost(), workspace, pvc, session, pod, service).Build() + r := configuredSessionReconciler(c, scheme) + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(session)}); err != nil { + t.Fatal(err) + } + assertObjectCounts(t, c, 0, 0) + var failed clusterv1alpha1.T4Session + if err := c.Get(ctx, client.ObjectKeyFromObject(session), &failed); err != nil { + t.Fatal(err) + } + condition := findCondition(failed.Status.Conditions, "WorkspaceReady") + if failed.Status.PodName != "" || failed.Status.ServiceName != "" || condition == nil || condition.Status != metav1.ConditionFalse || condition.ObservedGeneration != failed.Generation { + t.Fatalf("untrusted Workspace PVC retained session authority: status=%#v WorkspaceReady=%#v", failed.Status, condition) + } + }) + } +} + +func TestSessionFailureAndPendingRefreshEveryCondition(t *testing.T) { + t.Run("failure", func(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + session := testSession() + session.Generation = 8 + session.Status.ObservedGeneration = 7 + session.Status.Conditions = staleSessionConditions(7) + c := fake.NewClientBuilder().WithScheme(scheme).WithStatusSubresource(&clusterv1alpha1.T4Session{}).WithObjects(session).Build() + r := configuredSessionReconciler(c, scheme) + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(session)}); err != nil { + t.Fatal(err) + } + var failed clusterv1alpha1.T4Session + if err := c.Get(ctx, client.ObjectKeyFromObject(session), &failed); err != nil { + t.Fatal(err) + } + assertCurrentSessionConditions(t, &failed, map[string]metav1.ConditionStatus{ + "HostReady": metav1.ConditionFalse, "WorkspaceReady": metav1.ConditionUnknown, "RuntimeConfigured": metav1.ConditionUnknown, "Available": metav1.ConditionFalse, + }) + }) + + t.Run("pending", func(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + workspace := testWorkspace(clusterv1alpha1.RetentionPolicyDelete) + workspace.UID = "workspace-uid" + workspace.Status.PVCName = controllers.WorkspacePVCName(workspace) + pvc := ownedWorkspacePVC(workspace) + session := testSession() + session.UID = "session-uid" + session.Generation = 8 + session.Status.ObservedGeneration = 7 + session.Status.Conditions = staleSessionConditions(7) + _, service := ownedSessionResources(session) + service.Spec.Type = corev1.ServiceTypeNodePort + service.Spec.Ports = []corev1.ServicePort{{Name: "host", Port: 8787, NodePort: 32080}} + c := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&clusterv1alpha1.T4Session{}, &corev1.PersistentVolumeClaim{}, &corev1.Pod{}). + WithObjects(testHost(), workspace, pvc, session, service).Build() + r := configuredSessionReconciler(c, scheme) + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(session)}); err != nil { + t.Fatal(err) + } + var pending clusterv1alpha1.T4Session + if err := c.Get(ctx, client.ObjectKeyFromObject(session), &pending); err != nil { + t.Fatal(err) + } + assertCurrentSessionConditions(t, &pending, map[string]metav1.ConditionStatus{ + "HostReady": metav1.ConditionTrue, "WorkspaceReady": metav1.ConditionTrue, "RuntimeConfigured": metav1.ConditionTrue, "Available": metav1.ConditionFalse, + }) + }) +} + +func TestSessionResourcesWithAnyForeignOwnerFailClosed(t *testing.T) { + for _, path := range []string{"normal", "dependency-cleanup"} { + for _, kind := range []string{"Pod", "Service"} { + t.Run(path+"/"+kind, func(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + session := testSession() + session.UID = "session-uid" + pod, service := ownedSessionResources(session) + foreignOwner := metav1.OwnerReference{APIVersion: "example.test/v1", Kind: "Foreign", Name: "foreign", UID: "foreign-uid"} + if kind == "Pod" { + pod.OwnerReferences = append(pod.OwnerReferences, foreignOwner) + } else { + service.OwnerReferences = append(service.OwnerReferences, foreignOwner) + } + objects := []client.Object{session, pod, service} + if path == "normal" { + workspace := testWorkspace(clusterv1alpha1.RetentionPolicyDelete) + workspace.UID = "workspace-uid" + workspace.Status.PVCName = controllers.WorkspacePVCName(workspace) + objects = append(objects, testHost(), workspace, ownedWorkspacePVC(workspace)) + } + c := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&clusterv1alpha1.T4Session{}, &corev1.PersistentVolumeClaim{}, &corev1.Pod{}). + WithObjects(objects...).Build() + r := configuredSessionReconciler(c, scheme) + beforePod := pod.DeepCopy() + beforeService := service.DeepCopy() + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(session)}); err != nil { + t.Fatal(err) + } + if kind == "Pod" { + var got corev1.Pod + if err := c.Get(ctx, client.ObjectKeyFromObject(pod), &got); err != nil { + t.Fatalf("foreign-owned Pod was deleted: %v", err) + } + if !reflect.DeepEqual(got.ObjectMeta, beforePod.ObjectMeta) || !reflect.DeepEqual(got.Spec, beforePod.Spec) { + t.Fatalf("Pod with foreign OwnerReference was mutated: %#v", got) + } + } else { + var got corev1.Service + if err := c.Get(ctx, client.ObjectKeyFromObject(service), &got); err != nil { + t.Fatalf("foreign-owned Service was deleted: %v", err) + } + if !reflect.DeepEqual(got.ObjectMeta, beforeService.ObjectMeta) || !reflect.DeepEqual(got.Spec, beforeService.Spec) { + t.Fatalf("Service with foreign OwnerReference was mutated: %#v", got) + } + } + var sibling client.Object + if kind == "Pod" { + sibling = &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: controllers.SessionServiceName(session), Namespace: session.Namespace}} + } else { + sibling = &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: controllers.SessionPodName(session), Namespace: session.Namespace}} + } + if err := c.Get(ctx, client.ObjectKeyFromObject(sibling), sibling); !apierrors.IsNotFound(err) { + t.Fatalf("exclusively owned sibling remained after ownership conflict: %v", err) + } + var failed clusterv1alpha1.T4Session + if err := c.Get(ctx, client.ObjectKeyFromObject(session), &failed); err != nil { + t.Fatal(err) + } + available := findCondition(failed.Status.Conditions, "Available") + if available == nil || available.Status != metav1.ConditionFalse || !strings.Contains(available.Reason, "OwnershipConflict") { + t.Fatalf("foreign owner did not produce stable ownership conflict: %#v", available) + } + if path == "normal" { + assertCurrentSessionConditions(t, &failed, map[string]metav1.ConditionStatus{ + "HostReady": metav1.ConditionTrue, "WorkspaceReady": metav1.ConditionTrue, "RuntimeConfigured": metav1.ConditionTrue, "Available": metav1.ConditionFalse, + }) + } + }) + } + } +} + +func TestSessionCreateAlreadyExistsRefetchesForeignWinner(t *testing.T) { + for _, kind := range []string{"Pod", "Service"} { + t.Run(kind, func(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + workspace := testWorkspace(clusterv1alpha1.RetentionPolicyDelete) + workspace.UID = "workspace-uid" + workspace.Status.PVCName = controllers.WorkspacePVCName(workspace) + session := testSession() + session.UID = "session-uid" + pod, service := ownedSessionResources(session) + objects := []client.Object{testHost(), workspace, ownedWorkspacePVC(workspace), session} + if kind == "Pod" { + objects = append(objects, service) + } else { + objects = append(objects, pod) + } + base := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&clusterv1alpha1.T4Session{}, &corev1.PersistentVolumeClaim{}, &corev1.Pod{}). + WithObjects(objects...).Build() + c := &createAlreadyExistsClient{Client: base, raceKind: kind, hideWinnerFromCache: true} + r := configuredSessionReconciler(c, scheme) + r.APIReader = base + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(session)}); err != nil { + t.Fatal(err) + } + if c.winner == nil { + t.Fatalf("%s create race was not exercised", kind) + } + if kind == "Pod" { + var got corev1.Pod + if err := base.Get(ctx, client.ObjectKeyFromObject(c.winner), &got); err != nil { + t.Fatalf("foreign Pod winner was deleted: %v", err) + } + want := c.winner.(*corev1.Pod) + if !reflect.DeepEqual(got.ObjectMeta, want.ObjectMeta) || !reflect.DeepEqual(got.Spec, want.Spec) { + t.Fatalf("foreign Pod winner was mutated: %#v", got) + } + assertObjectCounts(t, base, 1, 0) + } else { + var got corev1.Service + if err := base.Get(ctx, client.ObjectKeyFromObject(c.winner), &got); err != nil { + t.Fatalf("foreign Service winner was deleted: %v", err) + } + want := c.winner.(*corev1.Service) + if !reflect.DeepEqual(got.ObjectMeta, want.ObjectMeta) || !reflect.DeepEqual(got.Spec, want.Spec) { + t.Fatalf("foreign Service winner was mutated: %#v", got) + } + assertObjectCounts(t, base, 0, 1) + } + var failed clusterv1alpha1.T4Session + if err := base.Get(ctx, client.ObjectKeyFromObject(session), &failed); err != nil { + t.Fatal(err) + } + assertCurrentSessionConditions(t, &failed, map[string]metav1.ConditionStatus{ + "HostReady": metav1.ConditionTrue, "WorkspaceReady": metav1.ConditionTrue, "RuntimeConfigured": metav1.ConditionTrue, "Available": metav1.ConditionFalse, + }) + available := findCondition(failed.Status.Conditions, "Available") + if available.Reason != kind+"OwnershipConflict" { + t.Fatalf("Available reason = %q, want %sOwnershipConflict", available.Reason, kind) + } + }) + } +} + +func TestWorkspaceTerminalPVCFailuresRevokePublishedAuthority(t *testing.T) { + for _, test := range []struct { + name string + mutate func(*corev1.PersistentVolumeClaim) + reason string + }{ + {name: "Bound without RWX", reason: "PVCNotRWX", mutate: func(pvc *corev1.PersistentVolumeClaim) { pvc.Spec.AccessModes = []corev1.PersistentVolumeAccessMode{corev1.ReadWriteOnce} }}, + {name: "Lost", reason: "PVCLost", mutate: func(pvc *corev1.PersistentVolumeClaim) { pvc.Status.Phase = corev1.ClaimLost }}, + } { + t.Run(test.name, func(t *testing.T) { + ctx := context.Background() + scheme := testScheme(t) + workspace := testWorkspace(clusterv1alpha1.RetentionPolicyDelete) + workspace.UID = "workspace-uid" + pvc := ownedWorkspacePVC(workspace) + test.mutate(pvc) + c := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&clusterv1alpha1.T4Workspace{}, &corev1.PersistentVolumeClaim{}). + WithObjects(testHost(), rwxStorageClass(), workspace, pvc).Build() + r := &controllers.WorkspaceReconciler{Client: c, Scheme: scheme} + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(workspace)}); err != nil { + t.Fatal(err) + } + var failed clusterv1alpha1.T4Workspace + if err := c.Get(ctx, client.ObjectKeyFromObject(workspace), &failed); err != nil { + t.Fatal(err) + } + storageReady := findCondition(failed.Status.Conditions, "StorageReady") + ready := findCondition(failed.Status.Conditions, "Ready") + if failed.Status.PVCName != "" || failed.Status.PVCPhase != "" || !failed.Status.Capacity.IsZero() || + storageReady == nil || storageReady.Status != metav1.ConditionFalse || storageReady.Reason != test.reason || storageReady.ObservedGeneration != failed.Generation || + ready == nil || ready.Status != metav1.ConditionFalse || ready.Reason != test.reason || ready.ObservedGeneration != failed.Generation { + t.Fatalf("terminal PVC failure retained authority: status=%#v StorageReady=%#v Ready=%#v", failed.Status, storageReady, ready) + } + }) + } +} + + +func ownedWorkspacePVC(workspace *clusterv1alpha1.T4Workspace) *corev1.PersistentVolumeClaim { + return &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{ + Name: controllers.WorkspacePVCName(workspace), Namespace: workspace.Namespace, + Annotations: map[string]string{clusterv1alpha1.WorkspaceUIDAnnotation: string(workspace.UID)}, + OwnerReferences: []metav1.OwnerReference{{APIVersion: clusterv1alpha1.GroupVersion.String(), Kind: "T4Workspace", Name: workspace.Name, UID: workspace.UID, Controller: ptr(true)}}, + }, + Spec: corev1.PersistentVolumeClaimSpec{StorageClassName: ptr("portable-rwx"), AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}}, + Status: corev1.PersistentVolumeClaimStatus{Phase: corev1.ClaimBound}, + } +} + +func ownedSessionResources(session *clusterv1alpha1.T4Session) (*corev1.Pod, *corev1.Service) { + owner := metav1.OwnerReference{APIVersion: clusterv1alpha1.GroupVersion.String(), Kind: "T4Session", Name: session.Name, UID: session.UID, Controller: ptr(true)} + pod := &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: controllers.SessionPodName(session), Namespace: session.Namespace, OwnerReferences: []metav1.OwnerReference{owner}}} + service := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: controllers.SessionServiceName(session), Namespace: session.Namespace, OwnerReferences: []metav1.OwnerReference{owner}}, Spec: corev1.ServiceSpec{Type: corev1.ServiceTypeClusterIP}} + return pod, service +} + +func staleSessionConditions(generation int64) []metav1.Condition { + return []metav1.Condition{ + {Type: "HostReady", Status: metav1.ConditionTrue, Reason: "Stale", ObservedGeneration: generation}, + {Type: "WorkspaceReady", Status: metav1.ConditionTrue, Reason: "Stale", ObservedGeneration: generation}, + {Type: "RuntimeConfigured", Status: metav1.ConditionTrue, Reason: "Stale", ObservedGeneration: generation}, + {Type: "Available", Status: metav1.ConditionTrue, Reason: "Stale", ObservedGeneration: generation}, + } +} + +func assertCurrentSessionConditions(t *testing.T, session *clusterv1alpha1.T4Session, want map[string]metav1.ConditionStatus) { + t.Helper() + for conditionType, status := range want { + condition := findCondition(session.Status.Conditions, conditionType) + if condition == nil || condition.Status != status || condition.ObservedGeneration != session.Generation { + t.Fatalf("%s = %#v, want %s at generation %d", conditionType, condition, status, session.Generation) + } + } +} + +func configuredSessionReconciler(c client.Client, scheme *runtime.Scheme) *controllers.SessionReconciler { + for _, object := range []client.Object{ + &corev1.ConfigMap{ObjectMeta: metav1.ObjectMeta{Name: "omp-runtime-config", Namespace: "team"}, Data: map[string]string{ + "provider-models": testOMPModels, "agent-settings": testOMPSettings, "other-models": otherTestOMPModels, "other-settings": otherTestOMPSettings, + }}, + &corev1.ConfigMap{ObjectMeta: metav1.ObjectMeta{Name: "other-omp-config", Namespace: "team"}, Data: map[string]string{ + "provider-models": testOMPModels, "agent-settings": testOMPSettings, "other-models": otherTestOMPModels, "other-settings": otherTestOMPSettings, + }}, + rwxStorageClass(), + } { + if err := c.Create(context.Background(), object); err != nil && !apierrors.IsAlreadyExists(err) { + panic(err) + } + } + return &controllers.SessionReconciler{ + Client: c, + APIReader: c, + Scheme: scheme, + RuntimeImage: testRuntimeImage, + OMPConfig: controllers.SessionOMPConfig{ + ConfigMapName: "omp-runtime-config", + ModelsKey: "provider-models", + SettingsKey: "agent-settings", + }, + } +} + +func testScheme(t *testing.T) *runtime.Scheme { + t.Helper() + scheme := runtime.NewScheme() + for _, add := range []func(*runtime.Scheme) error{corev1.AddToScheme, storagev1.AddToScheme, clusterv1alpha1.AddToScheme} { + if err := add(scheme); err != nil { + t.Fatal(err) + } + } + return scheme +} + +func testHost() *clusterv1alpha1.T4ClusterHost { + return &clusterv1alpha1.T4ClusterHost{ + ObjectMeta: metav1.ObjectMeta{Name: "host-a", Namespace: "team", UID: "host-uid"}, + Spec: clusterv1alpha1.T4ClusterHostSpec{StorageClassName: "portable-rwx", RuntimeProfiles: []string{"default"}}, + } +} + func rwxStorageClass() *storagev1.StorageClass { return &storagev1.StorageClass{ ObjectMeta: metav1.ObjectMeta{Name: "portable-rwx", Annotations: map[string]string{clusterv1alpha1.RWXStorageClassAnnotation: string(corev1.ReadWriteMany)}}, @@ -1326,4 +2426,112 @@ func hasReadOnlyMount(mounts []corev1.VolumeMount, name, path string) bool { return false } +type pvcOverrideReader struct { + client.Reader + pvc *corev1.PersistentVolumeClaim +} + +func (r *pvcOverrideReader) Get(ctx context.Context, key client.ObjectKey, object client.Object, options ...client.GetOption) error { + if pvc, ok := object.(*corev1.PersistentVolumeClaim); ok && key == client.ObjectKeyFromObject(r.pvc) { + r.pvc.DeepCopyInto(pvc) + return nil + } + return r.Reader.Get(ctx, key, object, options...) +} + +type createAlreadyExistsClient struct { + client.Client + raceKind string + winner client.Object + hideWinnerFromCache bool +} + +type replaceBeforeDeleteClient struct { + client.Client + raceKind string + replacementUID types.UID + raced bool + observedUID types.UID + observedResourceVersion string +} + +func (c *replaceBeforeDeleteClient) Delete(ctx context.Context, object client.Object, options ...client.DeleteOption) error { + if c.raced { + return c.Client.Delete(ctx, object, options...) + } + if _, isPod := object.(*corev1.Pod); !isPod || c.raceKind != "Pod" { + return c.Client.Delete(ctx, object, options...) + } + deleteOptions := (&client.DeleteOptions{}).ApplyOptions(options) + if deleteOptions.Preconditions != nil { + if deleteOptions.Preconditions.UID != nil { + c.observedUID = *deleteOptions.Preconditions.UID + } + if deleteOptions.Preconditions.ResourceVersion != nil { + c.observedResourceVersion = *deleteOptions.Preconditions.ResourceVersion + } + } + c.raced = true + if err := c.Client.Delete(ctx, object); err != nil { + return err + } + replacement := object.DeepCopyObject().(client.Object) + replacement.SetUID(c.replacementUID) + replacement.SetResourceVersion("") + replacement.SetDeletionTimestamp(nil) + replacement.SetFinalizers(nil) + replacement.SetOwnerReferences(nil) + if err := c.Client.Create(ctx, replacement); err != nil { + return err + } + return apierrors.NewConflict(schema.GroupResource{Resource: "pods"}, object.GetName(), errors.New("delete preconditions no longer match replacement")) +} + +func (c *createAlreadyExistsClient) Get(ctx context.Context, key client.ObjectKey, object client.Object, options ...client.GetOption) error { + if c.hideWinnerFromCache && c.winner != nil && key == client.ObjectKeyFromObject(c.winner) { + switch object.(type) { + case *corev1.Pod: + if c.raceKind == "Pod" { + return apierrors.NewNotFound(schema.GroupResource{Resource: "pods"}, key.Name) + } + case *corev1.Service: + if c.raceKind == "Service" { + return apierrors.NewNotFound(schema.GroupResource{Resource: "services"}, key.Name) + } + case *corev1.PersistentVolumeClaim: + if c.raceKind == "PVC" { + return apierrors.NewNotFound(schema.GroupResource{Resource: "persistentvolumeclaims"}, key.Name) + } + } + } + return c.Client.Get(ctx, key, object, options...) +} + +func (c *createAlreadyExistsClient) Create(ctx context.Context, object client.Object, options ...client.CreateOption) error { + var winner client.Object + switch object := object.(type) { + case *corev1.Pod: + if c.raceKind == "Pod" { + winner = object.DeepCopy() + } + case *corev1.Service: + if c.raceKind == "Service" { + winner = object.DeepCopy() + } + case *corev1.PersistentVolumeClaim: + if c.raceKind == "PVC" { + winner = object.DeepCopy() + } + } + if winner == nil || c.winner != nil { + return c.Client.Create(ctx, object, options...) + } + winner.SetOwnerReferences(nil) + if err := c.Client.Create(ctx, winner, options...); err != nil { + return err + } + c.winner = winner.DeepCopyObject().(client.Object) + return apierrors.NewAlreadyExists(schema.GroupResource{Resource: strings.ToLower(c.raceKind) + "s"}, object.GetName()) +} + func ptr[T any](value T) *T { return &value } diff --git a/packages/cluster-operator/controllers/session_controller.go b/packages/cluster-operator/controllers/session_controller.go index 681f4f0c..49caac41 100644 --- a/packages/cluster-operator/controllers/session_controller.go +++ b/packages/cluster-operator/controllers/session_controller.go @@ -4,6 +4,7 @@ import ( "context" "crypto/sha256" "encoding/json" + "errors" "fmt" "net/url" "reflect" @@ -25,6 +26,7 @@ import ( ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" + "sigs.k8s.io/controller-runtime/pkg/handler" clusterv1alpha1 "github.com/LycaonLLC/t4-code/packages/cluster-operator/api/v1alpha1" ) @@ -36,6 +38,13 @@ const ( SessionReviewerTokenExpirationSeconds int64 = 3600 ) +const ( + sessionHostRefIndexField = "t4.session.spec.hostRef" + sessionWorkspaceRefIndexField = "t4.session.spec.workspaceRef" +) + +var errSessionResourceOwnershipConflict = errors.New("session resource ownership conflict") + var ( configMapKeyPattern = regexp.MustCompile(`^[-._A-Za-z0-9]+$`) runtimeImagePattern = regexp.MustCompile(`^(?:(?:[A-Za-z0-9](?:[A-Za-z0-9.-]*[A-Za-z0-9])?|\[[A-Fa-f0-9:]+\])(?::[0-9]+)?/)?[a-z0-9]+(?:(?:[._]|__|-+)[a-z0-9]+)*(?:/[a-z0-9]+(?:(?:[._]|__|-+)[a-z0-9]+)*)*@sha256:[a-f0-9]{64}$`) @@ -271,6 +280,10 @@ func (r *SessionReconciler) Reconcile(ctx context.Context, request ctrl.Request) var session clusterv1alpha1.T4Session found := false defer func() { + if errors.Is(err, errSessionResourceOwnershipConflict) { + result = ctrl.Result{RequeueAfter: 30 * time.Second} + err = nil + } observeReconcile(metricKindSession, request.NamespacedName, session.Status.Conditions, conditionObjectPresent(&session, found, err), err) }() if err := r.Get(ctx, request.NamespacedName, &session); err != nil { @@ -289,19 +302,22 @@ func (r *SessionReconciler) Reconcile(ctx context.Context, request ctrl.Request) if err := r.deleteOwnedSessionResources(ctx, &session); err != nil { return ctrl.Result{}, err } - return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, "RuntimeConfigured", reason, message) + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, false, false, "RuntimeConfigured", reason, message) } if reason, message := r.OMPConfig.validationFailure(); reason != "" { if err := r.deleteOwnedSessionResources(ctx, &session); err != nil { return ctrl.Result{}, err } - return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, "RuntimeConfigured", reason, message) + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, false, false, "RuntimeConfigured", reason, message) } var host clusterv1alpha1.T4ClusterHost if err := r.Get(ctx, types.NamespacedName{Namespace: session.Namespace, Name: session.Spec.HostRef}, &host); err != nil { if apierrors.IsNotFound(err) { - return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, "HostReady", "HostNotFound", "referenced T4ClusterHost does not exist") + if err := r.deleteOwnedSessionResources(ctx, &session); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, false, false, "HostReady", "HostNotFound", "referenced T4ClusterHost does not exist") } return ctrl.Result{}, err } @@ -309,7 +325,7 @@ func (r *SessionReconciler) Reconcile(ctx context.Context, request ctrl.Request) if err := r.deleteOwnedSessionResources(ctx, &session); err != nil { return ctrl.Result{}, err } - return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, "RuntimeConfigured", "RuntimeProfileNotAllowed", "runtime profile is not allowed by the referenced T4ClusterHost") + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, true, false, "RuntimeConfigured", "RuntimeProfileNotAllowed", "runtime profile is not allowed by the referenced T4ClusterHost") } var storageClass storagev1.StorageClass if err := r.Get(ctx, types.NamespacedName{Name: host.Spec.StorageClassName}, &storageClass); err != nil { @@ -317,7 +333,7 @@ func (r *SessionReconciler) Reconcile(ctx context.Context, request ctrl.Request) if err := r.deleteOwnedSessionResources(ctx, &session); err != nil { return ctrl.Result{}, err } - return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, "WorkspaceReady", ReasonStorageClassNotFound, fmt.Sprintf("StorageClass %q selected by the referenced T4ClusterHost does not exist", host.Spec.StorageClassName)) + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, true, false, "WorkspaceReady", ReasonStorageClassNotFound, fmt.Sprintf("StorageClass %q selected by the referenced T4ClusterHost does not exist", host.Spec.StorageClassName)) } return ctrl.Result{}, err } @@ -325,30 +341,63 @@ func (r *SessionReconciler) Reconcile(ctx context.Context, request ctrl.Request) if err := r.deleteOwnedSessionResources(ctx, &session); err != nil { return ctrl.Result{}, err } - return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, "WorkspaceReady", ReasonStorageClassNotRWX, fmt.Sprintf("StorageClass %q selected by the referenced T4ClusterHost is not administrator-declared ReadWriteMany", host.Spec.StorageClassName)) + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, true, false, "WorkspaceReady", ReasonStorageClassNotRWX, fmt.Sprintf("StorageClass %q selected by the referenced T4ClusterHost is not administrator-declared ReadWriteMany", host.Spec.StorageClassName)) } var workspace clusterv1alpha1.T4Workspace if err := r.Get(ctx, types.NamespacedName{Namespace: session.Namespace, Name: session.Spec.WorkspaceRef}, &workspace); err != nil { if apierrors.IsNotFound(err) { - return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, "WorkspaceReady", "WorkspaceNotFound", "referenced T4Workspace does not exist") + if err := r.deleteOwnedSessionResources(ctx, &session); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, true, false, "WorkspaceReady", "WorkspaceNotFound", "referenced T4Workspace does not exist") } return ctrl.Result{}, err } if workspace.Spec.HostRef != session.Spec.HostRef { - return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, "WorkspaceReady", "HostMismatch", "session and workspace must reference the same T4ClusterHost") + if err := r.deleteOwnedSessionResources(ctx, &session); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, true, false, "WorkspaceReady", "HostMismatch", "session and workspace must reference the same T4ClusterHost") } if workspace.Status.PVCName == "" { - return ctrl.Result{RequeueAfter: 5 * time.Second}, r.updateSessionFailure(ctx, &session, "WorkspaceReady", "PVCNotDeclared", "workspace controller has not declared a PVC") + if err := r.deleteOwnedSessionResources(ctx, &session); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: 5 * time.Second}, r.updateSessionFailure(ctx, &session, true, false, "WorkspaceReady", "PVCNotDeclared", "workspace controller has not declared a PVC") + } + if workspace.UID != "" && workspace.Status.PVCName != WorkspacePVCName(&workspace) { + if err := r.deleteOwnedSessionResources(ctx, &session); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, true, false, "WorkspaceReady", "PVCIdentityMismatch", "workspace status does not reference its deterministic PVC") } var pvc corev1.PersistentVolumeClaim if err := r.Get(ctx, types.NamespacedName{Namespace: session.Namespace, Name: workspace.Status.PVCName}, &pvc); err != nil { if apierrors.IsNotFound(err) { - return ctrl.Result{RequeueAfter: 5 * time.Second}, r.updateSessionFailure(ctx, &session, "WorkspaceReady", "PVCNotFound", "workspace PVC does not exist") + if err := r.deleteOwnedSessionResources(ctx, &session); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: 5 * time.Second}, r.updateSessionFailure(ctx, &session, true, false, "WorkspaceReady", "PVCNotFound", "workspace PVC does not exist") } return ctrl.Result{}, err } + if workspace.UID != "" && !workspaceOwnsPVC(&workspace, &pvc) { + if err := r.deleteOwnedSessionResources(ctx, &session); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, true, false, "WorkspaceReady", "PVCOwnershipConflict", "workspace PVC identity or ownership is not authoritative") + } + if pvcStorageClassName(&pvc) != host.Spec.StorageClassName { + if err := r.deleteOwnedSessionResources(ctx, &session); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, true, false, "WorkspaceReady", ReasonStorageClassMismatch, fmt.Sprintf("workspace PVC uses StorageClass %q instead of host-selected %q", pvcStorageClassName(&pvc), host.Spec.StorageClassName)) + } if pvc.Status.Phase != corev1.ClaimBound || !pvcHasRWX(&pvc) { - return ctrl.Result{RequeueAfter: 5 * time.Second}, r.updateSessionFailure(ctx, &session, "WorkspaceReady", "PVCNotBoundRWX", "workspace PVC must be Bound and ReadWriteMany before a session starts") + if err := r.deleteOwnedSessionResources(ctx, &session); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: 5 * time.Second}, r.updateSessionFailure(ctx, &session, true, false, "WorkspaceReady", "PVCNotBoundRWX", "workspace PVC must be Bound and ReadWriteMany before a session starts") } runtimeVersions, reason, message, err := r.loadOMPResourceVersions(ctx, session.Namespace) if err != nil { @@ -358,9 +407,20 @@ func (r *SessionReconciler) Reconcile(ctx context.Context, request ctrl.Request) if err := r.deleteOwnedSessionResources(ctx, &session); err != nil { return ctrl.Result{}, err } - return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, "RuntimeConfigured", reason, message) + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, true, true, "RuntimeConfigured", reason, message) + } + reason, message, err = r.authoritativePVCValidation(ctx, &workspace, &pvc, host.Spec.StorageClassName) + if err != nil { + return ctrl.Result{}, err + } + if reason != "" { + if err := r.deleteOwnedSessionResources(ctx, &session); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, true, false, "WorkspaceReady", reason, message) } + serviceName := SessionServiceName(&session) podName := SessionPodName(&session) labels := map[string]string{ @@ -380,17 +440,31 @@ func (r *SessionReconciler) Reconcile(ctx context.Context, request ctrl.Request) return ctrl.Result{}, err } var service corev1.Service - if err := r.Get(ctx, types.NamespacedName{Namespace: session.Namespace, Name: serviceName}, &service); apierrors.IsNotFound(err) { + serviceKey := types.NamespacedName{Namespace: session.Namespace, Name: serviceName} + if err := r.Get(ctx, serviceKey, &service); apierrors.IsNotFound(err) { service = desiredService - if err := r.Create(ctx, &service); err != nil && !apierrors.IsAlreadyExists(err) { - return ctrl.Result{}, err + if err := r.Create(ctx, &service); err != nil { + if !apierrors.IsAlreadyExists(err) { + return ctrl.Result{}, err + } + reader := r.APIReader + if reader == nil { + reader = r.Client + } + if err := reader.Get(ctx, serviceKey, &service); err != nil { + return ctrl.Result{}, err + } } } else if err != nil { return ctrl.Result{}, err - } else if !metav1.IsControlledBy(&service, &session) { - return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, "Available", "ServiceOwnershipConflict", "deterministic session Service is not controlled by this session") + } + if !sessionExclusivelyOwnsResource(&service, &session) { + if err := r.deleteOwnedSessionResourcesAfterVerifiedDependencies(ctx, &session, "ServiceOwnershipConflict", "deterministic session Service has an unexpected owner"); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: 30 * time.Second}, nil } else if !serviceExposureIsInternal(&service) { - if err := r.Delete(ctx, &service); err != nil && !apierrors.IsNotFound(err) { + if err := deleteWithPreconditions(ctx, r.Client, &service); err != nil && !apierrors.IsNotFound(err) { return ctrl.Result{}, err } if err := r.updateSessionPending(ctx, &session, podName, serviceName, "ServiceExposureChanged", "session Service is being recreated with ClusterIP-only exposure"); err != nil { @@ -414,15 +488,39 @@ func (r *SessionReconciler) Reconcile(ctx context.Context, request ctrl.Request) return ctrl.Result{}, err } var pod corev1.Pod - if err := r.Get(ctx, types.NamespacedName{Namespace: session.Namespace, Name: podName}, &pod); apierrors.IsNotFound(err) { - pod = desiredPod - if err := r.Create(ctx, &pod); err != nil && !apierrors.IsAlreadyExists(err) { + podKey := types.NamespacedName{Namespace: session.Namespace, Name: podName} + if err := r.Get(ctx, podKey, &pod); apierrors.IsNotFound(err) { + reason, message, err := r.authoritativePVCValidation(ctx, &workspace, &pvc, host.Spec.StorageClassName) + if err != nil { return ctrl.Result{}, err } + if reason != "" { + if err := r.deleteOwnedSessionResources(ctx, &session); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, true, false, "WorkspaceReady", reason, message) + } + pod = desiredPod + if err := r.Create(ctx, &pod); err != nil { + if !apierrors.IsAlreadyExists(err) { + return ctrl.Result{}, err + } + reader := r.APIReader + if reader == nil { + reader = r.Client + } + if err := reader.Get(ctx, podKey, &pod); err != nil { + return ctrl.Result{}, err + } + } } else if err != nil { return ctrl.Result{}, err - } else if !metav1.IsControlledBy(&pod, &session) { - return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, "Available", "PodOwnershipConflict", "deterministic session Pod is not controlled by this session") + } + if !sessionExclusivelyOwnsResource(&pod, &session) { + if err := r.deleteOwnedSessionResourcesAfterVerifiedDependencies(ctx, &session, "PodOwnershipConflict", "deterministic session Pod has an unexpected owner"); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: 30 * time.Second}, nil } else if !labelsContain(pod.Labels, desiredPod.Labels) { if pod.Labels == nil { pod.Labels = map[string]string{} @@ -438,7 +536,7 @@ func (r *SessionReconciler) Reconcile(ctx context.Context, request ctrl.Request) } return ctrl.Result{RequeueAfter: time.Second}, nil } else if pod.Annotations[clusterv1alpha1.SessionPodSpecHashAnnotation] != desiredPod.Annotations[clusterv1alpha1.SessionPodSpecHashAnnotation] { - if err := r.Delete(ctx, &pod); err != nil && !apierrors.IsNotFound(err) { + if err := deleteWithPreconditions(ctx, r.Client, &pod); err != nil && !apierrors.IsNotFound(err) { return ctrl.Result{}, err } if err := r.updateSessionPending(ctx, &session, podName, serviceName, "PodSpecChanged", "session Pod is being recreated to apply immutable desired state"); err != nil { @@ -447,6 +545,17 @@ func (r *SessionReconciler) Reconcile(ctx context.Context, request ctrl.Request) return ctrl.Result{RequeueAfter: time.Second}, nil } + reason, message, err = r.authoritativePVCValidation(ctx, &workspace, &pvc, host.Spec.StorageClassName) + if err != nil { + return ctrl.Result{}, err + } + if reason != "" { + if err := r.deleteOwnedSessionResources(ctx, &session); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateSessionFailure(ctx, &session, true, false, "WorkspaceReady", reason, message) + } + original := session.Status if session.Status.Conditions != nil { original.Conditions = append([]metav1.Condition(nil), session.Status.Conditions...) @@ -454,6 +563,7 @@ func (r *SessionReconciler) Reconcile(ctx context.Context, request ctrl.Request) session.Status.ObservedGeneration = session.Generation session.Status.PodName = podName session.Status.ServiceName = serviceName + meta.SetStatusCondition(&session.Status.Conditions, condition("HostReady", metav1.ConditionTrue, "HostResolved", "referenced T4ClusterHost is available", session.Generation)) meta.SetStatusCondition(&session.Status.Conditions, condition("WorkspaceReady", metav1.ConditionTrue, "PVCBoundRWX", "workspace PVC is Bound and ReadWriteMany", session.Generation)) meta.SetStatusCondition(&session.Status.Conditions, condition("RuntimeConfigured", metav1.ConditionTrue, "OMPReferencesReady", "administrator-owned OMP runtime references are configured", session.Generation)) if podReady(&pod) { @@ -477,6 +587,36 @@ func (r *SessionReconciler) Reconcile(ctx context.Context, request ctrl.Request) return ctrl.Result{RequeueAfter: 30 * time.Second}, nil } +func (r *SessionReconciler) authoritativePVCValidation(ctx context.Context, workspace *clusterv1alpha1.T4Workspace, cachedPVC *corev1.PersistentVolumeClaim, storageClassName string) (string, string, error) { + reader := r.APIReader + if reader == nil { + reader = r.Client + } + var authoritativePVC corev1.PersistentVolumeClaim + if err := reader.Get(ctx, client.ObjectKeyFromObject(cachedPVC), &authoritativePVC); err != nil { + if apierrors.IsNotFound(err) { + return "PVCAuthorityChanged", "workspace PVC does not exist in authoritative API state", nil + } + return "", "", err + } + if authoritativePVC.UID != cachedPVC.UID { + return "PVCAuthorityChanged", "authoritative workspace PVC UID differs from the validated cached PVC", nil + } + if workspace.UID != "" && !workspaceOwnsPVC(workspace, &authoritativePVC) { + return "PVCAuthorityChanged", "authoritative workspace PVC owner reference does not belong to the workspace", nil + } + if pvcStorageClassName(&authoritativePVC) != storageClassName { + return "PVCAuthorityChanged", fmt.Sprintf("authoritative workspace PVC uses StorageClass %q instead of host-selected %q", pvcStorageClassName(&authoritativePVC), storageClassName), nil + } + if !pvcHasRWX(&authoritativePVC) { + return "PVCAuthorityChanged", "authoritative workspace PVC does not request ReadWriteMany", nil + } + if authoritativePVC.Status.Phase != corev1.ClaimBound { + return "PVCAuthorityChanged", "authoritative workspace PVC is not Bound", nil + } + return "", "", nil +} + func (r *SessionReconciler) desiredPod(session *clusterv1alpha1.T4Session, pvcName, podName string, labels map[string]string, runtimeVersions ompResourceVersions) (corev1.Pod, error) { falseValue := false trueValue := true @@ -609,6 +749,9 @@ func (r *SessionReconciler) updateSessionPending(ctx context.Context, session *c original.Conditions = append([]metav1.Condition(nil), session.Status.Conditions...) } session.Status.ObservedGeneration = session.Generation + meta.SetStatusCondition(&session.Status.Conditions, condition("HostReady", metav1.ConditionTrue, "HostResolved", "referenced T4ClusterHost is available", session.Generation)) + meta.SetStatusCondition(&session.Status.Conditions, condition("WorkspaceReady", metav1.ConditionTrue, "PVCBoundRWX", "workspace PVC is Bound and ReadWriteMany", session.Generation)) + meta.SetStatusCondition(&session.Status.Conditions, condition("RuntimeConfigured", metav1.ConditionTrue, "OMPReferencesReady", "administrator-owned OMP runtime references are configured", session.Generation)) session.Status.PodName = podName session.Status.ServiceName = serviceName session.Status.Phase = clusterv1alpha1.InfrastructurePending @@ -640,35 +783,46 @@ func (r *SessionReconciler) reconcileDelete(ctx context.Context, session *cluste &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: SessionServiceName(session), Namespace: session.Namespace}}, } existing := make([]client.Object, 0, len(objects)) + var ownershipConflict client.Object + reader := r.APIReader + if reader == nil { + reader = r.Client + } for _, object := range objects { - err := r.Get(ctx, client.ObjectKeyFromObject(object), object) + err := reader.Get(ctx, client.ObjectKeyFromObject(object), object) if apierrors.IsNotFound(err) { continue } if err != nil { return ctrl.Result{}, err } - if !metav1.IsControlledBy(object, session) { - before := session.Status - if session.Status.Conditions != nil { - before.Conditions = append([]metav1.Condition(nil), session.Status.Conditions...) - } - meta.SetStatusCondition(&session.Status.Conditions, condition("Available", metav1.ConditionFalse, "CleanupOwnershipConflict", fmt.Sprintf("deterministic %T is not controlled by this session", object), session.Generation)) - if !reflect.DeepEqual(before, session.Status) { - if err := r.Status().Update(ctx, session); err != nil { - return ctrl.Result{}, err - } + if !sessionExclusivelyOwnsResource(object, session) { + if ownershipConflict == nil { + ownershipConflict = object } - return ctrl.Result{RequeueAfter: 30 * time.Second}, nil + continue } existing = append(existing, object) } for _, object := range existing { if object.GetDeletionTimestamp().IsZero() { - if err := r.Delete(ctx, object); err != nil && !apierrors.IsNotFound(err) { + if err := deleteWithPreconditions(ctx, r.Client, object); err != nil && !apierrors.IsNotFound(err) { + return ctrl.Result{}, err + } + } + } + if ownershipConflict != nil { + before := session.Status + if session.Status.Conditions != nil { + before.Conditions = append([]metav1.Condition(nil), session.Status.Conditions...) + } + meta.SetStatusCondition(&session.Status.Conditions, condition("Available", metav1.ConditionFalse, "CleanupOwnershipConflict", fmt.Sprintf("deterministic %T is not controlled by this session", ownershipConflict), session.Generation)) + if !reflect.DeepEqual(before, session.Status) { + if err := r.Status().Update(ctx, session); err != nil { return ctrl.Result{}, err } } + return ctrl.Result{RequeueAfter: 30 * time.Second}, nil } if len(existing) > 0 { return ctrl.Result{RequeueAfter: time.Second}, nil @@ -678,46 +832,122 @@ func (r *SessionReconciler) reconcileDelete(ctx context.Context, session *cluste } func (r *SessionReconciler) deleteOwnedSessionResources(ctx context.Context, session *clusterv1alpha1.T4Session) error { + reader := r.APIReader + if reader == nil { + reader = r.Client + } + return r.deleteOwnedSessionResourcesWithFailure(ctx, reader, session, true, false, false, "ResourceOwnershipConflict", "one or more deterministic session resources have an unexpected owner") +} + +func (r *SessionReconciler) deleteOwnedSessionResourcesAfterVerifiedDependencies(ctx context.Context, session *clusterv1alpha1.T4Session, reason, message string) error { + reader := r.APIReader + if reader == nil { + reader = r.Client + } + return r.deleteOwnedSessionResourcesWithFailure(ctx, reader, session, false, true, true, reason, message) +} + +func (r *SessionReconciler) deleteOwnedSessionResourcesWithFailure(ctx context.Context, reader client.Reader, session *clusterv1alpha1.T4Session, deleteWithoutConflict, hostReady, workspaceReady bool, reason, message string) error { objects := []client.Object{ &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: SessionPodName(session), Namespace: session.Namespace}}, &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: SessionServiceName(session), Namespace: session.Namespace}}, } + owned := make([]client.Object, 0, len(objects)) + ownershipConflict := false for _, object := range objects { - if err := r.Get(ctx, client.ObjectKeyFromObject(object), object); err != nil { + if err := reader.Get(ctx, client.ObjectKeyFromObject(object), object); err != nil { if err := client.IgnoreNotFound(err); err != nil { return err } continue } - if !metav1.IsControlledBy(object, session) { + if !sessionExclusivelyOwnsResource(object, session) { + ownershipConflict = true continue } - if err := r.Delete(ctx, object); err != nil && !apierrors.IsNotFound(err) { + owned = append(owned, object) + } + if ownershipConflict || deleteWithoutConflict { + for _, object := range owned { + if err := deleteWithPreconditions(ctx, r.Client, object); err != nil && !apierrors.IsNotFound(err) { + return err + } + } + } + if ownershipConflict { + if err := r.updateSessionFailure(ctx, session, hostReady, workspaceReady, "Available", reason, message); err != nil { return err } + return errSessionResourceOwnershipConflict } return nil } -func (r *SessionReconciler) updateSessionFailure(ctx context.Context, session *clusterv1alpha1.T4Session, conditionType, reason, message string) error { +func deleteWithPreconditions(ctx context.Context, writer client.Client, object client.Object) error { + preconditions := metav1.Preconditions{} + if uid := object.GetUID(); uid != "" { + preconditions.UID = &uid + } + if resourceVersion := object.GetResourceVersion(); resourceVersion != "" { + preconditions.ResourceVersion = &resourceVersion + } + options := &client.DeleteOptions{} + if preconditions.UID != nil || preconditions.ResourceVersion != nil { + options.Preconditions = &preconditions + } + return writer.Delete(ctx, object, options) +} + +func (r *SessionReconciler) updateSessionFailure(ctx context.Context, session *clusterv1alpha1.T4Session, hostReady, workspaceReady bool, conditionType, reason, message string) error { original := session.Status if session.Status.Conditions != nil { original.Conditions = append([]metav1.Condition(nil), session.Status.Conditions...) } + if conditionType == "HostReady" { + meta.SetStatusCondition(&session.Status.Conditions, condition("HostReady", metav1.ConditionFalse, reason, message, session.Generation)) + } else if hostReady { + meta.SetStatusCondition(&session.Status.Conditions, condition("HostReady", metav1.ConditionTrue, "HostResolved", "referenced T4ClusterHost is available", session.Generation)) + } else { + meta.SetStatusCondition(&session.Status.Conditions, condition("HostReady", metav1.ConditionUnknown, "NotEvaluated", "host dependency was not evaluated", session.Generation)) + } + if conditionType == "WorkspaceReady" { + meta.SetStatusCondition(&session.Status.Conditions, condition("WorkspaceReady", metav1.ConditionFalse, reason, message, session.Generation)) + } else if workspaceReady { + meta.SetStatusCondition(&session.Status.Conditions, condition("WorkspaceReady", metav1.ConditionTrue, "PVCBoundRWX", "workspace PVC is Bound and ReadWriteMany", session.Generation)) + } else { + meta.SetStatusCondition(&session.Status.Conditions, condition("WorkspaceReady", metav1.ConditionUnknown, "NotEvaluated", "workspace dependency was not evaluated", session.Generation)) + } + if conditionType == "RuntimeConfigured" { + meta.SetStatusCondition(&session.Status.Conditions, condition("RuntimeConfigured", metav1.ConditionFalse, reason, message, session.Generation)) + } else if workspaceReady { + meta.SetStatusCondition(&session.Status.Conditions, condition("RuntimeConfigured", metav1.ConditionTrue, "OMPReferencesReady", "administrator-owned OMP runtime references are configured", session.Generation)) + } else { + meta.SetStatusCondition(&session.Status.Conditions, condition("RuntimeConfigured", metav1.ConditionUnknown, "NotEvaluated", "runtime configuration was not evaluated", session.Generation)) + } session.Status.ObservedGeneration = session.Generation session.Status.PodName = "" session.Status.ServiceName = "" session.Status.Phase = clusterv1alpha1.InfrastructureFailed - meta.SetStatusCondition(&session.Status.Conditions, condition(conditionType, metav1.ConditionFalse, reason, message, session.Generation)) - if conditionType != "Available" { - meta.SetStatusCondition(&session.Status.Conditions, condition("Available", metav1.ConditionFalse, reason, message, session.Generation)) - } + meta.SetStatusCondition(&session.Status.Conditions, condition("Available", metav1.ConditionFalse, reason, message, session.Generation)) if reflect.DeepEqual(original, session.Status) { return nil } return r.Status().Update(ctx, session) } +func sessionExclusivelyOwnsResource(object metav1.Object, session *clusterv1alpha1.T4Session) bool { + controller := metav1.GetControllerOf(object) + if controller == nil || controller.APIVersion != clusterv1alpha1.GroupVersion.String() || controller.Kind != "T4Session" || controller.Name != session.Name || controller.UID != session.UID { + return false + } + for _, reference := range object.GetOwnerReferences() { + if reference.APIVersion != clusterv1alpha1.GroupVersion.String() || reference.Kind != "T4Session" || reference.Name != session.Name || reference.UID != session.UID { + return false + } + } + return true +} + func serviceExposureIsInternal(service *corev1.Service) bool { if service.Spec.Type != corev1.ServiceTypeClusterIP || service.Spec.ClusterIP == corev1.ClusterIPNone || service.Spec.ExternalName != "" || len(service.Spec.ExternalIPs) != 0 || service.Spec.LoadBalancerIP != "" || len(service.Spec.LoadBalancerSourceRanges) != 0 || @@ -741,9 +971,62 @@ func labelsContain(actual, required map[string]string) bool { return true } +func indexSessionByHostRef(object client.Object) []string { + session, ok := object.(*clusterv1alpha1.T4Session) + if !ok || session.Spec.HostRef == "" { + return nil + } + return []string{session.Spec.HostRef} +} + +func indexSessionByWorkspaceRef(object client.Object) []string { + session, ok := object.(*clusterv1alpha1.T4Session) + if !ok || session.Spec.WorkspaceRef == "" { + return nil + } + return []string{session.Spec.WorkspaceRef} +} + +func (r *SessionReconciler) sessionRequestsForHost(ctx context.Context, object client.Object) []ctrl.Request { + host, ok := object.(*clusterv1alpha1.T4ClusterHost) + if !ok || host.Name == "" || host.Namespace == "" { + return nil + } + return r.sessionRequestsForReference(ctx, host.Namespace, sessionHostRefIndexField, host.Name, "clusterHost", client.ObjectKeyFromObject(host)) +} + +func (r *SessionReconciler) sessionRequestsForWorkspace(ctx context.Context, object client.Object) []ctrl.Request { + workspace, ok := object.(*clusterv1alpha1.T4Workspace) + if !ok || workspace.Name == "" || workspace.Namespace == "" { + return nil + } + return r.sessionRequestsForReference(ctx, workspace.Namespace, sessionWorkspaceRefIndexField, workspace.Name, "workspace", client.ObjectKeyFromObject(workspace)) +} + +func (r *SessionReconciler) sessionRequestsForReference(ctx context.Context, namespace, field, value, dependencyKind string, dependencyKey types.NamespacedName) []ctrl.Request { + var sessions clusterv1alpha1.T4SessionList + if err := r.List(ctx, &sessions, client.InNamespace(namespace), client.MatchingFields{field: value}); err != nil { + ctrl.LoggerFrom(ctx).Error(err, "unable to map dependency to sessions", dependencyKind, dependencyKey) + return nil + } + requests := make([]ctrl.Request, 0, len(sessions.Items)) + for i := range sessions.Items { + requests = append(requests, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(&sessions.Items[i])}) + } + return requests +} + func (r *SessionReconciler) SetupWithManager(manager ctrl.Manager) error { + if err := manager.GetFieldIndexer().IndexField(context.Background(), &clusterv1alpha1.T4Session{}, sessionHostRefIndexField, indexSessionByHostRef); err != nil { + return fmt.Errorf("index T4Session by host reference: %w", err) + } + if err := manager.GetFieldIndexer().IndexField(context.Background(), &clusterv1alpha1.T4Session{}, sessionWorkspaceRefIndexField, indexSessionByWorkspaceRef); err != nil { + return fmt.Errorf("index T4Session by workspace reference: %w", err) + } return ctrl.NewControllerManagedBy(manager). For(&clusterv1alpha1.T4Session{}). + Watches(&clusterv1alpha1.T4ClusterHost{}, handler.EnqueueRequestsFromMapFunc(r.sessionRequestsForHost)). + Watches(&clusterv1alpha1.T4Workspace{}, handler.EnqueueRequestsFromMapFunc(r.sessionRequestsForWorkspace)). Owns(&corev1.Pod{}). Owns(&corev1.Service{}). Complete(r) diff --git a/packages/cluster-operator/controllers/session_controller_test.go b/packages/cluster-operator/controllers/session_controller_test.go new file mode 100644 index 00000000..cba15fef --- /dev/null +++ b/packages/cluster-operator/controllers/session_controller_test.go @@ -0,0 +1,62 @@ +package controllers + +import ( + "context" + "testing" + + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + + clusterv1alpha1 "github.com/LycaonLLC/t4-code/packages/cluster-operator/api/v1alpha1" +) + +func TestSessionRequestsForHostOnlyEnqueuesAffectedSessions(t *testing.T) { + scheme := runtime.NewScheme() + if err := clusterv1alpha1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } + objects := []client.Object{ + &clusterv1alpha1.T4Session{ObjectMeta: metav1.ObjectMeta{Name: "session-a", Namespace: "team"}, Spec: clusterv1alpha1.T4SessionSpec{HostRef: "host-a", WorkspaceRef: "workspace-a"}}, + &clusterv1alpha1.T4Session{ObjectMeta: metav1.ObjectMeta{Name: "session-b", Namespace: "team"}, Spec: clusterv1alpha1.T4SessionSpec{HostRef: "host-a", WorkspaceRef: "workspace-b"}}, + &clusterv1alpha1.T4Session{ObjectMeta: metav1.ObjectMeta{Name: "other-host", Namespace: "team"}, Spec: clusterv1alpha1.T4SessionSpec{HostRef: "host-b", WorkspaceRef: "workspace-a"}}, + &clusterv1alpha1.T4Session{ObjectMeta: metav1.ObjectMeta{Name: "other-namespace", Namespace: "other"}, Spec: clusterv1alpha1.T4SessionSpec{HostRef: "host-a", WorkspaceRef: "workspace-a"}}, + } + c := fake.NewClientBuilder().WithScheme(scheme). + WithIndex(&clusterv1alpha1.T4Session{}, sessionHostRefIndexField, indexSessionByHostRef). + WithIndex(&clusterv1alpha1.T4Session{}, sessionWorkspaceRefIndexField, indexSessionByWorkspaceRef). + WithObjects(objects...).Build() + r := &SessionReconciler{Client: c, Scheme: scheme} + + requests := r.sessionRequestsForHost(context.Background(), &clusterv1alpha1.T4ClusterHost{ObjectMeta: metav1.ObjectMeta{Name: "host-a", Namespace: "team"}}) + assertRequestSet(t, requests, []types.NamespacedName{ + {Namespace: "team", Name: "session-a"}, + {Namespace: "team", Name: "session-b"}, + }) +} + +func TestSessionRequestsForWorkspaceOnlyEnqueuesAffectedSessions(t *testing.T) { + scheme := runtime.NewScheme() + if err := clusterv1alpha1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } + objects := []client.Object{ + &clusterv1alpha1.T4Session{ObjectMeta: metav1.ObjectMeta{Name: "session-a", Namespace: "team"}, Spec: clusterv1alpha1.T4SessionSpec{HostRef: "host-a", WorkspaceRef: "workspace-a"}}, + &clusterv1alpha1.T4Session{ObjectMeta: metav1.ObjectMeta{Name: "session-b", Namespace: "team"}, Spec: clusterv1alpha1.T4SessionSpec{HostRef: "host-b", WorkspaceRef: "workspace-a"}}, + &clusterv1alpha1.T4Session{ObjectMeta: metav1.ObjectMeta{Name: "other-workspace", Namespace: "team"}, Spec: clusterv1alpha1.T4SessionSpec{HostRef: "host-a", WorkspaceRef: "workspace-b"}}, + &clusterv1alpha1.T4Session{ObjectMeta: metav1.ObjectMeta{Name: "other-namespace", Namespace: "other"}, Spec: clusterv1alpha1.T4SessionSpec{HostRef: "host-a", WorkspaceRef: "workspace-a"}}, + } + c := fake.NewClientBuilder().WithScheme(scheme). + WithIndex(&clusterv1alpha1.T4Session{}, sessionHostRefIndexField, indexSessionByHostRef). + WithIndex(&clusterv1alpha1.T4Session{}, sessionWorkspaceRefIndexField, indexSessionByWorkspaceRef). + WithObjects(objects...).Build() + r := &SessionReconciler{Client: c, Scheme: scheme} + + requests := r.sessionRequestsForWorkspace(context.Background(), &clusterv1alpha1.T4Workspace{ObjectMeta: metav1.ObjectMeta{Name: "workspace-a", Namespace: "team"}}) + assertRequestSet(t, requests, []types.NamespacedName{ + {Namespace: "team", Name: "session-a"}, + {Namespace: "team", Name: "session-b"}, + }) +} diff --git a/packages/cluster-operator/controllers/workspace_controller.go b/packages/cluster-operator/controllers/workspace_controller.go index fbc60f45..e1b616dc 100644 --- a/packages/cluster-operator/controllers/workspace_controller.go +++ b/packages/cluster-operator/controllers/workspace_controller.go @@ -10,6 +10,7 @@ import ( storagev1 "k8s.io/api/storage/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" meta "k8s.io/apimachinery/pkg/api/meta" + apiresource "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/types" @@ -23,7 +24,8 @@ import ( type WorkspaceReconciler struct { client.Client - Scheme *runtime.Scheme + APIReader client.Reader + Scheme *runtime.Scheme } const ( @@ -73,8 +75,9 @@ func (r *WorkspaceReconciler) Reconcile(ctx context.Context, request ctrl.Reques } pvcName := WorkspacePVCName(&workspace) + pvcKey := types.NamespacedName{Namespace: workspace.Namespace, Name: pvcName} var pvc corev1.PersistentVolumeClaim - err = r.Get(ctx, types.NamespacedName{Namespace: workspace.Namespace, Name: pvcName}, &pvc) + err = r.Get(ctx, pvcKey, &pvc) if apierrors.IsNotFound(err) { volumeMode := corev1.PersistentVolumeFilesystem pvc = corev1.PersistentVolumeClaim{ @@ -98,13 +101,27 @@ func (r *WorkspaceReconciler) Reconcile(ctx context.Context, request ctrl.Reques return ctrl.Result{}, err } } - if err := r.Create(ctx, &pvc); err != nil && !apierrors.IsAlreadyExists(err) { - return ctrl.Result{}, err + if err := r.Create(ctx, &pvc); err != nil { + if !apierrors.IsAlreadyExists(err) { + return ctrl.Result{}, err + } + reader := r.APIReader + if reader == nil { + reader = r.Client + } + if err := reader.Get(ctx, pvcKey, &pvc); err != nil { + return ctrl.Result{}, err + } } } else if err != nil { return ctrl.Result{}, err - } else if !workspaceOwnsPVC(&workspace, &pvc) { + } + if !workspaceOwnsPVC(&workspace, &pvc) { return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateWorkspaceFailure(ctx, &workspace, "StorageReady", "PVCOwnershipConflict", "deterministic workspace PVC does not belong to this workspace") + } else if !pvcHasRWX(&pvc) { + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateWorkspaceFailure(ctx, &workspace, "StorageReady", "PVCNotRWX", "workspace PVC does not request ReadWriteMany") + } else if pvcStorageClassName(&pvc) != storageClassName { + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateWorkspaceFailure(ctx, &workspace, "StorageReady", ReasonStorageClassMismatch, fmt.Sprintf("workspace PVC uses StorageClass %q instead of host-selected %q; data-bearing PVCs are never recreated automatically", pvcStorageClassName(&pvc), storageClassName)) } else if workspace.Spec.RetentionPolicy == clusterv1alpha1.RetentionPolicyRetain && metav1.IsControlledBy(&pvc, &workspace) { before := pvc.DeepCopy() pvc.OwnerReferences = removeWorkspaceOwnerReference(pvc.OwnerReferences, workspace.UID) @@ -115,6 +132,31 @@ func (r *WorkspaceReconciler) Reconcile(ctx context.Context, request ctrl.Reques return ctrl.Result{Requeue: true}, nil } } + reader := r.APIReader + if reader == nil { + reader = r.Client + } + var authoritativePVC corev1.PersistentVolumeClaim + if err := reader.Get(ctx, pvcKey, &authoritativePVC); err != nil { + if apierrors.IsNotFound(err) { + return ctrl.Result{RequeueAfter: 5 * time.Second}, r.updateWorkspaceFailure(ctx, &workspace, "StorageReady", "PVCNotFound", "workspace PVC does not exist in authoritative API state") + } + return ctrl.Result{}, err + } + if authoritativePVC.UID != pvc.UID || !workspaceOwnsPVC(&workspace, &authoritativePVC) { + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateWorkspaceFailure(ctx, &workspace, "StorageReady", "PVCOwnershipConflict", "authoritative workspace PVC identity or ownership does not belong to this workspace") + } + if !pvcHasRWX(&authoritativePVC) { + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateWorkspaceFailure(ctx, &workspace, "StorageReady", "PVCNotRWX", "authoritative workspace PVC does not request ReadWriteMany") + } + if pvcStorageClassName(&authoritativePVC) != storageClassName { + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateWorkspaceFailure(ctx, &workspace, "StorageReady", ReasonStorageClassMismatch, fmt.Sprintf("authoritative workspace PVC uses StorageClass %q instead of host-selected %q; data-bearing PVCs are never recreated automatically", pvcStorageClassName(&authoritativePVC), storageClassName)) + } + pvc = authoritativePVC + + if pvc.Status.Phase == corev1.ClaimLost { + return ctrl.Result{RequeueAfter: 30 * time.Second}, r.updateWorkspaceFailure(ctx, &workspace, "StorageReady", "PVCLost", "workspace PVC lost its volume") + } original := workspace.Status original.Capacity = workspace.Status.Capacity.DeepCopy() @@ -126,19 +168,12 @@ func (r *WorkspaceReconciler) Reconcile(ctx context.Context, request ctrl.Reques workspace.Status.PVCPhase = pvc.Status.Phase capacity := pvc.Status.Capacity[corev1.ResourceStorage] workspace.Status.Capacity = capacity.DeepCopy() + meta.SetStatusCondition(&workspace.Status.Conditions, condition("HostReady", metav1.ConditionTrue, "HostResolved", "referenced T4ClusterHost is available", workspace.Generation)) meta.SetStatusCondition(&workspace.Status.Conditions, condition("StorageReady", metav1.ConditionTrue, ReasonStorageReady, "RWX StorageClass and workspace PVC are accepted", workspace.Generation)) switch pvc.Status.Phase { case corev1.ClaimBound: - if !pvcHasRWX(&pvc) { - workspace.Status.Phase = clusterv1alpha1.InfrastructureFailed - meta.SetStatusCondition(&workspace.Status.Conditions, condition("Ready", metav1.ConditionFalse, "PVCNotRWX", "bound workspace PVC does not request ReadWriteMany", workspace.Generation)) - } else { - workspace.Status.Phase = clusterv1alpha1.InfrastructureReady - meta.SetStatusCondition(&workspace.Status.Conditions, condition("Ready", metav1.ConditionTrue, "PVCBound", "workspace PVC is bound with ReadWriteMany access", workspace.Generation)) - } - case corev1.ClaimLost: - workspace.Status.Phase = clusterv1alpha1.InfrastructureFailed - meta.SetStatusCondition(&workspace.Status.Conditions, condition("Ready", metav1.ConditionFalse, "PVCLost", "workspace PVC lost its volume", workspace.Generation)) + workspace.Status.Phase = clusterv1alpha1.InfrastructureReady + meta.SetStatusCondition(&workspace.Status.Conditions, condition("Ready", metav1.ConditionTrue, "PVCBound", "workspace PVC is bound with ReadWriteMany access", workspace.Generation)) default: workspace.Status.Phase = clusterv1alpha1.InfrastructurePending meta.SetStatusCondition(&workspace.Status.Conditions, condition("Ready", metav1.ConditionFalse, "PVCBinding", "workspace PVC is waiting to bind", workspace.Generation)) @@ -158,6 +193,11 @@ func workspaceOwnsPVC(workspace *clusterv1alpha1.T4Workspace, pvc *corev1.Persis if pvc.Annotations[clusterv1alpha1.WorkspaceUIDAnnotation] != string(workspace.UID) { return false } + for _, reference := range pvc.OwnerReferences { + if reference.APIVersion != clusterv1alpha1.GroupVersion.String() || reference.Kind != "T4Workspace" || reference.Name != workspace.Name || reference.UID != workspace.UID { + return false + } + } controller := metav1.GetControllerOf(pvc) if workspace.Spec.RetentionPolicy == clusterv1alpha1.RetentionPolicyDelete { return controller != nil && controller.UID == workspace.UID @@ -192,7 +232,11 @@ func (r *WorkspaceReconciler) reconcileDelete(ctx context.Context, workspace *cl } } var sessions clusterv1alpha1.T4SessionList - if err := r.List(ctx, &sessions, client.InNamespace(workspace.Namespace)); err != nil { + sessionReader := r.APIReader + if sessionReader == nil { + sessionReader = r.Client + } + if err := sessionReader.List(ctx, &sessions, client.InNamespace(workspace.Namespace)); err != nil { return ctrl.Result{}, err } remainingSessions := 0 @@ -216,7 +260,11 @@ func (r *WorkspaceReconciler) reconcileDelete(ctx context.Context, workspace *cl } pvcKey := types.NamespacedName{Namespace: workspace.Namespace, Name: WorkspacePVCName(workspace)} var pvc corev1.PersistentVolumeClaim - err := r.Get(ctx, pvcKey, &pvc) + reader := r.APIReader + if reader == nil { + reader = r.Client + } + err := reader.Get(ctx, pvcKey, &pvc) if err == nil && !workspaceOwnsPVC(workspace, &pvc) { before := workspace.Status if workspace.Status.Conditions != nil { @@ -248,7 +296,7 @@ func (r *WorkspaceReconciler) reconcileDelete(ctx context.Context, workspace *cl } } else { if err == nil { - if err := r.Delete(ctx, &pvc); err != nil && !apierrors.IsNotFound(err) { + if err := deleteWithPreconditions(ctx, r.Client, &pvc); err != nil && !apierrors.IsNotFound(err) { return ctrl.Result{}, err } return ctrl.Result{RequeueAfter: time.Second}, nil @@ -268,11 +316,18 @@ func (r *WorkspaceReconciler) updateWorkspaceFailure(ctx context.Context, worksp original.Conditions = append([]metav1.Condition(nil), workspace.Status.Conditions...) } workspace.Status.ObservedGeneration = workspace.Generation + workspace.Status.PVCName = "" + workspace.Status.PVCPhase = "" + workspace.Status.Capacity = apiresource.Quantity{} workspace.Status.Phase = clusterv1alpha1.InfrastructureFailed - meta.SetStatusCondition(&workspace.Status.Conditions, condition(conditionType, metav1.ConditionFalse, reason, message, workspace.Generation)) - if conditionType != "Ready" { - meta.SetStatusCondition(&workspace.Status.Conditions, condition("Ready", metav1.ConditionFalse, reason, message, workspace.Generation)) + if conditionType == "HostReady" { + meta.SetStatusCondition(&workspace.Status.Conditions, condition("HostReady", metav1.ConditionFalse, reason, message, workspace.Generation)) + meta.SetStatusCondition(&workspace.Status.Conditions, condition("StorageReady", metav1.ConditionUnknown, "NotEvaluated", "storage dependency was not evaluated because the referenced host is unavailable", workspace.Generation)) + } else { + meta.SetStatusCondition(&workspace.Status.Conditions, condition("HostReady", metav1.ConditionTrue, "HostResolved", "referenced T4ClusterHost is available", workspace.Generation)) + meta.SetStatusCondition(&workspace.Status.Conditions, condition("StorageReady", metav1.ConditionFalse, reason, message, workspace.Generation)) } + meta.SetStatusCondition(&workspace.Status.Conditions, condition("Ready", metav1.ConditionFalse, reason, message, workspace.Generation)) if reflect.DeepEqual(original, workspace.Status) { return nil } @@ -337,6 +392,23 @@ func (r *WorkspaceReconciler) workspaceRequestsForStorageClass(ctx context.Conte return requests } +func (r *WorkspaceReconciler) workspaceRequestsForHost(ctx context.Context, object client.Object) []ctrl.Request { + host, ok := object.(*clusterv1alpha1.T4ClusterHost) + if !ok || host.Name == "" || host.Namespace == "" { + return nil + } + var workspaces clusterv1alpha1.T4WorkspaceList + if err := r.List(ctx, &workspaces, client.InNamespace(host.Namespace), client.MatchingFields{workspaceHostRefIndexField: host.Name}); err != nil { + ctrl.LoggerFrom(ctx).Error(err, "unable to map cluster host to workspaces", "clusterHost", client.ObjectKeyFromObject(host)) + return nil + } + requests := make([]ctrl.Request, 0, len(workspaces.Items)) + for i := range workspaces.Items { + requests = append(requests, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(&workspaces.Items[i])}) + } + return requests +} + func (r *WorkspaceReconciler) SetupWithManager(manager ctrl.Manager) error { if err := manager.GetFieldIndexer().IndexField(context.Background(), &clusterv1alpha1.T4ClusterHost{}, hostStorageClassIndexField, indexHostByStorageClass); err != nil { return fmt.Errorf("index T4ClusterHost by StorageClass: %w", err) @@ -346,6 +418,7 @@ func (r *WorkspaceReconciler) SetupWithManager(manager ctrl.Manager) error { } return ctrl.NewControllerManagedBy(manager). For(&clusterv1alpha1.T4Workspace{}). + Watches(&clusterv1alpha1.T4ClusterHost{}, handler.EnqueueRequestsFromMapFunc(r.workspaceRequestsForHost)). Watches(&corev1.PersistentVolumeClaim{}, handler.EnqueueRequestsFromMapFunc(workspaceRequestsForPVC)). Watches(&storagev1.StorageClass{}, handler.EnqueueRequestsFromMapFunc(r.workspaceRequestsForStorageClass)). Complete(r) diff --git a/packages/cluster-operator/controllers/workspace_controller_test.go b/packages/cluster-operator/controllers/workspace_controller_test.go index 412705d0..27d07759 100644 --- a/packages/cluster-operator/controllers/workspace_controller_test.go +++ b/packages/cluster-operator/controllers/workspace_controller_test.go @@ -9,6 +9,7 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/types" + ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/client/fake" @@ -107,3 +108,42 @@ func TestWorkspaceRequestsForStorageClassOnlyEnqueuesAffectedWorkspaces(t *testi } } } + +func TestWorkspaceRequestsForHostOnlyEnqueuesAffectedWorkspaces(t *testing.T) { + scheme := runtime.NewScheme() + if err := clusterv1alpha1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } + objects := []client.Object{ + &clusterv1alpha1.T4Workspace{ObjectMeta: metav1.ObjectMeta{Name: "workspace-a", Namespace: "team"}, Spec: clusterv1alpha1.T4WorkspaceSpec{HostRef: "host-a"}}, + &clusterv1alpha1.T4Workspace{ObjectMeta: metav1.ObjectMeta{Name: "workspace-b", Namespace: "team"}, Spec: clusterv1alpha1.T4WorkspaceSpec{HostRef: "host-a"}}, + &clusterv1alpha1.T4Workspace{ObjectMeta: metav1.ObjectMeta{Name: "other-host", Namespace: "team"}, Spec: clusterv1alpha1.T4WorkspaceSpec{HostRef: "host-b"}}, + &clusterv1alpha1.T4Workspace{ObjectMeta: metav1.ObjectMeta{Name: "other-namespace", Namespace: "other"}, Spec: clusterv1alpha1.T4WorkspaceSpec{HostRef: "host-a"}}, + } + c := fake.NewClientBuilder().WithScheme(scheme). + WithIndex(&clusterv1alpha1.T4Workspace{}, workspaceHostRefIndexField, indexWorkspaceByHostRef). + WithObjects(objects...).Build() + r := &WorkspaceReconciler{Client: c, Scheme: scheme} + + requests := r.workspaceRequestsForHost(context.Background(), &clusterv1alpha1.T4ClusterHost{ObjectMeta: metav1.ObjectMeta{Name: "host-a", Namespace: "team"}}) + assertRequestSet(t, requests, []types.NamespacedName{ + {Namespace: "team", Name: "workspace-a"}, + {Namespace: "team", Name: "workspace-b"}, + }) +} + +func assertRequestSet(t *testing.T, requests []ctrl.Request, want []types.NamespacedName) { + t.Helper() + got := make(map[types.NamespacedName]int, len(requests)) + for _, request := range requests { + got[request.NamespacedName]++ + } + if len(got) != len(want) { + t.Fatalf("requests = %#v, want exactly %v", requests, want) + } + for _, key := range want { + if got[key] != 1 { + t.Fatalf("requests = %#v, want %v exactly once", requests, key) + } + } +}