From 38df08c32647649fbccad0a1522e6e8039fae8ea Mon Sep 17 00:00:00 2001 From: sam Date: Tue, 6 Oct 2026 17:39:46 +0800 Subject: [PATCH] Add execution latency observations and trace correlation --- docs/getting-started/operations.md | 8 + docs/zh/getting-started/operations.md | 10 +- services/core/internal/execution/delivery.go | 8 +- .../execution/directory_preparation.go | 2 +- .../execution/environment_admission.go | 23 ++- .../execution/execution_observation.go | 46 +++++ .../execution/execution_observation_test.go | 178 ++++++++++++++++++ .../execution/executor_observation.go | 8 +- .../execution/executor_preparation.go | 2 +- .../core/internal/execution/preparation.go | 5 +- .../internal/execution/prepared_dispatch.go | 20 +- .../internal/execution/runtime_compute.go | 3 + .../internal/execution/runtime_connections.go | 4 + .../execution/runtime_initialization.go | 2 + .../internal/execution/runtime_lifecycle.go | 9 +- .../internal/execution/runtime_pending.go | 2 + services/core/internal/execution/worker.go | 2 + 17 files changed, 317 insertions(+), 15 deletions(-) create mode 100644 services/core/internal/execution/execution_observation.go create mode 100644 services/core/internal/execution/execution_observation_test.go diff --git a/docs/getting-started/operations.md b/docs/getting-started/operations.md index b692794be..8b210c993 100644 --- a/docs/getting-started/operations.md +++ b/docs/getting-started/operations.md @@ -193,3 +193,11 @@ Mutating `oac` commands hold `.oac.lock`. If another command holds it, retry aft Web signs administrators in with the Core key, checks the origin of every request, and forwards signed-in `/core/v1` requests to Core with the Core key, which stays on the server. It forwards `/v1` and `/api/v1` to Core unchanged, with the caller's own credential, serves only the non-secret node payload at `/node-install/`, and has no Docker or KVM access. Machine routes under `/api/v1` use their own enrollment and connection credentials. No service receives a Docker socket. Sandboxes are the isolation boundary ([Runtime and outer isolation](../concepts.md#runtime-and-outer-isolation)). Docker sandboxes share the node's kernel, and a Docker node is [root-equivalent](./nodes.md#what-the-installer-sets-up) on its host; microsandbox gives each sandbox a microVM with an explicit [network policy](./nodes.md#what-the-installer-sets-up). Core itself has no Docker socket or KVM access. + +## Execution latency + +Use the `environment input reserved` log to connect the submitting HTTP trace to `reservation_id` and `execution_trace_id`. The worker derives its diagnostic execution trace from the durable reservation identity, so observer disconnection and owner restart do not require a process-local trace map. This is an explicit link between the HTTP and execution traces, not continuation of the original HTTP span. Follow `environment input selected`, `execution preparation requested` and `environment input admitted` to connect the reservation to the Session, Environment, device, preparation request, Executor and Turn. `preparation_request_id` identifies the control request and is distinct from an HTTP `request_id`. For inputs included in Session creation, join `session initial input origin` to the worker reservation with the same `session_id` and `is_initial=true`; this covers both regular and streaming creation without an extra store query. + +`execution stage` records `stage`, `duration_ms` and `status` (`ok`, `error`, `cancelled` or `timeout`). Stages cover input validation/ownership, reservation and admission wait, execution configuration and admission promotion, lifecycle gate wait, Runtime observation/provisioning, provider create/resume and Environment initialization. Runtime connection logs record confirmed observed connection transitions, not the exact socket connection instant. Use Environment, allocation and node IDs to join background Runtime records to an input; a background record without the input trace is not an unbroken request span. Stage durations use a local monotonic clock. `queue_age_ms` instead compares selection time to the database reservation creation time and may reflect clock skew; it is not a measured scheduler-only interval. + +`control_ready_ms` starts after scheduling, Runtime readiness and configuration assembly. `start_control_ms` describes the control acknowledgement. `input_to_first_text_ms` starts immediately before execution delivery, includes start-control time and is not model-only TTFT. The admission, control and first-text intervals overlap; do not sum them as independent durations. A successful provider call describes that operation, not completed daemon connection or a completed Turn. Model request start, response headers, first model text and retries remain unavailable unless the selected native adapter provides those observations. Logs contain correlation IDs and finite status categories, not input content, credentials or raw provider errors. diff --git a/docs/zh/getting-started/operations.md b/docs/zh/getting-started/operations.md index facb13092..4d4e07649 100644 --- a/docs/zh/getting-started/operations.md +++ b/docs/zh/getting-started/operations.md @@ -1,7 +1,7 @@ --- title: "管理你的安装" source: docs/getting-started/operations.md -source_hash: e60a6e96cd61692b6adc7664874ba0ed2ce489982e7da2fe2347631ee2a94c0a +source_hash: 263a48e467c9aeec1f8357f0e11b40c6f1b329dd6162f65eda3a1038acd4f551 --- 安装运维人员负责 Core 主机、存储和可用性。节点主机运行各自的服务;参阅[节点](nodes.md)。设置见[配置参考](../configuration.md)。 @@ -196,3 +196,11 @@ cd && rm -rf ~/.oac/core Web 使用 Core 密钥让管理员登录,检查每个请求来源,并用保留在服务器上的 Core 密钥将已登录的 `/core/v1` 请求转发到 Core。它把 `/v1` 和 `/api/v1` 原样转发给 Core,使用调用方自己的凭据;Web 仅在 `/node-install/` 提供不含密钥的节点文件,没有 Docker 或 KVM 访问权限。`/api/v1` 机器路由使用独立注册和连接凭据。没有服务持有 Docker 套接字。 沙箱是隔离边界([Runtime 与外层隔离](../concepts.md#runtime-and-outer-isolation))。Docker 沙箱共享节点内核,Docker 节点在主机上[等同于 root 权限](nodes.md#what-the-installer-sets-up);microsandbox 为每个沙箱提供具有显式[网络策略](nodes.md#what-the-installer-sets-up)的 microVM。Core 自身无 Docker 套接字或 KVM 访问权限。 + +## 执行延迟 {#execution-latency} + +通过 `environment input reserved` 日志将提交输入的 HTTP Trace 与 `reservation_id`、`execution_trace_id` 关联。worker 从持久化预留标识派生执行诊断 Trace,观察连接断开或 owner 重启都不依赖进程内 Trace 映射。这是 HTTP Trace 与执行 Trace 之间的明确关联,不是原 HTTP span 的延续。沿 `environment input selected`、`execution preparation requested`、`environment input admitted` 将预留与 Session、Environment、device、准备请求、Executor、Turn 关联。`preparation_request_id` 标识控制准备请求,与 HTTP `request_id` 不同。创建 Session 时附带输入的情况,通过相同 `session_id` 和 `is_initial=true` 将 `session initial input origin` 与 worker 的预留关联;普通和流式创建都支持,不增加 store 查询。 + +`execution stage` 记录 `stage`、`duration_ms` 与 `status`(`ok`、`error`、`cancelled`、`timeout`)。阶段覆盖输入验证/ownership、预留和等待 admission、执行配置与 admission 提升、生命周期 gate 等待、Runtime 观察和供给、provider create/resume、Environment 初始化。Runtime 连接日志记录已确认观察到的连接变化,不是 socket 建立的精确时间。通过 Environment、allocation、node ID 将后台 Runtime 记录与输入关联;缺少输入 Trace 的后台记录不能视为连续的请求 span。阶段耗时使用本进程单调时钟。`queue_age_ms` 则比较选择时间与数据库预留创建时间,可能受到时钟偏差影响,不能当作单独测量的调度等待。 + +`control_ready_ms` 从调度、Runtime 就绪和配置组装完成后开始。`start_control_ms` 描述控制确认。`input_to_first_text_ms` 从执行投递前开始,包含 start-control 时间,不是纯模型 TTFT。admission、control、首字区间存在重叠,不要当作独立耗时相加。provider 调用成功表示该操作成功,不等于 daemon 连接完成或 Turn 完成。除非原生适配器提供相应观测,模型请求开始、响应 headers、模型首字、重试仍为不可观测。日志仅记录关联 ID 与有限状态分类,不记录输入内容、凭据或 provider 原始错误。 diff --git a/services/core/internal/execution/delivery.go b/services/core/internal/execution/delivery.go index a9c2ea6d6..babd31a2a 100644 --- a/services/core/internal/execution/delivery.go +++ b/services/core/internal/execution/delivery.go @@ -7,6 +7,8 @@ import ( "strconv" "time" + obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" @@ -37,7 +39,11 @@ func requestCancellation(ctx context.Context, peer *runtimegateway.Session, runI } func send(ctx context.Context, peer *runtimegateway.Session, kind, runID string, payload any) error { - env, err := proto.NewEnvelope(kind, runID, payload) + trace := "" + if carrier, ok := obslog.TraceFromContext(ctx); ok { + trace = carrier.String() + } + env, err := proto.NewEnvelopeWithTrace(kind, runID, payload, trace) if err != nil { return err } diff --git a/services/core/internal/execution/directory_preparation.go b/services/core/internal/execution/directory_preparation.go index 427e400cb..4fa0725fb 100644 --- a/services/core/internal/execution/directory_preparation.go +++ b/services/core/internal/execution/directory_preparation.go @@ -29,7 +29,7 @@ func (d *Dispatcher) withPreparedWorkspace(owner context.Context, peer *runtimeg if err := d.configurePreparedEnvironment(session, environment, bound, &req); err != nil { return ErrExecutionUnavailable } - prepared, err := newPreparedStart(peer) + prepared, err := newPreparedStart(owner, peer) if err != nil { return ErrExecutionUnavailable } diff --git a/services/core/internal/execution/environment_admission.go b/services/core/internal/execution/environment_admission.go index 835027d90..318fa2c34 100644 --- a/services/core/internal/execution/environment_admission.go +++ b/services/core/internal/execution/environment_admission.go @@ -7,6 +7,8 @@ import ( "slices" "time" + obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" ) @@ -72,11 +74,17 @@ func (w *Worker) validateCreation(ctx context.Context, input sessions.CreateSess return nil } -func (w *Worker) submitEnvironmentInputs(ctx context.Context, session sessions.Session, key string, inputs []sessions.Input) ([]sessions.InputReceipt, error) { - if err := w.validateEnvironmentAdmission(ctx, session.Engine, session.Configuration); err != nil { +func (w *Worker) submitEnvironmentInputs(ctx context.Context, session sessions.Session, key string, inputs []sessions.Input) (receipts []sessions.InputReceipt, err error) { + checkAt := time.Now() + err = w.validateEnvironmentAdmission(ctx, session.Engine, session.Configuration) + observeExecutionStage(ctx, "input_validate", checkAt, err, "session_id", session.ID) + if err != nil { return nil, err } - if err := w.checkAdmissionOwnership(ctx); err != nil { + checkAt = time.Now() + err = w.checkAdmissionOwnership(ctx) + observeExecutionStage(ctx, "input_ownership", checkAt, err, "session_id", session.ID) + if err != nil { return nil, err } kind := "" @@ -96,12 +104,21 @@ func (w *Worker) submitEnvironmentInputs(ctx context.Context, session sessions.S } changed, unsubscribe := w.dispatcher.notifications.subscribe(session.TenantID, session.ID) defer unsubscribe() + reservedAt := time.Now() reserve, cancel := context.WithTimeout(ctx, 5*time.Second) reservation, err := w.admission.ReserveEnvironmentInput(reserve, session.TenantID, session.ID, key, inputs) cancel() + observeExecutionStage(ctx, "input_reserve", reservedAt, err, "session_id", session.ID) if err != nil { return nil, err } + obslog.Info(ctx, "environment input reserved", "session_id", session.ID, + "reservation_id", reservation.ID, "execution_trace_id", reservationTraceID(reservation.ID).String(), "state", reservation.State) + admittedAt := time.Now() + defer func() { + observeExecutionStage(ctx, "input_admission_wait", admittedAt, err, + "session_id", session.ID, "reservation_id", reservation.ID, "state", reservation.State) + }() w.wakeScheduler() if reservation.State == sessions.EnvironmentInputPending && !reservation.IsInitial { w.hintRuntimeWake(ctx, session) diff --git a/services/core/internal/execution/execution_observation.go b/services/core/internal/execution/execution_observation.go new file mode 100644 index 000000000..8fc0e967e --- /dev/null +++ b/services/core/internal/execution/execution_observation.go @@ -0,0 +1,46 @@ +package execution + +import ( + "context" + "crypto/sha256" + "errors" + "time" + + obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" +) + +// A reservation survives observer disconnects and owner restarts. Its opaque +// identity anchors the execution trace; the submitting HTTP trace is linked in +// the reservation log rather than retained in a process-local map. +func reservationTraceID(reservation string) obslog.TraceID { + digest := sha256.Sum256([]byte("oac.environment-input:" + reservation)) + var id obslog.TraceID + copy(id[:], digest[:len(id)]) + return id +} + +func reservationTrace(ctx context.Context, reservation string) context.Context { + return obslog.WithTrace(ctx, obslog.Carrier{Trace: reservationTraceID(reservation), Span: obslog.NewSpanID()}) +} + +// Stage intervals use a local monotonic clock. Error prose and request content +// are deliberately excluded: provider errors may contain private configuration. +func observeExecutionStage(ctx context.Context, stage string, started time.Time, err error, attrs ...any) { + status := "ok" + switch { + case errors.Is(err, context.Canceled): + status = "cancelled" + case errors.Is(err, context.DeadlineExceeded): + status = "timeout" + case err != nil: + status = "error" + } + fields := []any{"stage", stage, "duration_ms", float64(time.Since(started).Microseconds()) / 1000, "status", status} + obslog.Info(ctx, "execution stage", append(fields, attrs...)...) +} + +// Initial input is committed by Session creation. Joining its origin to the +// worker's is_initial reservation avoids an additional diagnostic store query. +func recordInitialInputOrigin(ctx context.Context, session string) { + obslog.Info(ctx, "session initial input origin", "session_id", session) +} diff --git a/services/core/internal/execution/execution_observation_test.go b/services/core/internal/execution/execution_observation_test.go new file mode 100644 index 000000000..bbefdd3fe --- /dev/null +++ b/services/core/internal/execution/execution_observation_test.go @@ -0,0 +1,178 @@ +package execution + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "io" + "log/slog" + "strings" + "sync" + "testing" + "time" + + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" +) + +func captureExecutionLogs(t *testing.T) *bytes.Buffer { + t.Helper() + var out bytes.Buffer + previous := slog.Default() + slog.SetDefault(slog.New(obslog.NewContextHandler(slog.NewJSONHandler(&out, nil)))) + t.Cleanup(func() { slog.SetDefault(previous) }) + return &out +} + +func TestReservationTraceSurvivesRecoveryWithoutRetainingObserver(t *testing.T) { + parent, cancel := context.WithCancel(t.Context()) + first := reservationTrace(parent, "reservation-one") + recovered := reservationTrace(t.Context(), "reservation-one") + a, _ := obslog.TraceFromContext(first) + b, _ := obslog.TraceFromContext(recovered) + if a.Trace != b.Trace || a.Span == b.Span || a.Trace == reservationTraceID("reservation-two") { + t.Fatal("reservation recovery did not retain identity with a distinct attempt span") + } + if _, err := obslog.ParseTraceparent(a.String()); err != nil { + t.Fatal(err) + } + cancel() + if !errors.Is(first.Err(), context.Canceled) || recovered.Err() != nil { + t.Fatal("diagnostic context changed cancellation ownership") + } +} + +func TestExecutionStageStatusAndCredentialRedaction(t *testing.T) { + out := captureExecutionLogs(t) + ctx := reservationTrace(t.Context(), "reservation-one") + for _, sample := range []struct { + err error + status string + }{ + {nil, "ok"}, {context.Canceled, "cancelled"}, {context.DeadlineExceeded, "timeout"}, + {errors.New("private provider URL and credential must never appear"), "error"}, + } { + out.Reset() + observeExecutionStage(ctx, "provider_create", time.Now().Add(-time.Millisecond), sample.err, "allocation_id", "allocation-one") + var entry map[string]any + if err := json.Unmarshal(out.Bytes(), &entry); err != nil { + t.Fatal(err) + } + if entry["status"] != sample.status || entry["stage"] != "provider_create" || entry["trace_id"] != reservationTraceID("reservation-one").String() || entry["duration_ms"].(float64) < 1 { + t.Fatal("stage identity/status/duration lost", entry) + } + if strings.Contains(out.String(), "credential") || strings.Contains(out.String(), "private provider") { + t.Fatal("raw provider error leaked") + } + } +} + +func TestExecutorControlLogsDistinguishPreparationFromHTTP(t *testing.T) { + out := captureExecutionLogs(t) + ctx := obslog.WithRequestID(reservationTrace(t.Context(), "reservation-one"), "http-one") + prepared := &preparedStart{ctx: ctx, requestID: "prepare-one", executorID: "executor-one", createdAt: time.Now(), startSentAt: time.Now()} + recordExecutorReadiness(prepared, proto.PreparationStatusPayload{ExecutorID: "executor-one", Reused: true}) + recordExecutorStart(prepared, "turn-one") + for _, line := range bytes.Split(bytes.TrimSpace(out.Bytes()), []byte("\n")) { + var entry map[string]any + if err := json.Unmarshal(line, &entry); err != nil { + t.Fatal(err) + } + if entry["request_id"] != "http-one" || entry["preparation_request_id"] != "prepare-one" || entry["trace_id"] != reservationTraceID("reservation-one").String() { + t.Fatal("control log confused HTTP and preparation identity", entry) + } + } +} + +func TestLifecycleGateObservationPreservesCancellation(t *testing.T) { + out := captureExecutionLogs(t) + r := &runtimeLifecycle{ctx: t.Context(), gate: make(chan struct{}, 1), nodeID: "node-one"} + r.gate <- struct{}{} + ctx, cancel := context.WithCancel(t.Context()) + cancel() + if err := r.lock(ctx); !errors.Is(err, context.Canceled) { + t.Fatal(err) + } + if len(r.gate) != 1 || !strings.Contains(out.String(), `"status":"cancelled"`) { + t.Fatal("cancelled waiter changed gate or log status") + } +} + +type executionTraceConn struct { + writes chan []byte + closed chan struct{} + once sync.Once +} + +func (c *executionTraceConn) ReadMessage() (int, []byte, error) { <-c.closed; return 0, nil, io.EOF } +func (c *executionTraceConn) WriteMessage(_ int, data []byte) error { + c.writes <- append([]byte(nil), data...) + return nil +} +func (c *executionTraceConn) WriteControl(int, []byte, time.Time) error { return nil } +func (c *executionTraceConn) SetReadLimit(int64) {} +func (c *executionTraceConn) SetReadDeadline(time.Time) error { return nil } +func (c *executionTraceConn) SetWriteDeadline(time.Time) error { return nil } +func (c *executionTraceConn) Close() error { c.once.Do(func() { close(c.closed) }); return nil } + +func TestExecutionSendPreservesPayloadAndCarriesTrace(t *testing.T) { + for _, traced := range []bool{false, true} { + conn := &executionTraceConn{writes: make(chan []byte, 2), closed: make(chan struct{})} + peer := runtimegateway.NewSession(conn, "device-one", "workspace-one", "test", nil, nil) + peer.Start() + ctx := t.Context() + if traced { + ctx = reservationTrace(ctx, "reservation-one") + } + payload := proto.ExecutionStartPayload{Handle: "handle-one", ExecutorID: "executor-one", RunID: "turn-one", Input: proto.TextInput("private input preserved on wire")} + if err := send(ctx, peer, proto.TypeExecutionStart, "prepare-one", payload); err != nil { + peer.Close("test done") + t.Fatal(err) + } + select { + case data := <-conn.writes: + var env proto.Envelope + if err := json.Unmarshal(data, &env); err != nil { + t.Fatal(err) + } + var got proto.ExecutionStartPayload + if err := env.DecodePayload(&got); err != nil { + t.Fatal(err) + } + expected, _ := json.Marshal(payload) + actual, _ := json.Marshal(got) + if env.ID != "prepare-one" || env.Type != proto.TypeExecutionStart || !bytes.Equal(expected, actual) { + t.Fatal("trace changed execution identity or payload") + } + if traced { + carrier, _ := obslog.TraceFromContext(ctx) + if env.Trace != carrier.String() { + t.Fatal("execution trace did not reach actual gateway wire") + } + } else if env.Trace != "" { + t.Fatal("missing trace changed into a fabricated caller trace") + } + case <-time.After(time.Second): + t.Fatal("gateway did not send") + } + peer.Close("test done") + } +} + +func TestInitialInputOriginKeepsHTTPTraceAndSessionIdentity(t *testing.T) { + out := captureExecutionLogs(t) + ctx, origin := obslog.StartBackgroundTrace(t.Context(), "http") + recordInitialInputOrigin(ctx, "session-one") + var entry map[string]any + if err := json.Unmarshal(out.Bytes(), &entry); err != nil { + t.Fatal(err) + } + if entry["session_id"] != "session-one" || entry["trace_id"] != origin.Trace.String() { + t.Fatal("initial input origin lost its HTTP/session bridge", entry) + } + if _, present := entry["reservation_id"]; present { + t.Fatal("creation fabricated a reservation identity") + } +} diff --git a/services/core/internal/execution/executor_observation.go b/services/core/internal/execution/executor_observation.go index 9d1fd8132..b823447c0 100644 --- a/services/core/internal/execution/executor_observation.go +++ b/services/core/internal/execution/executor_observation.go @@ -11,20 +11,20 @@ import ( // These durations use one Core process monotonic clock. They describe control // readiness and input-to-first-text latency, not model-only or native tool time. func recordExecutorReadiness(p *preparedStart, status proto.PreparationStatusPayload) { - obslog.Bg().Info("executor ready", "request_id", p.requestID, + obslog.Info(p.ctx, "executor ready", "preparation_request_id", p.requestID, "executor_id", status.ExecutorID, "reused", status.Reused, "control_ready_ms", time.Since(p.createdAt).Milliseconds()) } // The started receipt confirms adapter ownership, not model consumption. func recordExecutorStart(p *preparedStart, turn string) { - obslog.Bg().Info("executor start acknowledged", "turn_id", turn, - "request_id", p.requestID, "executor_id", p.executorID, + obslog.Info(p.ctx, "executor start acknowledged", "turn_id", turn, + "preparation_request_id", p.requestID, "executor_id", p.executorID, "start_control_ms", time.Since(p.startSentAt).Milliseconds()) } func recordFirstText(ctx context.Context, session, turn string, submitted time.Time) { - obslog.Bg().InfoContext(ctx, "executor first text", "session_id", session, + obslog.Info(ctx, "executor first text", "session_id", session, "turn_id", turn, "input_to_first_text_ms", time.Since(submitted).Milliseconds()) } diff --git a/services/core/internal/execution/executor_preparation.go b/services/core/internal/execution/executor_preparation.go index fad7e9c36..5b17f084e 100644 --- a/services/core/internal/execution/executor_preparation.go +++ b/services/core/internal/execution/executor_preparation.go @@ -14,7 +14,7 @@ import ( // Core does not cache native ownership. A replacement requires the Runtime to // confirm cleanup and recover the exact Session history before returning ready. func (d *Dispatcher) prepareTurnExecutor(ctx context.Context, peer *runtimegateway.Session, tenant, session, turn string, request proto.PromptRequestPayload, expectedStatus string) (*preparedStart, error) { - prepared, err := newPreparedStart(peer) + prepared, err := newPreparedStart(ctx, peer) if err != nil { return nil, err } diff --git a/services/core/internal/execution/preparation.go b/services/core/internal/execution/preparation.go index 59ea9bfb4..6e55bfb23 100644 --- a/services/core/internal/execution/preparation.go +++ b/services/core/internal/execution/preparation.go @@ -21,6 +21,7 @@ type preparationRejection struct { func (e *preparationRejection) Error() string { return "preparation control rejected: " + e.code } type preparedStart struct { + ctx context.Context createdAt time.Time startSentAt time.Time startObserved bool @@ -32,13 +33,13 @@ type preparedStart struct { sub *runtimegateway.Subscription } -func newPreparedStart(peer *runtimegateway.Session) (*preparedStart, error) { +func newPreparedStart(ctx context.Context, peer *runtimegateway.Session) (*preparedStart, error) { id := uuid.NewString() sub, err := peer.SubscribePreparation(id) if err != nil { return nil, err } - return &preparedStart{peer: peer, requestID: id, sub: sub, createdAt: time.Now()}, nil + return &preparedStart{ctx: ctx, peer: peer, requestID: id, sub: sub, createdAt: time.Now()}, nil } func (p *preparedStart) close() { diff --git a/services/core/internal/execution/prepared_dispatch.go b/services/core/internal/execution/prepared_dispatch.go index 204a3d11b..70202be38 100644 --- a/services/core/internal/execution/prepared_dispatch.go +++ b/services/core/internal/execution/prepared_dispatch.go @@ -5,6 +5,9 @@ import ( "encoding/json" "errors" "strings" + "time" + + obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" @@ -19,6 +22,9 @@ type EnvironmentRun struct { // RunEnvironmentInput reserves a Turn on the Session-owned Runtime Executor. It // checks lease, the lease d.Store was built on, before any Runtime preparation. func (d *Dispatcher) RunEnvironmentInput(ctx context.Context, lease Ownership, tenantID, sessionID, reservationID string) (run EnvironmentRun, err error) { + ctx = reservationTrace(ctx, reservationID) + selectedAt := time.Now() + obslog.Info(ctx, "environment input selected", "session_id", sessionID, "reservation_id", reservationID) if err = lease.CheckOwnership(ctx); err != nil { return run, err } @@ -26,6 +32,10 @@ func (d *Dispatcher) RunEnvironmentInput(ctx context.Context, lease Ownership, t if err != nil || run.Reservation.State != sessions.EnvironmentInputPending { return run, err } + if !run.Reservation.CreatedAt.IsZero() { + obslog.Info(ctx, "environment input queue age", "session_id", sessionID, "reservation_id", reservationID, + "queue_age_ms", time.Since(run.Reservation.CreatedAt).Milliseconds(), "is_initial", run.Reservation.IsInitial) + } session, err := d.Store.GetSession(ctx, tenantID, sessionID) if err != nil { return run, err @@ -77,11 +87,15 @@ func (d *Dispatcher) RunEnvironmentInput(ctx context.Context, lease Ownership, t if err := d.configurePreparedEnvironment(session, environment, bound.Device, &req); err != nil { return run, err } - prepared, err := newPreparedStart(peer) + observeExecutionStage(owner, "execution_configuration", selectedAt, nil, + "session_id", sessionID, "reservation_id", reservationID, "environment_id", environment.ID, "device_id", bound.Device.ID) + prepared, err := newPreparedStart(owner, peer) if err != nil { return run, err } defer prepared.close() + obslog.Info(owner, "execution preparation requested", "session_id", sessionID, "reservation_id", reservationID, + "environment_id", environment.ID, "device_id", bound.Device.ID, "preparation_request_id", prepared.requestID) if err = send(owner, peer, proto.TypeExecutionPrepare, prepared.requestID, proto.ExecutionPreparePayload{SessionID: sessionID, Configuration: req}); err != nil { return run, err } @@ -92,7 +106,9 @@ func (d *Dispatcher) RunEnvironmentInput(ctx context.Context, lease Ownership, t if err := d.messageInputSupport(peer, session.Engine, snapshot, messages); err != nil { return run, err } + promoteAt := time.Now() promoted, err := d.Store.PromoteEnvironmentInput(owner, tenantID, sessionID, reservationID) + observeExecutionStage(owner, "input_promote", promoteAt, err, "session_id", sessionID, "reservation_id", reservationID) if errors.Is(err, sessions.ErrTurnConflict) { // A rejected claim leaves the reservation pending for a later attempt. return run, err @@ -108,6 +124,8 @@ func (d *Dispatcher) RunEnvironmentInput(ctx context.Context, lease Ownership, t return run, nil } req.RunID = run.Reservation.Receipts[0].TurnID + obslog.Info(owner, "environment input admitted", "session_id", sessionID, "reservation_id", reservationID, + "turn_id", req.RunID, "preparation_request_id", prepared.requestID, "executor_id", prepared.executorID) req.Input = messages through := run.Reservation.Receipts[len(run.Reservation.Receipts)-1].Sequence releaseDelivery, err := peer.TrackExecutionDelivery(req.RunID) diff --git a/services/core/internal/execution/runtime_compute.go b/services/core/internal/execution/runtime_compute.go index 2f7cc36bc..9a2d42b22 100644 --- a/services/core/internal/execution/runtime_compute.go +++ b/services/core/internal/execution/runtime_compute.go @@ -256,7 +256,10 @@ func (r *runtimeLifecycle) restoreCompute(ctx context.Context, p sandbox.Suspens if state.Target == nil || state.Retained == nil || state.Rollback { return sandbox.ErrOwnership } + resumeAt := time.Now() result, err := p.Resume(ctx, sandbox.ResumeRequest{Reference: runtimeReference(owner), OperationID: state.RestoreID, Retained: *state.Retained, Target: *state.Target, ReconcileOnly: observeOnly}) + observeExecutionStage(ctx, "provider_resume", resumeAt, err, "allocation_id", owner.ID, + "session_id", owner.SessionID, "environment_id", owner.EnvironmentID, "node_id", r.nodeID, "reconcile_only", observeOnly) if err != nil { return err } diff --git a/services/core/internal/execution/runtime_connections.go b/services/core/internal/execution/runtime_connections.go index 10c4bfd1a..5c3a4ff44 100644 --- a/services/core/internal/execution/runtime_connections.go +++ b/services/core/internal/execution/runtime_connections.go @@ -4,6 +4,8 @@ import ( "context" "errors" + obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" + "github.com/google/uuid" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment" @@ -62,6 +64,8 @@ func observeRuntimeConnection(ctx context.Context, operations *sessions.Executio return err } current.connected = connected + obslog.Info(ctx, "runtime connection observed", "environment_id", environment, "connected", connected, + "connection_generation", current.generation, "connection_revision", current.revision) return nil } diff --git a/services/core/internal/execution/runtime_initialization.go b/services/core/internal/execution/runtime_initialization.go index 4f94d4d52..b9a7ab23e 100644 --- a/services/core/internal/execution/runtime_initialization.go +++ b/services/core/internal/execution/runtime_initialization.go @@ -94,10 +94,12 @@ func (w *Worker) initializeEnvironment(ctx context.Context, owner sessions.Envir operation, cancel := context.WithTimeout(ctx, 30*time.Minute) defer cancel() failure := sessions.ProvisioningFailure{} + initializeAt := time.Now() err := w.prepareEnvironment(operation, owner, &failure) if err == nil { err = w.dispatcher.sessionExecution.CompleteEnvironmentInitialization(operation, owner) } + observeExecutionStage(ctx, "environment_initialize", initializeAt, err, "environment_id", owner.EnvironmentID, "session_id", owner.SessionID, "device_id", owner.DeviceID) if err != nil { log.Warn(ctx, "Environment preparation failed", "environment_id", owner.EnvironmentID, "session_id", owner.SessionID) // A later scan settles an unrecorded failure; it never retries the setup. diff --git a/services/core/internal/execution/runtime_lifecycle.go b/services/core/internal/execution/runtime_lifecycle.go index 28802b315..c36bca71e 100644 --- a/services/core/internal/execution/runtime_lifecycle.go +++ b/services/core/internal/execution/runtime_lifecycle.go @@ -119,7 +119,9 @@ func validatedRuntimeProvider(config *RuntimeProvider, registry *runtimegateway. return copied, nil } -func (r *runtimeLifecycle) lock(ctx context.Context) error { +func (r *runtimeLifecycle) lock(ctx context.Context) (err error) { + started := time.Now() + defer func() { observeExecutionStage(ctx, "runtime_gate_wait", started, err, "node_id", r.nodeID) }() select { case r.gate <- struct{}{}: if err := r.ctx.Err(); err != nil { @@ -226,10 +228,13 @@ func (r *runtimeLifecycle) provision(ctx context.Context, tenant, environment, p if err := r.lease.CheckOwnership(ctx); err != nil { return owner, err } + createAt := time.Now() info, err := provider.Create(ctx, sandbox.Bootstrap{ Reference: runtimeReference(owner), SessionID: owner.SessionID, DeviceID: owner.DeviceID, CoreURL: r.config.CoreURL, Credential: token, NetworkAccess: placement.NetworkAccess, AllowedDomains: placement.AllowedDomains, }) + observeExecutionStage(ctx, "provider_create", createAt, err, "allocation_id", owner.ID, + "session_id", owner.SessionID, "environment_id", environment, "device_id", owner.DeviceID, "node_id", r.nodeID) if info.Reference == runtimeReference(owner) && info.CreateSettled && info.State == "absent" { record, cancel := context.WithTimeout(context.WithoutCancel(ctx), 5*time.Second) released, releaseErr := r.deployment.ReleaseAbsentCreation(record, owner) @@ -301,7 +306,9 @@ func (r *runtimeLifecycle) reconcile(ctx context.Context) error { for _, owner := range rows { r.cursor = owner.ID operation, stop := context.WithTimeout(ctx, 30*time.Second) + observeAt := time.Now() err := r.observe(operation, owner) + observeExecutionStage(ctx, "runtime_observe", observeAt, err, "allocation_id", owner.ID, "session_id", owner.SessionID, "node_id", r.nodeID) r.recordObservation(ctx, owner, err) stop() if err != nil { diff --git a/services/core/internal/execution/runtime_pending.go b/services/core/internal/execution/runtime_pending.go index 254b33877..870d0d53a 100644 --- a/services/core/internal/execution/runtime_pending.go +++ b/services/core/internal/execution/runtime_pending.go @@ -25,7 +25,9 @@ func (r *runtimeLifecycle) provisionPending(ctx context.Context) error { r.pendingCursor = environment.ID provider := r.config.InstallationID operation, cancel := context.WithTimeout(ctx, 30*time.Second) + provisionAt := time.Now() _, err := r.provision(operation, environment.TenantID, environment.ID, provider) + observeExecutionStage(ctx, "runtime_provision", provisionAt, err, "environment_id", environment.ID, "node_id", r.nodeID) cancel() if err != nil { if ownership := r.lease.CheckOwnership(ctx); ownership != nil { diff --git a/services/core/internal/execution/worker.go b/services/core/internal/execution/worker.go index ba7660c5d..4f2113be4 100644 --- a/services/core/internal/execution/worker.go +++ b/services/core/internal/execution/worker.go @@ -149,6 +149,7 @@ func (w *Worker) CreateSession(ctx context.Context, tenant string, input session } session, err := w.admission.CreateSession(ctx, tenant, input) if err == nil && len(input.InitialInputs) > 0 { + recordInitialInputOrigin(ctx, session.ID) w.wakeScheduler() } return session, err @@ -161,6 +162,7 @@ func (w *Worker) CreateSessionStream(ctx context.Context, tenant string, input s } creation, err := w.admission.CreateSessionStream(ctx, tenant, input) if err == nil && len(input.InitialInputs) > 0 { + recordInitialInputOrigin(ctx, creation.Session.ID) w.wakeScheduler() } return creation, err