Add the crash reason and timestamp to the ActorStatus (#1867)

Added a new field to actor status that has the crash reason (e.g., the
atelet response error) and the timestamp of the crash for better UX.
This commit is contained in:
Luiz Oliveira
2026-09-25 15:33:08 +00:00
committed by GitHub
parent d72edfbb3c
commit e87c55fc2e
19 changed files with 1365 additions and 625 deletions
@@ -28,6 +28,7 @@ import (
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
"google.golang.org/protobuf/testing/protocmp"
"google.golang.org/protobuf/types/known/timestamppb"
"k8s.io/apimachinery/pkg/util/validation/field"
)
@@ -568,6 +569,89 @@ func TestValidateActorUpdate(t *testing.T) {
}
}
func TestValidateActorStatusCrash(t *testing.T) {
crashPath := field.NewPath("status", "crash")
withCrash := func(mutate ...func(*ateapipb.ActorCrash)) func(*ateapipb.Actor) {
return withActorStatus(func(s *ateapipb.ActorStatus) {
s.State = ateapipb.ActorState_ACTOR_STATE_CRASHED
s.Crash = &ateapipb.ActorCrash{
Message: crashMessageWorkerGone,
CrashTime: &timestamppb.Timestamp{Seconds: 867},
}
for _, m := range mutate {
m(s.Crash)
}
})
}
tests := []struct {
name string
oldVal *ateapipb.Actor
newVal *ateapipb.Actor
want field.ErrorList
}{
{
name: "set crash",
oldVal: validActor(withActorStatus()),
newVal: validActor(withCrash()),
},
{
name: "set empty crash",
oldVal: validActor(withActorStatus()),
newVal: validActor(withCrash(func(c *ateapipb.ActorCrash) { *c = ateapipb.ActorCrash{} })),
},
{
name: "clear crash",
oldVal: validActor(withCrash()),
newVal: validActor(withActorStatus()),
},
{
name: "replace crash",
oldVal: validActor(withCrash()),
newVal: validActor(withCrash(func(c *ateapipb.ActorCrash) {
c.Message = crashMessageWorkerDraining
c.CrashTime = &timestamppb.Timestamp{Seconds: 5309}
})),
},
{
name: "message at max length",
oldVal: validActor(withActorStatus()),
newVal: validActor(withCrash(func(c *ateapipb.ActorCrash) { c.Message = strings.Repeat("x", 4096) })),
},
{
name: "message too long",
oldVal: validActor(withActorStatus()),
newVal: validActor(withCrash(func(c *ateapipb.ActorCrash) { c.Message = strings.Repeat("x", 4097) })),
want: field.ErrorList{field.TooLong(crashPath.Child("message"), nil, 4096).WithOrigin("maxLength")},
},
{
name: "truncated crash message fits",
oldVal: validActor(withActorStatus()),
newVal: validActor(withCrash(func(c *ateapipb.ActorCrash) {
c.Message = newActorCrash("pause", strings.Repeat("x", 2*maxCrashMessageBytes)).GetMessage()
})),
},
{
// Unchanged fields are not revalidated on update, so a crash stored
// before a limit was tightened does not block later status writes.
name: "unchanged overlong crash is not revalidated",
oldVal: validActor(withCrash(func(c *ateapipb.ActorCrash) { c.Message = strings.Repeat("x", 4097) })),
newVal: validActor(withCrash(func(c *ateapipb.ActorCrash) { c.Message = strings.Repeat("x", 4097) })),
},
{
name: "changed overlong crash is revalidated",
oldVal: validActor(withCrash(func(c *ateapipb.ActorCrash) { c.Message = strings.Repeat("x", 4097) })),
newVal: validActor(withCrash(func(c *ateapipb.ActorCrash) { c.Message = strings.Repeat("y", 4097) })),
want: field.ErrorList{field.TooLong(crashPath.Child("message"), nil, 4096).WithOrigin("maxLength")},
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
assertValidateErr(t, validateActorUpdate(context.Background(), nil, tt.newVal, tt.oldVal, true), tt.want)
})
}
}
func TestValidateGetActorRequest(t *testing.T) {
tests := []struct {
name string
+44 -2
View File
@@ -19,17 +19,43 @@ import (
"errors"
"fmt"
"log/slog"
"strings"
"github.com/agent-substrate/substrate/cmd/ateapi/internal/store"
"github.com/agent-substrate/substrate/internal/actorevent"
"github.com/agent-substrate/substrate/internal/ateattr"
"github.com/agent-substrate/substrate/internal/resources"
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
"google.golang.org/grpc/status"
"google.golang.org/protobuf/types/known/timestamppb"
)
// Fixed messages recorded in ActorCrash.message for crashes due to non-atelet
// reasons. A failed atelet call builds its message with ateletCrashMessage instead.
const (
crashMessageWorkerAssignmentMissing = "actor has no worker assignment in a state that requires one"
crashMessageLocalSnapshotNodeUnknown = "node holding the actor's local snapshot is unknown"
crashMessageWorkerGone = "assigned worker no longer exists"
crashMessageWorkerDraining = "assigned worker is draining"
crashMessageWorkerReassigned = "assigned worker no longer hosts the actor"
crashMessageWorkerIneligible = "assigned worker no longer satisfies the actor's placement constraints"
crashMessageWorkerPodGone = "worker pod went away while hosting the actor"
)
// maxCrashMessageBytes matches the maxLength on ActorCrash.message.
const maxCrashMessageBytes = 4096
// ateletCrashMessage describes a failed atelet call with the error text atelet
// returned, since its gRPC code is Unknown for most failures.
// TODO: consider sanizing the error message returned by atelet.
func ateletCrashMessage(rpc string, err error) string {
return fmt.Sprintf("atelet %s: %s", rpc, status.Convert(err).Message())
}
// crashActor moves the actor to CRASHED state and frees the worker it was
// assigned to, if any, so the worker can host other actors.
func crashActor(ctx context.Context, st crashActorStore, actorRef resources.ActorRef, opName string) error {
// assigned to, if any, so the worker can host other actors. message is
// recorded in the actor's status.
func crashActor(ctx context.Context, st crashActorStore, actorRef resources.ActorRef, opName, message string) error {
actor, err := st.GetActor(ctx, actorRef)
if err != nil {
return fmt.Errorf("while loading actor to crash: %w", err)
@@ -57,6 +83,10 @@ func crashActor(ctx context.Context, st crashActorStore, actorRef resources.Acto
_, err = st.UpdateActor(ctx, actorRef, store.PreconditionFrom(actor), func(toUpdate *ateapipb.Actor) error {
toUpdate.Status.State = ateapipb.ActorState_ACTOR_STATE_CRASHED
// An actor crashed concurrently, e.g. by worker deletion, keeps its first crash.
if !wasAlreadyCrashed {
toUpdate.Status.Crash = newActorCrash(opName, message)
}
// InProgressSnapshotUri and InProgressLocalSnapshotName are kept so a
// later DeleteActor or RevertActor can delete what they name: each is
@@ -79,6 +109,18 @@ func crashActor(ctx context.Context, st crashActorStore, actorRef resources.Acto
return nil
}
// newActorCrash records a crash that happens now, prefixing message with the
// operation that failed when it is known.
func newActorCrash(opName, message string) *ateapipb.ActorCrash {
if opName = ateattr.NormalizeOperationName(opName); opName != ateattr.OperationUnknown {
message = opName + " failed: " + message
}
return &ateapipb.ActorCrash{
Message: truncateUTF8(strings.ToValidUTF8(message, "�"), maxCrashMessageBytes),
CrashTime: timestamppb.Now(),
}
}
// logActorCrashed carries the identity ate.actor.crashes cannot: actor identity
// is barred from metric labels, so this record is the only way to attribute a
// crash to one agent. Call it beside recordActorCrash, under the same guard.
+119 -5
View File
@@ -31,11 +31,15 @@ import (
"github.com/agent-substrate/substrate/internal/ateattr"
"github.com/agent-substrate/substrate/internal/resources"
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
"github.com/google/go-cmp/cmp"
otellog "go.opentelemetry.io/otel/log"
"go.opentelemetry.io/otel/log/global"
sdklog "go.opentelemetry.io/otel/sdk/log"
sdkmetric "go.opentelemetry.io/otel/sdk/metric"
"go.opentelemetry.io/otel/sdk/metric/metricdata"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
"google.golang.org/protobuf/testing/protocmp"
)
// seedActor stores a running actor with all worker-binding fields populated, so
@@ -272,13 +276,123 @@ func TestCrashActor(t *testing.T) {
tt.setup(t, ctx, st)
}
err := crashActor(ctx, st, actorRef, ateattr.OperationUnknown)
err := crashActor(ctx, st, actorRef, ateattr.OperationUnknown, "test crash")
tt.check(t, ctx, st, err)
})
}
}
func TestCrashActor_RecordsCrash(t *testing.T) {
ctx := context.Background()
st, cleanup := storetest.SetupTestStore(t)
defer cleanup()
actorRef := resources.ActorRef{Atespace: "team-a", Name: "actor-1"}
seedActor(t, ctx, st, actorRef)
before := time.Now().Truncate(time.Microsecond)
if err := crashActor(ctx, st, actorRef, ateattr.OperationResume, crashMessageWorkerDraining); err != nil {
t.Fatalf("crashActor() = %v, want nil", err)
}
first, err := st.GetActor(ctx, actorRef)
if err != nil {
t.Fatalf("GetActor: %v", err)
}
crash := first.GetStatus().GetCrash()
if want := "resume failed: " + crashMessageWorkerDraining; crash.GetMessage() != want {
t.Errorf("Crash.Message = %q, want %q", crash.GetMessage(), want)
}
if got := crash.GetCrashTime().AsTime(); got.Before(before) || got.After(time.Now()) {
t.Errorf("Crash.CrashTime = %v, want between %v and now", got, before)
}
// Crashing an already-crashed actor, as a concurrent crash does, keeps the first crash.
if err := crashActor(ctx, st, actorRef, ateattr.OperationResume, crashMessageWorkerGone); err != nil {
t.Fatalf("second crashActor() = %v, want nil", err)
}
second, err := st.GetActor(ctx, actorRef)
if err != nil {
t.Fatalf("GetActor: %v", err)
}
if diff := cmp.Diff(crash, second.GetStatus().GetCrash(), protocmp.Transform()); diff != "" {
t.Errorf("Crash after re-crash differs from the first crash (-want +got):\n%s", diff)
}
}
func TestAteletCrashMessage(t *testing.T) {
tests := []struct {
name string
err error
want string
}{
{
name: "status error keeps its text",
err: status.Error(codes.Unknown, "while uploading external snapshot: googleapi: Error 403: forbidden"),
want: "atelet Restore: while uploading external snapshot: googleapi: Error 403: forbidden",
},
{
name: "plain error keeps its text",
err: errors.New("connection refused"),
want: "atelet Restore: connection refused",
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if got := ateletCrashMessage("Restore", tt.err); got != tt.want {
t.Errorf("ateletCrashMessage() = %q, want %q", got, tt.want)
}
})
}
}
func TestNewActorCrash(t *testing.T) {
const resumeOpPrefix = "resume failed: "
tests := []struct {
name string
opName string
message string
want string
}{
{
name: "known operation prefixes the message",
opName: ateattr.OperationResume,
message: crashMessageWorkerGone,
want: resumeOpPrefix + crashMessageWorkerGone,
},
{
name: "unknown operation leaves the message bare",
opName: ateattr.OperationUnknown,
message: crashMessageWorkerPodGone,
want: crashMessageWorkerPodGone,
},
{
name: "invalid UTF-8 is replaced",
opName: ateattr.OperationResume,
message: "bad \xff byte",
want: resumeOpPrefix + "bad \uFFFD byte",
},
{
name: "long message is truncated to the limit",
opName: ateattr.OperationResume,
message: strings.Repeat("x", maxCrashMessageBytes),
want: resumeOpPrefix + strings.Repeat("x", maxCrashMessageBytes-len(resumeOpPrefix)),
},
{
name: "truncation does not split a rune",
opName: ateattr.OperationResume,
message: strings.Repeat("x", maxCrashMessageBytes-len(resumeOpPrefix)-1) + "é",
want: resumeOpPrefix + strings.Repeat("x", maxCrashMessageBytes-len(resumeOpPrefix)-1),
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if got := newActorCrash(tt.opName, tt.message).GetMessage(); got != tt.want {
t.Errorf("newActorCrash().Message = %q, want %q", got, tt.want)
}
})
}
}
func TestCrashActor_Metrics(t *testing.T) {
reader := sdkmetric.NewManualReader()
mp := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader))
@@ -325,7 +439,7 @@ func TestCrashActor_Metrics(t *testing.T) {
}
storetest.MustCreateActor(t, ctx, st, actor)
if err := crashActor(ctx, st, actorRef, ateattr.OperationResume); err != nil {
if err := crashActor(ctx, st, actorRef, ateattr.OperationResume, "test crash"); err != nil {
t.Fatalf("crashActor: %v", err)
}
@@ -424,7 +538,7 @@ func TestCrashActorReleaseFailureLeavesWorkerReclaimable(t *testing.T) {
seedWorker(t, ctx, st, actorRef)
releaseErr := errors.New("state store unavailable")
err := crashActor(ctx, failingReleaseStore{Interface: st, err: releaseErr}, actorRef, ateattr.OperationUnknown)
err := crashActor(ctx, failingReleaseStore{Interface: st, err: releaseErr}, actorRef, ateattr.OperationUnknown, "test crash")
if err == nil {
t.Fatal("crashActor() = nil, want error")
@@ -611,7 +725,7 @@ func TestCrashActor_RecordAndCounterAgree(t *testing.T) {
Status: &ateapipb.ActorStatus{State: ateapipb.ActorState_ACTOR_STATE_RUNNING},
})
if err := crashActor(ctx, st, actorRef, ateattr.OperationResume); err != nil {
if err := crashActor(ctx, st, actorRef, ateattr.OperationResume, "test crash"); err != nil {
t.Fatalf("crashActor: %v", err)
}
if len(*records) != 1 {
@@ -654,7 +768,7 @@ func TestCrashActor_RecordAndCounterAgree(t *testing.T) {
assertCopiesAgree(t, (*records)[0], gotEvents[0], actorevent.Crashed)
// Re-crashing an already-crashed actor must move neither signal.
if err := crashActor(ctx, st, actorRef, ateattr.OperationResume); err != nil {
if err := crashActor(ctx, st, actorRef, ateattr.OperationResume, "test crash"); err != nil {
t.Fatalf("second crashActor: %v", err)
}
if len(*records) != 1 {
@@ -3084,6 +3084,62 @@ func TestResumeActor_AteletFailureCrashesActor(t *testing.T) {
if actor.GetStatus().GetWorkerAssignment() != nil {
t.Errorf("expected worker assignment to be cleared, got %v", actor.GetStatus().GetWorkerAssignment())
}
assertActorCrashStatus(t, tc, name, "resume failed: atelet Restore: mock atelet failure")
worker, err := tc.persistence.GetWorker(context.Background(), podUID)
if err != nil {
t.Fatalf("GetWorker(%s) failed: %v", podUID, err)
}
if n := worker.GetStatus().GetAllocated().GetActors(); n != 0 {
t.Errorf("expected worker to be released after crash, still holds %d actors", n)
}
}
// TestResumeActor_LocalRestoreFailureCrashesActor: an atelet error restoring a
// paused actor from its local snapshot crashes it and releases its worker.
func TestResumeActor_LocalRestoreFailureCrashesActor(t *testing.T) {
ns := namespaceForTest("ns-resume-local-crash")
tc := setupTest(t, ns)
defer tc.cleanup()
createTemplate(t, tc, ns)
podUID := createWorkerPod(t, tc, ns, "worker-1", "node1", "pool1")
name := "id1"
ref := &ateapipb.ObjectRef{Atespace: testAtespace, Name: name}
if _, err := tc.client.CreateActor(context.Background(), &ateapipb.CreateActorRequest{Actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{Atespace: testAtespace, Name: name},
ActorTemplate: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "tmpl1"},
}}); err != nil {
t.Fatalf("CreateActor failed: %v", err)
}
if _, err := tc.client.ResumeActor(context.Background(), &ateapipb.ResumeActorRequest{Actor: ref}); err != nil {
t.Fatalf("ResumeActor failed: %v", err)
}
if _, err := tc.client.PauseActor(context.Background(), &ateapipb.PauseActorRequest{Actor: ref}); err != nil {
t.Fatalf("PauseActor failed: %v", err)
}
tc.fakeAtelet.Reset()
tc.fakeAtelet.FailRestore = status.Error(codes.Unavailable, "injected restore failure")
if _, err := tc.client.ResumeActor(context.Background(), &ateapipb.ResumeActorRequest{Actor: ref}); err == nil {
t.Fatal("ResumeActor succeeded despite failing restore")
}
if got := tc.fakeAtelet.RestoreRequest.GetType(); got != ateletpb.CheckpointType_CHECKPOINT_TYPE_LOCAL {
t.Fatalf("restore type = %v, want LOCAL", got)
}
actor, err := tc.client.GetActor(context.Background(), &ateapipb.GetActorRequest{Actor: ref})
if err != nil {
t.Fatalf("GetActor failed: %v", err)
}
if actor.GetStatus().GetState() != ateapipb.ActorState_ACTOR_STATE_CRASHED {
t.Fatalf("state after failed restore = %v, want CRASHED", actor.GetStatus().GetState())
}
if actor.GetStatus().GetWorkerAssignment() != nil {
t.Errorf("expected worker assignment to be cleared, got %v", actor.GetStatus().GetWorkerAssignment())
}
assertActorCrashStatus(t, tc, name, "resume failed: atelet Restore: injected restore failure")
worker, err := tc.persistence.GetWorker(context.Background(), podUID)
if err != nil {
@@ -3801,6 +3857,9 @@ func TestResumeActor_ReleasesStaleWorkerWhenPoolBecomesIneligible(t *testing.T)
if got := getResp.GetStatus().GetState(); got != ateapipb.ActorState_ACTOR_STATE_CRASHED {
t.Errorf("expected actor state CRASHED, got %v", got)
}
// The failed Run already crashed the actor, so the crash it records is
// that one, not the later eligibility check.
assertActorCrashStatus(t, tc, name, "resume failed: atelet Run: mock atelet failure")
listResp, err := tc.client.ListWorkers(context.Background(), &ateapipb.ListWorkersRequest{})
if err != nil {
@@ -3940,6 +3999,7 @@ func TestResumeActor_CrashesIfAssignedWorkerIsDraining(t *testing.T) {
if got := getResp.GetStatus().GetWorkerAssignment().GetWorkerPod(); got != "" {
t.Errorf("expected actor pod name to be empty, got %q", got)
}
assertActorCrashStatus(t, tc, id, "resume failed: assigned worker is draining")
// The draining worker must have been released.
listResp, err := tc.client.ListWorkers(context.Background(), &ateapipb.ListWorkersRequest{})
@@ -4367,6 +4427,7 @@ func TestSuspendActor_FromPaused_UploadFailureCrashes(t *testing.T) {
if crashed.GetStatus().GetWorkerAssignment() != nil {
t.Errorf("expected worker assignment to be cleared, got %v", crashed.GetStatus().GetWorkerAssignment())
}
assertActorCrashStatus(t, tc, name, "suspend failed: atelet UploadPausedCheckpoint: injected upload failure")
worker, err := tc.persistence.GetWorker(context.Background(), podUID)
if err != nil {
@@ -4388,6 +4449,86 @@ func TestSuspendActor_FromPaused_UploadFailureCrashes(t *testing.T) {
}
}
// TestCheckpointFailureCrashes: an atelet Checkpoint error while pausing or
// suspending a running actor crashes it, records atelet's error text, and
// releases its worker.
func TestCheckpointFailureCrashes(t *testing.T) {
tests := []struct {
name string
call func(tc *testContext, ref *ateapipb.ObjectRef) error
wantMessage string
}{
{
name: "pause",
call: func(tc *testContext, ref *ateapipb.ObjectRef) error {
_, err := tc.client.PauseActor(context.Background(), &ateapipb.PauseActorRequest{Actor: ref})
return err
},
wantMessage: "pause failed: atelet Checkpoint: injected checkpoint failure",
},
{
name: "suspend",
call: func(tc *testContext, ref *ateapipb.ObjectRef) error {
_, err := tc.client.SuspendActor(context.Background(), &ateapipb.SuspendActorRequest{Actor: ref})
return err
},
wantMessage: "suspend failed: atelet Checkpoint: injected checkpoint failure",
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
ns := namespaceForTest("ns-checkpoint-crash-" + tt.name)
tc := setupTest(t, ns)
defer tc.cleanup()
createTemplate(t, tc, ns)
podUID := createWorkerPod(t, tc, ns, "worker-1", "node1", "pool1")
name := "id1"
ref := &ateapipb.ObjectRef{Atespace: testAtespace, Name: name}
if _, err := tc.client.CreateActor(context.Background(), &ateapipb.CreateActorRequest{Actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{Atespace: testAtespace, Name: name},
ActorTemplate: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "tmpl1"},
}}); err != nil {
t.Fatalf("CreateActor failed: %v", err)
}
if _, err := tc.client.ResumeActor(context.Background(), &ateapipb.ResumeActorRequest{Actor: ref}); err != nil {
t.Fatalf("ResumeActor failed: %v", err)
}
tc.fakeAtelet.Reset()
tc.fakeAtelet.FailCheckpoint = status.Error(codes.Unavailable, "injected checkpoint failure")
err := tt.call(tc, ref)
if err == nil {
t.Fatalf("%s succeeded despite failing checkpoint", tt.name)
}
if !tc.fakeAtelet.CheckpointCalled {
t.Error("expected atelet Checkpoint to be called")
}
crashed, err := tc.client.GetActor(context.Background(), &ateapipb.GetActorRequest{Actor: ref})
if err != nil {
t.Fatalf("GetActor failed: %v", err)
}
if crashed.GetStatus().GetState() != ateapipb.ActorState_ACTOR_STATE_CRASHED {
t.Fatalf("state after failed checkpoint = %v, want CRASHED", crashed.GetStatus().GetState())
}
if crashed.GetStatus().GetWorkerAssignment() != nil {
t.Errorf("expected worker assignment to be cleared, got %v", crashed.GetStatus().GetWorkerAssignment())
}
assertActorCrashStatus(t, tc, name, tt.wantMessage)
worker, err := tc.persistence.GetWorker(context.Background(), podUID)
if err != nil {
t.Fatalf("GetWorker(%s) failed: %v", podUID, err)
}
if n := worker.GetStatus().GetAllocated().GetActors(); n != 0 {
t.Errorf("expected worker to be released after crash, still holds %d actors", n)
}
})
}
}
// TestResumeActor_RelocatesAfterSuspendFromPaused covers the capacity-recovery
// flow: a PAUSED actor is pinned to the node holding its local snapshot, so it
// cannot resume while that node is full. Suspending it uploads the snapshot and
@@ -4921,6 +5062,72 @@ func TestRevertActor_FromCrashed(t *testing.T) {
}
}
// TestRevertActor_TerminateFailureCrashes verifies that an atelet error while
// terminating the discarded execution crashes the actor, and that a later
// successful revert clears the recorded crash.
func TestRevertActor_TerminateFailureCrashes(t *testing.T) {
ns := namespaceForTest("ns-revert-terminate-crash")
tc := setupTest(t, ns)
defer tc.cleanup()
createTemplate(t, tc, ns)
createWorkerPod(t, tc, ns, "worker-1", "node1", "pool1")
ctx := context.Background()
const name = "id1"
actorRef := &ateapipb.ObjectRef{Atespace: testAtespace, Name: name}
if _, err := tc.client.CreateActor(ctx, &ateapipb.CreateActorRequest{Actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{Atespace: testAtespace, Name: name},
ActorTemplate: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "tmpl1"},
}}); err != nil {
t.Fatalf("CreateActor failed: %v", err)
}
if _, err := tc.client.ResumeActor(ctx, &ateapipb.ResumeActorRequest{Actor: actorRef}); err != nil {
t.Fatalf("ResumeActor failed: %v", err)
}
tc.fakeAtelet.Lock.Lock()
tc.fakeAtelet.FailTerminate = status.Error(codes.Internal, "injected terminate failure at /var/lib/node-path")
tc.fakeAtelet.Lock.Unlock()
if _, err := tc.client.RevertActor(ctx, &ateapipb.RevertActorRequest{Actor: actorRef}); err == nil {
t.Fatal("RevertActor succeeded despite failing terminate")
}
assertActorCrashStatus(t, tc, name, "revert failed: atelet Terminate: workflow failed at step CallAteletTerminate: while terminating actor on atelet: rpc error: code = Internal desc = injected terminate failure at /var/lib/node-path")
tc.fakeAtelet.Lock.Lock()
tc.fakeAtelet.FailTerminate = nil
tc.fakeAtelet.Lock.Unlock()
reverted, err := tc.client.RevertActor(ctx, &ateapipb.RevertActorRequest{Actor: actorRef})
if err != nil {
t.Fatalf("second RevertActor failed: %v", err)
}
if got := reverted.GetActor().GetStatus().GetCrash(); got != nil {
t.Errorf("Crash = %v, want cleared by the successful revert", got)
}
}
// assertActorCrashStatus reads the actor through the API and checks the crash it records.
func assertActorCrashStatus(t *testing.T, tc *testContext, name, wantMessage string) {
t.Helper()
actor, err := tc.client.GetActor(context.Background(), &ateapipb.GetActorRequest{
Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: name},
})
if err != nil {
t.Fatalf("GetActor failed: %v", err)
}
crash := actor.GetStatus().GetCrash()
if crash == nil {
t.Fatalf("Crash = nil, want message %q", wantMessage)
}
if crash.GetMessage() != wantMessage {
t.Errorf("Crash.Message = %q, want %q", crash.GetMessage(), wantMessage)
}
if crash.GetCrashTime() == nil {
t.Error("Crash.CrashTime = nil, want set")
}
}
// TestRevertActor_RejectsSuspended pins the rejection a lost race produces: a
// suspend that won left the actor SUSPENDED, and reporting success there would
// claim the opposite of what the revert asked for.
@@ -117,6 +117,7 @@ type FakeAteletServer struct {
CheckpointCalled bool
CheckpointRequest *ateletpb.CheckpointRequest
FailCheckpoint error
RestoreCalled bool
RestoreRequest *ateletpb.RestoreRequest
@@ -168,6 +169,7 @@ func (f *FakeAteletServer) Reset() {
f.CheckpointCalled = false
f.CheckpointRequest = nil
f.FailCheckpoint = nil
f.RestoreCalled = false
f.RestoreRequest = nil
@@ -219,6 +221,9 @@ func (f *FakeAteletServer) Checkpoint(ctx context.Context, req *ateletpb.Checkpo
f.CheckpointCalled = true
f.CheckpointRequest = proto.Clone(req).(*ateletpb.CheckpointRequest)
if f.FailCheckpoint != nil {
return nil, f.FailCheckpoint
}
if err := f.writeSnapshot(req.GetExternalConfig().GetSnapshotUri()); err != nil {
return nil, err
@@ -160,7 +160,7 @@ func (w *ActorWorkflow) ensureAteletPaused(ctx context.Context, actorRef resourc
assignment := actor.GetStatus().GetWorkerAssignment()
if assignment == nil {
// Missing active worker pod reference in PAUSING state indicates corrupted store state.
if err := crashActor(ctx, w.store, actorRef, ateattr.OperationPause); err != nil {
if err := crashActor(ctx, w.store, actorRef, ateattr.OperationPause, crashMessageWorkerAssignmentMissing); err != nil {
slog.ErrorContext(ctx, "Failed to crash actor", slog.String("err", err.Error()))
}
return "", status.Errorf(codes.FailedPrecondition, "CallAteletPause prerequisite not met for Actor: %s. No worker assignment", actorRef)
@@ -201,7 +201,7 @@ func (w *ActorWorkflow) ensureAteletPaused(ctx context.Context, actorRef resourc
if _, err = client.Checkpoint(ctx, req); err != nil {
slog.LogAttrs(ctx, slog.LevelError, "Setting Actor to crashed due to error",
append(ateattr.ActorRefLogAttrs(actorRef), slog.Any("err", err))...)
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationPause); cerr != nil {
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationPause, ateletCrashMessage("Checkpoint", err)); cerr != nil {
return wireSnapshotScope, cerr
}
return wireSnapshotScope, fmt.Errorf("actor %s crashed: %w", actorRef, err)
@@ -254,6 +254,7 @@ func (w *ActorWorkflow) ensurePausedFinalized(ctx context.Context, actorRef reso
}
wasAlreadyCrashed := latestActor.GetStatus().GetState() == ateapipb.ActorState_ACTOR_STATE_CRASHED
newState := ateapipb.ActorState_ACTOR_STATE_PAUSED
var crashStatus *ateapipb.ActorCrash
if nodeName == "" {
// Without a node name we cannot record where the local snapshot lives,
// so the actor can never be resumed (the scheduler would search for a
@@ -262,6 +263,7 @@ func (w *ActorWorkflow) ensurePausedFinalized(ctx context.Context, actorRef reso
slog.LogAttrs(ctx, slog.LevelError, "Node name not found during finalize pause, crashing actor",
ateattr.ActorRefLogAttrs(actorRef)...)
newState = ateapipb.ActorState_ACTOR_STATE_CRASHED
crashStatus = newActorCrash(ateattr.OperationPause, crashMessageLocalSnapshotNodeUnknown)
}
contentScope := actorTemplate.GetSnapshotConfig().GetOnPause()
sandboxClass := ""
@@ -274,6 +276,9 @@ func (w *ActorWorkflow) ensurePausedFinalized(ctx context.Context, actorRef reso
storedActor, err := w.store.UpdateActor(ctx, actorRef, store.PreconditionFrom(latestActor), func(toUpdate *ateapipb.Actor) error {
toUpdate.Status.State = newState
if newState == ateapipb.ActorState_ACTOR_STATE_CRASHED && !wasAlreadyCrashed {
toUpdate.Status.Crash = crashStatus
}
// TODO(dberkov) - what if InProgressLocalSnapshotName is empty? That shouldn't be possible.
if toUpdate.GetStatus().GetInProgressLocalSnapshotName() != "" {
localSnapshot := &ateapipb.LocalSnapshot{
@@ -74,6 +74,9 @@ func TestEnsurePausedFinalized_WorkerGone(t *testing.T) {
if got.GetStatus().GetState() != ateapipb.ActorState_ACTOR_STATE_CRASHED {
t.Errorf("state = %v, want CRASHED (node name unknown, cannot resume safely)", got.GetStatus().GetState())
}
if msg, want := got.GetStatus().GetCrash().GetMessage(), "pause failed: "+crashMessageLocalSnapshotNodeUnknown; msg != want {
t.Errorf("crash message = %q, want %q", msg, want)
}
for _, n := range got.GetStatus().GetLocalSnapshot().GetNodeVmsWithLocalSnapshots() {
if n == "" {
t.Errorf("BUG: empty string in NodeVmsWithLocalSnapshots, the scheduler's node restriction would never match a real worker")
@@ -343,6 +346,9 @@ func TestPauseActor_CrashesWhenPausingActorMissingWorkerPod(t *testing.T) {
if got.GetStatus().GetState() != ateapipb.ActorState_ACTOR_STATE_CRASHED {
t.Errorf("stored state = %v, want %v", got.GetStatus().GetState(), ateapipb.ActorState_ACTOR_STATE_CRASHED)
}
if msg, want := got.GetStatus().GetCrash().GetMessage(), "pause failed: "+crashMessageWorkerAssignmentMissing; msg != want {
t.Errorf("crash message = %q, want %q", msg, want)
}
}
// TestEnsureMarkedPausing_GoldenAtespaceRejected verifies golden actors
@@ -350,7 +350,7 @@ func (w *ActorWorkflow) validateAssignedWorker(ctx context.Context, actorRef res
slog.ErrorContext(ctx, "expected a worker assignment on a RESUMING actor, found none")
// Crash the actor if its worker assignment is missing. We should never be in this state.
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationResume); cerr != nil {
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationResume, crashMessageWorkerAssignmentMissing); cerr != nil {
return nil, cerr
}
return nil, status.Errorf(codes.Aborted, "actor %s crashed", actorRef)
@@ -360,7 +360,7 @@ func (w *ActorWorkflow) validateAssignedWorker(ctx context.Context, actorRef res
if err != nil {
// Crash the actor if it was assigned to a deleted pod.
if errors.Is(err, store.ErrNotFound) {
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationResume); cerr != nil {
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationResume, crashMessageWorkerGone); cerr != nil {
return nil, cerr
}
return nil, status.Errorf(codes.Aborted, "actor %s crashed", actorRef)
@@ -371,7 +371,7 @@ func (w *ActorWorkflow) validateAssignedWorker(ctx context.Context, actorRef res
slog.InfoContext(ctx, "Assigned worker is draining; crashing actor",
slog.String("actor", actorRef.String()),
slog.String("worker", worker.GetWorkerNamespace()+"/"+worker.GetWorkerPod()))
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationResume); cerr != nil {
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationResume, crashMessageWorkerDraining); cerr != nil {
return nil, cerr
}
return nil, status.Errorf(codes.Aborted, "actor %s crashed", actorRef.String())
@@ -384,7 +384,7 @@ func (w *ActorWorkflow) validateAssignedWorker(ctx context.Context, actorRef res
if !hosted {
slog.ErrorContext(ctx, "crashing actor because its assigned worker no longer hosts it",
slog.String("worker", worker.GetWorkerPod()))
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationResume); cerr != nil {
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationResume, crashMessageWorkerReassigned); cerr != nil {
return nil, fmt.Errorf("while crashing actor: %w", cerr)
}
return nil, status.Errorf(codes.Aborted, "actor %s crashed", actorRef)
@@ -402,7 +402,7 @@ func (w *ActorWorkflow) validateAssignedWorker(ctx context.Context, actorRef res
if _, err := w.store.ReleaseActorFromWorker(ctx, worker.GetMetadata().GetName(), actor.GetMetadata().GetUid()); err != nil {
return nil, fmt.Errorf("while releasing stale worker assignment: %w", err)
}
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationResume); cerr != nil {
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationResume, crashMessageWorkerIneligible); cerr != nil {
return nil, fmt.Errorf("while crashing actor: %w", cerr)
}
return nil, status.Errorf(codes.Aborted, "actor %s crashed", actorRef)
@@ -707,7 +707,7 @@ func (w *ActorWorkflow) ensureAteletRestored(ctx context.Context, actorRef resou
if _, err = client.Restore(ctx, req); err != nil {
slog.LogAttrs(ctx, slog.LevelError, "Setting Actor to crashed due to error",
append(ateattr.ActorRefLogAttrs(actorRef), slog.Any("err", err))...)
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationResume); cerr != nil {
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationResume, ateletCrashMessage("Restore", err)); cerr != nil {
return tele, cerr
}
return tele, fmt.Errorf("actor %s crashed: %w", actorRef, err)
@@ -757,7 +757,7 @@ func (w *ActorWorkflow) ensureAteletRestored(ctx context.Context, actorRef resou
if _, err = client.Restore(ctx, req); err != nil {
slog.LogAttrs(ctx, slog.LevelError, "Setting Actor to crashed due to error",
append(ateattr.ActorRefLogAttrs(actorRef), slog.Any("err", err))...)
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationResume); cerr != nil {
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationResume, ateletCrashMessage("Restore", err)); cerr != nil {
return tele, cerr
}
return tele, fmt.Errorf("actor %s crashed: %w", actorRef, err)
@@ -792,7 +792,7 @@ func (w *ActorWorkflow) ensureAteletRestored(ctx context.Context, actorRef resou
if _, err = client.Run(ctx, req); err != nil {
slog.LogAttrs(ctx, slog.LevelError, "Setting Actor to crashed due to error",
append(ateattr.ActorRefLogAttrs(actorRef), slog.Any("err", err))...)
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationResume); cerr != nil {
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationResume, ateletCrashMessage("Run", err)); cerr != nil {
return tele, cerr
}
return tele, fmt.Errorf("actor %s crashed: %w", actorRef, err)
@@ -739,15 +739,18 @@ func TestResumeActor_CrashesOnMissingWorkerAssignment(t *testing.T) {
if got.GetStatus().GetState() != ateapipb.ActorState_ACTOR_STATE_CRASHED {
t.Errorf("stored state = %v, want %v", got.GetStatus().GetState(), ateapipb.ActorState_ACTOR_STATE_CRASHED)
}
if msg, want := got.GetStatus().GetCrash().GetMessage(), "resume failed: "+crashMessageWorkerAssignmentMissing; msg != want {
t.Errorf("crash message = %q, want %q", msg, want)
}
}
// TestValidateAssignedWorker_WorkerOwnership verifies that RESUMING recovery
// only proceeds on a worker whose assignment still names this actor: the
// TestValidateAssignedWorker verifies that RESUMING recovery only proceeds on
// a live, non-draining worker whose assignment still names this actor: the
// recovery path loads the worker by pod name only, so the assignment may have
// been cleared and the worker re-claimed by another actor in the meantime. On
// a mismatch the actor is crashed and the worker — which is not ours — must
// not be written.
func TestValidateAssignedWorker_WorkerOwnership(t *testing.T) {
func TestValidateAssignedWorker(t *testing.T) {
ownAssignment := &ateapipb.ActorAssignment{
Actor: &ateapipb.ObjectRef{Atespace: "team-a", Name: "shared"},
ActorUid: "own-actor-uid",
@@ -760,14 +763,21 @@ func TestValidateAssignedWorker_WorkerOwnership(t *testing.T) {
Actor: &ateapipb.ObjectRef{Atespace: "team-a", Name: "shared"},
ActorUid: "stale-incarnation-uid",
}
activeStatus := &ateapipb.WorkerStatus{State: ateapipb.WorkerState_WORKER_STATE_ACTIVE, Capacity: &ateapipb.WorkerResources{Actors: 1}}
drainingStatus := &ateapipb.WorkerStatus{State: ateapipb.WorkerState_WORKER_STATE_DRAINING, Capacity: &ateapipb.WorkerResources{Actors: 1}}
tests := []struct {
name string
name string
// workerStatus is the stored worker's status, nil for a worker that
// is gone.
workerStatus *ateapipb.WorkerStatus
sandboxClass string
assignment *ateapipb.ActorAssignment
// wantCode is codes.OK when validateAssignedWorker must return nil.
wantCode codes.Code
wantActorState ateapipb.ActorState
// wantCrashMessage is the crash recorded when the actor is crashed.
wantCrashMessage string
// wantAssignment is the assignment expected on the stored worker
// afterwards; wantWorkerWrite false additionally asserts the worker
// version did not move (no write at all).
@@ -775,31 +785,54 @@ func TestValidateAssignedWorker_WorkerOwnership(t *testing.T) {
wantWorkerWrite bool
}{
{
name: "crashes actor and leaves worker untouched when assigned to another actor",
sandboxClass: "gvisor",
assignment: otherAssignment,
wantCode: codes.Aborted,
wantActorState: ateapipb.ActorState_ACTOR_STATE_CRASHED,
wantAssignment: otherAssignment,
name: "crashes actor when worker is gone",
wantCode: codes.Aborted,
wantActorState: ateapipb.ActorState_ACTOR_STATE_CRASHED,
wantCrashMessage: "resume failed: " + crashMessageWorkerGone,
},
{
name: "crashes actor and leaves worker untouched when assigned to previous incarnation of same actor",
sandboxClass: "gvisor",
assignment: staleIncarnationAssignment,
wantCode: codes.Aborted,
wantActorState: ateapipb.ActorState_ACTOR_STATE_CRASHED,
wantAssignment: staleIncarnationAssignment,
name: "crashes actor and leaves worker untouched when worker is draining",
workerStatus: drainingStatus,
sandboxClass: "gvisor",
assignment: ownAssignment,
wantCode: codes.Aborted,
wantActorState: ateapipb.ActorState_ACTOR_STATE_CRASHED,
wantCrashMessage: "resume failed: " + crashMessageWorkerDraining,
wantAssignment: ownAssignment,
},
{
name: "crashes actor and leaves worker untouched when assignment is cleared",
sandboxClass: "gvisor",
assignment: nil,
wantCode: codes.Aborted,
wantActorState: ateapipb.ActorState_ACTOR_STATE_CRASHED,
wantAssignment: nil,
name: "crashes actor and leaves worker untouched when assigned to another actor",
workerStatus: activeStatus,
sandboxClass: "gvisor",
assignment: otherAssignment,
wantCode: codes.Aborted,
wantActorState: ateapipb.ActorState_ACTOR_STATE_CRASHED,
wantCrashMessage: "resume failed: " + crashMessageWorkerReassigned,
wantAssignment: otherAssignment,
},
{
name: "crashes actor and leaves worker untouched when assigned to previous incarnation of same actor",
workerStatus: activeStatus,
sandboxClass: "gvisor",
assignment: staleIncarnationAssignment,
wantCode: codes.Aborted,
wantActorState: ateapipb.ActorState_ACTOR_STATE_CRASHED,
wantCrashMessage: "resume failed: " + crashMessageWorkerReassigned,
wantAssignment: staleIncarnationAssignment,
},
{
name: "crashes actor and leaves worker untouched when assignment is cleared",
workerStatus: activeStatus,
sandboxClass: "gvisor",
assignment: nil,
wantCode: codes.Aborted,
wantActorState: ateapipb.ActorState_ACTOR_STATE_CRASHED,
wantCrashMessage: "resume failed: " + crashMessageWorkerReassigned,
wantAssignment: nil,
},
{
name: "passes for own eligible worker",
workerStatus: activeStatus,
sandboxClass: "gvisor",
assignment: ownAssignment,
wantCode: codes.OK,
@@ -807,13 +840,15 @@ func TestValidateAssignedWorker_WorkerOwnership(t *testing.T) {
wantAssignment: ownAssignment,
},
{
name: "releases own ineligible worker and crashes actor",
sandboxClass: "microvm",
assignment: ownAssignment,
wantCode: codes.Aborted,
wantActorState: ateapipb.ActorState_ACTOR_STATE_CRASHED,
wantAssignment: nil,
wantWorkerWrite: true,
name: "releases own ineligible worker and crashes actor",
workerStatus: activeStatus,
sandboxClass: "microvm",
assignment: ownAssignment,
wantCode: codes.Aborted,
wantActorState: ateapipb.ActorState_ACTOR_STATE_CRASHED,
wantCrashMessage: "resume failed: " + crashMessageWorkerIneligible,
wantAssignment: nil,
wantWorkerWrite: true,
},
}
@@ -822,23 +857,26 @@ func TestValidateAssignedWorker_WorkerOwnership(t *testing.T) {
ctx := context.Background()
persistence := newTestPersistence(t)
if _, err := persistence.CreateWorker(ctx, &ateapipb.Worker{
Metadata: &ateapipb.ResourceMetadata{Name: testWorkerUID("pod-1")},
WorkerNamespace: "worker-ns",
WorkerPool: "pool",
WorkerPod: "pod-1",
WorkerPodUid: testWorkerUID("pod-1"),
SandboxClass: tt.sandboxClass,
Status: &ateapipb.WorkerStatus{State: ateapipb.WorkerState_WORKER_STATE_ACTIVE, Capacity: &ateapipb.WorkerResources{Actors: 1}},
}); err != nil {
t.Fatalf("CreateWorker: %v", err)
}
seedAssignment(t, persistence, testWorkerUID("pod-1"), tt.assignment)
// Fetch the stored version so the no-write assertion below can
// detect any optimistic update.
seeded, err := persistence.GetWorker(ctx, testWorkerUID("pod-1"))
if err != nil {
t.Fatalf("GetWorker: %v", err)
var seeded *ateapipb.Worker
if tt.workerStatus != nil {
if _, err := persistence.CreateWorker(ctx, &ateapipb.Worker{
Metadata: &ateapipb.ResourceMetadata{Name: testWorkerUID("pod-1")},
WorkerNamespace: "worker-ns",
WorkerPool: "pool",
WorkerPod: "pod-1",
WorkerPodUid: testWorkerUID("pod-1"),
SandboxClass: tt.sandboxClass,
Status: tt.workerStatus,
}); err != nil {
t.Fatalf("CreateWorker: %v", err)
}
seedAssignment(t, persistence, testWorkerUID("pod-1"), tt.assignment)
// Fetch the stored version so the no-write assertion below can
// detect any optimistic update.
var err error
if seeded, err = persistence.GetWorker(ctx, testWorkerUID("pod-1")); err != nil {
t.Fatalf("GetWorker: %v", err)
}
}
seedWorkflowActor(t, ctx, persistence, resources.ActorRef{Atespace: "team-a", Name: "shared"}, "ns", "tmpl1", ateapipb.ActorState_ACTOR_STATE_RESUMING)
@@ -858,7 +896,7 @@ func TestValidateAssignedWorker_WorkerOwnership(t *testing.T) {
},
}
tmpl := &ateapipb.ActorTemplate{SandboxConfig: &ateapipb.SandboxConfig{SandboxClass: ateapipb.SandboxClass_SANDBOX_CLASS_GVISOR}}
_, err = w.validateAssignedWorker(ctx, resources.ActorRef{Atespace: "team-a", Name: "shared"}, resumingActor, tmpl)
_, err := w.validateAssignedWorker(ctx, resources.ActorRef{Atespace: "team-a", Name: "shared"}, resumingActor, tmpl)
if got := status.Code(err); got != tt.wantCode {
t.Fatalf("status.Code(err) = %v, want %v (err: %v)", got, tt.wantCode, err)
}
@@ -870,7 +908,13 @@ func TestValidateAssignedWorker_WorkerOwnership(t *testing.T) {
if actor.GetStatus().GetState() != tt.wantActorState {
t.Errorf("stored actor state = %v, want %v", actor.GetStatus().GetState(), tt.wantActorState)
}
if msg := actor.GetStatus().GetCrash().GetMessage(); msg != tt.wantCrashMessage {
t.Errorf("crash message = %q, want %q", msg, tt.wantCrashMessage)
}
if tt.workerStatus == nil {
return
}
stored, err := persistence.GetWorker(ctx, testWorkerUID("pod-1"))
if err != nil {
t.Fatalf("GetWorker: %v", err)
@@ -185,7 +185,7 @@ func (w *ActorWorkflow) ensureWorkerDiscarded(ctx context.Context, actorRef reso
// can revert again — that retry needs no live worker.
slog.LogAttrs(ctx, slog.LevelError, "Setting Actor to crashed due to error",
append(ateattr.ActorRefLogAttrs(actorRef), slog.Any("err", terr))...)
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationRevert); cerr != nil {
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationRevert, ateletCrashMessage("Terminate", terr)); cerr != nil {
return cerr
}
return fmt.Errorf("actor %s crashed: %w", actorRef, terr)
@@ -256,6 +256,7 @@ func (w *ActorWorkflow) ensureRevertedFinalized(ctx context.Context, actorRef re
toUpdate.Status.InProgressSnapshotUri = ""
toUpdate.Status.InProgressLocalSnapshotName = ""
toUpdate.Status.LocalSnapshot = nil
toUpdate.Status.Crash = nil
return nil
})
if err != nil {
@@ -203,7 +203,9 @@ func TestRevertActor_NoSnapshotToRevertTo(t *testing.T) {
w := newTestActorWorkflow(t, st, "ns", "tmpl1")
actorRef := resources.ActorRef{Atespace: "team-a", Name: "id1"}
seedWorkflowActor(t, ctx, st, actorRef, "ns", "tmpl1", ateapipb.ActorState_ACTOR_STATE_CRASHED)
seedWorkflowActor(t, ctx, st, actorRef, "ns", "tmpl1", ateapipb.ActorState_ACTOR_STATE_CRASHED, func(a *ateapipb.Actor) {
a.Status.Crash = newActorCrash("resume", crashMessageWorkerGone)
})
reverted, err := w.RevertActor(ctx, actorRef)
if err != nil {
@@ -212,6 +214,9 @@ func TestRevertActor_NoSnapshotToRevertTo(t *testing.T) {
if got := reverted.GetStatus().GetState(); got != ateapipb.ActorState_ACTOR_STATE_SUSPENDED {
t.Errorf("state = %v, want SUSPENDED", got)
}
if got := reverted.GetStatus().GetCrash(); got != nil {
t.Errorf("Crash = %v, want cleared by the revert", got)
}
if got := reverted.GetStatus().GetExternalSnapshot().GetSnapshotUri(); got != "" {
t.Errorf("external snapshot = %q, want none", got)
}
@@ -213,7 +213,7 @@ func (w *ActorWorkflow) ensureAteletSuspended(ctx context.Context, actorRef reso
assignment := actor.GetStatus().GetWorkerAssignment()
if assignment == nil {
// Missing active worker pod reference in SUSPENDING state indicates corrupted store state.
if err := crashActor(ctx, w.store, actorRef, ateattr.OperationSuspend); err != nil {
if err := crashActor(ctx, w.store, actorRef, ateattr.OperationSuspend, crashMessageWorkerAssignmentMissing); err != nil {
slog.ErrorContext(ctx, "Failed to crash actor", slog.String("err", err.Error()))
}
return "", fmt.Errorf("actor is CRASHED because it was in SUSPENDING state but has no active worker")
@@ -254,7 +254,7 @@ func (w *ActorWorkflow) ensureAteletSuspended(ctx context.Context, actorRef reso
if _, err = client.Checkpoint(ctx, req); err != nil {
slog.LogAttrs(ctx, slog.LevelError, "Setting Actor to crashed due to error",
append(ateattr.ActorRefLogAttrs(actorRef), slog.Any("err", err))...)
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationSuspend); cerr != nil {
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationSuspend, ateletCrashMessage("Checkpoint", err)); cerr != nil {
return wireSnapshotScope, cerr
}
return wireSnapshotScope, fmt.Errorf("actor %s crashed: %w", actorRef, err)
@@ -276,7 +276,7 @@ func (w *ActorWorkflow) ensurePausedSnapshotUploaded(ctx context.Context, actorR
if len(local.GetNodeVmsWithLocalSnapshots()) == 0 {
// Without the node the snapshot can never be found (mirrors
// FinalizePaused, which crashes rather than record an unknown node).
if err := crashActor(ctx, w.store, actorRef, ateattr.OperationSuspend); err != nil {
if err := crashActor(ctx, w.store, actorRef, ateattr.OperationSuspend, crashMessageLocalSnapshotNodeUnknown); err != nil {
slog.ErrorContext(ctx, "Failed to crash actor", slog.String("err", err.Error()))
}
return "", fmt.Errorf("actor is CRASHED because it was suspending a paused snapshot with no node recorded")
@@ -308,7 +308,7 @@ func (w *ActorWorkflow) ensurePausedSnapshotUploaded(ctx context.Context, actorR
if _, err = client.UploadPausedCheckpoint(ctx, req); err != nil {
slog.LogAttrs(ctx, slog.LevelError, "Setting Actor to crashed due to error",
append(ateattr.ActorRefLogAttrs(actorRef), slog.Any("err", err))...)
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationSuspend); cerr != nil {
if cerr := crashActor(ctx, w.store, actorRef, ateattr.OperationSuspend, ateletCrashMessage("UploadPausedCheckpoint", err)); cerr != nil {
return wireSnapshotScope, cerr
}
return wireSnapshotScope, fmt.Errorf("actor %s crashed: %w", actorRef, err)
@@ -192,6 +192,9 @@ func TestSuspendActor_CrashesWhenSuspendingActorMissingWorkerPod(t *testing.T) {
if got.GetStatus().GetState() != ateapipb.ActorState_ACTOR_STATE_CRASHED {
t.Errorf("stored state = %v, want %v", got.GetStatus().GetState(), ateapipb.ActorState_ACTOR_STATE_CRASHED)
}
if msg, want := got.GetStatus().GetCrash().GetMessage(), "suspend failed: "+crashMessageWorkerAssignmentMissing; msg != want {
t.Errorf("crash message = %q, want %q", msg, want)
}
}
// newTestPersistence returns an isolated PostgreSQL-backed store.
@@ -735,4 +738,7 @@ func TestSuspendActor_PausedWithoutLocalSnapshotCrashes(t *testing.T) {
if got.GetStatus().GetState() != ateapipb.ActorState_ACTOR_STATE_CRASHED {
t.Errorf("stored state = %v, want %v", got.GetStatus().GetState(), ateapipb.ActorState_ACTOR_STATE_CRASHED)
}
if msg, want := got.GetStatus().GetCrash().GetMessage(), "suspend failed: "+crashMessageLocalSnapshotNodeUnknown; msg != want {
t.Errorf("crash message = %q, want %q", msg, want)
}
}
@@ -214,6 +214,9 @@ func (w *WorkerWorkflow) releaseBoundActor(ctx context.Context, worker *ateapipb
slog.String("worker", name))...)
_, err = w.store.UpdateActor(ctx, actorRef, store.PreconditionFrom(actor), func(toUpdate *ateapipb.Actor) error {
toUpdate.Status.State = ateapipb.ActorState_ACTOR_STATE_CRASHED
if !wasAlreadyCrashed {
toUpdate.Status.Crash = newActorCrash(opName, crashMessageWorkerPodGone)
}
toUpdate.Status.WorkerAssignment = nil
// Local in-progress checkpoint dies with the worker: it lived on the node
// that went away. The external in-progress checkpoint is kept so delete
@@ -126,6 +126,9 @@ func TestDeleteWorkerWorkflow_ReleasesBoundActor(t *testing.T) {
if got.GetStatus().GetState() != ateapipb.ActorState_ACTOR_STATE_CRASHED {
t.Errorf("actor state = %v, want CRASHED: it never suspended cleanly", got.GetStatus().GetState())
}
if msg := got.GetStatus().GetCrash().GetMessage(); msg != crashMessageWorkerPodGone {
t.Errorf("crash message = %q, want %q", msg, crashMessageWorkerPodGone)
}
if got.GetStatus().GetWorkerAssignment() != nil {
t.Errorf("actor worker assignment = %v, want it cleared", got.GetStatus().GetWorkerAssignment())
}
@@ -245,6 +245,74 @@ func Validate_Actor(
return errs
}
// Validate_ActorCrash validates an instance of ActorCrash according
// to declarative validation rules in the API schema.
func Validate_ActorCrash(
ctx context.Context, op operation.Operation, fldPath *field.Path,
obj, oldObj *ateapipb.ActorCrash) (errs field.ErrorList) {
{ // field ateapipb.ActorCrash.Message
fn := func(
fldPath *field.Path,
obj, oldObj *string,
oldValueCorrelated bool) (errs field.ErrorList) {
// don't revalidate unchanged data
if oldValueCorrelated && op.Type == operation.Update {
if obj == oldObj || (obj != nil && oldObj != nil && *obj == *oldObj) {
return nil
}
}
// call field-attached validations
earlyReturn := false
if e := validate.OptionalValue(ctx, op, fldPath, obj, oldObj).MarkShortCircuit(); len(e) != 0 {
earlyReturn = true
}
if earlyReturn {
return // do not proceed
}
if e := validate.MaxLength(ctx, op, fldPath, obj, oldObj, 4096); len(e) != 0 {
errs = append(errs, e...)
}
return
}
oldVal := safe.Field(oldObj,
func(oldObj *ateapipb.ActorCrash) *string {
return &oldObj.Message
})
errs = append(errs, fn(fldPath.Child("message"), &obj.Message, oldVal, oldObj != nil)...)
}
{ // field ateapipb.ActorCrash.CrashTime
fn := func(
fldPath *field.Path,
obj, oldObj *timestamppb.Timestamp,
oldValueCorrelated bool) (errs field.ErrorList) {
// don't revalidate unchanged data
if oldValueCorrelated && op.Type == operation.Update {
if ateDeepEqual(obj, oldObj) {
return nil
}
}
// call field-attached validations
earlyReturn := false
if e := validate.OptionalPointer(ctx, op, fldPath, obj, oldObj).MarkShortCircuit(); len(e) != 0 {
earlyReturn = true
}
if earlyReturn {
return // do not proceed
}
return
}
oldVal := safe.Field(oldObj,
func(oldObj *ateapipb.ActorCrash) *timestamppb.Timestamp {
return oldObj.CrashTime
})
errs = append(errs, fn(fldPath.Child("crash_time"), obj.CrashTime, oldVal, oldObj != nil)...)
}
return errs
}
// Validate_ActorMetadataDataSource validates an instance of ActorMetadataDataSource according
// to declarative validation rules in the API schema.
func Validate_ActorMetadataDataSource(
@@ -630,6 +698,36 @@ func Validate_ActorStatus(
errs = append(errs, fn(fldPath.Child("in_progress_local_snapshot_name"), &obj.InProgressLocalSnapshotName, oldVal, oldObj != nil)...)
}
{ // field ateapipb.ActorStatus.Crash
fn := func(
fldPath *field.Path,
obj, oldObj *ateapipb.ActorCrash,
oldValueCorrelated bool) (errs field.ErrorList) {
// don't revalidate unchanged data
if oldValueCorrelated && op.Type == operation.Update {
if ateDeepEqual(obj, oldObj) {
return nil
}
}
// call field-attached validations
earlyReturn := false
if e := validate.OptionalPointer(ctx, op, fldPath, obj, oldObj).MarkShortCircuit(); len(e) != 0 {
earlyReturn = true
}
if earlyReturn {
return // do not proceed
}
// call the type's validation function
errs = append(errs, Validate_ActorCrash(ctx, op, fldPath, obj, oldObj)...)
return
}
oldVal := safe.Field(oldObj,
func(oldObj *ateapipb.ActorStatus) *ateapipb.ActorCrash {
return oldObj.Crash
})
errs = append(errs, fn(fldPath.Child("crash"), obj.Crash, oldVal, oldObj != nil)...)
}
return errs
}
+16
View File
@@ -1606,6 +1606,22 @@ func TestWorkerPodDeletion(t *testing.T) {
t.Logf("Waiting for actor %q to transition to CRASHED...", actorName)
waitForActorState(ctx, t, clients, actorName, ateapipb.ActorState_ACTOR_STATE_CRASHED)
crashed, err := clients.SubstrateAPI.GetActor(ctx, &ateapipb.GetActorRequest{
Actor: &ateapipb.ObjectRef{Atespace: demoAtespace, Name: actorName},
})
if err != nil {
t.Fatalf("failed to get crashed actor %q: %v", actorName, err)
}
crash := crashed.GetStatus().GetCrash()
// The actor was RUNNING, not mid-operation, so the crash message carries no
// "<op> failed: " prefix.
if want := "worker pod went away while hosting the actor"; crash.GetMessage() != want {
t.Errorf("Crash.Message = %q, want %q", crash.GetMessage(), want)
}
if crash.GetCrashTime() == nil {
t.Errorf("Crash.CrashTime = nil, want set")
}
// Verify the worker is cleaned up (deleted) from store
t.Logf("Verifying worker for pod %s/%s is removed from store...", podNamespace, podName)
deadline := time.Now().Add(30 * time.Second)
File diff suppressed because it is too large Load Diff
+20
View File
@@ -591,6 +591,26 @@ message ActorStatus {
// +k8s:optional
// +k8s:format=k8s-short-name
string in_progress_local_snapshot_name = 8;
// crash records why and when the Actor entered CRASHED. It is set with the
// CRASHED state and cleared when a revert returns the Actor to SUSPENDED.
//
// +k8s:optional
ActorCrash crash = 9;
}
// ActorCrash describes the failure that moved an Actor to CRASHED.
message ActorCrash {
// message is a human-readable description of the failure.
//
// +k8s:optional
// +k8s:maxLength=4096
string message = 1;
// crash_time is when the Actor entered CRASHED.
//
// +k8s:optional
google.protobuf.Timestamp crash_time = 2;
}
// WorkerAssignment points at the Worker currently hosting an Actor.