Skip to content
Open
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
17 changes: 10 additions & 7 deletions cmd/atelet/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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()
}
}()

Expand Down Expand Up @@ -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
}
Expand Down
44 changes: 38 additions & 6 deletions cmd/atelet/systeminfovolume.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 := &registeredActor{uid: actorUID, ref: ref, volumes: volumes}
// Held until the initial write finishes so a refresh cannot interleave.
actor.mu.Lock()
Expand All @@ -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 &registration{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]
Expand Down
125 changes: 111 additions & 14 deletions cmd/atelet/systeminfovolume_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -98,16 +98,18 @@ func metadataVolumeSpec() *ateletpb.SystemInfoVolume {

// registerTrustVolume registers actorUID with one volume projecting the
// egress bundle at <dir>/<uid>/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 {
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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)
}

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

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

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

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