From 51caf14ac27a8ba9dc0d6b8877b8d50ceb6ef620 Mon Sep 17 00:00:00 2001 From: Balaji Ganesan Date: Mon, 28 Sep 2026 17:21:58 -0700 Subject: [PATCH 1/5] fix(nvca): let a pod steal a capture claim whose owner is gone The capture-once claim in NvSnapFunctionState protects an in-flight capture with a lease of about 50 minutes, and only lease expiry made it stealable. A pod that dies mid-capture leaves the function version in Capturing with nobody working on it; seen on dev1 2026-09-28 when the nvsnap agent DaemonSet rolled while the checkpoint request was in flight and the pod was then replaced. The new pod reconciled, saw the live claim, and stood down for the full lease. tryClaimCapture now takes an owner liveness check. When the claim is live but its owner pod no longer exists or is terminating, the claim is taken over at once. Unknown liveness (no check, or an API error other than NotFound) keeps the claim, so a transient failure cannot let two pods capture at the same time. Co-Authored-By: Balaji Ganesan --- .../nvsnap/reconciler/captureonce_test.go | 64 +++++++++++++++++++ .../pkg/nvca/nvsnap/reconciler/reconciler.go | 21 +++++- .../nvca/pkg/nvca/nvsnap/reconciler/state.go | 24 ++++++- 3 files changed, 106 insertions(+), 3 deletions(-) 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..4f6541fc7a 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 @@ -150,6 +150,40 @@ 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) bool { return false } + alive := func(context.Context, string) bool { return true } + if claimed, _ := tryClaimCaptureLive(ctx, dyn, fvID, "ns1/podB", 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", lease, now, nil); claimed { + t.Fatal("unknown liveness (nil) must keep the claim") + } + claimed, err := tryClaimCaptureLive(ctx, dyn, fvID, "ns1/podB", 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" { + t.Errorf("owner after steal = %q, want ns1/podB", st.CaptureOwner) + } +} + func TestWriteStatusReleasesClaim(t *testing.T) { const fvID = "fv-release" ctx := context.Background() @@ -203,6 +237,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 +254,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/reconciler.go b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/reconciler.go index 85d99d1259..a0fd2ff7fd 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" @@ -414,7 +415,7 @@ 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) + claimed, err := tryClaimCaptureLive(ctx, r.DynClient, fvID, owner, claimNow.Add(r.CaptureLeaseTTL), claimNow, r.claimOwnerAlive) if err != nil { return fmt.Errorf("capture-once claim for %s: %w", fvID, err) } @@ -766,3 +767,21 @@ 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 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) + } + 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/state.go b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/state.go index bd16d450a0..9f15e90e6b 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 @@ -331,6 +331,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 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 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) @@ -343,9 +359,13 @@ func tryClaimCapture(ctx context.Context, dc dynamic.Interface, fvID, owner stri 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 + if ownerAlive == nil || ownerAlive(ctx, st.CaptureOwner) { + 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. From 0b206041aa9e0701692db1fe9c9c440354da9acd Mon Sep 17 00:00:00 2001 From: Balaji Ganesan Date: Tue, 29 Sep 2026 11:55:55 -0700 Subject: [PATCH 2/5] fix(nvca): compare the claimant pod UID and fence terminal claim writes Two gaps in the capture-once claim after the dead-owner steal: Inference pods have deterministic names, so a later reconciliation can recreate a dead claimant's namespace/name. claimOwnerAlive found the replacement and treated the original owner as alive, pinning the function version until lease expiry. The claim now records the owner pod's UID (status.captureOwnerUID); a same-name pod with another UID reads as owner gone, and a re-entrant claim must match both. A reconcile whose pod began terminating could still be polling its checkpoint when a peer stole the claim. Its terminal writeStatus then cleared the new owner's claim or replaced the new owner's result. Every terminal write now carries the writer's claim token and is rejected with ErrClaimSuperseded when the object's claim belongs to someone else; the superseded reconcile logs and returns without requeue. Writes from paths that never held a claim (recovery before claiming, the sweep) stay unfenced. Relates to #2154 Co-Authored-By: Balaji Ganesan --- .../v1alpha1/nvsnapfunctionstate_types.go | 6 ++ .../nvsnap/reconciler/captureonce_test.go | 78 +++++++++++++++++-- .../reconciler/coldstart_pioneer_test.go | 2 +- .../pkg/nvca/nvsnap/reconciler/reconciler.go | 41 +++++++--- .../nvca/nvsnap/reconciler/reconciler_test.go | 8 +- .../nvca/pkg/nvca/nvsnap/reconciler/state.go | 65 ++++++++++++++-- .../nvca/pkg/nvca/nvsnap/reconciler/sweep.go | 2 +- 7 files changed, 168 insertions(+), 34 deletions(-) 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/captureonce_test.go b/src/compute-plane-services/nvca/pkg/nvca/nvsnap/reconciler/captureonce_test.go index 4f6541fc7a..b5493248f3 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" @@ -163,15 +166,15 @@ func TestTryClaimCapture_DeadOwnerIsStealable(t *testing.T) { 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) bool { return false } - alive := func(context.Context, string) bool { return true } - if claimed, _ := tryClaimCaptureLive(ctx, dyn, fvID, "ns1/podB", lease, now, alive); claimed { + 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", lease, now, nil); claimed { + 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", lease, now, dead) + 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) } @@ -179,8 +182,67 @@ func TestTryClaimCapture_DeadOwnerIsStealable(t *testing.T) { if err != nil { t.Fatal(err) } - if st := readStatus(cur); st.CaptureOwner != "ns1/podB" { - t.Errorf("owner after steal = %q, want ns1/podB", st.CaptureOwner) + 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) } } @@ -193,7 +255,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{}) 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 a0fd2ff7fd..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 @@ -358,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) @@ -415,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 := tryClaimCaptureLive(ctx, r.DynClient, fvID, owner, claimNow.Add(r.CaptureLeaseTTL), claimNow, r.claimOwnerAlive) + 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) } @@ -445,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) @@ -454,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 @@ -499,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") @@ -540,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": @@ -554,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, @@ -578,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 } @@ -706,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, @@ -723,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() @@ -771,7 +782,7 @@ func (r *Reconciler) removeCheckpointOnWarm(ctx context.Context, pod *corev1.Pod // 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 string) bool { +func (r *Reconciler) claimOwnerAlive(ctx context.Context, owner, ownerUID string) bool { ns, name, ok := strings.Cut(owner, "/") if !ok || r.KubeClient == nil { return true @@ -780,6 +791,12 @@ func (r *Reconciler) claimOwnerAlive(ctx context.Context, owner string) bool { 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 } 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 9f15e90e6b..542a5dd352 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,44 @@ 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 +} + +func (t claimToken) empty() bool { return t.Owner == "" } + +// 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() || st.CaptureOwner == "" { + return false + } + 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 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 +314,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,13 +372,13 @@ 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) + 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 string) bool +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 @@ -346,7 +387,7 @@ type OwnerAliveFunc func(ctx context.Context, owner string) bool // 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 string, leaseExpiry, now time.Time, ownerAlive OwnerAliveFunc) (bool, error) { +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) @@ -358,8 +399,11 @@ func tryClaimCaptureLive(ctx context.Context, dc dynamic.Interface, fvID, owner return false, nil case nvsnapv1alpha1.LocalCacheStateCapturing: leaseLive := st.CaptureLeaseExpiry != nil && now.Before(st.CaptureLeaseExpiry.Time) - if leaseLive && st.CaptureOwner != owner { - if ownerAlive == nil || ownerAlive(ctx, st.CaptureOwner) { + // 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 @@ -375,6 +419,11 @@ func tryClaimCaptureLive(ctx context.Context, dc dynamic.Interface, fvID, owner } 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) @@ -505,7 +554,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, @@ -519,5 +568,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..e23d1b2e82 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 @@ -91,7 +91,7 @@ func (r *Reconciler) SweepOnce(ctx context.Context) { AttemptCount: 0, LastError: "", LastAttemptAt: now, - }); err != nil { + }, claimToken{}); err != nil { log.WithError(err).WithField("functionVersionID", fvID). Warn("sweep: writeStatus Warm failed; will retry next tick") continue From 1f0e3f1054dbbc0c649c54f1da13a8a08b2821dd Mon Sep 17 00:00:00 2001 From: Balaji Ganesan Date: Tue, 29 Sep 2026 12:05:20 -0700 Subject: [PATCH 3/5] fix(nvca): fence the durable-warm sweep against an in-flight capture The sweep flipped a Capturing function version to Warm with an unfenced write: it cleared the live owner's claim, and a later terminal write from that owner (which no longer matched any claim) could overwrite the recovered Warm status. The sweep now leaves a live claim whose owner pod is alive alone, takes over a dead or expired one fenced on the claim it observed, and writes with an expect-unclaimed token otherwise, so a claim that appeared after the list is never released. A terminal write whose claim has already been closed is superseded rather than applied. Relates to #2154 Co-Authored-By: Balaji Ganesan --- .../nvsnap/reconciler/captureonce_test.go | 14 ++++ .../nvca/pkg/nvca/nvsnap/reconciler/state.go | 18 ++++- .../nvca/pkg/nvca/nvsnap/reconciler/sweep.go | 17 ++++- .../pkg/nvca/nvsnap/reconciler/sweep_test.go | 67 +++++++++++++++++++ 4 files changed, 113 insertions(+), 3 deletions(-) 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 b5493248f3..0a6c143146 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 @@ -244,6 +244,20 @@ func TestWriteStatusRejectsSupersededClaim(t *testing.T) { 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) + } } func TestWriteStatusReleasesClaim(t *testing.T) { 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 542a5dd352..15155eda19 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 @@ -259,9 +259,14 @@ type statusUpdate struct { 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 } -func (t claimToken) empty() bool { return t.Owner == "" } +func (t claimToken) empty() bool { return t.Owner == "" && !t.ExpectUnclaimed } // ErrClaimSuperseded is returned by writeStatus when the caller's claim // token no longer matches the claim on the object: another reconcile owns @@ -271,9 +276,18 @@ 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() || st.CaptureOwner == "" { + if token.empty() { return false } + if token.ExpectUnclaimed { + return st.CaptureOwner != "" + } + 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 } 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 e23d1b2e82..8b3f80794d 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,21 @@ 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); a dead or expired claim is taken over, fenced on the + // claim as observed so a newer claimant cannot be released. + token := claimToken{ExpectUnclaimed: true} + 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} + } + now := time.Now() if err := writeStatus(ctx, r.DynClient, fvID, statusUpdate{ CheckpointHash: hash, @@ -91,7 +106,7 @@ func (r *Reconciler) SweepOnce(ctx context.Context) { AttemptCount: 0, LastError: "", LastAttemptAt: now, - }, claimToken{}); 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..37ccc6a4d0 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,15 @@ 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" + k8sfake "k8s.io/client-go/kubernetes/fake" "net/http" "net/http/httptest" "strings" "testing" + "time" "github.com/sirupsen/logrus" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -228,3 +233,65 @@ 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) + } +} From 33ffb787fc96187aba0066ec14d47df34186dd5c Mon Sep 17 00:00:00 2001 From: Balaji Ganesan Date: Tue, 29 Sep 2026 12:14:26 -0700 Subject: [PATCH 4/5] fix(nvca): fence sweep writes on the observed resourceVersion State comparison cannot see a claim that opened and closed between the sweep's list and its write, or a lease the same owner refreshed. Every sweep write now carries the resourceVersion it listed and is rejected when the object changed at all since; the next tick re-evaluates from fresh state. Relates to #2154 Co-Authored-By: Balaji Ganesan --- .../nvsnap/reconciler/captureonce_test.go | 12 +++++++ .../nvca/pkg/nvca/nvsnap/reconciler/state.go | 15 ++++++++- .../nvca/pkg/nvca/nvsnap/reconciler/sweep.go | 12 ++++--- .../pkg/nvca/nvsnap/reconciler/sweep_test.go | 31 +++++++++++++++++++ 4 files changed, 65 insertions(+), 5 deletions(-) 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 0a6c143146..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 @@ -258,6 +258,18 @@ func TestWriteStatusRejectsSupersededClaim(t *testing.T) { 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) { 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 15155eda19..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 @@ -264,9 +264,16 @@ type claimToken struct { // 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 } +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 @@ -282,6 +289,9 @@ func (st cfsStatus) supersedes(token claimToken) bool { 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 @@ -299,6 +309,9 @@ func writeStatus(ctx context.Context, dc dynamic.Interface, fvID string, upd sta 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 } 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 8b3f80794d..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 @@ -84,9 +84,13 @@ func (r *Reconciler) SweepOnce(ctx context.Context) { // An in-flight capture is not recovered over. A live claim whose // owner pod is alive is left alone (unknown liveness counts as - // alive); a dead or expired claim is taken over, fenced on the - // claim as observed so a newer claimant cannot be released. - token := claimToken{ExpectUnclaimed: true} + // 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) { @@ -94,7 +98,7 @@ func (r *Reconciler) SweepOnce(ctx context.Context) { Debug("sweep: capture in flight with a live owner; not recovering over it") continue } - token = claimToken{Owner: st.CaptureOwner, UID: st.CaptureOwnerUID} + token = claimToken{Owner: st.CaptureOwner, UID: st.CaptureOwnerUID, ObservedResourceVersion: cfs.GetResourceVersion()} } now := time.Now() 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 37ccc6a4d0..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 @@ -26,7 +26,9 @@ import ( "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" @@ -295,3 +297,32 @@ func TestSweep_RecoversOverDeadOwner(t *testing.T) { 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) + } +} From ce8a314228fdafb86a411cbe9822360ec121a3c4 Mon Sep 17 00:00:00 2001 From: Balaji Ganesan Date: Tue, 29 Sep 2026 12:45:22 -0700 Subject: [PATCH 5/5] build(nvca): add client-go/testing to the reconciler test target The sweep test uses a fake client reactor. Co-Authored-By: Balaji Ganesan --- .../nvca/pkg/nvca/nvsnap/reconciler/BUILD.bazel | 1 + 1 file changed, 1 insertion(+) 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", ], )