From 313ae1efd909565d88f438e6e929d6c4d96c4e01 Mon Sep 17 00:00:00 2001 From: Angela Hu Date: Wed, 29 Jul 2026 09:28:30 -0700 Subject: [PATCH 1/5] feat(ateapi): add resumed boolean to ResumeActor response --- cmd/ateapi/internal/controlapi/resume_actor.go | 4 ++-- cmd/ateapi/internal/controlapi/workflow.go | 8 ++++---- cmd/ateapi/internal/controlapi/workflow_resume.go | 2 ++ cmd/ateapi/internal/controlapi/workflow_resume_test.go | 7 +++++-- pkg/proto/ateapipb/ateapi.pb.go | 8 ++++++++ pkg/proto/ateapipb/ateapi.proto | 4 ++++ 6 files changed, 25 insertions(+), 8 deletions(-) diff --git a/cmd/ateapi/internal/controlapi/resume_actor.go b/cmd/ateapi/internal/controlapi/resume_actor.go index bc99be203..eac4081d5 100644 --- a/cmd/ateapi/internal/controlapi/resume_actor.go +++ b/cmd/ateapi/internal/controlapi/resume_actor.go @@ -33,7 +33,7 @@ func (s *Service) ResumeActor(ctx context.Context, req *ateapipb.ResumeActorRequ actorRef := resources.ActorRefFromObjectRef(req.GetActor()) setSpanActorRefAttributes(ctx, actorRef) - actor, err := s.actorWorkflow.ResumeActor(ctx, actorRef, req.GetBoot()) + actor, resumed, err := s.actorWorkflow.ResumeActor(ctx, actorRef, req.GetBoot()) if err != nil { if errors.Is(err, store.ErrVersionConflict) { return nil, status.Error(codes.Aborted, "concurrent update conflict, please retry") @@ -45,7 +45,7 @@ func (s *Service) ResumeActor(ctx context.Context, req *ateapipb.ResumeActorRequ } setSpanActorAttributes(ctx, actor) - return &ateapipb.ResumeActorResponse{Actor: actor}, nil + return &ateapipb.ResumeActorResponse{Actor: actor, Resumed: resumed}, nil } func validateResumeActorRequest(req *ateapipb.ResumeActorRequest) error { diff --git a/cmd/ateapi/internal/controlapi/workflow.go b/cmd/ateapi/internal/controlapi/workflow.go index f06efa7b6..e8fc4292a 100644 --- a/cmd/ateapi/internal/controlapi/workflow.go +++ b/cmd/ateapi/internal/controlapi/workflow.go @@ -165,7 +165,7 @@ func NewActorWorkflow( } // ResumeActor executes the workflow to resume a suspended actor. Idempotent. -func (w *ActorWorkflow) ResumeActor(ctx context.Context, actorRef resources.ActorRef, boot bool) (*ateapipb.Actor, error) { +func (w *ActorWorkflow) ResumeActor(ctx context.Context, actorRef resources.ActorRef, boot bool) (*ateapipb.Actor, bool, error) { input := &ResumeInput{ ActorRef: actorRef, Boot: boot, @@ -174,7 +174,7 @@ func (w *ActorWorkflow) ResumeActor(ctx context.Context, actorRef resources.Acto ctx, lock, err := w.acquireActorLock(ctx, actorRef) if err != nil { - return nil, err + return nil, false, err } defer lock.Close() @@ -187,10 +187,10 @@ func (w *ActorWorkflow) ResumeActor(ctx context.Context, actorRef resources.Acto } if err := RunWorkflow(ctx, input, state, steps); err != nil { - return nil, err + return nil, false, err } - return state.Actor, nil + return state.Actor, !state.WasRunning, nil } // SuspendActor executes the workflow to suspend a running actor. Idempotent. diff --git a/cmd/ateapi/internal/controlapi/workflow_resume.go b/cmd/ateapi/internal/controlapi/workflow_resume.go index 40f88a126..80add4f85 100644 --- a/cmd/ateapi/internal/controlapi/workflow_resume.go +++ b/cmd/ateapi/internal/controlapi/workflow_resume.go @@ -49,6 +49,7 @@ type ResumeState struct { Actor *ateapipb.Actor Worker *ateapipb.Worker ActorTemplate *atev1alpha1.ActorTemplate + WasRunning bool } type LoadActorForResumeStep struct { @@ -73,6 +74,7 @@ func (s *LoadActorForResumeStep) Execute(ctx context.Context, input *ResumeInput return fmt.Errorf("while getting actor from DB: %w", err) } state.Actor = actor + state.WasRunning = (actor.GetStatus() == ateapipb.Actor_STATUS_RUNNING) actorTemplate, err := s.actorTemplateLister.ActorTemplates(actor.GetActorTemplateNamespace()).Get(actor.GetActorTemplateName()) if err != nil { diff --git a/cmd/ateapi/internal/controlapi/workflow_resume_test.go b/cmd/ateapi/internal/controlapi/workflow_resume_test.go index 5a4e47a6a..568aaca08 100644 --- a/cmd/ateapi/internal/controlapi/workflow_resume_test.go +++ b/cmd/ateapi/internal/controlapi/workflow_resume_test.go @@ -442,7 +442,7 @@ func TestResumeActorWorkflow_RejectedAndIdempotentPaths(t *testing.T) { a.WorkerPoolName = "pool1" }) - actor, err := w.ResumeActor(ctx, resources.ActorRef{Atespace: "team-a", Name: "id1"}, false) + actor, resumed, err := w.ResumeActor(ctx, resources.ActorRef{Atespace: "team-a", Name: "id1"}, false) if tc.wantErr { if got := status.Code(err); got != codes.FailedPrecondition { t.Fatalf("status.Code(err) = %v, want %v (err: %v)", got, codes.FailedPrecondition, err) @@ -454,6 +454,9 @@ func TestResumeActorWorkflow_RejectedAndIdempotentPaths(t *testing.T) { if actor.GetStatus() != tc.wantStatus { t.Errorf("returned status = %v, want %v", actor.GetStatus(), tc.wantStatus) } + if tc.seedStatus == ateapipb.Actor_STATUS_RUNNING && resumed { + t.Errorf("expected resumed = false for already running actor, got true") + } } got, err := st.GetActor(ctx, resources.ActorRef{Atespace: "team-a", Name: "id1"}) @@ -548,7 +551,7 @@ func TestResumeActor_CrashesOnCorruptWorkerAssignment(t *testing.T) { a.AteomPodName = "worker-1" // AteomPodUid and WorkerPoolName left empty }) - _, err := w.ResumeActor(ctx, resources.ActorRef{Atespace: "team-a", Name: "id1"}, false) + _, _, err := w.ResumeActor(ctx, resources.ActorRef{Atespace: "team-a", Name: "id1"}, false) if got := status.Code(err); got != codes.Aborted { t.Fatalf("status.Code(err) = %v, want %v (err: %v)", got, codes.Aborted, err) } diff --git a/pkg/proto/ateapipb/ateapi.pb.go b/pkg/proto/ateapipb/ateapi.pb.go index 6d5a6e70f..6a7d2cdb8 100644 --- a/pkg/proto/ateapipb/ateapi.pb.go +++ b/pkg/proto/ateapipb/ateapi.pb.go @@ -1464,6 +1464,7 @@ func (x *ResumeActorRequest) GetBoot() bool { type ResumeActorResponse struct { state protoimpl.MessageState `protogen:"open.v1"` Actor *Actor `protobuf:"bytes,1,opt,name=actor,proto3" json:"actor,omitempty"` + Resumed bool `protobuf:"varint,2,opt,name=resumed,proto3" json:"resumed,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } @@ -1505,6 +1506,13 @@ func (x *ResumeActorResponse) GetActor() *Actor { return nil } +func (x *ResumeActorResponse) GetResumed() bool { + if x != nil { + return x.Resumed + } + return false +} + type DeleteActorRequest struct { state protoimpl.MessageState `protogen:"open.v1"` Actor *ObjectRef `protobuf:"bytes,1,opt,name=actor,proto3" json:"actor,omitempty"` diff --git a/pkg/proto/ateapipb/ateapi.proto b/pkg/proto/ateapipb/ateapi.proto index 27c87e552..1d48165de 100644 --- a/pkg/proto/ateapipb/ateapi.proto +++ b/pkg/proto/ateapipb/ateapi.proto @@ -276,6 +276,10 @@ message ResumeActorRequest { message ResumeActorResponse { Actor actor = 1; + + // True if a resume workflow was executed to activate the actor. + // False if the actor was already RUNNING. + bool resumed = 2; } message DeleteActorRequest { From 052f942c3df9c8b4349e2e290a23e0ee1270b44c Mon Sep 17 00:00:00 2001 From: Angela Hu Date: Wed, 29 Jul 2026 10:02:27 -0700 Subject: [PATCH 2/5] feat(atenet): disambiguate singleflight resume outcomes --- cmd/atenet/internal/router/resumer.go | 64 +++++++++++++++--- cmd/atenet/internal/router/resumer_test.go | 79 ++++++++++++++++++---- 2 files changed, 120 insertions(+), 23 deletions(-) diff --git a/cmd/atenet/internal/router/resumer.go b/cmd/atenet/internal/router/resumer.go index 61f1aeda4..151eb0467 100644 --- a/cmd/atenet/internal/router/resumer.go +++ b/cmd/atenet/internal/router/resumer.go @@ -17,6 +17,7 @@ package router import ( "context" "math" + "sync/atomic" "time" "github.com/agent-substrate/substrate/internal/ateattr" @@ -64,6 +65,25 @@ type budgetExhaustedError struct{ lastErr error } func (e *budgetExhaustedError) Error() string { return e.lastErr.Error() } func (e *budgetExhaustedError) Unwrap() error { return e.lastErr } +// ResumeOutcome indicates the singleflight execution state of an actor resumption request. +type ResumeOutcome string + +const ( + // ResumeOutcomeNone indicates the actor was already running (steady-state warm route). + ResumeOutcomeNone ResumeOutcome = "none" + // ResumeOutcomeTriggered indicates this request won the singleflight lock and initiated cold activation. + ResumeOutcomeTriggered ResumeOutcome = "triggered" + // ResumeOutcomeJoined indicates this request parked on an in-flight singleflight resume. + ResumeOutcomeJoined ResumeOutcome = "joined" +) + +type resumeCallResult struct { + actor *ateapipb.Actor + resumed bool + leaderID uint64 + err error +} + // ActorResumer coordinates safe, deduplicated resumption of actors. type ActorResumer struct { apiClient ateapipb.ControlClient @@ -78,6 +98,7 @@ type ActorResumer struct { budget time.Duration // backoff paces the retries within the budget. backoff wait.Backoff + nextID uint64 } // resumerOption configures an ActorResumer. @@ -134,11 +155,13 @@ func (r *ActorResumer) retryable(err error) bool { // ResumeActor ensures the requested actor is running. It deduplicates concurrent // requests within the process and, when parking is enabled, holds the request // while retrying transient failures until the budget elapses. -func (r *ActorResumer) ResumeActor(ctx context.Context, actorRef resources.ActorRef) (*ateapipb.Actor, error) { +func (r *ActorResumer) ResumeActor(ctx context.Context, actorRef resources.ActorRef) (*ateapipb.Actor, ResumeOutcome, error) { ctx, span := otel.Tracer(routerServiceName).Start(ctx, "ResumeActor", trace.WithAttributes(ateattr.ActorRefAttributes(actorRef)...)) defer span.End() + reqID := atomic.AddUint64(&r.nextID, 1) + ch := r.flight.DoChan(actorRef.String(), func() (interface{}, error) { // We detach the context from the first caller using a fixed background budget. // This guarantees that if Caller 1 disconnects or times out, the underlying @@ -167,8 +190,8 @@ func (r *ActorResumer) ResumeActor(ctx context.Context, actorRef resources.Actor } if r.retryable(err) { - lastRetryErr = err // remember it in case the budget elapses - return false, nil // park: retry until the budget elapses + lastRetryErr = err + return false, nil } return false, err }) @@ -187,21 +210,42 @@ func (r *ActorResumer) ResumeActor(ctx context.Context, actorRef resources.Actor // as a 504. bgCtx is this loop's only deadline source, so checking it // covers both landing spots (mid-RPC and between retries). if lastRetryErr != nil && (bgCtx.Err() != nil || wait.Interrupted(err)) { - return nil, &budgetExhaustedError{lastErr: lastRetryErr} + return &resumeCallResult{leaderID: reqID, err: &budgetExhaustedError{lastErr: lastRetryErr}}, nil } - return nil, err + return &resumeCallResult{leaderID: reqID, err: err}, nil } - return resumeResp.GetActor(), nil + return &resumeCallResult{ + actor: resumeResp.GetActor(), + resumed: resumeResp.GetResumed(), + leaderID: reqID, + }, nil }) select { case <-ctx.Done(): - return nil, ctx.Err() + return nil, ResumeOutcomeNone, ctx.Err() case res := <-ch: - if res.Err != nil { - return nil, res.Err + callRes, _ := res.Val.(*resumeCallResult) + if callRes == nil { + if res.Err != nil { + return nil, ResumeOutcomeNone, res.Err + } + return nil, ResumeOutcomeNone, status.Error(codes.Internal, "resume call returned nil result") + } + if callRes.err != nil { + return nil, ResumeOutcomeNone, callRes.err + } + + outcome := ResumeOutcomeNone + if callRes.resumed { + if callRes.leaderID == reqID { + outcome = ResumeOutcomeTriggered + } else { + outcome = ResumeOutcomeJoined + } } - return res.Val.(*ateapipb.Actor), nil + + return callRes.actor, outcome, nil } } diff --git a/cmd/atenet/internal/router/resumer_test.go b/cmd/atenet/internal/router/resumer_test.go index 1ab6308aa..641bc2793 100644 --- a/cmd/atenet/internal/router/resumer_test.go +++ b/cmd/atenet/internal/router/resumer_test.go @@ -58,23 +58,51 @@ func TestActorResumer_ResumeActor(t *testing.T) { Status: ateapipb.Actor_STATUS_RUNNING, AteomPodIp: expectedIP, }, + Resumed: true, }, nil }, } resumer := NewActorResumer(mock) - actor, err := resumer.ResumeActor(context.Background(), testActorRef) + actor, outcome, err := resumer.ResumeActor(context.Background(), testActorRef) if err != nil { t.Fatalf("unexpected error: %v", err) } if actor.GetAteomPodIp() != expectedIP { t.Errorf("expected IP %q, got %q", expectedIP, actor.GetAteomPodIp()) } + if outcome != ResumeOutcomeTriggered { + t.Errorf("expected outcome %q, got %q", ResumeOutcomeTriggered, outcome) + } if resumeCalled != 1 { t.Errorf("expected ResumeActor called 1 time, got %d", resumeCalled) } }) + t.Run("WarmRouting_Disambiguation", func(t *testing.T) { + mock := &resumerMockClient{ + resumeFn: func(ctx context.Context, in *ateapipb.ResumeActorRequest, opts ...grpc.CallOption) (*ateapipb.ResumeActorResponse, error) { + return &ateapipb.ResumeActorResponse{ + Actor: &ateapipb.Actor{ + Metadata: &ateapipb.ResourceMetadata{Name: testActorName}, + Status: ateapipb.Actor_STATUS_RUNNING, + AteomPodIp: expectedIP, + }, + Resumed: false, + }, nil + }, + } + + resumer := NewActorResumer(mock) + _, outcome, err := resumer.ResumeActor(context.Background(), testActorRef) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if outcome != ResumeOutcomeNone { + t.Errorf("expected outcome %q for warm routing, got %q", ResumeOutcomeNone, outcome) + } + }) + t.Run("RetryOnAbortedConflict", func(t *testing.T) { var resumeCalled int mock := &resumerMockClient{ @@ -89,18 +117,22 @@ func TestActorResumer_ResumeActor(t *testing.T) { Status: ateapipb.Actor_STATUS_RUNNING, AteomPodIp: expectedIP, }, + Resumed: true, }, nil }, } resumer := NewActorResumer(mock) - actor, err := resumer.ResumeActor(context.Background(), testActorRef) + actor, outcome, err := resumer.ResumeActor(context.Background(), testActorRef) if err != nil { t.Fatalf("unexpected error: %v", err) } if actor.GetAteomPodIp() != expectedIP { t.Errorf("expected IP %q, got %q", expectedIP, actor.GetAteomPodIp()) } + if outcome != ResumeOutcomeTriggered { + t.Errorf("expected outcome %q, got %q", ResumeOutcomeTriggered, outcome) + } if resumeCalled != 3 { t.Errorf("expected ResumeActor called 3 times, got %d", resumeCalled) } @@ -114,13 +146,16 @@ func TestActorResumer_ResumeActor(t *testing.T) { } resumer := NewActorResumer(mock) - _, err := resumer.ResumeActor(context.Background(), testActorRef) + _, outcome, err := resumer.ResumeActor(context.Background(), testActorRef) if got := status.Code(err); got != codes.NotFound { t.Errorf("expected gRPC code NotFound, got %v (err=%v)", got, err) } + if outcome != ResumeOutcomeNone { + t.Errorf("expected outcome %q on error, got %q", ResumeOutcomeNone, outcome) + } }) - t.Run("SingleflightDeduplication", func(t *testing.T) { + t.Run("SingleflightDeduplication_Disambiguation", func(t *testing.T) { var resumeCalled int var mu sync.Mutex @@ -136,6 +171,7 @@ func TestActorResumer_ResumeActor(t *testing.T) { Status: ateapipb.Actor_STATUS_RUNNING, AteomPodIp: expectedIP, }, + Resumed: true, }, nil }, } @@ -145,17 +181,19 @@ func TestActorResumer_ResumeActor(t *testing.T) { var wg sync.WaitGroup const concurrentRequests = 10 results := make([]*ateapipb.Actor, concurrentRequests) + outcomes := make([]ResumeOutcome, concurrentRequests) errs := make([]error, concurrentRequests) wg.Add(concurrentRequests) for i := 0; i < concurrentRequests; i++ { go func(idx int) { defer wg.Done() - results[idx], errs[idx] = resumer.ResumeActor(context.Background(), testActorRef) + results[idx], outcomes[idx], errs[idx] = resumer.ResumeActor(context.Background(), testActorRef) }(i) } wg.Wait() + var triggeredCount, joinedCount int for i := 0; i < concurrentRequests; i++ { if errs[i] != nil { t.Fatalf("request %d failed: %v", i, errs[i]) @@ -163,6 +201,21 @@ func TestActorResumer_ResumeActor(t *testing.T) { if results[i].GetAteomPodIp() != expectedIP { t.Errorf("request %d expected IP %q, got %q", i, expectedIP, results[i].GetAteomPodIp()) } + switch outcomes[i] { + case ResumeOutcomeTriggered: + triggeredCount++ + case ResumeOutcomeJoined: + joinedCount++ + default: + t.Errorf("unexpected outcome for request %d: %q", i, outcomes[i]) + } + } + + if triggeredCount != 1 { + t.Errorf("expected exactly 1 request to have outcome 'triggered', got %d", triggeredCount) + } + if joinedCount != concurrentRequests-1 { + t.Errorf("expected %d requests to have outcome 'joined', got %d", concurrentRequests-1, joinedCount) } mu.Lock() @@ -199,7 +252,7 @@ func TestActorResumer_Parking(t *testing.T) { } resumer := NewActorResumer(mock, withParking(ParkedRequestConfig{Max: 1, Budget: 5 * time.Second})) - actor, err := resumer.ResumeActor(context.Background(), testActorRef) + actor, _, err := resumer.ResumeActor(context.Background(), testActorRef) if err != nil { t.Fatalf("unexpected error: %v", err) } @@ -228,7 +281,7 @@ func TestActorResumer_Parking(t *testing.T) { // Budget large enough for a few ~100ms-spaced retries before it elapses; // the pool never frees up. resumer := NewActorResumer(mock, withParking(ParkedRequestConfig{Max: 1, Budget: 1500 * time.Millisecond})) - _, err := resumer.ResumeActor(context.Background(), testActorRef) + _, _, err := resumer.ResumeActor(context.Background(), testActorRef) // The client must see the meaningful capacity error, not a generic // timeout: status.Code must unwrap through the budget-exhaustion marker. if got := status.Code(err); got != codes.FailedPrecondition { @@ -265,7 +318,7 @@ func TestActorResumer_Parking(t *testing.T) { } resumer := NewActorResumer(mock, withParking(ParkedRequestConfig{Max: 1, Budget: 5 * time.Second})) - actor, err := resumer.ResumeActor(context.Background(), testActorRef) + actor, _, err := resumer.ResumeActor(context.Background(), testActorRef) if err != nil { t.Fatalf("unexpected error: %v", err) } @@ -292,7 +345,7 @@ func TestActorResumer_Parking(t *testing.T) { } resumer := NewActorResumer(mock) - _, err := resumer.ResumeActor(context.Background(), testActorRef) + _, _, err := resumer.ResumeActor(context.Background(), testActorRef) if got := status.Code(err); got != codes.Unavailable { t.Errorf("expected Unavailable, got %v (err=%v)", got, err) } @@ -327,7 +380,7 @@ func TestActorResumer_Parking(t *testing.T) { } resumer := NewActorResumer(mock, withParking(ParkedRequestConfig{Max: 1, Budget: 300 * time.Millisecond})) - _, err := resumer.ResumeActor(context.Background(), testActorRef) + _, _, err := resumer.ResumeActor(context.Background(), testActorRef) // The deadline landed mid-RPC; the client must still see the capacity // error (503 "no free workers available"), not a generic timeout (504). if got := status.Code(err); got != codes.FailedPrecondition { @@ -353,7 +406,7 @@ func TestActorResumer_Parking(t *testing.T) { // Default constructor => parking disabled => fail-fast. resumer := NewActorResumer(mock) - _, err := resumer.ResumeActor(context.Background(), testActorRef) + _, _, err := resumer.ResumeActor(context.Background(), testActorRef) if got := status.Code(err); got != codes.FailedPrecondition { t.Errorf("expected FailedPrecondition, got %v (err=%v)", got, err) } @@ -403,7 +456,7 @@ func TestActorResumer_CallerCancelDoesNotAbortFlight(t *testing.T) { ctx1, cancel := context.WithCancel(context.Background()) errCh := make(chan error, 1) go func() { - _, err := resumer.ResumeActor(ctx1, testActorRef) + _, _, err := resumer.ResumeActor(ctx1, testActorRef) errCh <- err }() <-started @@ -425,7 +478,7 @@ func TestActorResumer_CallerCancelDoesNotAbortFlight(t *testing.T) { } resCh := make(chan result, 1) go func() { - a, rerr := resumer.ResumeActor(context.Background(), testActorRef) + a, _, rerr := resumer.ResumeActor(context.Background(), testActorRef) resCh <- result{a, rerr} }() // Give caller 2 a moment to join before releasing the flight, so the From 42342ee221879c57dc555640f613b7a53034dd94 Mon Sep 17 00:00:00 2001 From: Angela Hu Date: Wed, 29 Jul 2026 10:08:01 -0700 Subject: [PATCH 3/5] feat(atenet): extend route duration metric with resume and outcome labels --- cmd/atenet/internal/router/extproc.go | 59 +++++++++++++------ cmd/atenet/internal/router/extproc_test.go | 68 +++++++++++++++++++++- 2 files changed, 106 insertions(+), 21 deletions(-) diff --git a/cmd/atenet/internal/router/extproc.go b/cmd/atenet/internal/router/extproc.go index 2463930f5..4d22b489e 100644 --- a/cmd/atenet/internal/router/extproc.go +++ b/cmd/atenet/internal/router/extproc.go @@ -31,6 +31,8 @@ import ( "go.opentelemetry.io/otel/metric" "go.opentelemetry.io/otel/propagation" "google.golang.org/grpc" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" "github.com/agent-substrate/substrate/pkg/proto/ateapipb" ) @@ -94,8 +96,10 @@ func (s *ExtProcServer) Process(stream extprocv3.ExternalProcessor_ProcessServer switch reqType := req.Request.(type) { case *extprocv3.ProcessingRequest_RequestHeaders: start := time.Now() - hResponse, rqm, target, tmplNs, tmplName, err := s.handleRequestHeaders(stream.Context(), reqType.RequestHeaders) + hResponse, rqm, target, tmplNs, tmplName, resumeOutcome, err := s.handleRequestHeaders(stream.Context(), reqType.RequestHeaders) elapsed := time.Since(start) + outcomeStr := classifyOutcome(err) + resumeStr := string(resumeOutcome) if err != nil { slog.ErrorContext(stream.Context(), "Error during ext_proc RequestHeaders processing", slog.String("err", err.Error())) var reqErr *reqError @@ -104,11 +108,11 @@ func (s *ExtProcServer) Process(stream extprocv3.ExternalProcessor_ProcessServer } else { resp = immediateResponse(envoy_type.StatusCode_InternalServerError, err.Error()) } - s.recordRouteDuration(stream.Context(), elapsed, tmplNs, tmplName, classifyOutcome(err)) + s.recordRouteDuration(stream.Context(), elapsed, tmplNs, tmplName, outcomeStr, resumeStr) s.recorder.AddRouterRequest(start, elapsed, "Error", "-", rqm) } else { resp.Response = &extprocv3.ProcessingResponse_RequestHeaders{RequestHeaders: hResponse} - s.recordRouteDuration(stream.Context(), elapsed, tmplNs, tmplName, "ok") + s.recordRouteDuration(stream.Context(), elapsed, tmplNs, tmplName, outcomeStr, resumeStr) s.recorder.AddRouterRequest(start, elapsed, "Route ok", target, rqm) } @@ -132,7 +136,7 @@ func (s *ExtProcServer) Process(stream extprocv3.ExternalProcessor_ProcessServer func (s *ExtProcServer) handleRequestHeaders( ctx context.Context, reqHeaders *extprocv3.HttpHeaders, -) (*extprocv3.HeadersResponse, *requestMetadata, string, string, string, error) { +) (*extprocv3.HeadersResponse, *requestMetadata, string, string, string, ResumeOutcome, error) { metadata := newRequestMetadata(reqHeaders.Headers.GetHeaders()) slog.InfoContext(ctx, "Request", slog.String("host", metadata.host)) @@ -147,7 +151,7 @@ func (s *ExtProcServer) handleRequestHeaders( actorRef, err := parseActorRef(metadata.host) if err != nil { // Host is invalid, respond with 404. - return nil, metadata, "", "", "", invalidHostErr(metadata.host, err) + return nil, metadata, "", "", "", ResumeOutcomeNone, invalidHostErr(metadata.host, err) } // Admit the request to the parking lot before resuming. While resume is @@ -157,15 +161,14 @@ func (s *ExtProcServer) handleRequestHeaders( // backpressure instead of queueing without bound. release, ok := s.parking.enter(ctx) if !ok { - return nil, metadata, "", "", "", parkingFullErr(actorRef.String()) + return nil, metadata, "", "", "", ResumeOutcomeNone, parkingFullErr(actorRef.String()) } slog.InfoContext(ctx, "ResumeActor", slog.Any("actor", actorRef)) - actor, err := s.resumer.ResumeActor(ctx, actorRef) + actor, resumeOutcome, err := s.resumer.ResumeActor(ctx, actorRef) release(parkOutcomeFor(err)) - if err != nil { - return nil, metadata, "", "", "", mapResumeError(actorRef, err) + return nil, metadata, "", "", "", resumeOutcome, mapResumeError(actorRef, err) } // Actor template identity, used as low-cardinality route-latency metric @@ -180,7 +183,7 @@ func (s *ExtProcServer) handleRequestHeaders( slog.String("workerIP", workerIP)) if ip := net.ParseIP(workerIP); ip == nil { - return nil, metadata, "", tmplNs, tmplName, newReqError(envoy_type.StatusCode_InternalServerError, + return nil, metadata, "", tmplNs, tmplName, resumeOutcome, newReqError(envoy_type.StatusCode_InternalServerError, "actor %s routing failed", actorRef) } @@ -197,10 +200,10 @@ func (s *ExtProcServer) handleRequestHeaders( Response: &extprocv3.CommonResponse{ HeaderMutation: mutation, }, - }, metadata, targetAddr, tmplNs, tmplName, nil + }, metadata, targetAddr, tmplNs, tmplName, resumeOutcome, nil } -func (s *ExtProcServer) recordRouteDuration(ctx context.Context, d time.Duration, tmplNs, tmplName, outcome string) { +func (s *ExtProcServer) recordRouteDuration(ctx context.Context, d time.Duration, tmplNs, tmplName, outcome, resume string) { if s.routeDuration == nil { return } @@ -208,18 +211,38 @@ func (s *ExtProcServer) recordRouteDuration(ctx context.Context, d time.Duration attribute.String("actor_template_namespace", tmplNs), attribute.String("actor_template_name", tmplName), attribute.String("outcome", outcome), + attribute.String("resume", resume), )) } func classifyOutcome(err error) string { - switch { - case errors.Is(err, context.Canceled), errors.Is(err, context.DeadlineExceeded): + if err == nil { + return "ok" + } + if errors.Is(err, context.Canceled) || status.Code(err) == codes.Canceled { return "cancelled" - default: - var re *reqError - if errors.As(err, &re) && re.statusCode == int(envoy_type.StatusCode_NotFound) { + } + if errors.Is(err, context.DeadlineExceeded) || status.Code(err) == codes.DeadlineExceeded { + return "timeout" + } + switch status.Code(err) { + case codes.FailedPrecondition: + return "no_capacity" + case codes.Aborted: + return "lock_conflict" + case codes.NotFound: + return "not_found" + } + var re *reqError + if errors.As(err, &re) { + switch envoy_type.StatusCode(re.statusCode) { + case envoy_type.StatusCode_NotFound: return "not_found" + case envoy_type.StatusCode_ServiceUnavailable: + return "no_capacity" + case envoy_type.StatusCode_GatewayTimeout: + return "timeout" } - return "error" } + return "resume_error" } diff --git a/cmd/atenet/internal/router/extproc_test.go b/cmd/atenet/internal/router/extproc_test.go index b6a889c11..2da0fd207 100644 --- a/cmd/atenet/internal/router/extproc_test.go +++ b/cmd/atenet/internal/router/extproc_test.go @@ -69,7 +69,7 @@ func TestHandleRequestHeadersDoesNotLogSensitiveData(t *testing.T) { }, } - _, metadata, target, _, _, err := s.handleRequestHeaders(context.Background(), reqHeaders) + _, metadata, target, _, _, _, err := s.handleRequestHeaders(context.Background(), reqHeaders) if err != nil { t.Fatalf("unexpected error: %v", err) } @@ -204,7 +204,7 @@ func TestExtProcHeadersEvaluation(t *testing.T) { }, } - res, metadata, target, _, _, err := s.handleRequestHeaders(context.Background(), reqHeaders) + res, metadata, target, _, _, _, err := s.handleRequestHeaders(context.Background(), reqHeaders) if tc.expectErr { if err == nil { t.Fatalf("expected error but got nil") @@ -286,7 +286,7 @@ func TestExtProc_ParkingLotFull(t *testing.T) { }, } - _, _, _, _, _, err := s.handleRequestHeaders(context.Background(), reqHeaders) + _, _, _, _, _, _, err := s.handleRequestHeaders(context.Background(), reqHeaders) if err == nil { t.Fatal("expected error when parking lot is full") } @@ -304,3 +304,65 @@ func TestExtProc_ParkingLotFull(t *testing.T) { t.Error("resume must not be attempted for a shed request") } } + +func TestClassifyOutcome(t *testing.T) { + tests := []struct { + name string + err error + expected string + }{ + { + name: "nil error maps to ok", + err: nil, + expected: "ok", + }, + { + name: "context Canceled maps to cancelled", + err: context.Canceled, + expected: "cancelled", + }, + { + name: "context DeadlineExceeded maps to timeout", + err: context.DeadlineExceeded, + expected: "timeout", + }, + { + name: "FailedPrecondition gRPC code maps to no_capacity", + err: status.Error(codes.FailedPrecondition, "capacity full"), + expected: "no_capacity", + }, + { + name: "Aborted gRPC code maps to lock_conflict", + err: status.Error(codes.Aborted, "lock conflict"), + expected: "lock_conflict", + }, + { + name: "NotFound gRPC code maps to not_found", + err: status.Error(codes.NotFound, "missing"), + expected: "not_found", + }, + { + name: "StatusCode_NotFound reqError maps to not_found", + err: newReqError(envoy_type.StatusCode_NotFound, "missing"), + expected: "not_found", + }, + { + name: "StatusCode_ServiceUnavailable reqError maps to no_capacity", + err: newReqError(envoy_type.StatusCode_ServiceUnavailable, "no free workers"), + expected: "no_capacity", + }, + { + name: "Unknown error maps to resume_error", + err: errors.New("internal storage glitch"), + expected: "resume_error", + }, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + if got := classifyOutcome(tc.err); got != tc.expected { + t.Errorf("classifyOutcome(%v) = %q, want %q", tc.err, got, tc.expected) + } + }) + } +} From 3497427f5961be4c8a1ef7ef4131a9f660eb55e9 Mon Sep 17 00:00:00 2001 From: Angela Hu Date: Wed, 29 Jul 2026 10:13:01 -0700 Subject: [PATCH 4/5] test(e2e): configure atenet-router OTLP export and add metrics to e2e --- .../controlapi/workflow_resume_test.go | 10 +++- cmd/atenet/internal/router/extproc.go | 16 ++++-- cmd/atenet/internal/router/extproc_test.go | 52 +++++++++++++++++++ cmd/atenet/internal/router/resumer.go | 36 +++++++++---- docs/observability.md | 6 ++- internal/ateattr/ateattr.go | 19 +++++-- internal/e2e/collector_metrics.go | 1 + internal/e2e/suites/metrics/metrics_test.go | 13 +++++ manifests/ate-install/kind/kustomization.yaml | 15 ++++++ pkg/proto/ateapipb/ateapi.pb.go | 13 +++-- 10 files changed, 154 insertions(+), 27 deletions(-) diff --git a/cmd/ateapi/internal/controlapi/workflow_resume_test.go b/cmd/ateapi/internal/controlapi/workflow_resume_test.go index 568aaca08..b466dc295 100644 --- a/cmd/ateapi/internal/controlapi/workflow_resume_test.go +++ b/cmd/ateapi/internal/controlapi/workflow_resume_test.go @@ -454,8 +454,14 @@ func TestResumeActorWorkflow_RejectedAndIdempotentPaths(t *testing.T) { if actor.GetStatus() != tc.wantStatus { t.Errorf("returned status = %v, want %v", actor.GetStatus(), tc.wantStatus) } - if tc.seedStatus == ateapipb.Actor_STATUS_RUNNING && resumed { - t.Errorf("expected resumed = false for already running actor, got true") + if tc.seedStatus == ateapipb.Actor_STATUS_RUNNING { + if resumed { + t.Errorf("expected resumed = false for already running actor, got true") + } + } else { + if !resumed { + t.Errorf("expected resumed = true for cold activation, got false") + } } } diff --git a/cmd/atenet/internal/router/extproc.go b/cmd/atenet/internal/router/extproc.go index 4d22b489e..f6e69ff61 100644 --- a/cmd/atenet/internal/router/extproc.go +++ b/cmd/atenet/internal/router/extproc.go @@ -23,11 +23,11 @@ import ( "net" "time" + "github.com/agent-substrate/substrate/internal/ateattr" extprocv3 "github.com/envoyproxy/go-control-plane/envoy/service/ext_proc/v3" envoy_type "github.com/envoyproxy/go-control-plane/envoy/type/v3" "go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc" "go.opentelemetry.io/otel" - "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/metric" "go.opentelemetry.io/otel/propagation" "google.golang.org/grpc" @@ -208,10 +208,10 @@ func (s *ExtProcServer) recordRouteDuration(ctx context.Context, d time.Duration return } s.routeDuration.Record(ctx, d.Seconds(), metric.WithAttributes( - attribute.String("actor_template_namespace", tmplNs), - attribute.String("actor_template_name", tmplName), - attribute.String("outcome", outcome), - attribute.String("resume", resume), + ateattr.TemplateNamespaceKey.String(tmplNs), + ateattr.TemplateNameKey.String(tmplName), + ateattr.RouterOutcomeKey.String(outcome), + ateattr.RouterResumeKey.String(resume), )) } @@ -232,6 +232,10 @@ func classifyOutcome(err error) string { return "lock_conflict" case codes.NotFound: return "not_found" + case codes.Unavailable: + return "unavailable" + case codes.ResourceExhausted: + return "rate_limited" } var re *reqError if errors.As(err, &re) { @@ -242,6 +246,8 @@ func classifyOutcome(err error) string { return "no_capacity" case envoy_type.StatusCode_GatewayTimeout: return "timeout" + case envoy_type.StatusCode_TooManyRequests: + return "rate_limited" } } return "resume_error" diff --git a/cmd/atenet/internal/router/extproc_test.go b/cmd/atenet/internal/router/extproc_test.go index 2da0fd207..8bb83d393 100644 --- a/cmd/atenet/internal/router/extproc_test.go +++ b/cmd/atenet/internal/router/extproc_test.go @@ -28,6 +28,9 @@ import ( corev3 "github.com/envoyproxy/go-control-plane/envoy/config/core/v3" extprocv3 "github.com/envoyproxy/go-control-plane/envoy/service/ext_proc/v3" envoy_type "github.com/envoyproxy/go-control-plane/envoy/type/v3" + "go.opentelemetry.io/otel/attribute" + sdkmetric "go.opentelemetry.io/otel/sdk/metric" + "go.opentelemetry.io/otel/sdk/metric/metricdata" "google.golang.org/grpc" "google.golang.org/grpc/codes" "google.golang.org/grpc/status" @@ -341,6 +344,16 @@ func TestClassifyOutcome(t *testing.T) { err: status.Error(codes.NotFound, "missing"), expected: "not_found", }, + { + name: "Unavailable gRPC code maps to unavailable", + err: status.Error(codes.Unavailable, "control-plane down"), + expected: "unavailable", + }, + { + name: "ResourceExhausted gRPC code maps to rate_limited", + err: status.Error(codes.ResourceExhausted, "rate limit exceeded"), + expected: "rate_limited", + }, { name: "StatusCode_NotFound reqError maps to not_found", err: newReqError(envoy_type.StatusCode_NotFound, "missing"), @@ -351,6 +364,11 @@ func TestClassifyOutcome(t *testing.T) { err: newReqError(envoy_type.StatusCode_ServiceUnavailable, "no free workers"), expected: "no_capacity", }, + { + name: "StatusCode_TooManyRequests reqError maps to rate_limited", + err: newReqError(envoy_type.StatusCode_TooManyRequests, "rate limited"), + expected: "rate_limited", + }, { name: "Unknown error maps to resume_error", err: errors.New("internal storage glitch"), @@ -366,3 +384,37 @@ func TestClassifyOutcome(t *testing.T) { }) } } + +func TestRecordRouteDuration_Attributes(t *testing.T) { + reader := sdkmetric.NewManualReader() + mp := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader)) + h, err := mp.Meter("atenet-router").Float64Histogram(routeDurationMetricName) + if err != nil { + t.Fatalf("failed to create histogram: %v", err) + } + + s := NewExtProcServer(50051, nil, h, ParkedRequestConfig{}, nil) + s.recordRouteDuration(context.Background(), 10*time.Millisecond, "team-a-ns", "tmpl-a", classifyOutcome(nil), string(ResumeOutcomeTriggered)) + + var rm metricdata.ResourceMetrics + if err := reader.Collect(context.Background(), &rm); err != nil { + t.Fatalf("Collect failed: %v", err) + } + + dp := rm.ScopeMetrics[0].Metrics[0].Data.(metricdata.Histogram[float64]).DataPoints[0] + wantAttrs := map[string]string{ + "ate.template.namespace": "team-a-ns", + "ate.template.name": "tmpl-a", + "ate.router.outcome": "ok", + "ate.router.resume": "triggered", + } + + for k, want := range wantAttrs { + val, exists := dp.Attributes.Value(attribute.Key(k)) + if !exists { + t.Errorf("missing metric attribute %q", k) + } else if val.AsString() != want { + t.Errorf("attribute %q = %q, want %q", k, val.AsString(), want) + } + } +} diff --git a/cmd/atenet/internal/router/resumer.go b/cmd/atenet/internal/router/resumer.go index 151eb0467..e6b01f705 100644 --- a/cmd/atenet/internal/router/resumer.go +++ b/cmd/atenet/internal/router/resumer.go @@ -69,17 +69,19 @@ func (e *budgetExhaustedError) Unwrap() error { return e.lastErr } type ResumeOutcome string const ( - // ResumeOutcomeNone indicates the actor was already running (steady-state warm route). - ResumeOutcomeNone ResumeOutcome = "none" - // ResumeOutcomeTriggered indicates this request won the singleflight lock and initiated cold activation. - ResumeOutcomeTriggered ResumeOutcome = "triggered" - // ResumeOutcomeJoined indicates this request parked on an in-flight singleflight resume. - ResumeOutcomeJoined ResumeOutcome = "joined" + ResumeOutcomeNone ResumeOutcome = ateattr.RouterResumeNone + ResumeOutcomeTriggered ResumeOutcome = ateattr.RouterResumeTriggered + ResumeOutcomeJoined ResumeOutcome = ateattr.RouterResumeJoined ) type resumeCallResult struct { - actor *ateapipb.Actor - resumed bool + actor *ateapipb.Actor + // resumed is true if ResumeActor call executed a cold activation + // false if the actor was already running + resumed bool + // leaderID is the unique request ID (reqID) of the leader that initiated + // the singleflight execution. It helps disambiguates the leader caller + // (ResumeOutcomeTriggered) from joiner callers (ResumeOutcomeJoined). leaderID uint64 err error } @@ -98,7 +100,10 @@ type ActorResumer struct { budget time.Duration // backoff paces the retries within the budget. backoff wait.Backoff - nextID uint64 + // nextID is a counter assigned to each incoming ResumeActor call. + // Used as a unique ID to identify requests (reqID) and disambiguate the + // leader vs joiners for singleflight outcome classification. + nextID uint64 } // resumerOption configures an ActorResumer. @@ -190,8 +195,8 @@ func (r *ActorResumer) ResumeActor(ctx context.Context, actorRef resources.Actor } if r.retryable(err) { - lastRetryErr = err - return false, nil + lastRetryErr = err // remember it in case the budget elapses + return false, nil // park: retry until the budget elapses } return false, err }) @@ -224,6 +229,8 @@ func (r *ActorResumer) ResumeActor(ctx context.Context, actorRef resources.Actor select { case <-ctx.Done(): + // The caller's request context was canceled before the singleflight resume completed. + // Return early with ResumeOutcomeNone ("none") return nil, ResumeOutcomeNone, ctx.Err() case res := <-ch: callRes, _ := res.Val.(*resumeCallResult) @@ -233,10 +240,17 @@ func (r *ActorResumer) ResumeActor(ctx context.Context, actorRef resources.Actor } return nil, ResumeOutcomeNone, status.Error(codes.Internal, "resume call returned nil result") } + + // On error, return ResumeOutcomeNone ("none") so the failure is tagged + // under the 'outcome' label rather than misreported as an activation. if callRes.err != nil { return nil, ResumeOutcomeNone, callRes.err } + // Disambiguate singleflight resume outcome: + // - ResumeOutcomeNone ("none"): resumed == false, actor was already active/running. + // - ResumeOutcomeTriggered ("triggered"): Cold activation leader (resumed == true, caller's reqID == leaderID). + // - ResumeOutcomeJoined ("joined"): Cold activation joiner (resumed == true, caller's reqID != leaderID). outcome := ResumeOutcomeNone if callRes.resumed { if callRes.leaderID == reqID { diff --git a/docs/observability.md b/docs/observability.md index 0c77eb4b5..78ab8e922 100644 --- a/docs/observability.md +++ b/docs/observability.md @@ -110,11 +110,15 @@ Agent Substrate emits foundational OpenTelemetry system and server metrics to mo | Metric | Emitted by | Type | Measures | |--------|------------|------|----------| | `rpc.server.call.duration` | ateapi & atelet (gRPC servers, via `otelgrpc`) | histogram | per-method gRPC latency, request rate, and errors (labels `rpc.method`, `rpc.response.status_code`) | -| `atenet.router.route.duration` | atenet-router | histogram | Substrate E2E — Envoy receiving a request to Envoy forwarding it to the resolved worker, excluding actor compute and the response | +| `atenet.router.route.duration` | atenet-router | histogram | Substrate E2E — Envoy receiving a request to Envoy forwarding it to the resolved worker, excluding actor compute and the response (labels `ate.template.namespace`, `ate.template.name`, `ate.router.outcome`, `ate.router.resume`) | | `atelet.snapshot.size` | atelet | histogram | uncompressed size in bytes of each gVisor snapshot image written during checkpoint (labels `kind`, `actor_template_namespace`, `actor_template_name`) | The table lists the OpenTelemetry instrument names. How a name appears in a query depends on the backend (Cloud Monitoring (GMP) / Kind collector). +For `atenet.router.route.duration`: +* `ate.router.outcome` categorizes the route attempt result: `ok`, `cancelled`, `timeout`, `no_capacity`, `lock_conflict`, `not_found`, `unavailable`, `rate_limited`, or `resume_error`. +* `ate.router.resume` indicates the singleflight execution state of actor resumption: `none` (actor already running), `triggered` (initiated cold activation), or `joined` (parked on in-flight activation). + ### Local Metrics with Prometheus (Kind Cluster) For local development inside a `kind` cluster, Agent Substrate automatically provisions a Prometheus server in the `otel-system` namespace. diff --git a/internal/ateattr/ateattr.go b/internal/ateattr/ateattr.go index 670342008..e6361d5ed 100644 --- a/internal/ateattr/ateattr.go +++ b/internal/ateattr/ateattr.go @@ -50,9 +50,22 @@ const ( // WorkerStateKey stays worker-rooted rather than nesting under the pool so it // can grow siblings. const ( - WorkerPoolNameKey = attribute.Key("ate.workerpool.name") - WorkerStateKey = attribute.Key("ate.worker.state") - SandboxClassKey = attribute.Key("ate.sandbox.class") + WorkerPoolNamespaceKey = attribute.Key("ate.workerpool.namespace") + WorkerPoolNameKey = attribute.Key("ate.workerpool.name") + WorkerStateKey = attribute.Key("ate.worker.state") + SandboxClassKey = attribute.Key("ate.sandbox.class") + RouterResumeKey = attribute.Key("ate.router.resume") + RouterOutcomeKey = attribute.Key("ate.router.outcome") +) + +// Values for RouterResumeKey. +const ( + // RouterResumeNone indicates the actor was already running (steady-state route). + RouterResumeNone = "none" + // RouterResumeTriggered indicates this request won the singleflight lock and initiated cold activation. + RouterResumeTriggered = "triggered" + // RouterResumeJoined indicates this request parked on an in-flight singleflight resume. + RouterResumeJoined = "joined" ) // Values for WorkerStateKey. Only idle and assigned are representable today; diff --git a/internal/e2e/collector_metrics.go b/internal/e2e/collector_metrics.go index 9733b50bf..64376012d 100644 --- a/internal/e2e/collector_metrics.go +++ b/internal/e2e/collector_metrics.go @@ -41,6 +41,7 @@ const ( // it pins the worker-count instrument introduced alongside this harness. var PlatformMetricPrefixes = []string{ "ate_workerpool_workers", + "atenet_router_route_duration", } // ScrapeCollectorMetrics port-forwards the kind stack's OTel Collector and reads diff --git a/internal/e2e/suites/metrics/metrics_test.go b/internal/e2e/suites/metrics/metrics_test.go index 50d4a9010..46a912eee 100644 --- a/internal/e2e/suites/metrics/metrics_test.go +++ b/internal/e2e/suites/metrics/metrics_test.go @@ -27,6 +27,7 @@ import ( "time" "github.com/agent-substrate/substrate/internal/e2e" + "github.com/agent-substrate/substrate/internal/resources" "github.com/agent-substrate/substrate/pkg/proto/ateapipb" ) @@ -71,6 +72,18 @@ func TestPlatformMetricsEmitted(t *testing.T) { // they add the drive steps their instruments need. resume(t, ctx, clients, actorID) + // Drive request through the router so Envoy ext_proc emits atenet_router_route_duration. + rClient, err := e2e.NewRouterClient(ctx) + if err != nil { + t.Fatalf("NewRouterClient: %v", err) + } + defer rClient.Close() + resp, err := rClient.Get(ctx, resources.ActorRef{Atespace: metricsAtespace, Name: actorID}, "/") + if err != nil { + t.Fatalf("rClient.Get: %v", err) + } + _ = resp.Body.Close() + deadline := time.Now().Add(2 * time.Minute) var missing []string var ateomSeen bool diff --git a/manifests/ate-install/kind/kustomization.yaml b/manifests/ate-install/kind/kustomization.yaml index a6dcca402..ba79b25d5 100644 --- a/manifests/ate-install/kind/kustomization.yaml +++ b/manifests/ate-install/kind/kustomization.yaml @@ -70,3 +70,18 @@ patches: env: - name: OTEL_EXPORTER_OTLP_ENDPOINT value: http://opentelemetry-collector.otel-system.svc:4317 + - patch: |- + apiVersion: apps/v1 + kind: Deployment + metadata: + name: atenet-router + namespace: ate-system + spec: + template: + spec: + containers: + - name: atenet-router + env: + - name: OTEL_EXPORTER_OTLP_ENDPOINT + value: http://opentelemetry-collector.otel-system.svc:4317 + diff --git a/pkg/proto/ateapipb/ateapi.pb.go b/pkg/proto/ateapipb/ateapi.pb.go index 6a7d2cdb8..974d674b2 100644 --- a/pkg/proto/ateapipb/ateapi.pb.go +++ b/pkg/proto/ateapipb/ateapi.pb.go @@ -1462,9 +1462,11 @@ func (x *ResumeActorRequest) GetBoot() bool { } type ResumeActorResponse struct { - state protoimpl.MessageState `protogen:"open.v1"` - Actor *Actor `protobuf:"bytes,1,opt,name=actor,proto3" json:"actor,omitempty"` - Resumed bool `protobuf:"varint,2,opt,name=resumed,proto3" json:"resumed,omitempty"` + state protoimpl.MessageState `protogen:"open.v1"` + Actor *Actor `protobuf:"bytes,1,opt,name=actor,proto3" json:"actor,omitempty"` + // True if a resume workflow was executed to activate the actor. + // False if the actor was already RUNNING. + Resumed bool `protobuf:"varint,2,opt,name=resumed,proto3" json:"resumed,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } @@ -2432,9 +2434,10 @@ const file_ateapi_proto_rawDesc = "" + "\x05actor\x18\x01 \x01(\v2\r.ateapi.ActorR\x05actor\"Q\n" + "\x12ResumeActorRequest\x12'\n" + "\x05actor\x18\x01 \x01(\v2\x11.ateapi.ObjectRefR\x05actor\x12\x12\n" + - "\x04boot\x18\x02 \x01(\bR\x04boot\":\n" + + "\x04boot\x18\x02 \x01(\bR\x04boot\"T\n" + "\x13ResumeActorResponse\x12#\n" + - "\x05actor\x18\x01 \x01(\v2\r.ateapi.ActorR\x05actor\"=\n" + + "\x05actor\x18\x01 \x01(\v2\r.ateapi.ActorR\x05actor\x12\x18\n" + + "\aresumed\x18\x02 \x01(\bR\aresumed\"=\n" + "\x12DeleteActorRequest\x12'\n" + "\x05actor\x18\x01 \x01(\v2\x11.ateapi.ObjectRefR\x05actor\"P\n" + "\x12ListWorkersRequest\x12\x1b\n" + From ef9183b5ddb948bb7cb006b6228ccf7d4f65085e Mon Sep 17 00:00:00 2001 From: Angela Hu Date: Fri, 31 Jul 2026 07:43:01 -0700 Subject: [PATCH 5/5] trigger ci