mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
Record actor state changes from ateapi (#1638)
# Record actor state changes from ateapi
ateapi now writes a log record every time an actor changes state.
## Before
Nothing tells you what state an actor is in. Metrics can't carry actor
identity, and the control plane samples traces at 10%, so most
transitions leave nothing behind at all. If you want to know whether an
agent is running, paused or suspended, there's nowhere to look.
## After
One filter:
```text
ate.actor.uid="8f2a1c4e6b0d47f1" AND ate.actor.state!=""
```
Last record wins. That's the state it's in now, and the timestamp is
when it got there.
The record:
```json
{
"time": "…",
"level": "INFO",
"msg": "Actor state changed",
"ate.atespace": "ate-demo-counter",
"ate.actor.name": "counter-1",
"ate.actor.uid": "8f2a…",
"ate.template.atespace": "ate-demo-counter",
"ate.template.name": "counter",
"ate.actor.operation.name": "suspend",
"ate.actor.state": "suspended"
}
```
## The two new keys
`ate.actor.state` is the `ateapipb.ActorState` values lowercased, so the
log wording and the state machine can't drift apart.
`ate.actor.operation.name` says what caused the change. The state on its
own doesn't tell you: an actor lands in `suspended` from a suspend and
`paused` from a pause, and only one of those gives the worker back.
`Actor crashed` picks up the same state key. A crash is a state change
too, and it's the one an actor can reach without any operation
finishing.
## Notes for review
**Why ateapi.** It owns the state machine. It also sees transitions that
never reach a worker, like deleting an actor that was already suspended,
or a resume that fails in the scheduler.
**It can't log a state the store never had.** The record goes out after
the store commit, and every state commit has a version precondition, so
a losing writer in a concurrent update writes nothing.
- [x] Tests pass
- [x] Appropriate changes to documentation are included in the PR
This commit is contained in:
@@ -142,6 +142,10 @@ func (s *ServiceImpl) CreateActor(ctx context.Context, inActor *ateapipb.Actor)
|
||||
return nil, fmt.Errorf("while recording actor: %w", err)
|
||||
}
|
||||
|
||||
// Without this an actor that is created and never resumed has no record at
|
||||
// all, at any retention.
|
||||
logActorStateChanged(ctx, stored, ateattr.OperationCreate)
|
||||
|
||||
return stored, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -116,9 +116,14 @@ func crashActor(ctx context.Context, st crashActorStore, actorRef resources.Acto
|
||||
// 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.
|
||||
//
|
||||
// It names ate.actor.state for the same reason ateom's lifecycle records do: a
|
||||
// crash is the one transition ateom never observes, so a consumer taking the
|
||||
// last state an actor reached has to see this record to reach "crashed" at all.
|
||||
func logActorCrashed(ctx context.Context, actor *ateapipb.Actor, opName, reason string) {
|
||||
attrs := ateattr.ActorLogAttrs(resources.ActorAttributionFromActor(actor))
|
||||
attrs = append(attrs, slog.String(string(ateattr.ActorOperationNameKey), opName))
|
||||
attrs = append(attrs, slog.String(string(ateattr.ActorStateKey), ateattr.ActorStateCrashed))
|
||||
attrs = append(attrs, ateattr.FailureLogAttrs(reason)...)
|
||||
slog.LogAttrs(ctx, slog.LevelError, "Actor crashed", attrs...)
|
||||
}
|
||||
|
||||
@@ -580,11 +580,17 @@ func TestCrashActorReleaseFailureLeavesWorkerReclaimable(t *testing.T) {
|
||||
// logs through the slog default, so this swaps it and the caller cannot be parallel.
|
||||
func crashRecords(t *testing.T) *[]map[string]string {
|
||||
t.Helper()
|
||||
return logRecords(t, "Actor crashed")
|
||||
}
|
||||
|
||||
// logRecords captures the attributes of every record with the given message.
|
||||
func logRecords(t *testing.T, msg string) *[]map[string]string {
|
||||
t.Helper()
|
||||
|
||||
var records []map[string]string
|
||||
prev := slog.Default()
|
||||
slog.SetDefault(slog.New(slogHandlerFunc(func(r slog.Record) {
|
||||
if r.Message != "Actor crashed" {
|
||||
if r.Message != msg {
|
||||
return
|
||||
}
|
||||
fields := map[string]string{}
|
||||
@@ -648,6 +654,7 @@ func TestCrashActor_RecordAndCounterAgree(t *testing.T) {
|
||||
string(ateattr.TemplateAtespaceKey): "demo-ns",
|
||||
string(ateattr.TemplateNameKey): "counter-template",
|
||||
string(ateattr.ActorOperationNameKey): ateattr.OperationResume,
|
||||
string(ateattr.ActorStateKey): ateattr.ActorStateCrashed,
|
||||
string(ateattr.FailureReasonKey): ateattr.ReasonWorkerPodGone,
|
||||
string(ateattr.FailureDomainKey): ateattr.FailureDomainInfrastructure,
|
||||
}
|
||||
|
||||
@@ -18,10 +18,12 @@ import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
|
||||
"github.com/agent-substrate/substrate/cmd/ateapi/internal/scheduling"
|
||||
"github.com/agent-substrate/substrate/cmd/ateapi/internal/store"
|
||||
"github.com/agent-substrate/substrate/cmd/ateapi/internal/workercache"
|
||||
"github.com/agent-substrate/substrate/internal/ateattr"
|
||||
"github.com/agent-substrate/substrate/internal/objectstore"
|
||||
"github.com/agent-substrate/substrate/internal/resources"
|
||||
listersv1alpha1 "github.com/agent-substrate/substrate/pkg/client/listers/api/v1alpha1"
|
||||
@@ -67,6 +69,35 @@ func markSkipped(ctx context.Context, reason string) {
|
||||
)
|
||||
}
|
||||
|
||||
// logActorStateChanged records an actor state change. The last record for an
|
||||
// actor's uid is the state it is in now.
|
||||
//
|
||||
// The state is read off the committed record, never passed in, so a record
|
||||
// cannot claim a state the store did not hold. Pass the actor the store returned
|
||||
// and call it after UpdateActor returns, not inside the mutate closure, which
|
||||
// can be retried. Every state commit carries a version precondition, so a call
|
||||
// that returns is the one that made the change.
|
||||
//
|
||||
// Crashes go through logActorCrashed instead, so read the state off
|
||||
// ate.actor.state rather than off the message.
|
||||
func logActorStateChanged(ctx context.Context, actor *ateapipb.Actor, opName string) {
|
||||
logActorState(ctx, actor, opName, ateattr.ActorStateValue(actor.GetStatus().GetState()))
|
||||
}
|
||||
|
||||
// logActorDeleted records the terminal transition. The row is gone, so there is
|
||||
// no committed state left to read and this is the one state named by hand.
|
||||
func logActorDeleted(ctx context.Context, actor *ateapipb.Actor, opName string) {
|
||||
logActorState(ctx, actor, opName, ateattr.ActorStateDeleted)
|
||||
}
|
||||
|
||||
func logActorState(ctx context.Context, actor *ateapipb.Actor, opName, state string) {
|
||||
attrs := ateattr.ActorLogAttrs(resources.ActorAttributionFromActor(actor))
|
||||
attrs = append(attrs,
|
||||
slog.String(string(ateattr.ActorOperationNameKey), ateattr.NormalizeOperationName(opName)),
|
||||
slog.String(string(ateattr.ActorStateKey), state))
|
||||
slog.LogAttrs(ctx, slog.LevelInfo, "Actor state changed", attrs...)
|
||||
}
|
||||
|
||||
// ActorWorkflow handles the workflows for actor's resume / suspend operations.
|
||||
type ActorWorkflow struct {
|
||||
store actorWorkflowStore
|
||||
|
||||
@@ -21,6 +21,7 @@ import (
|
||||
"log/slog"
|
||||
|
||||
"github.com/agent-substrate/substrate/cmd/ateapi/internal/store"
|
||||
"github.com/agent-substrate/substrate/internal/ateattr"
|
||||
"github.com/agent-substrate/substrate/internal/objectstore"
|
||||
"github.com/agent-substrate/substrate/internal/proto/ateletpb"
|
||||
"github.com/agent-substrate/substrate/internal/resources"
|
||||
@@ -346,6 +347,7 @@ func (w *ActorWorkflow) ensureMarkedDeleting(ctx context.Context, actorRef resou
|
||||
}
|
||||
return nil, fmt.Errorf("while setting actor state to DELETING: %w", err)
|
||||
}
|
||||
logActorStateChanged(ctx, storedActor, ateattr.OperationDelete)
|
||||
return storedActor, nil
|
||||
}
|
||||
|
||||
@@ -453,5 +455,6 @@ func (w *ActorWorkflow) finalizeDeleted(ctx context.Context, actorRef resources.
|
||||
}
|
||||
return nil, fmt.Errorf("while deleting actor from DB: %w", err)
|
||||
}
|
||||
logActorDeleted(ctx, deleted, ateattr.OperationDelete)
|
||||
return deleted, nil
|
||||
}
|
||||
|
||||
@@ -143,6 +143,7 @@ func (w *ActorWorkflow) ensureMarkedPausing(ctx context.Context, actorRef resour
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
logActorStateChanged(ctx, storedActor, ateattr.OperationPause)
|
||||
return storedActor, nil
|
||||
}
|
||||
|
||||
@@ -292,6 +293,9 @@ func (w *ActorWorkflow) ensurePausedFinalized(ctx context.Context, actorRef reso
|
||||
logActorCrashed(ctx, latestActor, ateattr.OperationPause, ateattr.ReasonCorruptedAssignment)
|
||||
recordActorCrash(ctx, crashAttrs)
|
||||
}
|
||||
if err == nil && storedActor.GetStatus().GetState() == ateapipb.ActorState_ACTOR_STATE_PAUSED {
|
||||
logActorStateChanged(ctx, storedActor, ateattr.OperationPause)
|
||||
}
|
||||
if err != nil {
|
||||
if errors.Is(err, store.ErrVersionConflict) {
|
||||
return nil, status.Error(codes.Aborted, "concurrent update conflict, please retry")
|
||||
|
||||
@@ -556,6 +556,7 @@ func (w *ActorWorkflow) assignWorkerAttempt(ctx context.Context, actorRef resour
|
||||
poolNamespace = assignedWorker.GetWorkerNamespace()
|
||||
pool = assignedWorker.GetWorkerPool()
|
||||
outcome = ateattr.SchedulerOutcomeAssigned
|
||||
logActorStateChanged(ctx, storedActor, ateattr.OperationResume)
|
||||
return storedActor, assignedWorker, nil
|
||||
}
|
||||
|
||||
@@ -813,5 +814,6 @@ func (w *ActorWorkflow) finalizeRunning(ctx context.Context, actorRef resources.
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
logActorStateChanged(ctx, storedActor, ateattr.OperationResume)
|
||||
return storedActor, nil
|
||||
}
|
||||
|
||||
@@ -160,6 +160,7 @@ func (w *ActorWorkflow) ensureMarkedSuspending(ctx context.Context, actorRef res
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
logActorStateChanged(ctx, storedActor, ateattr.OperationSuspend)
|
||||
return storedActor, nil
|
||||
}
|
||||
|
||||
@@ -442,6 +443,7 @@ func (w *ActorWorkflow) ensureSuspendedFinalized(ctx context.Context, actorRef r
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
logActorStateChanged(ctx, storedActor, ateattr.OperationSuspend)
|
||||
return storedActor, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,267 @@
|
||||
// Copyright 2026 Google LLC
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
package controlapi
|
||||
|
||||
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/ateattr"
|
||||
"github.com/agent-substrate/substrate/internal/resources"
|
||||
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
|
||||
)
|
||||
|
||||
// TestActorStateChangeRecords drives each transition that commits a new state
|
||||
// and pins the record it writes. The state a consumer reads has to be the state
|
||||
// the store now holds, so each case asserts both.
|
||||
func TestActorStateChangeRecords(t *testing.T) {
|
||||
const (
|
||||
tmplAtespace = "ns"
|
||||
tmplName = "tmpl1"
|
||||
)
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
seedState ateapipb.ActorState
|
||||
// transition runs the workflow step that commits the new state.
|
||||
transition func(t *testing.T, w *ActorWorkflow, ref resources.ActorRef, actor *ateapipb.Actor, tmpl *ateapipb.ActorTemplate)
|
||||
wantOp string
|
||||
wantState string
|
||||
wantStored ateapipb.ActorState
|
||||
}{
|
||||
{
|
||||
name: "suspend marks the actor suspending",
|
||||
seedState: ateapipb.ActorState_ACTOR_STATE_RUNNING,
|
||||
transition: func(t *testing.T, w *ActorWorkflow, ref resources.ActorRef, actor *ateapipb.Actor, tmpl *ateapipb.ActorTemplate) {
|
||||
if _, err := w.ensureMarkedSuspending(context.Background(), ref, actor, tmpl); err != nil {
|
||||
t.Fatalf("ensureMarkedSuspending: %v", err)
|
||||
}
|
||||
},
|
||||
wantOp: ateattr.OperationSuspend,
|
||||
wantState: ateattr.ActorStateSuspending,
|
||||
wantStored: ateapipb.ActorState_ACTOR_STATE_SUSPENDING,
|
||||
},
|
||||
{
|
||||
name: "pause marks the actor pausing",
|
||||
seedState: ateapipb.ActorState_ACTOR_STATE_RUNNING,
|
||||
transition: func(t *testing.T, w *ActorWorkflow, ref resources.ActorRef, actor *ateapipb.Actor, _ *ateapipb.ActorTemplate) {
|
||||
if _, err := w.ensureMarkedPausing(context.Background(), ref, actor); err != nil {
|
||||
t.Fatalf("ensureMarkedPausing: %v", err)
|
||||
}
|
||||
},
|
||||
wantOp: ateattr.OperationPause,
|
||||
wantState: ateattr.ActorStatePausing,
|
||||
wantStored: ateapipb.ActorState_ACTOR_STATE_PAUSING,
|
||||
},
|
||||
{
|
||||
name: "delete marks the actor deleting",
|
||||
seedState: ateapipb.ActorState_ACTOR_STATE_SUSPENDED,
|
||||
transition: func(t *testing.T, w *ActorWorkflow, ref resources.ActorRef, actor *ateapipb.Actor, _ *ateapipb.ActorTemplate) {
|
||||
if _, err := w.ensureMarkedDeleting(context.Background(), ref, actor, false); err != nil {
|
||||
t.Fatalf("ensureMarkedDeleting: %v", err)
|
||||
}
|
||||
},
|
||||
wantOp: ateattr.OperationDelete,
|
||||
wantState: ateattr.ActorStateDeleting,
|
||||
wantStored: ateapipb.ActorState_ACTOR_STATE_DELETING,
|
||||
},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
records := logRecords(t, "Actor state changed")
|
||||
|
||||
persistence := newTestPersistence(t)
|
||||
storetest.MustCreateAtespace(t, ctx, persistence, tmplAtespace)
|
||||
if _, err := persistence.CreateActorTemplate(ctx, &ateapipb.ActorTemplate{
|
||||
Metadata: &ateapipb.ResourceMetadata{Atespace: tmplAtespace, Name: tmplName},
|
||||
SnapshotsConfig: &ateapipb.SnapshotsConfig{
|
||||
StorageLocation: testStorageLocation,
|
||||
OnPause: ateapipb.SnapshotContentScope_SNAPSHOT_CONTENT_SCOPE_FULL,
|
||||
},
|
||||
}); err != nil {
|
||||
t.Fatalf("create template: %v", err)
|
||||
}
|
||||
|
||||
actorRef := resources.ActorRef{Atespace: "team-a", Name: "id1"}
|
||||
seedWorkflowActor(t, ctx, persistence, actorRef, tmplAtespace, tmplName, tt.seedState)
|
||||
actor, err := persistence.GetActor(ctx, actorRef)
|
||||
if err != nil {
|
||||
t.Fatalf("get actor: %v", err)
|
||||
}
|
||||
tmpl, err := persistence.GetActorTemplate(ctx, resources.ActorTemplateRef{Atespace: tmplAtespace, Name: tmplName})
|
||||
if err != nil {
|
||||
t.Fatalf("get template: %v", err)
|
||||
}
|
||||
|
||||
w := &ActorWorkflow{store: persistence}
|
||||
tt.transition(t, w, actorRef, actor, tmpl)
|
||||
|
||||
if len(*records) != 1 {
|
||||
t.Fatalf("got %d state records, want 1: %v", len(*records), *records)
|
||||
}
|
||||
got := (*records)[0]
|
||||
want := map[string]string{
|
||||
string(ateattr.AtespaceKey): actorRef.Atespace,
|
||||
string(ateattr.ActorNameKey): actorRef.Name,
|
||||
string(ateattr.ActorUIDKey): actor.GetMetadata().GetUid(),
|
||||
string(ateattr.TemplateAtespaceKey): tmplAtespace,
|
||||
string(ateattr.TemplateNameKey): tmplName,
|
||||
string(ateattr.ActorOperationNameKey): tt.wantOp,
|
||||
string(ateattr.ActorStateKey): tt.wantState,
|
||||
}
|
||||
for k, wv := range want {
|
||||
if got[k] != wv {
|
||||
t.Errorf("%s = %q, want %q", k, got[k], wv)
|
||||
}
|
||||
}
|
||||
if len(got) != len(want) {
|
||||
t.Errorf("got %d attributes, want %d: %v", len(got), len(want), got)
|
||||
}
|
||||
|
||||
stored, err := persistence.GetActor(ctx, actorRef)
|
||||
if err != nil {
|
||||
t.Fatalf("reload actor: %v", err)
|
||||
}
|
||||
if stored.GetStatus().GetState() != tt.wantStored {
|
||||
t.Errorf("stored state = %v, want %v", stored.GetStatus().GetState(), tt.wantStored)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestActorCreatedRecord covers creation. An actor created and never resumed
|
||||
// makes no other transition, so without this it has no record at all, at any
|
||||
// retention.
|
||||
func TestActorCreatedRecord(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
records := logRecords(t, "Actor state changed")
|
||||
|
||||
persistence := newTestPersistence(t)
|
||||
storetest.MustCreateAtespace(t, ctx, persistence, "ns")
|
||||
storetest.MustCreateActor(t, ctx, persistence, &ateapipb.Actor{
|
||||
Metadata: &ateapipb.ResourceMetadata{Atespace: "team-a", Name: "id1"},
|
||||
ActorTemplate: &ateapipb.ObjectRef{Atespace: "ns", Name: "tmpl1"},
|
||||
Status: &ateapipb.ActorStatus{State: ateapipb.ActorState_ACTOR_STATE_SUSPENDED},
|
||||
})
|
||||
stored, err := persistence.GetActor(ctx, resources.ActorRef{Atespace: "team-a", Name: "id1"})
|
||||
if err != nil {
|
||||
t.Fatalf("get actor: %v", err)
|
||||
}
|
||||
|
||||
logActorStateChanged(ctx, stored, ateattr.OperationCreate)
|
||||
|
||||
if len(*records) != 1 {
|
||||
t.Fatalf("got %d state records, want 1: %v", len(*records), *records)
|
||||
}
|
||||
got := (*records)[0]
|
||||
if got[string(ateattr.ActorOperationNameKey)] != ateattr.OperationCreate {
|
||||
t.Errorf("operation = %q, want %q", got[string(ateattr.ActorOperationNameKey)], ateattr.OperationCreate)
|
||||
}
|
||||
// A new actor is born suspended, and the record has to say so rather than
|
||||
// inventing a "created" state the store does not have.
|
||||
if got[string(ateattr.ActorStateKey)] != ateattr.ActorStateSuspended {
|
||||
t.Errorf("state = %q, want %q", got[string(ateattr.ActorStateKey)], ateattr.ActorStateSuspended)
|
||||
}
|
||||
}
|
||||
|
||||
// TestActorDeletedRecord covers the terminal record. Without it "deleting" is
|
||||
// the last thing a deleted actor ever reports, and a consumer cannot tell a
|
||||
// finished delete from one that is stuck.
|
||||
func TestActorDeletedRecord(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
records := logRecords(t, "Actor state changed")
|
||||
|
||||
persistence := newTestPersistence(t)
|
||||
storetest.MustCreateAtespace(t, ctx, persistence, "ns")
|
||||
if _, err := persistence.CreateActorTemplate(ctx, &ateapipb.ActorTemplate{
|
||||
Metadata: &ateapipb.ResourceMetadata{Atespace: "ns", Name: "tmpl1"},
|
||||
SnapshotsConfig: &ateapipb.SnapshotsConfig{StorageLocation: testStorageLocation},
|
||||
}); err != nil {
|
||||
t.Fatalf("create template: %v", err)
|
||||
}
|
||||
|
||||
actorRef := resources.ActorRef{Atespace: "team-a", Name: "id1"}
|
||||
seedWorkflowActor(t, ctx, persistence, actorRef, "ns", "tmpl1", ateapipb.ActorState_ACTOR_STATE_DELETING)
|
||||
actor, err := persistence.GetActor(ctx, actorRef)
|
||||
if err != nil {
|
||||
t.Fatalf("get actor: %v", err)
|
||||
}
|
||||
|
||||
w := &ActorWorkflow{store: persistence}
|
||||
if _, err := w.finalizeDeleted(ctx, actorRef); err != nil {
|
||||
t.Fatalf("finalizeDeleted: %v", err)
|
||||
}
|
||||
|
||||
if len(*records) != 1 {
|
||||
t.Fatalf("got %d state records, want 1: %v", len(*records), *records)
|
||||
}
|
||||
got := (*records)[0]
|
||||
if got[string(ateattr.ActorStateKey)] != ateattr.ActorStateDeleted {
|
||||
t.Errorf("state = %q, want %q", got[string(ateattr.ActorStateKey)], ateattr.ActorStateDeleted)
|
||||
}
|
||||
// The identity has to survive the row it described, or the terminal record
|
||||
// cannot be joined to the rest of the actor's history.
|
||||
if got[string(ateattr.ActorUIDKey)] != actor.GetMetadata().GetUid() {
|
||||
t.Errorf("uid = %q, want %q", got[string(ateattr.ActorUIDKey)], actor.GetMetadata().GetUid())
|
||||
}
|
||||
|
||||
if _, err := persistence.GetActor(ctx, actorRef); err == nil {
|
||||
t.Error("actor still readable after finalizeDeleted")
|
||||
}
|
||||
}
|
||||
|
||||
// TestActorStateChangeRecordSkippedOnConflict pins that a transition which did
|
||||
// not commit writes nothing. A losing writer that still logged would put a state
|
||||
// in the stream the store never held.
|
||||
func TestActorStateChangeRecordSkippedOnConflict(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
records := logRecords(t, "Actor state changed")
|
||||
|
||||
persistence := newTestPersistence(t)
|
||||
storetest.MustCreateAtespace(t, ctx, persistence, "ns")
|
||||
if _, err := persistence.CreateActorTemplate(ctx, &ateapipb.ActorTemplate{
|
||||
Metadata: &ateapipb.ResourceMetadata{Atespace: "ns", Name: "tmpl1"},
|
||||
SnapshotsConfig: &ateapipb.SnapshotsConfig{StorageLocation: testStorageLocation},
|
||||
}); err != nil {
|
||||
t.Fatalf("create template: %v", err)
|
||||
}
|
||||
|
||||
actorRef := resources.ActorRef{Atespace: "team-a", Name: "id1"}
|
||||
seedWorkflowActor(t, ctx, persistence, actorRef, "ns", "tmpl1", ateapipb.ActorState_ACTOR_STATE_SUSPENDED)
|
||||
stale, err := persistence.GetActor(ctx, actorRef)
|
||||
if err != nil {
|
||||
t.Fatalf("get actor: %v", err)
|
||||
}
|
||||
|
||||
// Bump the stored version so the workflow's precondition is stale.
|
||||
if _, err := persistence.UpdateActor(ctx, actorRef, store.PreconditionFrom(stale), func(toUpdate *ateapipb.Actor) error {
|
||||
toUpdate.Status.InProgressSnapshotName = "someone-else"
|
||||
return nil
|
||||
}); err != nil {
|
||||
t.Fatalf("bump version: %v", err)
|
||||
}
|
||||
|
||||
w := &ActorWorkflow{store: persistence}
|
||||
if _, err := w.ensureMarkedDeleting(ctx, actorRef, stale, false); err == nil {
|
||||
t.Fatal("ensureMarkedDeleting on a stale actor = nil, want a conflict")
|
||||
}
|
||||
if len(*records) != 0 {
|
||||
t.Errorf("got %d state records from a losing writer, want 0: %v", len(*records), *records)
|
||||
}
|
||||
}
|
||||
+24
-2
@@ -141,13 +141,35 @@ The duration keys are the [`ate.actor.restore.duration`](#the-metric-registry) i
|
||||
|
||||
This is the record to use for a per-actor wake-up distribution. The histogram cannot answer that question at all, because actor identity is barred from metric labels; traces can, but the data plane is head-sampled at 1%.
|
||||
|
||||
ateapi's `Actor crashed` is the other one. It is written once per committed transition into `ACTOR_STATE_CRASHED`, beside the [`ate.actor.crashes`](#the-metric-registry) increment and under the same already-crashed guard, so the two can never disagree about how many crashes happened:
|
||||
ateapi's `Actor state changed` is written once per committed actor state transition. ateapi owns the state machine, so this is where an actor's state and the time it reached it come from:
|
||||
|
||||
```json
|
||||
{"time":"…","level":"INFO","msg":"Actor state changed",
|
||||
"ate.atespace":"ate-demo-counter","ate.actor.name":"counter-1","ate.actor.uid":"8f2a…",
|
||||
"ate.template.atespace":"ate-demo-counter","ate.template.name":"counter",
|
||||
"ate.actor.operation.name":"suspend","ate.actor.state":"suspended",
|
||||
"trace_id":"4bf92f…","span_id":"00f067…","trace_flags":"01"}
|
||||
```
|
||||
|
||||
`ate.actor.state` takes the `ateapipb.ActorState` values lowercased, so the log vocabulary and the state machine cannot fork. The last record for an actor's uid is the state it is in now, and its timestamp is when that state began. Query it per actor, not in aggregate: neither key is a metric label, because both only ever appear beside actor identity, which [the cardinality rules](#the-metric-registry) keep off metrics entirely.
|
||||
|
||||
`ate.actor.operation.name` says which operation drove the transition, which the state alone does not: an actor reaches `suspended` from a suspend and `paused` from a pause, and the two differ in whether the worker was released.
|
||||
|
||||
The record goes out after the store commit, never before, and the state is read straight off the committed record rather than named by the caller. Every state commit has a version check too, so if two writers race, the one that lost writes nothing. You will never see a state here that the store did not actually hold.
|
||||
|
||||
Creating an actor counts as a change. A new actor is born suspended, so it gets a record saying so, with `ate.actor.operation.name` set to `create`. Otherwise an actor that is created and never resumed would have no record at all, no matter how long you keep your logs.
|
||||
|
||||
`deleted` is the only state with no `ateapipb.ActorState` behind it. It is written once the actor row is gone, so there is nothing left to read back. It is also the last record an actor ever gets. Without it, `deleting` would be the end of the story, and a delete that finished would look just like one that got stuck.
|
||||
|
||||
**What this stream won't tell you.** It only writes when something changes. So an actor that has been sitting in the same state since before your logs roll over has no record, and no state. Ask the control plane what state something is in right now. Use this stream to see how it got there and when. Records can also go missing, like any other log, and a gap looks the same as an actor that just sat still. If you want to count activations, use the router's access log instead.
|
||||
|
||||
`Actor crashed` is the exception, and carries the same two keys with `ate.actor.state="crashed"`. It is written once per committed transition into `ACTOR_STATE_CRASHED`, beside the [`ate.actor.crashes`](#the-metric-registry) increment and under the same already-crashed guard, so the two can never disagree about how many crashes happened. A consumer deriving state therefore selects on `ate.actor.state`, not on the message:
|
||||
|
||||
```json
|
||||
{"time":"…","level":"ERROR","msg":"Actor crashed",
|
||||
"ate.atespace":"ate-demo-counter","ate.actor.name":"counter-1","ate.actor.uid":"8f2a…",
|
||||
"ate.template.atespace":"ate-demo-counter","ate.template.name":"counter",
|
||||
"ate.actor.operation.name":"resume",
|
||||
"ate.actor.operation.name":"resume","ate.actor.state":"crashed",
|
||||
"ate.failure.reason":"WORKER_POD_GONE","ate.failure.domain":"infrastructure",
|
||||
"trace_id":"4bf92f…","span_id":"00f067…","trace_flags":"01"}
|
||||
```
|
||||
|
||||
@@ -77,6 +77,52 @@ const (
|
||||
// Only the components that have a relay to take or miss carry it.
|
||||
const OTLPRelayKey = attribute.Key("ate.otlp.relay")
|
||||
|
||||
// ActorStateKey is log-only. It is bounded, but it is only ever recorded beside
|
||||
// actor identity, which no metric may carry.
|
||||
const ActorStateKey = attribute.Key("ate.actor.state")
|
||||
|
||||
// Values for ActorStateKey, mirroring ateapipb.ActorState with the
|
||||
// ACTOR_STATE_ prefix dropped so the two cannot fork. ActorStateDeleted is the
|
||||
// exception and has no enum counterpart: the record is gone, so no stored state
|
||||
// can stand for it.
|
||||
const (
|
||||
ActorStateResuming = "resuming"
|
||||
ActorStateRunning = "running"
|
||||
ActorStateSuspending = "suspending"
|
||||
ActorStateSuspended = "suspended"
|
||||
ActorStatePausing = "pausing"
|
||||
ActorStatePaused = "paused"
|
||||
ActorStateCrashed = "crashed"
|
||||
ActorStateDeleting = "deleting"
|
||||
ActorStateDeleted = "deleted"
|
||||
ActorStateUnknown = "unknown"
|
||||
)
|
||||
|
||||
// ActorStateValue maps the committed state onto its label value, so a producer
|
||||
// reports the state the store holds rather than one it names by hand.
|
||||
func ActorStateValue(state ateapipb.ActorState) string {
|
||||
switch state {
|
||||
case ateapipb.ActorState_ACTOR_STATE_RESUMING:
|
||||
return ActorStateResuming
|
||||
case ateapipb.ActorState_ACTOR_STATE_RUNNING:
|
||||
return ActorStateRunning
|
||||
case ateapipb.ActorState_ACTOR_STATE_SUSPENDING:
|
||||
return ActorStateSuspending
|
||||
case ateapipb.ActorState_ACTOR_STATE_SUSPENDED:
|
||||
return ActorStateSuspended
|
||||
case ateapipb.ActorState_ACTOR_STATE_PAUSING:
|
||||
return ActorStatePausing
|
||||
case ateapipb.ActorState_ACTOR_STATE_PAUSED:
|
||||
return ActorStatePaused
|
||||
case ateapipb.ActorState_ACTOR_STATE_CRASHED:
|
||||
return ActorStateCrashed
|
||||
case ateapipb.ActorState_ACTOR_STATE_DELETING:
|
||||
return ActorStateDeleting
|
||||
default:
|
||||
return ActorStateUnknown
|
||||
}
|
||||
}
|
||||
|
||||
// Metric-label keys: the only ate.* attributes allowed on metric datapoints,
|
||||
// each with a small bounded value set. High-cardinality identity (actor
|
||||
// name/uid, atespace) is absent by design; it belongs on spans and logs.
|
||||
|
||||
@@ -19,6 +19,7 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"maps"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"go.opentelemetry.io/otel/attribute"
|
||||
@@ -158,6 +159,7 @@ func TestKeySpellings(t *testing.T) {
|
||||
{TemplateNameKey, "ate.template.name"},
|
||||
{TemplateAtespaceKey, "ate.template.atespace"},
|
||||
{ActorVersionKey, "ate.actor.version"},
|
||||
{ActorStateKey, "ate.actor.state"},
|
||||
{ActorOperationNameKey, "ate.actor.operation.name"},
|
||||
{WorkerPoolNamespaceKey, "ate.workerpool.namespace"},
|
||||
{WorkerPoolNameKey, "ate.workerpool.name"},
|
||||
@@ -603,6 +605,105 @@ func TestSnapshotScopeValue(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestActorStateValues pins the spelling of every state a record can report.
|
||||
func TestActorStateValues(t *testing.T) {
|
||||
tests := []struct {
|
||||
got string
|
||||
want string
|
||||
}{
|
||||
{ActorStateResuming, "resuming"},
|
||||
{ActorStateRunning, "running"},
|
||||
{ActorStateSuspending, "suspending"},
|
||||
{ActorStateSuspended, "suspended"},
|
||||
{ActorStatePausing, "pausing"},
|
||||
{ActorStatePaused, "paused"},
|
||||
{ActorStateCrashed, "crashed"},
|
||||
{ActorStateDeleting, "deleting"},
|
||||
{ActorStateDeleted, "deleted"},
|
||||
{ActorStateUnknown, "unknown"},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.want, func(t *testing.T) {
|
||||
if tt.got != tt.want {
|
||||
t.Errorf("got %q, want %q", tt.got, tt.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestActorStateValue pins the mapping producers read the state through, so a
|
||||
// record reports the state the store holds rather than one named by hand.
|
||||
func TestActorStateValue(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
state ateapipb.ActorState
|
||||
want string
|
||||
}{
|
||||
{name: "resuming", state: ateapipb.ActorState_ACTOR_STATE_RESUMING, want: ActorStateResuming},
|
||||
{name: "running", state: ateapipb.ActorState_ACTOR_STATE_RUNNING, want: ActorStateRunning},
|
||||
{name: "suspending", state: ateapipb.ActorState_ACTOR_STATE_SUSPENDING, want: ActorStateSuspending},
|
||||
{name: "suspended", state: ateapipb.ActorState_ACTOR_STATE_SUSPENDED, want: ActorStateSuspended},
|
||||
{name: "pausing", state: ateapipb.ActorState_ACTOR_STATE_PAUSING, want: ActorStatePausing},
|
||||
{name: "paused", state: ateapipb.ActorState_ACTOR_STATE_PAUSED, want: ActorStatePaused},
|
||||
{name: "crashed", state: ateapipb.ActorState_ACTOR_STATE_CRASHED, want: ActorStateCrashed},
|
||||
{name: "deleting", state: ateapipb.ActorState_ACTOR_STATE_DELETING, want: ActorStateDeleting},
|
||||
{name: "unspecified", state: ateapipb.ActorState_ACTOR_STATE_UNSPECIFIED, want: ActorStateUnknown},
|
||||
{name: "value outside the enum", state: ateapipb.ActorState(9999), want: ActorStateUnknown},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
if got := ActorStateValue(tt.state); got != tt.want {
|
||||
t.Errorf("ActorStateValue(%v) = %q, want %q", tt.state, got, tt.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestActorStateValuesMirrorActorState holds the log vocabulary to the control
|
||||
// plane's state machine. A state added to the enum without a value here would
|
||||
// leave the stream unable to name the state an actor is in.
|
||||
//
|
||||
// ActorStateDeleted is excluded on purpose: the actor row is gone by the time it
|
||||
// is reported, so no enum value can stand for it. A second such value has to be
|
||||
// added to this list deliberately.
|
||||
func TestActorStateValuesMirrorActorState(t *testing.T) {
|
||||
// deleted has no enum value because the row is gone; unknown is the
|
||||
// mapper's fallback for UNSPECIFIED and anything off the enum.
|
||||
noEnumCounterpart := map[string]bool{ActorStateDeleted: true, ActorStateUnknown: true}
|
||||
|
||||
got := map[string]bool{
|
||||
ActorStateResuming: true,
|
||||
ActorStateRunning: true,
|
||||
ActorStateSuspending: true,
|
||||
ActorStateSuspended: true,
|
||||
ActorStatePausing: true,
|
||||
ActorStatePaused: true,
|
||||
ActorStateCrashed: true,
|
||||
ActorStateDeleting: true,
|
||||
ActorStateDeleted: true,
|
||||
ActorStateUnknown: true,
|
||||
}
|
||||
|
||||
want := map[string]bool{}
|
||||
for value, name := range ateapipb.ActorState_name {
|
||||
if ateapipb.ActorState(value) == ateapipb.ActorState_ACTOR_STATE_UNSPECIFIED {
|
||||
continue
|
||||
}
|
||||
want[strings.ToLower(strings.TrimPrefix(name, "ACTOR_STATE_"))] = true
|
||||
}
|
||||
|
||||
for state := range want {
|
||||
if !got[state] {
|
||||
t.Errorf("ateapipb.ActorState has %q with no ateattr constant", state)
|
||||
}
|
||||
}
|
||||
for state := range got {
|
||||
if !want[state] && !noEnumCounterpart[state] {
|
||||
t.Errorf("ateattr has state %q that ateapipb.ActorState does not", state)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestNormalizeSandboxClass covers the cardinality guard: atelet reads the class
|
||||
// out of a snapshot manifest nothing validates, so anything unrecognized has to
|
||||
// collapse onto a single value.
|
||||
|
||||
Reference in New Issue
Block a user