diff --git a/src/compute-plane-services/nvca/pkg/apis/nvsnap/v1alpha1/nvsnapfunctionstate_types.go b/src/compute-plane-services/nvca/pkg/apis/nvsnap/v1alpha1/nvsnapfunctionstate_types.go index c5fa8b0c04..e112494549 100644 --- a/src/compute-plane-services/nvca/pkg/apis/nvsnap/v1alpha1/nvsnapfunctionstate_types.go +++ b/src/compute-plane-services/nvca/pkg/apis/nvsnap/v1alpha1/nvsnapfunctionstate_types.go @@ -179,6 +179,12 @@ type NvSnapFunctionStateStatus struct { // observers attribute the in-flight claim. CaptureOwner string `json:"captureOwner,omitempty"` + // CaptureOwnerUID is the UID of the CaptureOwner pod. Inference pod + // names are deterministic, so a replacement pod can reuse a dead + // claimant's namespace/name; the UID tells the two apart when the + // claim's liveness is checked. + CaptureOwnerUID string `json:"captureOwnerUID,omitempty"` + // CaptureLeaseExpiry bounds how long a Capturing claim is honored. // If the owning reconcile dies mid-capture (controller restart, // lost leader election, ctx cancel) the claim would otherwise pin diff --git a/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/BUILD.bazel b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/BUILD.bazel index b2bb250ac7..06f4b7f1b6 100644 --- a/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/BUILD.bazel +++ b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/BUILD.bazel @@ -63,5 +63,6 @@ go_test( "//src/compute-plane-services/nvca/vendor/k8s.io/client-go/dynamic", "//src/compute-plane-services/nvca/vendor/k8s.io/client-go/dynamic/fake", "//src/compute-plane-services/nvca/vendor/k8s.io/client-go/kubernetes/fake", + "//src/compute-plane-services/nvca/vendor/k8s.io/client-go/testing", ], ) diff --git a/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/captureonce_test.go b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/captureonce_test.go index 96d4bd002d..441aa602ad 100644 --- a/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/captureonce_test.go +++ b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/captureonce_test.go @@ -17,6 +17,9 @@ package reconciler import ( "context" "encoding/json" + "errors" + corev1 "k8s.io/api/core/v1" + k8sfake "k8s.io/client-go/kubernetes/fake" "net/http" "net/http/httptest" "strings" @@ -150,6 +153,125 @@ func TestTryClaimCapture(t *testing.T) { // the Capturing claim — without this a completed/failed capture would // leave captureOwner/leaseExpiry set and either pin the function or let // readStatus mis-report a stale owner. +// A live lease held by a pod that no longer exists is stealable at once: +// the owner died mid-capture (dev1 2026-09-28: the node agent restarted +// under it and the pod was replaced) and nothing would release the claim +// before the lease expired. +func TestTryClaimCapture_DeadOwnerIsStealable(t *testing.T) { + ctx := context.Background() + fvID := "fv-dead" + now := time.Now() + lease := now.Add(50 * time.Minute) + dyn := newFakeDynamic(coldCFS(fvID)) + if claimed, err := tryClaimCapture(ctx, dyn, fvID, "ns1/podA", lease, now); err != nil || !claimed { + t.Fatalf("podA claim: claimed=%v err=%v", claimed, err) + } + dead := func(context.Context, string, string) bool { return false } + alive := func(context.Context, string, string) bool { return true } + if claimed, _ := tryClaimCaptureLive(ctx, dyn, fvID, "ns1/podB", "uid-b", lease, now, alive); claimed { + t.Fatal("podB must not steal a live claim from a live owner") + } + if claimed, _ := tryClaimCaptureLive(ctx, dyn, fvID, "ns1/podB", "uid-b", lease, now, nil); claimed { + t.Fatal("unknown liveness (nil) must keep the claim") + } + claimed, err := tryClaimCaptureLive(ctx, dyn, fvID, "ns1/podB", "uid-b", lease, now, dead) + if err != nil || !claimed { + t.Fatalf("podB must steal the claim from a dead owner: claimed=%v err=%v", claimed, err) + } + cur, err := dyn.Resource(CFSResource).Get(ctx, fvID, metav1.GetOptions{}) + if err != nil { + t.Fatal(err) + } + if st := readStatus(cur); st.CaptureOwner != "ns1/podB" || st.CaptureOwnerUID != "uid-b" { + t.Errorf("owner after steal = %q/%q, want ns1/podB uid-b", st.CaptureOwner, st.CaptureOwnerUID) + } +} + +// A pod that reuses the claimant's namespace/name with a new UID is a +// replacement, not the owner: its own claim is a foreign steal, subject to +// the liveness check like any other, and the liveness check itself treats +// the UID mismatch as "owner gone". +func TestTryClaimCapture_SameNameDifferentUIDIsNotTheOwner(t *testing.T) { + ctx := context.Background() + fvID := "fv-same-name" + now := time.Now() + lease := now.Add(50 * time.Minute) + dyn := newFakeDynamic(coldCFS(fvID)) + if claimed, err := tryClaimCaptureLive(ctx, dyn, fvID, "ns1/podA", "uid-old", lease, now, nil); err != nil || !claimed { + t.Fatalf("first claim: %v %v", claimed, err) + } + alive := func(context.Context, string, string) bool { return true } + if claimed, _ := tryClaimCaptureLive(ctx, dyn, fvID, "ns1/podA", "uid-new", lease, now, alive); claimed { + t.Fatal("same name with a different UID must not be treated as the re-entrant owner") + } + if claimed, _ := tryClaimCaptureLive(ctx, dyn, fvID, "ns1/podA", "uid-old", lease, now, alive); !claimed { + t.Fatal("the real owner re-enters its own claim") + } + // claimOwnerAlive compares the UID of the pod it finds. + r := &Reconciler{KubeClient: k8sfake.NewSimpleClientset(&corev1.Pod{ + ObjectMeta: metav1.ObjectMeta{Name: "podA", Namespace: "ns1", UID: "uid-new"}, + Status: corev1.PodStatus{Phase: corev1.PodRunning}, + })} + if r.claimOwnerAlive(ctx, "ns1/podA", "uid-old") { + t.Error("a replacement pod under the old name must read as owner gone") + } + if !r.claimOwnerAlive(ctx, "ns1/podA", "uid-new") { + t.Error("the pod with the recorded UID is alive") + } + if !r.claimOwnerAlive(ctx, "ns1/podA", "") { + t.Error("without a recorded UID the name alone decides, as before") + } +} + +// Terminal writes are fenced by the claim: a reconcile whose claim was taken +// over must not release or overwrite the new owner's claim. +func TestWriteStatusRejectsSupersededClaim(t *testing.T) { + const fvID = "fv-fenced" + ctx := context.Background() + dyn := newFakeDynamic(capturingCFS(fvID, "ns1/podB", time.Now().Add(30*time.Minute))) + warm := statusUpdate{CheckpointHash: "cafe", CapturedHere: true, LocalCacheState: nvsnapv1alpha1.LocalCacheStateWarm} + err := writeStatus(ctx, dyn, fvID, warm, claimToken{Owner: "ns1/podA", UID: "uid-a"}) + if !errors.Is(err, ErrClaimSuperseded) { + t.Fatalf("podA's write must be rejected once podB owns the claim: %v", err) + } + cur, _ := dyn.Resource(CFSResource).Get(ctx, fvID, metav1.GetOptions{}) + if st := readStatus(cur); st.CaptureOwner != "ns1/podB" || st.LocalCacheState != nvsnapv1alpha1.LocalCacheStateCapturing { + t.Errorf("status must be untouched by the superseded writer: %+v", st) + } + if err := writeStatus(ctx, dyn, fvID, warm, claimToken{Owner: "ns1/podB"}); err != nil { + t.Fatalf("the owner's write goes through: %v", err) + } + if err := writeStatus(ctx, dyn, fvID, warm, claimToken{}); err != nil { + t.Fatalf("an unfenced write is never rejected: %v", err) + } + // After a terminal write the claim is gone: a late write from the + // former owner is superseded, not applied. + if err := writeStatus(ctx, dyn, fvID, statusUpdate{LocalCacheState: nvsnapv1alpha1.LocalCacheStateFailed, LastError: "late"}, claimToken{Owner: "ns1/podB"}); !errors.Is(err, ErrClaimSuperseded) { + t.Fatalf("a writer whose claim was already closed must be rejected: %v", err) + } + // A path that observed no claim is rejected once a claim appeared. + dyn2 := newFakeDynamic(capturingCFS("fv-claimed", "ns1/podC", time.Now().Add(30*time.Minute))) + if err := writeStatus(ctx, dyn2, "fv-claimed", warm, claimToken{ExpectUnclaimed: true}); !errors.Is(err, ErrClaimSuperseded) { + t.Fatalf("ExpectUnclaimed must reject a write over a live claim: %v", err) + } + cur2, _ := dyn2.Resource(CFSResource).Get(ctx, "fv-claimed", metav1.GetOptions{}) + if st := readStatus(cur2); st.CaptureOwner != "ns1/podC" { + t.Errorf("claim must survive: %+v", st) + } + // The version fence catches what state comparison cannot: a claim that + // opened and closed, or a lease refresh by the same owner, between the + // observation and the write. + obj := coldCFS("fv-versioned") + obj.SetResourceVersion("41") + dyn3 := newFakeDynamic(obj) + if err := writeStatus(ctx, dyn3, "fv-versioned", warm, claimToken{ExpectUnclaimed: true, ObservedResourceVersion: "40"}); !errors.Is(err, ErrClaimSuperseded) { + t.Fatalf("a write observed at an older version must be rejected: %v", err) + } + if err := writeStatus(ctx, dyn3, "fv-versioned", warm, claimToken{ExpectUnclaimed: true, ObservedResourceVersion: "41"}); err != nil { + t.Fatalf("a write at the observed version goes through: %v", err) + } +} + func TestWriteStatusReleasesClaim(t *testing.T) { const fvID = "fv-release" ctx := context.Background() @@ -159,7 +281,7 @@ func TestWriteStatusReleasesClaim(t *testing.T) { CheckpointHash: "feedface", CapturedHere: true, LocalCacheState: nvsnapv1alpha1.LocalCacheStateWarm, - }); err != nil { + }, claimToken{Owner: "ns1/podA"}); err != nil { t.Fatalf("writeStatus: %v", err) } cur, _ := dyn.Resource(CFSResource).Get(ctx, fvID, metav1.GetOptions{}) @@ -203,6 +325,10 @@ func TestReconcileSkipsWhenCaptureInFlight(t *testing.T) { // podA already holds a live Capturing claim on this fvID. dyn := newFakeDynamic(capturingCFS(fvID, "ns1/podA", time.Now().Add(30*time.Minute))) r := newTestReconciler(t, pod, srv, dyn) + // podA, the claim owner, is alive: the in-flight guard must hold. + if _, err := r.KubeClient.CoreV1().Pods("ns1").Create(context.Background(), inferencePod("podA", "ns1", fvID), metav1.CreateOptions{}); err != nil { + t.Fatal(err) + } if err := r.Reconcile(context.Background(), pod); err != nil { t.Fatalf("Reconcile: %v", err) @@ -216,3 +342,29 @@ func TestReconcileSkipsWhenCaptureInFlight(t *testing.T) { t.Errorf("claim owner = %q, want unchanged ns1/podA", st.CaptureOwner) } } + +// Same scenario, but the claim owner no longer exists: the claim is dead +// and podB must capture instead of waiting out the lease. +func TestReconcileStealsClaimFromDeadOwner(t *testing.T) { + fvID := "fv-dead-owner" + posted := 0 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch { + case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/api/v1/checkpoints/lookup"): + _ = json.NewEncoder(w).Encode(map[string]any{"matches": []any{}}) + case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/api/v1/checkpoints"): + posted++ + w.WriteHeader(http.StatusInternalServerError) // stop the flow here; the claim is what is under test + default: + w.WriteHeader(http.StatusNotFound) + } + })) + defer srv.Close() + pod := inferencePod("podB", "ns1", fvID) + dyn := newFakeDynamic(capturingCFS(fvID, "ns1/podA", time.Now().Add(30*time.Minute))) + r := newTestReconciler(t, pod, srv, dyn) + _ = r.Reconcile(context.Background(), pod) + if posted != 1 { + t.Fatalf("podB must take over a claim whose owner pod is gone; POSTs=%d", posted) + } +} diff --git a/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/coldstart_pioneer_test.go b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/coldstart_pioneer_test.go index ab453c25d2..4cf2409369 100644 --- a/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/coldstart_pioneer_test.go +++ b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/coldstart_pioneer_test.go @@ -167,7 +167,7 @@ func TestWriteStatusReleasesPioneerClaim(t *testing.T) { CheckpointHash: "feedface", CapturedHere: true, LocalCacheState: nvsnapv1alpha1.LocalCacheStateWarm, - }); err != nil { + }, claimToken{}); err != nil { t.Fatalf("writeStatus: %v", err) } st := readStatus(mustGet(t, dyn, fvID)) diff --git a/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/reconciler.go b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/reconciler.go index 85d99d1259..e7305c8431 100644 --- a/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/reconciler.go +++ b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/reconciler.go @@ -29,6 +29,7 @@ import ( "errors" "fmt" "net/http" + "strings" "sync" "time" @@ -357,7 +358,7 @@ func (r *Reconciler) Reconcile(ctx context.Context, pod *corev1.Pod) error { AttemptCount: 0, LastError: "", LastAttemptAt: now, - }); err != nil { + }, claimToken{}); err != nil { return fmt.Errorf("recovery: writeStatus Warm: %w", err) } return r.removeCheckpointOnWarm(ctx, pod, log) @@ -414,7 +415,8 @@ func (r *Reconciler) Reconcile(ctx context.Context, pod *corev1.Pod) error { // buffer so losers skip immediately rather than each dwelling. claimNow := time.Now() owner := pod.Namespace + "/" + pod.Name - claimed, err := tryClaimCapture(ctx, r.DynClient, fvID, owner, claimNow.Add(r.CaptureLeaseTTL), claimNow) + token := claimToken{Owner: owner, UID: string(pod.UID)} + claimed, err := tryClaimCaptureLive(ctx, r.DynClient, fvID, owner, token.UID, claimNow.Add(r.CaptureLeaseTTL), claimNow, r.claimOwnerAlive) if err != nil { return fmt.Errorf("capture-once claim for %s: %w", fvID, err) } @@ -444,7 +446,7 @@ func (r *Reconciler) Reconcile(ctx context.Context, pod *corev1.Pod) error { LeaveRunning: true, }) if err != nil { - return r.recordFailure(ctx, fvID, "create_failed", prev, fmt.Errorf("nvsnap CreateCheckpoint: %w", err), log) + return r.recordFailure(ctx, fvID, "create_failed", prev, fmt.Errorf("nvsnap CreateCheckpoint: %w", err), token, log) } log = log.WithField("checkpointID", ckpt.ID) @@ -453,11 +455,11 @@ func (r *Reconciler) Reconcile(ctx context.Context, pod *corev1.Pod) error { defer pollCancel() final, err := r.pollCheckpointTerminal(pollCtx, ckpt.ID, log) if err != nil { - return r.recordFailure(ctx, fvID, "poll_failed", prev, err, log) + return r.recordFailure(ctx, fvID, "poll_failed", prev, err, token, log) } if final.Phase == nvsnap.PhaseFailed { return r.recordFailure(ctx, fvID, "completed_failed", prev, - fmt.Errorf("nvsnap-server reported checkpoint %s Phase=Failed: %s", ckpt.ID, final.Error), log) + fmt.Errorf("nvsnap-server reported checkpoint %s Phase=Failed: %s", ckpt.ID, final.Error), token, log) } // Defense in depth (nvca#14 + nvca#15). nvsnap-server should always @@ -498,7 +500,7 @@ func (r *Reconciler) Reconcile(ctx context.Context, pod *corev1.Pod) error { // an attempt. This branch's contract is "treat it as if the capture // never happened" -- prior state, no CFS-level error surface -- so it // restores prev verbatim and lets writeStatus drop the claim fields. - if err := releaseCaptureClaim(ctx, r.DynClient, fvID, prev); err != nil { + if err := releaseCaptureClaim(ctx, r.DynClient, fvID, prev, token); err != nil { // Non-fatal: the lease still expires on its own, just later. log.WithError(err).Warn("failed to release capture claim after empty hash; " + "peers will be gated until the lease expires") @@ -539,7 +541,7 @@ func (r *Reconciler) Reconcile(ctx context.Context, pod *corev1.Pod) error { log.WithError(err).WithField("hash", final.Hash). Warn("L2 promote-state poll did not reach terminal — falling back to cold start (capture itself succeeded)") return r.recordFailure(ctx, fvID, "promote_poll_failed", prev, - fmt.Errorf("L2 promote poll: %w", err), log) + fmt.Errorf("L2 promote poll: %w", err), token, log) } switch promote.State { case "failed": @@ -553,7 +555,7 @@ func (r *Reconciler) Reconcile(ctx context.Context, pod *corev1.Pod) error { "pvc_name": promote.PVCName, }).Warn("L2 promote terminally failed — pods will cold-start until next successful capture") return r.recordFailure(ctx, fvID, "promote_failed", prev, - fmt.Errorf("L2 promote state=failed for hash %s", final.Hash), log) + fmt.Errorf("L2 promote state=failed for hash %s", final.Hash), token, log) case "ready": log.WithFields(logrus.Fields{ "hash": final.Hash, @@ -577,7 +579,13 @@ func (r *Reconciler) Reconcile(ctx context.Context, pod *corev1.Pod) error { AttemptCount: 0, LastError: "", LastAttemptAt: now, - }); err != nil { + }, token); err != nil { + if errors.Is(err, ErrClaimSuperseded) { + // Another pod took the claim while this capture finished (this + // pod was terminating). Its result stands; ours is not recorded. + log.WithField("hash", final.Hash).Warn("capture finished after the claim was taken over; leaving the new owner's status untouched") + return nil + } log.WithError(err).Error("status update failed after successful checkpoint; will retry on next reconcile") return err } @@ -705,7 +713,7 @@ func isNonTransientAPIError(err error) bool { // // `reason` is a small enum that lands in the metric label — keep the // set bounded (see metrics.go). -func (r *Reconciler) recordFailure(ctx context.Context, fvID, reason string, prev cfsStatus, cause error, log logrus.FieldLogger) error { +func (r *Reconciler) recordFailure(ctx context.Context, fvID, reason string, prev cfsStatus, cause error, token claimToken, log logrus.FieldLogger) error { now := time.Now() upd := statusUpdate{ CheckpointHash: prev.CheckpointHash, @@ -722,7 +730,11 @@ func (r *Reconciler) recordFailure(ctx context.Context, fvID, reason string, pre if prev.CapturedAt != nil { upd.CapturedAt = prev.CapturedAt.Time } - if err := writeStatus(ctx, r.DynClient, fvID, upd); err != nil { + if err := writeStatus(ctx, r.DynClient, fvID, upd, token); err != nil { + if errors.Is(err, ErrClaimSuperseded) { + log.WithError(cause).Warn("attempt failed after the claim was taken over; the new owner's status is left untouched") + return nil + } log.WithError(err).Error("failed to write failure status; original error preserved in log") } checkpointAttemptFailures.WithLabelValues(reason).Inc() @@ -766,3 +778,27 @@ func (r *Reconciler) removeCheckpointOnWarm(ctx context.Context, pod *corev1.Pod } return nil } + +// claimOwnerAlive reports whether the pod holding a capture claim still +// exists and is not terminating. Errors other than NotFound count as alive: +// a transient API failure must not let two pods capture at once. +func (r *Reconciler) claimOwnerAlive(ctx context.Context, owner, ownerUID string) bool { + ns, name, ok := strings.Cut(owner, "/") + if !ok || r.KubeClient == nil { + return true + } + p, err := r.KubeClient.CoreV1().Pods(ns).Get(ctx, name, metav1.GetOptions{}) + if err != nil { + return !apierrors.IsNotFound(err) + } + // Inference pods have deterministic names: a later reconciliation + // can recreate the same name after the claimant died. A different + // UID under the same name is a replacement, not the owner. + if ownerUID != "" && string(p.UID) != ownerUID { + return false + } + if p.DeletionTimestamp != nil { + return false + } + return p.Status.Phase != corev1.PodSucceeded && p.Status.Phase != corev1.PodFailed +} diff --git a/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/reconciler_test.go b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/reconciler_test.go index ee7e9f2a7b..992f4ed889 100644 --- a/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/reconciler_test.go +++ b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/reconciler_test.go @@ -709,7 +709,7 @@ func TestRecordFailurePreservesCapturedAt(t *testing.T) { AttemptCount: 1, } _ = r.recordFailure(context.Background(), "fv-failtest", "completed_failed", prev, - fmtErrorf("simulated"), logrus.NewEntry(logrus.New())) + fmtErrorf("simulated"), claimToken{}, logrus.NewEntry(logrus.New())) got, err := dyn.Resource(CFSResource).Get(context.Background(), "fv-failtest", metav1.GetOptions{}) if err != nil { @@ -757,7 +757,7 @@ func TestWriteStatusPreservesUnmanagedKeys(t *testing.T) { if err := writeStatus(context.Background(), dyn, "fv-conditions", statusUpdate{ CheckpointHash: "deadbeef", LocalCacheState: nvsnapv1alpha1.LocalCacheStateWarm, - }); err != nil { + }, claimToken{}); err != nil { t.Fatalf("writeStatus: %v", err) } @@ -811,7 +811,7 @@ func TestRecordFailureReturnsNil_NoRequeue(t *testing.T) { r := &Reconciler{DynClient: dyn} got := r.recordFailure(context.Background(), "fv-noretry", "completed_failed", - cfsStatus{}, fmtErrorf("simulated capture failure"), logrus.NewEntry(logrus.New())) + cfsStatus{}, fmtErrorf("simulated capture failure"), claimToken{}, logrus.NewEntry(logrus.New())) if got != nil { t.Errorf("recordFailure returned %v; want nil so controller-runtime DOES NOT requeue (nvsnap-h100-a 2026-06-03 retry storm)", got) } @@ -834,7 +834,7 @@ func TestRecordFailureBumpsCounterByReason(t *testing.T) { before := testutil.ToFloat64(checkpointAttemptFailures.WithLabelValues(reason)) err := r.recordFailure(context.Background(), "fv-"+reason, reason, - cfsStatus{}, fmtErrorf("simulated"), logrus.NewEntry(logrus.New())) + cfsStatus{}, fmtErrorf("simulated"), claimToken{}, logrus.NewEntry(logrus.New())) if err != nil { t.Fatalf("recordFailure returned %v; must be nil", err) } diff --git a/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/state.go b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/state.go index bd16d450a0..fd449e8a26 100644 --- a/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/state.go +++ b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/state.go @@ -18,6 +18,7 @@ package reconciler import ( "context" + "errors" "fmt" "time" @@ -63,6 +64,11 @@ type cfsStatus struct { // crash recovery. Both zero-valued when not Capturing. CaptureOwner string CaptureLeaseExpiry *metav1.Time + // CaptureOwnerUID is the UID of the CaptureOwner pod. Inference pods + // have deterministic names (request name plus instance index), so a + // replacement pod can carry the same namespace/name as a dead + // claimant; the UID tells them apart. + CaptureOwnerUID string // ColdStartPioneer / ColdStartPioneerExpiry back the serialized-herd // cold-start pioneer election. ColdStartPioneer is the namespace/name // of the ICMSRequest that won the right to cold-start + capture while @@ -142,6 +148,7 @@ func readStatus(cfs *unstructured.Unstructured) cfsStatus { } } out.CaptureOwner, _, _ = unstructured.NestedString(cfs.Object, "status", "captureOwner") + out.CaptureOwnerUID, _, _ = unstructured.NestedString(cfs.Object, "status", "captureOwnerUID") if rfc, found, _ := unstructured.NestedString(cfs.Object, "status", "captureLeaseExpiry"); found { if t, err := time.Parse(time.RFC3339, rfc); err == nil { mt := metav1.NewTime(t) @@ -243,11 +250,71 @@ type statusUpdate struct { // /status, and the alternative — failing every status write — // would break Hook B entirely. Tracking removal: nvca-nvsnap // status-subresource-registration follow-up. -func writeStatus(ctx context.Context, dc dynamic.Interface, fvID string, upd statusUpdate) error { +// claimToken identifies the reconcile that holds the capture claim: the +// owner pod's namespace/name and UID. Terminal status writes carry it so +// a reconcile whose claim was taken over (its pod began terminating and +// a peer stole the claim while it was still polling) cannot release or +// overwrite the new owner's claim. The zero token writes unfenced, for +// paths that never held a claim (recovery before claiming, the sweep). +type claimToken struct { + Owner string + UID string + // ExpectUnclaimed fences a write from a path that holds no claim but + // observed none on the object (the sweep): the write is rejected if a + // claim appeared in between, so a fresh capture is never released by + // a recovery that read stale state. + ExpectUnclaimed bool + // ObservedResourceVersion, when set, rejects the write if the object + // changed at all since the writer observed it. Owner and unclaimed + // checks compare states and cannot see a claim that opened and closed + // in between, or a lease the same pod refreshed; the version can. + ObservedResourceVersion string +} + +func (t claimToken) empty() bool { + return t.Owner == "" && !t.ExpectUnclaimed && t.ObservedResourceVersion == "" +} + +// ErrClaimSuperseded is returned by writeStatus when the caller's claim +// token no longer matches the claim on the object: another reconcile owns +// the capture now and this one must not touch the status. +var ErrClaimSuperseded = errors.New("capture claim superseded by another owner") + +// supersedes reports whether the claim recorded on the object belongs to +// someone other than token. An object with no claim never supersedes. +func (st cfsStatus) supersedes(token claimToken) bool { + if token.empty() { + return false + } + if token.ExpectUnclaimed { + return st.CaptureOwner != "" + } + if token.Owner == "" { + return false // version-only fence, checked by the caller + } + if st.CaptureOwner == "" { + // The writer's claim is gone: another terminal write (a recovery + // or a takeover that finished) already closed this capture, so + // the writer's outcome must not replace it. + return true + } + if st.CaptureOwner != token.Owner { + return true + } + return st.CaptureOwnerUID != "" && token.UID != "" && st.CaptureOwnerUID != token.UID +} + +func writeStatus(ctx context.Context, dc dynamic.Interface, fvID string, upd statusUpdate, token claimToken) error { cur, err := dc.Resource(CFSResource).Get(ctx, fvID, metav1.GetOptions{}) if err != nil { return fmt.Errorf("get NvSnapFunctionState %s: %w", fvID, err) } + if token.ObservedResourceVersion != "" && cur.GetResourceVersion() != token.ObservedResourceVersion { + return ErrClaimSuperseded + } + if readStatus(cur).supersedes(token) { + return ErrClaimSuperseded + } // Read the existing status map and patch only the keys we manage. // SetNestedField(...,"status") would otherwise replace the whole // subtree, silently wiping unmanaged keys like `conditions` that @@ -274,6 +341,7 @@ func writeStatus(ctx context.Context, dc dynamic.Interface, fvID string, upd sta // would either pin the function until lease expiry or let readStatus // mis-report a stale owner. delete(status, "captureOwner") + delete(status, "captureOwnerUID") delete(status, "captureLeaseExpiry") // Release the cold-start pioneer claim too (serialized-herd). A // terminal status write means the pioneer's cold-start + capture is @@ -331,6 +399,22 @@ func writeStatus(ctx context.Context, dc dynamic.Interface, fvID string, upd sta // wall-clock deadline to stamp; now is injected for deterministic // tests. func tryClaimCapture(ctx context.Context, dc dynamic.Interface, fvID, owner string, leaseExpiry, now time.Time) (bool, error) { + return tryClaimCaptureLive(ctx, dc, fvID, owner, "", leaseExpiry, now, nil) +} + +// OwnerAliveFunc reports whether the pod named "namespace/name" that holds a +// capture claim still exists and is not terminating. nil means "unknown", +// which keeps the claim. +type OwnerAliveFunc func(ctx context.Context, owner, ownerUID string) bool + +// tryClaimCaptureLive is tryClaimCapture with a liveness check on the +// current owner. A live lease normally protects an in-flight capture from a +// second pioneer, but the lease outlives its owner: a pod that dies mid +// capture (evicted, replaced by a redeploy, its node agent restarted under +// it) leaves the version Capturing with nobody working on it until the +// lease expires, roughly 50 minutes. If the owner is gone, the claim is +// stealable at once. +func tryClaimCaptureLive(ctx context.Context, dc dynamic.Interface, fvID, owner, ownerUID string, leaseExpiry, now time.Time, ownerAlive OwnerAliveFunc) (bool, error) { cur, err := dc.Resource(CFSResource).Get(ctx, fvID, metav1.GetOptions{}) if err != nil { return false, fmt.Errorf("get NvSnapFunctionState %s: %w", fvID, err) @@ -342,10 +426,17 @@ func tryClaimCapture(ctx context.Context, dc dynamic.Interface, fvID, owner stri return false, nil case nvsnapv1alpha1.LocalCacheStateCapturing: leaseLive := st.CaptureLeaseExpiry != nil && now.Before(st.CaptureLeaseExpiry.Time) - if leaseLive && st.CaptureOwner != owner { - return false, nil // another pod holds a live claim + // The same namespace/name with a different UID is a replacement + // pod, not the claimant: treat it as foreign. + sameOwner := st.CaptureOwner == owner && (st.CaptureOwnerUID == "" || ownerUID == "" || st.CaptureOwnerUID == ownerUID) + if leaseLive && !sameOwner { + if ownerAlive == nil || ownerAlive(ctx, st.CaptureOwner, st.CaptureOwnerUID) { + return false, nil // another pod holds a live claim + } + // The owner is gone: nothing will finish or release this + // claim. Steal it now rather than at lease expiry. } - // expired lease (steal) or our own claim (re-entrant) → fall through + // expired lease, dead owner (steal) or our own claim (re-entrant) → fall through } // Patch only the keys we manage; preserve checkpointHash/attemptCount/etc. @@ -355,6 +446,11 @@ func tryClaimCapture(ctx context.Context, dc dynamic.Interface, fvID, owner stri } status["localCacheState"] = string(nvsnapv1alpha1.LocalCacheStateCapturing) status["captureOwner"] = owner + if ownerUID != "" { + status["captureOwnerUID"] = ownerUID + } else { + delete(status, "captureOwnerUID") + } status["captureLeaseExpiry"] = leaseExpiry.UTC().Format(time.RFC3339) if err := unstructured.SetNestedField(cur.Object, status, "status"); err != nil { return false, fmt.Errorf("set status: %w", err) @@ -485,7 +581,7 @@ func TryClaimColdStartPioneer(ctx context.Context, dc dynamic.Interface, fvID, o // Distinct from recordFailure on purpose: this does not set LastError or // increment AttemptCount, so a path whose contract is "treat it as if the // capture never happened" leaves no CFS-level error surface behind. -func releaseCaptureClaim(ctx context.Context, dc dynamic.Interface, fvID string, prev cfsStatus) error { +func releaseCaptureClaim(ctx context.Context, dc dynamic.Interface, fvID string, prev cfsStatus, token claimToken) error { upd := statusUpdate{ CheckpointHash: prev.CheckpointHash, CapturedHere: prev.CapturedHere, @@ -499,5 +595,5 @@ func releaseCaptureClaim(ctx context.Context, dc dynamic.Interface, fvID string, if prev.LastAttemptAt != nil { upd.LastAttemptAt = prev.LastAttemptAt.Time } - return writeStatus(ctx, dc, fvID, upd) + return writeStatus(ctx, dc, fvID, upd, token) } diff --git a/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/sweep.go b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/sweep.go index f239a2c0f6..3bc7b2278c 100644 --- a/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/sweep.go +++ b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/sweep.go @@ -82,6 +82,25 @@ func (r *Reconciler) SweepOnce(ctx context.Context) { continue // no usable capture yet — leave for a future tick } + // An in-flight capture is not recovered over. A live claim whose + // owner pod is alive is left alone (unknown liveness counts as + // alive). Every sweep write is fenced on the resourceVersion it + // listed: any change since then (a claim that opened, or opened + // and closed, a lease the owner refreshed, a result that landed) + // makes the write fail and the next tick re-evaluate. The owner + // and unclaimed checks stay as a second line for objects whose + // version the fake or an old server does not track. + token := claimToken{ExpectUnclaimed: true, ObservedResourceVersion: cfs.GetResourceVersion()} + if st.CaptureOwner != "" { + leaseLive := st.CaptureLeaseExpiry != nil && time.Now().Before(st.CaptureLeaseExpiry.Time) + if leaseLive && r.claimOwnerAlive(ctx, st.CaptureOwner, st.CaptureOwnerUID) { + log.WithFields(logrus.Fields{"functionVersionID": fvID, "owner": st.CaptureOwner}). + Debug("sweep: capture in flight with a live owner; not recovering over it") + continue + } + token = claimToken{Owner: st.CaptureOwner, UID: st.CaptureOwnerUID, ObservedResourceVersion: cfs.GetResourceVersion()} + } + now := time.Now() if err := writeStatus(ctx, r.DynClient, fvID, statusUpdate{ CheckpointHash: hash, @@ -91,7 +110,7 @@ func (r *Reconciler) SweepOnce(ctx context.Context) { AttemptCount: 0, LastError: "", LastAttemptAt: now, - }); err != nil { + }, token); err != nil { log.WithError(err).WithField("functionVersionID", fvID). Warn("sweep: writeStatus Warm failed; will retry next tick") continue diff --git a/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/sweep_test.go b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/sweep_test.go index 69d9011247..cf64e281ce 100644 --- a/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/sweep_test.go +++ b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/sweep_test.go @@ -23,10 +23,17 @@ package reconciler import ( "context" "encoding/json" + "errors" + nvsnapv1alpha1 "github.com/NVIDIA/nvcf/src/compute-plane-services/nvca/pkg/apis/nvsnap/v1alpha1" + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/runtime" + k8sfake "k8s.io/client-go/kubernetes/fake" + k8stesting "k8s.io/client-go/testing" "net/http" "net/http/httptest" "strings" "testing" + "time" "github.com/sirupsen/logrus" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -228,3 +235,94 @@ func TestWriteWorkloadLookup_PersistsAndIdempotent(t *testing.T) { t.Errorf("empty imageRef should be a no-op, got %v", err) } } + +// capturingWithLookup is cfsWithLookup with an in-flight capture claim. +func capturingWithLookup(fvID, image, owner, uid string, lease time.Time) *unstructured.Unstructured { + obj := cfsWithLookup(fvID, image, "", nil) + status, _, _ := unstructured.NestedMap(obj.Object, "status") + if status == nil { + status = map[string]any{} + } + status["localCacheState"] = string(nvsnapv1alpha1.LocalCacheStateCapturing) + status["captureOwner"] = owner + status["captureOwnerUID"] = uid + status["captureLeaseExpiry"] = lease.UTC().Format(time.RFC3339) + _ = unstructured.SetNestedField(obj.Object, status, "status") + return obj +} + +// The sweep never recovers over an in-flight capture whose owner pod is +// alive: the owner's own terminal write decides that capture. +func TestSweep_LeavesLiveCaptureAlone(t *testing.T) { + const fvID, image, hash = "fv-sweep-live", "ngc.io/fn:1", "deadbeefcafe1002" + srv := usableCaptureServer(image, hash) + defer srv.Close() + dyn := newFakeDynamic(capturingWithLookup(fvID, image, "ns1/podA", "uid-a", time.Now().Add(30*time.Minute))) + r := sweepReconciler(srv, dyn) + r.KubeClient = k8sfake.NewSimpleClientset(&corev1.Pod{ + ObjectMeta: metav1.ObjectMeta{Name: "podA", Namespace: "ns1", UID: "uid-a"}, + Status: corev1.PodStatus{Phase: corev1.PodRunning}, + }) + r.SweepOnce(context.Background()) + got, _ := dyn.Resource(CFSResource).Get(context.Background(), fvID, metav1.GetOptions{}) + if s := readStatus(got); s.LocalCacheState != nvsnapv1alpha1.LocalCacheStateCapturing || s.CaptureOwner != "ns1/podA" { + t.Errorf("live capture must be untouched by the sweep: %+v", s) + } +} + +// A claim whose owner is gone (or a replacement pod under the same name) is +// taken over by the sweep, fenced on the claim it observed. +func TestSweep_RecoversOverDeadOwner(t *testing.T) { + const fvID, image, hash = "fv-sweep-dead", "ngc.io/fn:1", "deadbeefcafe1003" + srv := usableCaptureServer(image, hash) + defer srv.Close() + dyn := newFakeDynamic(capturingWithLookup(fvID, image, "ns1/podA", "uid-old", time.Now().Add(30*time.Minute))) + r := sweepReconciler(srv, dyn) + r.KubeClient = k8sfake.NewSimpleClientset(&corev1.Pod{ // same name, new UID: a replacement + ObjectMeta: metav1.ObjectMeta{Name: "podA", Namespace: "ns1", UID: "uid-new"}, + Status: corev1.PodStatus{Phase: corev1.PodRunning}, + }) + r.SweepOnce(context.Background()) + got, _ := dyn.Resource(CFSResource).Get(context.Background(), fvID, metav1.GetOptions{}) + s := readStatus(got) + if s.LocalCacheState != nvsnapv1alpha1.LocalCacheStateWarm || s.CaptureOwner != "" || s.CheckpointHash != hash { + t.Errorf("dead-owner capture must be recovered to Warm: %+v", s) + } + // The former owner's late failure write is rejected: its claim is closed. + err := writeStatus(context.Background(), dyn, fvID, statusUpdate{LocalCacheState: nvsnapv1alpha1.LocalCacheStateFailed, LastError: "late"}, claimToken{Owner: "ns1/podA", UID: "uid-old"}) + if !errors.Is(err, ErrClaimSuperseded) { + t.Fatalf("late write from the recovered-over owner must be superseded: %v", err) + } + if s := readStatus(mustGet(t, dyn, fvID)); s.LocalCacheState != nvsnapv1alpha1.LocalCacheStateWarm { + t.Errorf("recovered Warm must stand: %+v", s) + } +} + +// The sweep fences on the version it listed: an object that changed after +// the list (here: the owner refreshed its lease) is not written. +func TestSweep_StaleObservationIsNotWritten(t *testing.T) { + const fvID, image, hash = "fv-sweep-stale", "ngc.io/fn:1", "deadbeefcafe1004" + srv := usableCaptureServer(image, hash) + defer srv.Close() + obj := capturingWithLookup(fvID, image, "ns1/podA", "uid-a", time.Now().Add(-time.Minute)) // expired: recoverable + obj.SetResourceVersion("7") + dyn := newFakeDynamic(obj) + // Between the sweep's list and its write the owner refreshes the lease: + // simulate by bumping the version the Get returns. + dyn.PrependReactor("get", "nvsnapfunctionstates", func(action k8stesting.Action) (bool, runtime.Object, error) { + cur, err := dyn.Tracker().Get(CFSResource, "", fvID) + if err != nil { + return true, nil, err + } + u := cur.(*unstructured.Unstructured).DeepCopy() + u.SetResourceVersion("8") + return true, u, nil + }) + r := sweepReconciler(srv, dyn) + r.KubeClient = k8sfake.NewSimpleClientset() + r.SweepOnce(context.Background()) + got, _ := dyn.Tracker().Get(CFSResource, "", fvID) + if s := readStatus(got.(*unstructured.Unstructured)); s.LocalCacheState != nvsnapv1alpha1.LocalCacheStateCapturing { + t.Errorf("an object that changed since the list must not be written: %+v", s) + } +}