From 3990b941b0f70f1b389d14db5102b4357c105e65 Mon Sep 17 00:00:00 2001 From: Max Thompson Date: Thu, 17 Sep 2026 12:30:42 -0700 Subject: [PATCH] atelet: undo only the system-info registration a failed Run or Restore made Register supersedes an existing entry for the same actor UID, but the failure paths in Run and Restore still undid their work with Deregister(actorUID), which removes whatever entry the map holds. When a retried Restore had already replaced the first attempt's entry, the first attempt's error path deleted the retry's live registration, and the running sandbox silently stopped receiving trust bundle refreshes. Register now returns a handle scoped to the entry it created. Undo removes that entry only if it is still the actor's current one and marks it stale, so a failed attempt cannot drop a later registration. A Register that fails while writing its volumes removes its own entry before returning. Deregister keeps its by-UID semantics for Terminate and Checkpoint, where the caller means whatever is registered. --- cmd/atelet/main.go | 17 ++-- cmd/atelet/systeminfovolume.go | 44 ++++++++-- cmd/atelet/systeminfovolume_test.go | 125 ++++++++++++++++++++++++---- 3 files changed, 159 insertions(+), 27 deletions(-) diff --git a/cmd/atelet/main.go b/cmd/atelet/main.go index 6155e49aa0..2fde126a77 100644 --- a/cmd/atelet/main.go +++ b/cmd/atelet/main.go @@ -504,14 +504,15 @@ func (s *AteomHerder) Run(ctx context.Context, req *ateletpb.RunRequest) (resp * return nil, fmt.Errorf("while recording sandbox assets: %w", err) } + systemInfoReg, err := s.systemInfoVolumes.Register(actorUID, actorRef, systemInfoVolumesFor(actorUID, req.GetSpec())) + if err != nil { + return nil, err + } defer func() { if err != nil { - s.systemInfoVolumes.Deregister(actorUID) + systemInfoReg.Undo() } }() - if err := s.systemInfoVolumes.Register(actorUID, actorRef, systemInfoVolumesFor(actorUID, req.GetSpec())); err != nil { - return nil, err - } if err := s.prepareOCIBundles(ctx, actorUID, actorRef, req.GetSpec(), sandboxRec.PauseImage, req.GetTargetAteomUid(), ); err != nil { @@ -1113,10 +1114,12 @@ func (s *AteomHerder) Restore(ctx context.Context, req *ateletpb.RestoreRequest) runtimeRec = goldenRec } - // Undo the Register if the restore fails. + // Undo the Register if the restore fails. Written by the prep leg below and + // read only after g.Wait returns. + var systemInfoReg *registration defer func() { if err != nil { - s.systemInfoVolumes.Deregister(actorUID) + systemInfoReg.Undo() } }() @@ -1188,7 +1191,7 @@ func (s *AteomHerder) Restore(ctx context.Context, req *ateletpb.RestoreRequest) prepFailedPhase = ateattr.SnapshotPhaseSandboxAssets return ateerrors.CrashIfReason(ctx, err, ateerrors.ReasonFailedGetExternalObject, ateerrors.ReasonInvalidObjectURL, ateerrors.ReasonTerminalFileSystemError, ateerrors.ReasonInvalidSandboxAsset) } - if err = s.systemInfoVolumes.Register(actorUID, actorRef, systemInfoVolumesFor(actorUID, req.GetSpec())); err != nil { + if systemInfoReg, err = s.systemInfoVolumes.Register(actorUID, actorRef, systemInfoVolumesFor(actorUID, req.GetSpec())); err != nil { prepFailedPhase = ateattr.SnapshotPhaseOCIUnpack return err } diff --git a/cmd/atelet/systeminfovolume.go b/cmd/atelet/systeminfovolume.go index 08516d6a83..396ed0ff7d 100644 --- a/cmd/atelet/systeminfovolume.go +++ b/cmd/atelet/systeminfovolume.go @@ -107,8 +107,9 @@ func newSystemInfoVolumeRefresher(lister certlisters.ClusterTrustBundleLister, i // Register records actorUID's system-info volumes and writes their contents // from current cluster state. If actorUID is already registered (for example // after a worker pod crash left a stale entry without Terminate), the previous -// registration is superseded. -func (r *systemInfoVolumeRefresher) Register(actorUID string, ref resources.ActorRef, volumes []*systemInfoVolume) error { +// registration is superseded. The returned registration undoes this call +// alone; a Register that fails leaves no entry behind. +func (r *systemInfoVolumeRefresher) Register(actorUID string, ref resources.ActorRef, volumes []*systemInfoVolume) (*registration, error) { actor := ®isteredActor{uid: actorUID, ref: ref, volumes: volumes} // Held until the initial write finishes so a refresh cannot interleave. actor.mu.Lock() @@ -129,14 +130,45 @@ func (r *systemInfoVolumeRefresher) Register(actorUID string, ref resources.Acto for _, v := range volumes { if err := r.write(ref, actorUID, v); err != nil { - return fmt.Errorf("while populating system-info volume %q: %w", v.Name, err) + r.unlink(actor) + actor.stale = true + return nil, fmt.Errorf("while populating system-info volume %q: %w", v.Name, err) } } - return nil + return ®istration{r: r, actor: actor}, nil +} + +// registration is the entry one Register call created. +type registration struct { + r *systemInfoVolumeRefresher + actor *registeredActor +} + +// Undo drops the entry unless a later Register for the same actor has +// replaced it, and stops writes to its volumes. Undo on a nil registration +// is a no-op. +func (h *registration) Undo() { + if h == nil { + return + } + h.r.unlink(h.actor) + h.actor.mu.Lock() + h.actor.stale = true + h.actor.mu.Unlock() +} + +// unlink removes actor from the registry if it is still the entry for its UID. +func (r *systemInfoVolumeRefresher) unlink(actor *registeredActor) { + r.mu.Lock() + defer r.mu.Unlock() + if r.actors[actor.uid] == actor { + delete(r.actors, actor.uid) + } } -// Deregister drops actorUID's registration. After Deregister returns, no more -// system-info volumes will be written for the actor. +// Deregister drops actorUID's current registration, whichever Register call +// made it. After Deregister returns, no more system-info volumes will be +// written for the actor. func (r *systemInfoVolumeRefresher) Deregister(actorUID string) { r.mu.Lock() actor := r.actors[actorUID] diff --git a/cmd/atelet/systeminfovolume_test.go b/cmd/atelet/systeminfovolume_test.go index 3ff2953a2f..9f951ca854 100644 --- a/cmd/atelet/systeminfovolume_test.go +++ b/cmd/atelet/systeminfovolume_test.go @@ -98,16 +98,18 @@ func metadataVolumeSpec() *ateletpb.SystemInfoVolume { // registerTrustVolume registers actorUID with one volume projecting the // egress bundle at //system-info/trust/ca.pem. -func registerTrustVolume(t *testing.T, r *systemInfoVolumeRefresher, dir, actorUID string) { +func registerTrustVolume(t *testing.T, r *systemInfoVolumeRefresher, dir, actorUID string) *registration { t.Helper() vol := &systemInfoVolume{ Name: "trust", Root: filepath.Join(dir, actorUID, "system-info", "trust"), Spec: trustVolumeSpec("ca.pem"), } - if err := r.Register(actorUID, resources.ActorRef{Atespace: "team-a", Name: actorUID}, []*systemInfoVolume{vol}); err != nil { + reg, err := r.Register(actorUID, resources.ActorRef{Atespace: "team-a", Name: actorUID}, []*systemInfoVolume{vol}) + if err != nil { t.Fatalf("Register(%s): %v", actorUID, err) } + return reg } func readProjected(t *testing.T, dir, actorUID, volume, relPath string) string { @@ -228,7 +230,7 @@ func TestSystemInfoVolumeRefresher_RegisterEmptyStopsRefreshing(t *testing.T) { // The actor comes back under a spec with no system-info volumes: it stays // tracked, but its former volumes stop refreshing. r.Deregister("uid-1") - if err := r.Register("uid-1", resources.ActorRef{Atespace: "team-a", Name: "uid-1"}, nil); err != nil { + if _, err := r.Register("uid-1", resources.ActorRef{Atespace: "team-a", Name: "uid-1"}, nil); err != nil { t.Fatalf("Register(empty): %v", err) } if r.actors["uid-1"] == nil { @@ -360,7 +362,7 @@ func TestSystemInfoVolumeRefresher_RotationLeavesUnchangedFilesAlone(t *testing. spec := &ateletpb.SystemInfoVolume{ DataSources: append(metadataVolumeSpec().GetDataSources(), trustVolumeSpec("trust/ca.pem").GetDataSources()...), } - if err := r.Register("uid-1", resources.ActorRef{Atespace: "team-a", Name: "actor-1"}, []*systemInfoVolume{{Name: "vol1", Root: root, Spec: spec}}); err != nil { + if _, err := r.Register("uid-1", resources.ActorRef{Atespace: "team-a", Name: "actor-1"}, []*systemInfoVolume{{Name: "vol1", Root: root, Spec: spec}}); err != nil { t.Fatalf("Register: %v", err) } @@ -393,14 +395,14 @@ func TestSystemInfoVolumeRegister_WritesActorMetadata(t *testing.T) { } golden := resources.ActorRef{Atespace: "ate-e2e-probe", Name: "golden-actor"} - if err := r.Register("uid-golden", golden, []*systemInfoVolume{vol()}); err != nil { + if _, err := r.Register("uid-golden", golden, []*systemInfoVolume{vol()}); err != nil { t.Fatalf("Register: %v", err) } // Overwrite with a different actor, as happens when a snapshot taken from // one actor seeds another on resume: files must carry the new values. alpha := resources.ActorRef{Atespace: "ate-e2e-probe", Name: "probe-alpha"} - if err := r.Register("uid-alpha", alpha, []*systemInfoVolume{vol()}); err != nil { + if _, err := r.Register("uid-alpha", alpha, []*systemInfoVolume{vol()}); err != nil { t.Fatalf("Register (rewrite): %v", err) } @@ -440,7 +442,7 @@ func TestSystemInfoVolumeRegister_StableRealPaths(t *testing.T) { } golden := resources.ActorRef{Atespace: "ate-e2e-probe", Name: "golden-actor"} - if err := r.Register("uid-golden", golden, []*systemInfoVolume{vol()}); err != nil { + if _, err := r.Register("uid-golden", golden, []*systemInfoVolume{vol()}); err != nil { t.Fatalf("Register: %v", err) } @@ -465,7 +467,7 @@ func TestSystemInfoVolumeRegister_StableRealPaths(t *testing.T) { // Regenerate for a different actor, as a restore from a shared golden // snapshot does. alpha := resources.ActorRef{Atespace: "ate-e2e-probe", Name: "probe-alpha"} - if err := r.Register("uid-alpha", alpha, []*systemInfoVolume{vol()}); err != nil { + if _, err := r.Register("uid-alpha", alpha, []*systemInfoVolume{vol()}); err != nil { t.Fatalf("Register (rewrite): %v", err) } @@ -552,7 +554,7 @@ func TestSystemInfoVolumeRefresher_RegisterTwiceSupersedes(t *testing.T) { registerTrustVolume(t, r, dir, "uid-1") first := r.actors["uid-1"] - if err := r.Register("uid-1", resources.ActorRef{Atespace: "team-a", Name: "uid-1"}, nil); err != nil { + if _, err := r.Register("uid-1", resources.ActorRef{Atespace: "team-a", Name: "uid-1"}, nil); err != nil { t.Fatalf("second Register: %v", err) } if !first.stale { @@ -563,6 +565,95 @@ func TestSystemInfoVolumeRefresher_RegisterTwiceSupersedes(t *testing.T) { } } +// A failed Run or Restore undoes only the registration it made; a later +// Register for the same actor keeps refreshing. +func TestSystemInfoVolumeRefresher_UndoLeavesSupersedingRegistration(t *testing.T) { + ctx := context.Background() + certA, certB := string(testCertPEM(t)), string(testCertPEM(t)) + store := newCTBStore(t) + store.set(t, certA) + r := newSystemInfoVolumeRefresher(store.lister, nil) + dir1, dir2 := t.TempDir(), t.TempDir() + first := registerTrustVolume(t, r, dir1, "uid-1") + registerTrustVolume(t, r, dir2, "uid-1") + + first.Undo() + if r.actors["uid-1"] == nil { + t.Fatal("undoing a superseded registration dropped the current one") + } + store.set(t, certB) + r.refreshBundle(ctx, EgressTrustBundleName) + if got := readProjected(t, dir2, "uid-1", "trust", "ca.pem"); got != certB { + t.Errorf("current registration's file = %q, want the rotation applied", got) + } + if got := readProjected(t, dir1, "uid-1", "trust", "ca.pem"); got != certA { + t.Errorf("undone registration's file = %q, refreshed after Undo", got) + } +} + +func TestSystemInfoVolumeRefresher_UndoDropsCurrentRegistration(t *testing.T) { + ctx := context.Background() + certA, certB := string(testCertPEM(t)), string(testCertPEM(t)) + store := newCTBStore(t) + store.set(t, certA) + r := newSystemInfoVolumeRefresher(store.lister, nil) + dir1, dir2 := t.TempDir(), t.TempDir() + registerTrustVolume(t, r, dir1, "uid-1") + second := registerTrustVolume(t, r, dir2, "uid-1") + + second.Undo() + if _, ok := r.actors["uid-1"]; ok { + t.Fatal("Undo left an entry for the actor") + } + if !second.actor.stale { + t.Error("Undo left the removed entry unfenced") + } + store.set(t, certB) + r.refreshBundle(ctx, EgressTrustBundleName) + for _, dir := range []string{dir1, dir2} { + if got := readProjected(t, dir, "uid-1", "trust", "ca.pem"); got != certA { + t.Errorf("file under %s = %q, refreshed after Undo", dir, got) + } + } +} + +func TestSystemInfoVolumeRefresher_FailedRegisterLeavesNoEntry(t *testing.T) { + register := func(t *testing.T, r *systemInfoVolumeRefresher) *registration { + t.Helper() + vol := &systemInfoVolume{Name: "trust", Root: filepath.Join(t.TempDir(), "trust"), Spec: trustVolumeSpec("ca.pem")} + reg, err := r.Register("uid-1", resources.ActorRef{Atespace: "team-a", Name: "uid-1"}, []*systemInfoVolume{vol}) + if err == nil { + t.Fatal("Register succeeded without a trust bundle to project") + } + if _, ok := r.actors["uid-1"]; ok { + t.Error("failed Register left an entry for the actor") + } + return reg + } + + t.Run("first registration", func(t *testing.T) { + r := newSystemInfoVolumeRefresher(newCTBStore(t).lister, nil) + reg := register(t, r) + if reg != nil { + t.Error("failed Register returned a registration") + } + // Restore defers Undo before Register has run. + reg.Undo() + }) + + t.Run("after superseding a live entry", func(t *testing.T) { + store := newCTBStore(t) + store.set(t, string(testCertPEM(t))) + r := newSystemInfoVolumeRefresher(store.lister, nil) + first := registerTrustVolume(t, r, t.TempDir(), "uid-1") + store.remove(t) + register(t, r) + if !first.actor.stale { + t.Error("superseded entry was not marked stale") + } + }) +} + func TestSystemInfoVolumesFor(t *testing.T) { spec := &ateletpb.WorkloadSpec{Volumes: []*ateletpb.Volume{ {Name: "data", Source: &ateletpb.Volume_DurableDir{DurableDir: &ateletpb.DurableDirVolume{}}}, @@ -583,7 +674,8 @@ func TestSystemInfoVolumesFor(t *testing.T) { } } -// Register, Deregister, and refresh run concurrently for the race detector. +// Register, Undo, Deregister, and refresh run concurrently for the race +// detector. func TestSystemInfoVolumeRefresher_ConcurrentLifecycle(t *testing.T) { ctx := context.Background() certA, certB := string(testCertPEM(t)), string(testCertPEM(t)) @@ -598,17 +690,22 @@ func TestSystemInfoVolumeRefresher_ConcurrentLifecycle(t *testing.T) { wg.Add(1) go func() { defer wg.Done() - for range 20 { + for n := range 20 { vol := &systemInfoVolume{ Name: "trust", Root: filepath.Join(dir, uid, "system-info", "trust"), Spec: trustVolumeSpec("ca.pem"), } - if err := r.Register(uid, resources.ActorRef{Atespace: "team-a", Name: uid}, []*systemInfoVolume{vol}); err != nil { + reg, err := r.Register(uid, resources.ActorRef{Atespace: "team-a", Name: uid}, []*systemInfoVolume{vol}) + if err != nil { t.Errorf("Register(%s): %v", uid, err) return } - r.Deregister(uid) + if n%2 == 0 { + reg.Undo() + } else { + r.Deregister(uid) + } } }() } @@ -671,7 +768,7 @@ func TestSystemInfoVolumeRegister_TrustBundle(t *testing.T) { t.Run("resolution failure fails the start rather than produce an empty trust file", func(t *testing.T) { r := newSystemInfoVolumeRefresher(ctbLister(t), nil) vol := &systemInfoVolume{Name: "trust", Root: filepath.Join(t.TempDir(), "trust"), Spec: trustVolumeSpec("ca.pem")} - err := r.Register("uid-2", resources.ActorRef{Atespace: "team-a", Name: "actor-2"}, []*systemInfoVolume{vol}) + _, err := r.Register("uid-2", resources.ActorRef{Atespace: "team-a", Name: "actor-2"}, []*systemInfoVolume{vol}) if err == nil || !strings.Contains(err.Error(), "not found") || !strings.Contains(err.Error(), `"trust"`) { t.Errorf("Register = %v, want not-found error naming the volume", err) }