mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
Update RPC error messages to include both actor name and atespace.
Now that we pass actorRef around, make sure we use it to report the full name (atespace + name) of the actor when reporting errors.
This commit is contained in:
committed by
Haven Xia
parent
2dc1fc6266
commit
f83b65d150
@@ -41,7 +41,7 @@ func maybeCrashActor(ctx context.Context, st store.Interface, actorRef resources
|
||||
slog.ErrorContext(ctx, "Failed to crash actor", slog.Any("cerr", cerr))
|
||||
return cerr
|
||||
}
|
||||
return status.Errorf(codes.DataLoss, "actor %s crashed", actorRef.Name)
|
||||
return status.Errorf(codes.DataLoss, "actor %s crashed", actorRef)
|
||||
}
|
||||
return fmt.Errorf("%s: %w", wrapMsg, err)
|
||||
}
|
||||
|
||||
@@ -37,7 +37,7 @@ func (s *Service) DeleteActor(ctx context.Context, req *ateapipb.DeleteActorRequ
|
||||
actor, err := s.persistence.GetActor(ctx, actorRef)
|
||||
if err != nil {
|
||||
if errors.Is(err, store.ErrNotFound) {
|
||||
return nil, status.Errorf(codes.NotFound, "Actor %s not found", actorRef.Name)
|
||||
return nil, status.Errorf(codes.NotFound, "Actor %s not found", actorRef)
|
||||
}
|
||||
return nil, fmt.Errorf("while fetching actor: %w", err)
|
||||
}
|
||||
@@ -50,14 +50,14 @@ func (s *Service) DeleteActor(ctx context.Context, req *ateapipb.DeleteActorRequ
|
||||
deleted, err := s.persistence.DeleteActor(ctx, actorRef)
|
||||
if err != nil {
|
||||
if errors.Is(err, store.ErrNotFound) {
|
||||
return nil, status.Errorf(codes.NotFound, "Actor %s not found", actorRef.Name)
|
||||
return nil, status.Errorf(codes.NotFound, "Actor %s not found", actorRef)
|
||||
}
|
||||
if errors.Is(err, store.ErrFailedPrecondition) {
|
||||
current, getErr := s.persistence.GetActor(ctx, actorRef)
|
||||
if getErr == nil {
|
||||
return nil, status.Errorf(codes.FailedPrecondition, "Actor %s is not suspended (status: %v)", actorRef.Name, current.GetStatus())
|
||||
return nil, status.Errorf(codes.FailedPrecondition, "Actor %s is not suspended (status: %v)", actorRef, current.GetStatus())
|
||||
}
|
||||
return nil, status.Errorf(codes.FailedPrecondition, "Actor %s is not suspended", actorRef.Name)
|
||||
return nil, status.Errorf(codes.FailedPrecondition, "Actor %s is not suspended", actorRef)
|
||||
}
|
||||
if errors.Is(err, store.ErrVersionConflict) {
|
||||
return nil, status.Error(codes.Aborted, "concurrent update conflict, please retry")
|
||||
|
||||
@@ -844,7 +844,7 @@ func TestGetActor_NotFound(t *testing.T) {
|
||||
_, err := tc.client.GetActor(context.Background(), &ateapipb.GetActorRequest{
|
||||
Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "non-existent"},
|
||||
})
|
||||
assertGrpcError(t, err, codes.NotFound, "Actor non-existent not found")
|
||||
assertGrpcError(t, err, codes.NotFound, "Actor test-atespace/non-existent not found")
|
||||
}
|
||||
|
||||
// TestListActors tests that all created actors can be listed.
|
||||
@@ -958,7 +958,7 @@ func TestListActors_ByAtespace(t *testing.T) {
|
||||
t.Errorf("GetActor(id1, team-a) failed: %v", err)
|
||||
}
|
||||
_, err = tc.client.GetActor(context.Background(), &ateapipb.GetActorRequest{Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "id1"}})
|
||||
assertGrpcError(t, err, codes.NotFound, "Actor id1 not found")
|
||||
assertGrpcError(t, err, codes.NotFound, "Actor test-atespace/id1 not found")
|
||||
}
|
||||
|
||||
// TestListActors_AllAtespaces verifies that an empty atespace lists actors across
|
||||
@@ -1735,7 +1735,7 @@ func TestUpdateActor_NotFound(t *testing.T) {
|
||||
defer tc.cleanup()
|
||||
|
||||
_, err := tc.client.UpdateActor(context.Background(), &ateapipb.UpdateActorRequest{Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "does-not-exist"}})
|
||||
assertGrpcError(t, err, codes.NotFound, "Actor does-not-exist not found")
|
||||
assertGrpcError(t, err, codes.NotFound, "Actor test-atespace/does-not-exist not found")
|
||||
}
|
||||
|
||||
// TestResumeActor_ReleasesStaleWorkerWhenPoolBecomesIneligible verifies that
|
||||
@@ -2578,7 +2578,7 @@ func TestDeleteActor_Success(t *testing.T) {
|
||||
_, err = tc.client.GetActor(context.Background(), &ateapipb.GetActorRequest{
|
||||
Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "id1"},
|
||||
})
|
||||
assertGrpcError(t, err, codes.NotFound, "Actor id1 not found")
|
||||
assertGrpcError(t, err, codes.NotFound, "Actor test-atespace/id1 not found")
|
||||
}
|
||||
|
||||
func TestDeleteActor_NotSuspended(t *testing.T) {
|
||||
@@ -2608,7 +2608,7 @@ func TestDeleteActor_NotSuspended(t *testing.T) {
|
||||
_, err = tc.client.DeleteActor(context.Background(), &ateapipb.DeleteActorRequest{
|
||||
Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "id1"},
|
||||
})
|
||||
assertGrpcError(t, err, codes.FailedPrecondition, "Actor id1 is not suspended (status: STATUS_RUNNING)")
|
||||
assertGrpcError(t, err, codes.FailedPrecondition, "Actor test-atespace/id1 is not suspended (status: STATUS_RUNNING)")
|
||||
}
|
||||
|
||||
func TestDeleteActor_Crashed(t *testing.T) {
|
||||
@@ -2649,7 +2649,7 @@ func TestDeleteActor_Crashed(t *testing.T) {
|
||||
_, err = tc.client.GetActor(context.Background(), &ateapipb.GetActorRequest{
|
||||
Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "id1"},
|
||||
})
|
||||
assertGrpcError(t, err, codes.NotFound, "Actor id1 not found")
|
||||
assertGrpcError(t, err, codes.NotFound, "Actor test-atespace/id1 not found")
|
||||
}
|
||||
|
||||
func TestDeleteActor_NotFound(t *testing.T) {
|
||||
@@ -2660,7 +2660,7 @@ func TestDeleteActor_NotFound(t *testing.T) {
|
||||
_, err := tc.client.DeleteActor(context.Background(), &ateapipb.DeleteActorRequest{
|
||||
Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "non-existent"},
|
||||
})
|
||||
assertGrpcError(t, err, codes.NotFound, "Actor non-existent not found")
|
||||
assertGrpcError(t, err, codes.NotFound, "Actor test-atespace/non-existent not found")
|
||||
}
|
||||
|
||||
func assertGrpcErrorRegex(t *testing.T, err error, wantCode codes.Code, wantMsg string) {
|
||||
|
||||
@@ -34,7 +34,7 @@ func (s *Service) GetActor(ctx context.Context, req *ateapipb.GetActorRequest) (
|
||||
actorRef := resources.ActorRefFromObjectRef(req.GetActor())
|
||||
actor, err := s.persistence.GetActor(ctx, actorRef)
|
||||
if errors.Is(err, store.ErrNotFound) {
|
||||
return nil, status.Errorf(codes.NotFound, "Actor %s not found", actorRef.Name)
|
||||
return nil, status.Errorf(codes.NotFound, "Actor %s not found", actorRef)
|
||||
} else if err != nil {
|
||||
return nil, fmt.Errorf("while getting actor from DB: %w", err)
|
||||
}
|
||||
|
||||
@@ -39,7 +39,7 @@ func (s *Service) PauseActor(ctx context.Context, req *ateapipb.PauseActorReques
|
||||
return nil, status.Error(codes.Aborted, "concurrent update conflict, please retry")
|
||||
}
|
||||
if errors.Is(err, store.ErrNotFound) {
|
||||
return nil, status.Errorf(codes.NotFound, "Actor %s not found", actorRef.Name)
|
||||
return nil, status.Errorf(codes.NotFound, "Actor %s not found", actorRef)
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -39,7 +39,7 @@ func (s *Service) ResumeActor(ctx context.Context, req *ateapipb.ResumeActorRequ
|
||||
return nil, status.Error(codes.Aborted, "concurrent update conflict, please retry")
|
||||
}
|
||||
if errors.Is(err, store.ErrNotFound) {
|
||||
return nil, status.Errorf(codes.NotFound, "Actor %s not found", actorRef.Name)
|
||||
return nil, status.Errorf(codes.NotFound, "Actor %s not found", actorRef)
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -39,7 +39,7 @@ func (s *Service) SuspendActor(ctx context.Context, req *ateapipb.SuspendActorRe
|
||||
return nil, status.Error(codes.Aborted, "concurrent update conflict, please retry")
|
||||
}
|
||||
if errors.Is(err, store.ErrNotFound) {
|
||||
return nil, status.Errorf(codes.NotFound, "Actor %s not found", actorRef.Name)
|
||||
return nil, status.Errorf(codes.NotFound, "Actor %s not found", actorRef)
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -36,7 +36,7 @@ func (s *Service) UpdateActor(ctx context.Context, req *ateapipb.UpdateActorRequ
|
||||
actor, err := s.persistence.GetActor(ctx, actorRef)
|
||||
if err != nil {
|
||||
if errors.Is(err, store.ErrNotFound) {
|
||||
return nil, status.Errorf(codes.NotFound, "Actor %s not found", actorRef.Name)
|
||||
return nil, status.Errorf(codes.NotFound, "Actor %s not found", actorRef)
|
||||
}
|
||||
return nil, fmt.Errorf("while getting actor: %w", err)
|
||||
}
|
||||
|
||||
@@ -87,7 +87,7 @@ func (s *MarkPausingStep) IsComplete(ctx context.Context, input *PauseInput, sta
|
||||
func (s *MarkPausingStep) CheckPrerequisite(ctx context.Context, input *PauseInput, state *PauseState) error {
|
||||
// The pause edge only exists from RUNNING; PAUSING/PAUSED are fast-forwarded by IsComplete.
|
||||
if state.Actor.GetStatus() != ateapipb.Actor_STATUS_RUNNING {
|
||||
return status.Errorf(codes.FailedPrecondition, "MarkPausingStep prerequisite not met for Actor: %s (got: %v, want %s)", input.ActorRef.Name, state.Actor.GetStatus(), ateapipb.Actor_STATUS_RUNNING)
|
||||
return status.Errorf(codes.FailedPrecondition, "MarkPausingStep prerequisite not met for Actor: %s (got: %v, want %s)", input.ActorRef, state.Actor.GetStatus(), ateapipb.Actor_STATUS_RUNNING)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -116,13 +116,13 @@ func (s *CallAteletPauseStep) IsComplete(ctx context.Context, input *PauseInput,
|
||||
}
|
||||
func (s *CallAteletPauseStep) CheckPrerequisite(ctx context.Context, input *PauseInput, state *PauseState) error {
|
||||
if state.Actor.GetStatus() != ateapipb.Actor_STATUS_PAUSING {
|
||||
return status.Errorf(codes.FailedPrecondition, "CallAteletPauseStep prerequisite not met for Actor: %s (got: %v, want %s)", input.ActorRef.Name, state.Actor.GetStatus(), ateapipb.Actor_STATUS_PAUSING)
|
||||
return status.Errorf(codes.FailedPrecondition, "CallAteletPauseStep prerequisite not met for Actor: %s (got: %v, want %s)", input.ActorRef, state.Actor.GetStatus(), ateapipb.Actor_STATUS_PAUSING)
|
||||
}
|
||||
if state.Actor.GetAteomPodNamespace() == "" || state.Actor.GetAteomPodName() == "" {
|
||||
if err := crashActor(ctx, s.store, input.ActorRef); err != nil {
|
||||
slog.ErrorContext(ctx, "Failed to crash actor", slog.String("err", err.Error()))
|
||||
}
|
||||
return status.Errorf(codes.FailedPrecondition, "CallAteletPauseStep prerequisite not met for Actor: %s. AteomPodNamespace: %s, GetAteomPodName %s", input.ActorRef.Name, state.Actor.GetAteomPodNamespace(), state.Actor.GetAteomPodName())
|
||||
return status.Errorf(codes.FailedPrecondition, "CallAteletPauseStep prerequisite not met for Actor: %s. AteomPodNamespace: %s, GetAteomPodName %s", input.ActorRef, state.Actor.GetAteomPodNamespace(), state.Actor.GetAteomPodName())
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -210,7 +210,7 @@ func (s *FinalizePausedStep) IsComplete(ctx context.Context, input *PauseInput,
|
||||
|
||||
func (s *FinalizePausedStep) CheckPrerequisite(ctx context.Context, input *PauseInput, state *PauseState) error {
|
||||
if state.Actor.GetStatus() != ateapipb.Actor_STATUS_PAUSING {
|
||||
return status.Errorf(codes.FailedPrecondition, "FinalizePausedStep prerequisite not met for Actor: %s (got: %v, want %s)", input.ActorRef.Name, state.Actor.GetStatus(), ateapipb.Actor_STATUS_PAUSING)
|
||||
return status.Errorf(codes.FailedPrecondition, "FinalizePausedStep prerequisite not met for Actor: %s (got: %v, want %s)", input.ActorRef, state.Actor.GetStatus(), ateapipb.Actor_STATUS_PAUSING)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -68,7 +68,7 @@ func (s *LoadActorForResumeStep) Execute(ctx context.Context, input *ResumeInput
|
||||
actor, err := s.store.GetActor(ctx, input.ActorRef)
|
||||
if err != nil {
|
||||
if errors.Is(err, store.ErrNotFound) {
|
||||
return status.Errorf(codes.NotFound, "Actor %s not found", input.ActorRef.Name)
|
||||
return status.Errorf(codes.NotFound, "Actor %s not found", input.ActorRef)
|
||||
}
|
||||
return fmt.Errorf("while getting actor from DB: %w", err)
|
||||
}
|
||||
@@ -94,7 +94,7 @@ func (s *LoadActorForResumeStep) Execute(ctx context.Context, input *ResumeInput
|
||||
if cerr := crashActor(ctx, s.store, input.ActorRef); cerr != nil {
|
||||
return cerr
|
||||
}
|
||||
return status.Errorf(codes.Aborted, "actor %s crashed", input.ActorRef.Name)
|
||||
return status.Errorf(codes.Aborted, "actor %s crashed", input.ActorRef)
|
||||
}
|
||||
|
||||
wk, err := s.store.GetWorker(ctx, actor.AteomPodNamespace, actor.WorkerPoolName, actor.AteomPodName)
|
||||
@@ -104,7 +104,7 @@ func (s *LoadActorForResumeStep) Execute(ctx context.Context, input *ResumeInput
|
||||
if cerr := crashActor(ctx, s.store, input.ActorRef); cerr != nil {
|
||||
return cerr
|
||||
}
|
||||
return status.Errorf(codes.Aborted, "actor %s crashed", input.ActorRef.Name)
|
||||
return status.Errorf(codes.Aborted, "actor %s crashed", input.ActorRef)
|
||||
}
|
||||
return fmt.Errorf("failed to get already assigned worker for actor %w", err)
|
||||
}
|
||||
@@ -133,7 +133,7 @@ func (s *AssignWorkerStep) CheckPrerequisite(ctx context.Context, input *ResumeI
|
||||
case ateapipb.Actor_STATUS_SUSPENDED, ateapipb.Actor_STATUS_PAUSED:
|
||||
return nil
|
||||
default:
|
||||
return status.Errorf(codes.FailedPrecondition, "AssignWorkerStep prerequisite not met for Actor: %s (got: %v, want %s or %s)", input.ActorRef.Name, state.Actor.GetStatus(), ateapipb.Actor_STATUS_SUSPENDED, ateapipb.Actor_STATUS_PAUSED)
|
||||
return status.Errorf(codes.FailedPrecondition, "AssignWorkerStep prerequisite not met for Actor: %s (got: %v, want %s or %s)", input.ActorRef, state.Actor.GetStatus(), ateapipb.Actor_STATUS_SUSPENDED, ateapipb.Actor_STATUS_PAUSED)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -234,7 +234,7 @@ func (s *AssignWorkerStep) Execute(ctx context.Context, input *ResumeInput, stat
|
||||
state.Actor = fresh
|
||||
return err
|
||||
default:
|
||||
return status.Errorf(codes.Aborted, "actor %s is %s and can no longer be resumed", input.ActorRef.Name, fresh.GetStatus())
|
||||
return status.Errorf(codes.Aborted, "actor %s is %s and can no longer be resumed", input.ActorRef, fresh.GetStatus())
|
||||
}
|
||||
}
|
||||
state.Actor = updatedActor
|
||||
@@ -326,7 +326,7 @@ func (s *CallAteletRestoreStep) IsComplete(ctx context.Context, input *ResumeInp
|
||||
}
|
||||
func (s *CallAteletRestoreStep) CheckPrerequisite(ctx context.Context, input *ResumeInput, state *ResumeState) error {
|
||||
if state.Actor.GetStatus() != ateapipb.Actor_STATUS_RESUMING {
|
||||
return status.Errorf(codes.FailedPrecondition, "CallAteletRestoreStep prerequisite not met for Actor: %s (got: %v, want %s)", input.ActorRef.Name, state.Actor.GetStatus(), ateapipb.Actor_STATUS_RESUMING)
|
||||
return status.Errorf(codes.FailedPrecondition, "CallAteletRestoreStep prerequisite not met for Actor: %s (got: %v, want %s)", input.ActorRef, state.Actor.GetStatus(), ateapipb.Actor_STATUS_RESUMING)
|
||||
}
|
||||
if state.Worker == nil {
|
||||
return status.Errorf(codes.FailedPrecondition, "Assigned worker is nil")
|
||||
@@ -340,7 +340,7 @@ func (s *CallAteletRestoreStep) CheckPrerequisite(ctx context.Context, input *Re
|
||||
if cerr := crashActor(ctx, s.store, input.ActorRef); cerr != nil {
|
||||
return fmt.Errorf("while crashing actor: %w", cerr)
|
||||
}
|
||||
return status.Errorf(codes.Aborted, "actor %s crashed", input.ActorRef.Name)
|
||||
return status.Errorf(codes.Aborted, "actor %s crashed", input.ActorRef)
|
||||
}
|
||||
constraints, err := schedulingConstraints(state.Actor, state.ActorTemplate)
|
||||
if err != nil {
|
||||
@@ -360,7 +360,7 @@ func (s *CallAteletRestoreStep) CheckPrerequisite(ctx context.Context, input *Re
|
||||
if cerr := crashActor(ctx, s.store, input.ActorRef); cerr != nil {
|
||||
return fmt.Errorf("while crashing actor: %w", cerr)
|
||||
}
|
||||
return status.Errorf(codes.Aborted, "actor %s crashed", input.ActorRef.Name)
|
||||
return status.Errorf(codes.Aborted, "actor %s crashed", input.ActorRef)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -473,7 +473,7 @@ func (s *FinalizeRunningStep) IsComplete(ctx context.Context, input *ResumeInput
|
||||
}
|
||||
func (s *FinalizeRunningStep) CheckPrerequisite(ctx context.Context, input *ResumeInput, state *ResumeState) error {
|
||||
if state.Actor.GetStatus() != ateapipb.Actor_STATUS_RESUMING {
|
||||
return status.Errorf(codes.FailedPrecondition, "FinalizeRunningStep prerequisite not met for Actor: %s (got: %v, want %s)", input.ActorRef.Name, state.Actor.GetStatus(), ateapipb.Actor_STATUS_RESUMING)
|
||||
return status.Errorf(codes.FailedPrecondition, "FinalizeRunningStep prerequisite not met for Actor: %s (got: %v, want %s)", input.ActorRef, state.Actor.GetStatus(), ateapipb.Actor_STATUS_RESUMING)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -350,7 +350,7 @@ func TestAssignWorkerStep_ConflictRefreshesActor(t *testing.T) {
|
||||
|
||||
var injected *ateapipb.Actor
|
||||
st := &conflictInjectingStore{Interface: persistence, inject: func() {
|
||||
fresh, err := persistence.GetActor(ctx, "team-a", "id1")
|
||||
fresh, err := persistence.GetActor(ctx, resources.ActorRef{Atespace: "team-a", Name: "id1"})
|
||||
if err != nil {
|
||||
t.Errorf("inject GetActor: %v", err)
|
||||
return
|
||||
@@ -369,7 +369,7 @@ func TestAssignWorkerStep_ConflictRefreshesActor(t *testing.T) {
|
||||
Spec: atev1alpha1.ActorTemplateSpec{SandboxClass: atev1alpha1.SandboxClassGvisor},
|
||||
},
|
||||
}
|
||||
err := step.Execute(ctx, &ResumeInput{ActorName: "id1", Atespace: "team-a"}, state)
|
||||
err := step.Execute(ctx, &ResumeInput{ActorRef: resources.ActorRef{Atespace: "team-a", Name: "id1"}}, state)
|
||||
|
||||
if tc.wantRetry {
|
||||
if !errors.Is(err, store.ErrVersionConflict) {
|
||||
@@ -387,7 +387,7 @@ func TestAssignWorkerStep_ConflictRefreshesActor(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
stored, err := persistence.GetActor(ctx, "team-a", "id1")
|
||||
stored, err := persistence.GetActor(ctx, resources.ActorRef{Atespace: "team-a", Name: "id1"})
|
||||
if err != nil {
|
||||
t.Fatalf("GetActor: %v", err)
|
||||
}
|
||||
|
||||
@@ -87,7 +87,7 @@ func (s *MarkSuspendingStep) IsComplete(ctx context.Context, input *SuspendInput
|
||||
}
|
||||
func (s *MarkSuspendingStep) CheckPrerequisite(ctx context.Context, input *SuspendInput, state *SuspendState) error {
|
||||
if state.Actor.GetStatus() != ateapipb.Actor_STATUS_RUNNING {
|
||||
return status.Errorf(codes.FailedPrecondition, "MarkSuspendingStep prerequisite not met for Actor: %s (got: %v, want %s)", input.ActorRef.Name, state.Actor.GetStatus(), ateapipb.Actor_STATUS_RUNNING)
|
||||
return status.Errorf(codes.FailedPrecondition, "MarkSuspendingStep prerequisite not met for Actor: %s (got: %v, want %s)", input.ActorRef, state.Actor.GetStatus(), ateapipb.Actor_STATUS_RUNNING)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -117,7 +117,7 @@ func (s *CallAteletSuspendStep) IsComplete(ctx context.Context, input *SuspendIn
|
||||
}
|
||||
func (s *CallAteletSuspendStep) CheckPrerequisite(ctx context.Context, input *SuspendInput, state *SuspendState) error {
|
||||
if state.Actor.GetStatus() != ateapipb.Actor_STATUS_SUSPENDING {
|
||||
return status.Errorf(codes.FailedPrecondition, "CallAteletSuspendStep prerequisite not met for Actor: %s (got: %v, want %s)", input.ActorRef.Name, state.Actor.GetStatus(), ateapipb.Actor_STATUS_SUSPENDING)
|
||||
return status.Errorf(codes.FailedPrecondition, "CallAteletSuspendStep prerequisite not met for Actor: %s (got: %v, want %s)", input.ActorRef, state.Actor.GetStatus(), ateapipb.Actor_STATUS_SUSPENDING)
|
||||
}
|
||||
if state.Actor.GetAteomPodNamespace() == "" || state.Actor.GetAteomPodName() == "" {
|
||||
if err := crashActor(ctx, s.store, input.ActorRef); err != nil {
|
||||
@@ -203,7 +203,7 @@ func (s *FinalizeSuspendedStep) IsComplete(ctx context.Context, input *SuspendIn
|
||||
}
|
||||
func (s *FinalizeSuspendedStep) CheckPrerequisite(ctx context.Context, input *SuspendInput, state *SuspendState) error {
|
||||
if state.Actor.GetStatus() != ateapipb.Actor_STATUS_SUSPENDING {
|
||||
return status.Errorf(codes.FailedPrecondition, "FinalizeSuspendedStep prerequisite not met for Actor: %s (got: %v, want %s)", input.ActorRef.Name, state.Actor.GetStatus(), ateapipb.Actor_STATUS_SUSPENDING)
|
||||
return status.Errorf(codes.FailedPrecondition, "FinalizeSuspendedStep prerequisite not met for Actor: %s (got: %v, want %s)", input.ActorRef, state.Actor.GetStatus(), ateapipb.Actor_STATUS_SUSPENDING)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -17,6 +17,7 @@ package router
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"github.com/agent-substrate/substrate/internal/resources"
|
||||
envoy_type "github.com/envoyproxy/go-control-plane/envoy/type/v3"
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
@@ -32,8 +33,8 @@ func newReqError(code envoy_type.StatusCode, format string, args ...any) error {
|
||||
}
|
||||
|
||||
// actorNotFoundErr returns a 404 reqError identifying the missing actor.
|
||||
func actorNotFoundErr(actorName string) error {
|
||||
return newReqError(envoy_type.StatusCode_NotFound, "actor %q not found", actorName)
|
||||
func actorNotFoundErr(actorRef resources.ActorRef) error {
|
||||
return newReqError(envoy_type.StatusCode_NotFound, "actor %s not found", actorRef)
|
||||
}
|
||||
|
||||
// invalidHostErr returns a 404 reqError explaining why the request host was
|
||||
@@ -53,7 +54,7 @@ func invalidHostErr(host string, cause error) error {
|
||||
//
|
||||
// Unrecognized errors collapse to 500 with a generic body to avoid leaking
|
||||
// server-side detail (stack traces, internal IDs) to clients.
|
||||
func mapResumeError(actorName string, err error) error {
|
||||
func mapResumeError(actorRef resources.ActorRef, err error) error {
|
||||
if err == nil {
|
||||
return nil
|
||||
}
|
||||
@@ -62,31 +63,31 @@ func mapResumeError(actorName string, err error) error {
|
||||
switch status.Code(err) {
|
||||
case codes.NotFound:
|
||||
re.statusCode = int(envoy_type.StatusCode_NotFound)
|
||||
re.msg = fmt.Sprintf("actor %q not found", actorName)
|
||||
re.msg = fmt.Sprintf("actor %s not found", actorRef)
|
||||
case codes.FailedPrecondition:
|
||||
// Preserve the gRPC description for FailedPrecondition only: it carries
|
||||
// actionable client-facing context (e.g. "no free workers available")
|
||||
// and is not security-sensitive.
|
||||
re.statusCode = int(envoy_type.StatusCode_ServiceUnavailable)
|
||||
re.msg = fmt.Sprintf("actor %q unavailable: %s", actorName, status.Convert(err).Message())
|
||||
re.msg = fmt.Sprintf("actor %s unavailable: %s", actorRef, status.Convert(err).Message())
|
||||
case codes.Unavailable:
|
||||
re.statusCode = int(envoy_type.StatusCode_ServiceUnavailable)
|
||||
re.msg = fmt.Sprintf("actor %q unavailable", actorName)
|
||||
re.msg = fmt.Sprintf("actor %s unavailable", actorRef)
|
||||
case codes.DeadlineExceeded:
|
||||
re.statusCode = int(envoy_type.StatusCode_GatewayTimeout)
|
||||
re.msg = fmt.Sprintf("actor %q request timed out", actorName)
|
||||
re.msg = fmt.Sprintf("actor %s request timed out", actorRef)
|
||||
case codes.PermissionDenied:
|
||||
re.statusCode = int(envoy_type.StatusCode_Forbidden)
|
||||
re.msg = fmt.Sprintf("actor %q access denied", actorName)
|
||||
re.msg = fmt.Sprintf("actor %s access denied", actorRef)
|
||||
case codes.Unauthenticated:
|
||||
re.statusCode = int(envoy_type.StatusCode_Unauthorized)
|
||||
re.msg = fmt.Sprintf("actor %q authentication required", actorName)
|
||||
re.msg = fmt.Sprintf("actor %s authentication required", actorRef)
|
||||
case codes.ResourceExhausted:
|
||||
re.statusCode = int(envoy_type.StatusCode_TooManyRequests)
|
||||
re.msg = fmt.Sprintf("actor %q rate limited", actorName)
|
||||
re.msg = fmt.Sprintf("actor %s rate limited", actorRef)
|
||||
default:
|
||||
re.statusCode = int(envoy_type.StatusCode_InternalServerError)
|
||||
re.msg = fmt.Sprintf("error resuming actor %q", actorName)
|
||||
re.msg = fmt.Sprintf("error resuming actor %s", actorRef)
|
||||
}
|
||||
return re
|
||||
}
|
||||
|
||||
@@ -18,6 +18,7 @@ import (
|
||||
"errors"
|
||||
"testing"
|
||||
|
||||
"github.com/agent-substrate/substrate/internal/resources"
|
||||
envoy_type "github.com/envoyproxy/go-control-plane/envoy/type/v3"
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
@@ -45,7 +46,7 @@ func TestNewReqError(t *testing.T) {
|
||||
func TestActorNotFoundErr(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
err := actorNotFoundErr("ctr6")
|
||||
err := actorNotFoundErr(resources.ActorRef{Atespace: "team-a", Name: "ctr6"})
|
||||
var reqErr *reqError
|
||||
if !errors.As(err, &reqErr) {
|
||||
t.Fatalf("errors.As(*reqError) = false, want true; err type = %T", err)
|
||||
@@ -53,7 +54,7 @@ func TestActorNotFoundErr(t *testing.T) {
|
||||
if reqErr.statusCode != int(envoy_type.StatusCode_NotFound) {
|
||||
t.Errorf("statusCode = %d, want %d", reqErr.statusCode, envoy_type.StatusCode_NotFound)
|
||||
}
|
||||
if got, want := err.Error(), `actor "ctr6" not found`; got != want {
|
||||
if got, want := err.Error(), `actor team-a/ctr6 not found`; got != want {
|
||||
t.Errorf("Error() = %q, want %q", got, want)
|
||||
}
|
||||
}
|
||||
@@ -82,7 +83,7 @@ func TestInvalidHostErr(t *testing.T) {
|
||||
func TestMapResumeError(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
const actorName = "ctr6"
|
||||
actorRef := resources.ActorRef{Atespace: "team-a", Name: "ctr6"}
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
@@ -94,55 +95,55 @@ func TestMapResumeError(t *testing.T) {
|
||||
name: "NotFound maps to 404",
|
||||
err: status.Error(codes.NotFound, "actor not found"),
|
||||
wantCode: envoy_type.StatusCode_NotFound,
|
||||
wantBody: `actor "ctr6" not found`,
|
||||
wantBody: `actor team-a/ctr6 not found`,
|
||||
},
|
||||
{
|
||||
name: "FailedPrecondition maps to 503 and preserves desc",
|
||||
err: status.Error(codes.FailedPrecondition, "no free workers available"),
|
||||
wantCode: envoy_type.StatusCode_ServiceUnavailable,
|
||||
wantBody: `actor "ctr6" unavailable: no free workers available`,
|
||||
wantBody: `actor team-a/ctr6 unavailable: no free workers available`,
|
||||
},
|
||||
{
|
||||
name: "Unavailable maps to 503",
|
||||
err: status.Error(codes.Unavailable, "control-plane down"),
|
||||
wantCode: envoy_type.StatusCode_ServiceUnavailable,
|
||||
wantBody: `actor "ctr6" unavailable`,
|
||||
wantBody: `actor team-a/ctr6 unavailable`,
|
||||
},
|
||||
{
|
||||
name: "DeadlineExceeded maps to 504",
|
||||
err: status.Error(codes.DeadlineExceeded, "context deadline exceeded"),
|
||||
wantCode: envoy_type.StatusCode_GatewayTimeout,
|
||||
wantBody: `actor "ctr6" request timed out`,
|
||||
wantBody: `actor team-a/ctr6 request timed out`,
|
||||
},
|
||||
{
|
||||
name: "PermissionDenied maps to 403",
|
||||
err: status.Error(codes.PermissionDenied, "denied"),
|
||||
wantCode: envoy_type.StatusCode_Forbidden,
|
||||
wantBody: `actor "ctr6" access denied`,
|
||||
wantBody: `actor team-a/ctr6 access denied`,
|
||||
},
|
||||
{
|
||||
name: "Unauthenticated maps to 401",
|
||||
err: status.Error(codes.Unauthenticated, "no creds"),
|
||||
wantCode: envoy_type.StatusCode_Unauthorized,
|
||||
wantBody: `actor "ctr6" authentication required`,
|
||||
wantBody: `actor team-a/ctr6 authentication required`,
|
||||
},
|
||||
{
|
||||
name: "ResourceExhausted maps to 429",
|
||||
err: status.Error(codes.ResourceExhausted, "quota"),
|
||||
wantCode: envoy_type.StatusCode_TooManyRequests,
|
||||
wantBody: `actor "ctr6" rate limited`,
|
||||
wantBody: `actor team-a/ctr6 rate limited`,
|
||||
},
|
||||
{
|
||||
name: "unknown gRPC code maps to 500 without leaking desc",
|
||||
err: status.Error(codes.Internal, "stack trace: foo bar"),
|
||||
wantCode: envoy_type.StatusCode_InternalServerError,
|
||||
wantBody: `error resuming actor "ctr6"`,
|
||||
wantBody: `error resuming actor team-a/ctr6`,
|
||||
},
|
||||
{
|
||||
name: "non-gRPC error maps to 500 without leaking message",
|
||||
err: errors.New("raw error with secret"),
|
||||
wantCode: envoy_type.StatusCode_InternalServerError,
|
||||
wantBody: `error resuming actor "ctr6"`,
|
||||
wantBody: `error resuming actor team-a/ctr6`,
|
||||
},
|
||||
}
|
||||
|
||||
@@ -150,7 +151,7 @@ func TestMapResumeError(t *testing.T) {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
got := mapResumeError(actorName, tc.err)
|
||||
got := mapResumeError(actorRef, tc.err)
|
||||
if got == nil {
|
||||
t.Fatal("mapResumeError returned nil")
|
||||
}
|
||||
@@ -176,7 +177,7 @@ func TestMapResumeError_NilError(t *testing.T) {
|
||||
|
||||
// Guard against accidental nil-error calls. Returning nil keeps the
|
||||
// happy path explicit at callsites instead of constructing a bogus 500.
|
||||
if got := mapResumeError("ctr6", nil); got != nil {
|
||||
if got := mapResumeError(resources.ActorRef{Atespace: "team-a", Name: "ctr6"}, nil); got != nil {
|
||||
t.Errorf("mapResumeError(_, nil) = %v, want nil", got)
|
||||
}
|
||||
}
|
||||
@@ -186,7 +187,7 @@ func TestMapResumeError_NilError(t *testing.T) {
|
||||
func TestMapResumeError_IsReqError(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
err := mapResumeError("x", status.Error(codes.NotFound, "x"))
|
||||
err := mapResumeError(resources.ActorRef{Atespace: "team-a", Name: "x"}, status.Error(codes.NotFound, "x"))
|
||||
var reqErr *reqError
|
||||
if !errors.As(err, &reqErr) {
|
||||
t.Fatalf("errors.As(*reqError) = false, want true; err type = %T", err)
|
||||
|
||||
@@ -151,7 +151,7 @@ func (s *ExtProcServer) handleRequestHeaders(
|
||||
slog.InfoContext(ctx, "ResumeActor", slog.Any("actor", actorRef))
|
||||
actor, err := s.resumer.ResumeActor(ctx, actorRef)
|
||||
if err != nil {
|
||||
return nil, metadata, "", "", "", mapResumeError(actorRef.Name, err)
|
||||
return nil, metadata, "", "", "", mapResumeError(actorRef, err)
|
||||
}
|
||||
|
||||
// Actor template identity, used as low-cardinality route-latency metric
|
||||
@@ -167,7 +167,7 @@ func (s *ExtProcServer) handleRequestHeaders(
|
||||
|
||||
if ip := net.ParseIP(workerIP); ip == nil {
|
||||
return nil, metadata, "", tmplNs, tmplName, newReqError(envoy_type.StatusCode_InternalServerError,
|
||||
"actor %q routing failed", actorRef.Name)
|
||||
"actor %s routing failed", actorRef)
|
||||
}
|
||||
|
||||
// TODO(bowei) -- handle more than port 80 on the actor.
|
||||
|
||||
@@ -115,7 +115,7 @@ func TestExtProcHeadersEvaluation(t *testing.T) {
|
||||
authority: testUUID + ".team-a.actors.resources.substrate.ate.dev",
|
||||
resumeErr: errors.New("resume failed with sensitive detail"),
|
||||
expectErr: true,
|
||||
expectedErrStr: `error resuming actor "123e4567-e89b-12d3-a456-426614174000"`,
|
||||
expectedErrStr: `error resuming actor team-a/123e4567-e89b-12d3-a456-426614174000`,
|
||||
expectedStatus: envoy_type.StatusCode_InternalServerError,
|
||||
},
|
||||
{
|
||||
@@ -123,7 +123,7 @@ func TestExtProcHeadersEvaluation(t *testing.T) {
|
||||
authority: testUUID + ".team-a.actors.resources.substrate.ate.dev",
|
||||
resumeErr: status.Error(codes.FailedPrecondition, "no free workers available"),
|
||||
expectErr: true,
|
||||
expectedErrStr: `actor "123e4567-e89b-12d3-a456-426614174000" unavailable: no free workers available`,
|
||||
expectedErrStr: `actor team-a/123e4567-e89b-12d3-a456-426614174000 unavailable: no free workers available`,
|
||||
expectedStatus: envoy_type.StatusCode_ServiceUnavailable,
|
||||
},
|
||||
{
|
||||
@@ -131,7 +131,7 @@ func TestExtProcHeadersEvaluation(t *testing.T) {
|
||||
authority: testUUID + ".team-a.actors.resources.substrate.ate.dev",
|
||||
resumeErr: status.Error(codes.NotFound, "actor missing"),
|
||||
expectErr: true,
|
||||
expectedErrStr: `actor "123e4567-e89b-12d3-a456-426614174000" not found`,
|
||||
expectedErrStr: `actor team-a/123e4567-e89b-12d3-a456-426614174000 not found`,
|
||||
expectedStatus: envoy_type.StatusCode_NotFound,
|
||||
},
|
||||
{
|
||||
@@ -139,7 +139,7 @@ func TestExtProcHeadersEvaluation(t *testing.T) {
|
||||
authority: testUUID + ".team-a.actors.resources.substrate.ate.dev",
|
||||
resumeErr: status.Error(codes.Unavailable, "control-plane down"),
|
||||
expectErr: true,
|
||||
expectedErrStr: `actor "123e4567-e89b-12d3-a456-426614174000" unavailable`,
|
||||
expectedErrStr: `actor team-a/123e4567-e89b-12d3-a456-426614174000 unavailable`,
|
||||
expectedStatus: envoy_type.StatusCode_ServiceUnavailable,
|
||||
},
|
||||
{
|
||||
@@ -147,7 +147,7 @@ func TestExtProcHeadersEvaluation(t *testing.T) {
|
||||
authority: testUUID + ".team-a.actors.resources.substrate.ate.dev",
|
||||
resumeErr: status.Error(codes.DeadlineExceeded, "deadline"),
|
||||
expectErr: true,
|
||||
expectedErrStr: `actor "123e4567-e89b-12d3-a456-426614174000" request timed out`,
|
||||
expectedErrStr: `actor team-a/123e4567-e89b-12d3-a456-426614174000 request timed out`,
|
||||
expectedStatus: envoy_type.StatusCode_GatewayTimeout,
|
||||
},
|
||||
{
|
||||
@@ -159,7 +159,7 @@ func TestExtProcHeadersEvaluation(t *testing.T) {
|
||||
},
|
||||
},
|
||||
expectErr: true,
|
||||
expectedErrStr: `actor "123e4567-e89b-12d3-a456-426614174000" routing failed`,
|
||||
expectedErrStr: `actor team-a/123e4567-e89b-12d3-a456-426614174000 routing failed`,
|
||||
expectedStatus: envoy_type.StatusCode_InternalServerError,
|
||||
},
|
||||
{
|
||||
|
||||
@@ -304,7 +304,7 @@ func TestLogsActorRunner_Run_OneShot_ActorNotRunning(t *testing.T) {
|
||||
t.Fatal("expected error, got nil")
|
||||
}
|
||||
|
||||
wantErrMsg := "actor act-123 is not currently running on any worker pod"
|
||||
wantErrMsg := "actor space-1/act-123 is not currently running on any worker pod"
|
||||
if !strings.Contains(err.Error(), wantErrMsg) {
|
||||
t.Errorf("unexpected error message: %v (expected substring %q)", err, wantErrMsg)
|
||||
}
|
||||
@@ -442,7 +442,7 @@ func TestLogsActorRunner_Run_Follow_NotFoundActor(t *testing.T) {
|
||||
t.Fatal("expected error, got nil")
|
||||
}
|
||||
|
||||
wantErrMsg := "actor act-notfound not found"
|
||||
wantErrMsg := "actor space-1/act-notfound not found"
|
||||
if !strings.Contains(err.Error(), wantErrMsg) {
|
||||
t.Errorf("unexpected error: %v (expected %q)", err, wantErrMsg)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user