mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
fix(atelet): Supersede stale system-info volume registrations in atelet (#1684)
When a worker pod crashes or is deleted out-of-band, `ateapi` transitions the actor to `CRASHED` and clears `worker_assignment` without calling `atelet.Terminate`. If the actor is later recovered and scheduled onto a new worker pod on the same node, `systemInfoVolumeRefresher.Register` panicked on the leftover entry in r.actors. Supersede any existing registration for the same actor UID by marking the previous entry stale under its mutex and replacing it in r.actors. Temporary fix for #1710 - [X] Tests pass - [X] Appropriate changes to documentation are included in the PR
This commit is contained in:
@@ -105,8 +105,9 @@ func newSystemInfoVolumeRefresher(lister certlisters.ClusterTrustBundleLister, i
|
||||
}
|
||||
|
||||
// Register records actorUID's system-info volumes and writes their contents
|
||||
// from current cluster state.
|
||||
// Double registering an actor causes a panic.
|
||||
// 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 {
|
||||
actor := ®isteredActor{uid: actorUID, ref: ref, volumes: volumes}
|
||||
// Held until the initial write finishes so a refresh cannot interleave.
|
||||
@@ -114,13 +115,16 @@ func (r *systemInfoVolumeRefresher) Register(actorUID string, ref resources.Acto
|
||||
defer actor.mu.Unlock()
|
||||
|
||||
r.mu.Lock()
|
||||
_, dup := r.actors[actorUID]
|
||||
if !dup {
|
||||
r.actors[actorUID] = actor
|
||||
}
|
||||
prev := r.actors[actorUID]
|
||||
r.actors[actorUID] = actor
|
||||
r.mu.Unlock()
|
||||
if dup {
|
||||
panic(fmt.Sprintf("system-info volumes: actor %s registered twice", actorUID))
|
||||
if prev != nil {
|
||||
prev.mu.Lock()
|
||||
prev.stale = true
|
||||
prev.mu.Unlock()
|
||||
slog.Info("Superseded a stale system-info volume registration",
|
||||
slog.String("actor_uid", actorUID),
|
||||
slog.Any("actor", ref))
|
||||
}
|
||||
|
||||
for _, v := range volumes {
|
||||
|
||||
@@ -544,19 +544,23 @@ func TestSystemInfoVolumeRefresher_DeregisterMarksStale(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestSystemInfoVolumeRefresher_RegisterTwicePanics(t *testing.T) {
|
||||
func TestSystemInfoVolumeRefresher_RegisterTwiceSupersedes(t *testing.T) {
|
||||
store := newCTBStore(t)
|
||||
store.set(t, string(testCertPEM(t)))
|
||||
r := newSystemInfoVolumeRefresher(store.lister, nil)
|
||||
dir := t.TempDir()
|
||||
registerTrustVolume(t, r, dir, "uid-1")
|
||||
first := r.actors["uid-1"]
|
||||
|
||||
defer func() {
|
||||
if recover() == nil {
|
||||
t.Error("Register of a still-registered UID did not panic")
|
||||
}
|
||||
}()
|
||||
_ = r.Register("uid-1", resources.ActorRef{Atespace: "team-a", Name: "uid-1"}, 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 {
|
||||
t.Error("superseded entry was not marked stale")
|
||||
}
|
||||
if r.actors["uid-1"] == first {
|
||||
t.Error("superseded entry was not replaced in r.actors")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSystemInfoVolumesFor(t *testing.T) {
|
||||
@@ -672,4 +676,15 @@ func TestSystemInfoVolumeRegister_TrustBundle(t *testing.T) {
|
||||
t.Errorf("Register = %v, want not-found error naming the volume", err)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("re-registration supersedes stale entry without panicking", func(t *testing.T) {
|
||||
r := newSystemInfoVolumeRefresher(store.lister, nil)
|
||||
dir1 := t.TempDir()
|
||||
dir2 := t.TempDir()
|
||||
registerTrustVolume(t, r, dir1, "uid-rereg")
|
||||
registerTrustVolume(t, r, dir2, "uid-rereg")
|
||||
if got := readProjected(t, dir2, "uid-rereg", "trust", "ca.pem"); got != string(certPEM) {
|
||||
t.Errorf("content = %q, want the sanitized bundle", got)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user