Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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",
],
)
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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()
Expand All @@ -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{})
Expand Down Expand Up @@ -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)
Expand All @@ -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)
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ import (
"errors"
"fmt"
"net/http"
"strings"
"sync"
"time"

Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)

Expand All @@ -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
Expand Down Expand Up @@ -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")
Expand Down Expand Up @@ -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":
Expand All @@ -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,
Expand All @@ -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
}
Expand Down Expand Up @@ -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,
Expand All @@ -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()
Expand Down Expand Up @@ -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{})
Comment thread
coderabbitai[bot] marked this conversation as resolved.
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
}
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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)
}

Expand Down Expand Up @@ -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)
}
Expand All @@ -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)
}
Expand Down
Loading
Loading