Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions cmd/ateapi/internal/controlapi/resume_actor.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand All @@ -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) field.ErrorList {
Expand Down
8 changes: 4 additions & 4 deletions cmd/ateapi/internal/controlapi/workflow.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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()

Expand All @@ -188,10 +188,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.
Expand Down
2 changes: 2 additions & 0 deletions cmd/ateapi/internal/controlapi/workflow_resume.go
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@ type ResumeState struct {
Actor *ateapipb.Actor
Worker *ateapipb.Worker
ActorTemplate *atev1alpha1.ActorTemplate
WasRunning bool
SnapshotLocation string
SnapshotScope ateapipb.SnapshotContentScope
}
Expand All @@ -75,6 +76,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 {
Expand Down
13 changes: 11 additions & 2 deletions cmd/ateapi/internal/controlapi/workflow_resume_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -448,7 +448,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)
Expand All @@ -460,6 +460,15 @@ 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 {
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")
}
}
}

got, err := st.GetActor(ctx, resources.ActorRef{Atespace: "team-a", Name: "id1"})
Expand Down Expand Up @@ -555,7 +564,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)
}
Expand Down
73 changes: 51 additions & 22 deletions cmd/atenet/internal/router/extproc.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,14 +23,16 @@ 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"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"

"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
)
Expand Down Expand Up @@ -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
Expand All @@ -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)
Comment thread
Angelawork marked this conversation as resolved.
s.recorder.AddRouterRequest(start, elapsed, "Route ok", target, rqm)
}

Expand All @@ -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))

Expand All @@ -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
Expand All @@ -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
Expand All @@ -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)
}

Expand All @@ -197,29 +200,55 @@ 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
}
s.routeDuration.Record(ctx, d.Seconds(), metric.WithAttributes(
attribute.String("actor_template_namespace", tmplNs),
attribute.String("actor_template_name", tmplName),
attribute.String("outcome", outcome),
ateattr.TemplateNamespaceKey.String(tmplNs),
ateattr.TemplateNameKey.String(tmplName),
ateattr.RouterOutcomeKey.String(outcome),
ateattr.RouterResumeKey.String(resume),
))
}

func classifyOutcome(err error) string {
Comment thread
Angelawork marked this conversation as resolved.
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"
case codes.Unavailable:
return "unavailable"
case codes.ResourceExhausted:
return "rate_limited"
}
Comment thread
Angelawork marked this conversation as resolved.
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"
case envoy_type.StatusCode_TooManyRequests:
return "rate_limited"
}
return "error"
}
return "resume_error"
}
Loading
Loading