mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
Stop leaking in-progress snapshots after worker deletion (#1541)
Fixes #1539 If an actor is `CRASHED`, it can't be resumed/paused/suspended anymore, it can only be deleted. An actor deletion will clean up everything under the prefix of `Actor.external_snapshot.snapshot_uri` and `Actor.in_progress_snapshot_name` (the latter will be empty if the underlying worker was deleted). Before, we risked leaking a snapshot after a worker deletion because we were clearing the `in_progress_snapshot_name` field. We don't need to clear it up because the CRASHED state is terminal. - [x] Tests pass - [x] Appropriate changes to documentation are included in the PR
This commit is contained in:
@@ -18,7 +18,9 @@ import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"github.com/agent-substrate/substrate/cmd/ateapi/internal/store"
|
||||
"github.com/agent-substrate/substrate/cmd/ateapi/internal/store/storetest"
|
||||
"github.com/agent-substrate/substrate/internal/objectstore/objectstoretest"
|
||||
"github.com/agent-substrate/substrate/internal/resources"
|
||||
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
|
||||
"google.golang.org/grpc/codes"
|
||||
@@ -326,3 +328,106 @@ func TestEnsureExternalSnapshotsReleased_CollectsStrandedSnapshots(t *testing.T)
|
||||
t.Errorf("another actor's external snapshot %v was released", otherActorsSnapshot)
|
||||
}
|
||||
}
|
||||
|
||||
// TestDeleteActor_CollectsSnapshotsAfterWorkerDelete verifies that
|
||||
// deleting an actor whose suspend a worker delete crashed mid-finalize reclaims
|
||||
// every object that suspend wrote. CRASHED is terminal, so the actor delete is
|
||||
// the only collector left: whatever it cannot name is leaked for good.
|
||||
func TestDeleteActor_CollectsSnapshotsAfterWorkerDelete(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
// hasBeenSuspended means the actor was replacing a snapshot it took
|
||||
// itself, rather than suspending for the first time. Only the first
|
||||
// suspend needs the retained in-progress name to find what it wrote; a
|
||||
// replacement is already reachable through the snapshot it owns.
|
||||
hasBeenSuspended bool
|
||||
}{
|
||||
{
|
||||
name: "replacing a snapshot the actor owned",
|
||||
hasBeenSuspended: true,
|
||||
},
|
||||
{
|
||||
name: "first suspend, nothing to replace",
|
||||
hasBeenSuspended: false,
|
||||
},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
persistence := newTestPersistence(t)
|
||||
template := seedSubstrateTemplate(t, ctx, persistence, "sub-tmpl")
|
||||
objects := objectstoretest.New()
|
||||
|
||||
actorRef := resources.ActorRef{Atespace: "team-a", Name: "actor-1"}
|
||||
workerName := testWorkerUID("pod-1")
|
||||
actor := storetest.MustCreateActor(t, ctx, persistence, &ateapipb.Actor{
|
||||
Metadata: &ateapipb.ResourceMetadata{Atespace: actorRef.Atespace, Name: actorRef.Name},
|
||||
ActorTemplate: &ateapipb.ObjectRef{Atespace: "team-a", Name: "sub-tmpl"},
|
||||
Status: &ateapipb.ActorStatus{
|
||||
State: ateapipb.ActorState_ACTOR_STATE_RUNNING,
|
||||
WorkerAssignment: &ateapipb.WorkerAssignment{
|
||||
Worker: &ateapipb.ObjectRef{Name: workerName},
|
||||
WorkerNamespace: "worker-ns",
|
||||
WorkerPool: "pool",
|
||||
WorkerPod: "pod-1",
|
||||
WorkerPodUid: workerName,
|
||||
},
|
||||
},
|
||||
})
|
||||
if _, err := persistence.CreateWorker(ctx, &ateapipb.Worker{
|
||||
Metadata: &ateapipb.ResourceMetadata{Name: workerName},
|
||||
WorkerNamespace: "worker-ns",
|
||||
WorkerPool: "pool",
|
||||
WorkerPod: "pod-1",
|
||||
WorkerPodUid: workerName,
|
||||
Status: &ateapipb.WorkerStatus{},
|
||||
}); err != nil {
|
||||
t.Fatalf("CreateWorker: %v", err)
|
||||
}
|
||||
seedAssignment(t, persistence, workerName, &ateapipb.ActorAssignment{
|
||||
Actor: &ateapipb.ObjectRef{Atespace: actorRef.Atespace, Name: actorRef.Name},
|
||||
ActorUid: actor.GetMetadata().GetUid(),
|
||||
})
|
||||
|
||||
if tt.hasBeenSuspended {
|
||||
previous := mustActorSnapshotURI(t, template, actor, "old")
|
||||
objects.PutSnapshot(t, previous, "manifest.json")
|
||||
actor = mustUpdateActorStatus(t, ctx, persistence, actor, func(s *ateapipb.ActorStatus) {
|
||||
s.ExternalSnapshot = &ateapipb.ExternalSnapshot{SnapshotUri: previous.String()}
|
||||
})
|
||||
}
|
||||
|
||||
actorWorkflow := NewActorWorkflow(persistence, nil, nil, nil, nil, nil, "", nil, objects)
|
||||
// Suspend the actor as far as it gets: MarkSuspending mints the
|
||||
// in-progress name, and the checkpoint writes under it
|
||||
actor, err := actorWorkflow.ensureMarkedSuspending(ctx, actorRef, actor, template)
|
||||
if err != nil {
|
||||
t.Fatalf("ensureMarkedSuspending: %v", err)
|
||||
}
|
||||
fresh := mustActorSnapshotURI(t, template, actor, actor.GetStatus().GetInProgressSnapshotName())
|
||||
objects.PutSnapshot(t, fresh, "manifest.json")
|
||||
|
||||
// The worker's pod goes away with the commit still outstanding, so
|
||||
// the suspend never gets to finish.
|
||||
if _, err := NewWorkerWorkflow(persistence).DeleteWorker(ctx, workerName, store.DeletePreconditions{}); err != nil {
|
||||
t.Fatalf("DeleteWorker: %v", err)
|
||||
}
|
||||
|
||||
// The actor is CRASHED and can only be deleted from here.
|
||||
stored, err := persistence.GetActor(ctx, actorRef)
|
||||
if err != nil {
|
||||
t.Fatalf("GetActor: %v", err)
|
||||
}
|
||||
if got := stored.GetStatus().GetState(); got != ateapipb.ActorState_ACTOR_STATE_CRASHED {
|
||||
t.Fatalf("state = %v, want CRASHED", got)
|
||||
}
|
||||
if _, err := actorWorkflow.DeleteActor(ctx, actorRef, true); err != nil {
|
||||
t.Fatalf("DeleteActor: %v", err)
|
||||
}
|
||||
if left := objects.Prefix(t, fresh.OwnerPrefix()); len(left) != 0 {
|
||||
t.Errorf("deleting the actor left %v under its own prefix, want everything it wrote collected", left)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -213,9 +213,10 @@ func (w *WorkerWorkflow) releaseBoundActor(ctx context.Context, worker *ateapipb
|
||||
_, err = w.store.UpdateActor(ctx, actorRef, store.PreconditionFrom(actor), func(toUpdate *ateapipb.Actor) error {
|
||||
toUpdate.Status.State = ateapipb.ActorState_ACTOR_STATE_CRASHED
|
||||
toUpdate.Status.WorkerAssignment = nil
|
||||
// Both in-progress checkpoints die with the worker: the durable one was
|
||||
// never uploaded, the local one lived on the node that went away.
|
||||
toUpdate.Status.InProgressSnapshotName = ""
|
||||
// Local in-progress checkpoint dies with the worker: it lived on the node
|
||||
// that went away. The external in-progress checkpoint is kept. It'll be deleted
|
||||
// with the actor when the actor is deleted (only possible outcome from CRASHED
|
||||
// state).
|
||||
toUpdate.Status.InProgressLocalSnapshotName = ""
|
||||
return nil
|
||||
})
|
||||
|
||||
@@ -129,10 +129,15 @@ func TestDeleteWorkerWorkflow_ReleasesBoundActor(t *testing.T) {
|
||||
if got.GetStatus().GetWorkerAssignment() != nil {
|
||||
t.Errorf("actor worker assignment = %v, want it cleared", got.GetStatus().GetWorkerAssignment())
|
||||
}
|
||||
// The durable checkpoint was never uploaded and the local one lived on the
|
||||
// node that went away, so both die with the worker.
|
||||
if got.GetStatus().GetInProgressSnapshotName() != "" || got.GetStatus().GetInProgressLocalSnapshotName() != "" {
|
||||
t.Errorf("in-progress checkpoints not cleared: %v", got.GetStatus())
|
||||
// The local checkpoint lived on the node that went away, so it dies with
|
||||
// the worker.
|
||||
if got.GetStatus().GetInProgressLocalSnapshotName() != "" {
|
||||
t.Errorf("in-progress local checkpoint not cleared: %v", got.GetStatus())
|
||||
}
|
||||
// The durable one is kept: it names the prefix whatever atelet already
|
||||
// uploaded lives under, which the actor's delete needs to collect it.
|
||||
if got.GetStatus().GetInProgressSnapshotName() != "partial-snapshot" {
|
||||
t.Errorf("in-progress external checkpoint not preserved: %v", got.GetStatus())
|
||||
}
|
||||
// The last completed snapshot is what makes the actor resumable, so it stays.
|
||||
if want := someActorSnapshotURI(t, testStorageLocation, apiActorRef.Atespace, "last"); got.GetStatus().GetExternalSnapshot().GetSnapshotUri() != want {
|
||||
|
||||
Reference in New Issue
Block a user