Skip to content
Merged
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
7 changes: 7 additions & 0 deletions apps/daemon/internal/cli/agent_discovery.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,10 +3,13 @@ package cli
import (
"context"
"fmt"
"time"

"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent"
"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent/claudesdk"
"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent/codex"
"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent/mcode"
obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log"
)

var harnessDeclarations = []agent.Declaration{codex.Declaration, mcode.Declaration, claudesdk.Declaration}
Expand All @@ -27,7 +30,11 @@ func discoverAgentCLIs(parent context.Context, rc *runContext, profile string, d
if rc.installedKinds != nil && !rc.installedKinds[declaration.Info.Kind] {
continue
}
started := time.Now()
runtime := declaration.Discover(parent, agent.DiscoveryOptions{Profile: profile, Stdout: rc.stdout, Stderr: rc.stderr}, declaration.Info)
obslog.Info(parent, "runtime startup stage", "stage", "harness_discovery",
"harness_kind", declaration.Info.Kind, "duration_ms", float64(time.Since(started))/float64(time.Millisecond),
"available", runtime != nil && runtime.Info.Available)
if runtime == nil {
continue
}
Expand Down
5 changes: 5 additions & 0 deletions apps/daemon/internal/cli/connect.go
Original file line number Diff line number Diff line change
Expand Up @@ -138,6 +138,7 @@ func runConnect(ctx *runContext, args []string) error {

// Self-check before pairing/loading credentials so a machine with
// no supported agent CLI fails before consuming a one-shot token.
initializeRuntimeObservations(ctx)
agentCLIs, err := preflightAgentCLIs(context.Background(), ctx, *profile)
if err != nil {
return err
Expand Down Expand Up @@ -308,6 +309,7 @@ func mainLoopRemote(parent context.Context, rc *runContext, profile string, prof
rootCtx, cancel := daemonize.NotifyContext(parent)
defer cancel()

bootstrapStarted := time.Now()
bootCtx, bootCancel := context.WithTimeout(rootCtx, bootstrapTimeout)
var boot *transport.BootstrapResponse
var err error
Expand All @@ -317,6 +319,7 @@ func mainLoopRemote(parent context.Context, rc *runContext, profile string, prof
boot, err = environmentBootstrap(bootCtx, prof, remote)
}
bootCancel()
observeRuntimeStartup(rootCtx, "bootstrap", bootstrapStarted, err)
if err != nil {
return fmt.Errorf("connect: bootstrap: %w", err)
}
Expand All @@ -337,6 +340,7 @@ func mainLoopRemote(parent context.Context, rc *runContext, profile string, prof
defer control.Close()
}
dial := func(ctx context.Context) (*transport.Conn, error) {
dialStarted := time.Now()
conn, err := transport.Dial(ctx, transport.DialOptions{
WSURL: wsURL,
DeviceID: boot.DeviceID,
Expand All @@ -347,6 +351,7 @@ func mainLoopRemote(parent context.Context, rc *runContext, profile string, prof
// in heartbeat's DaemonVersion field.
DaemonVersion: proto.Version,
})
observeRuntimeStartup(ctx, "transport_dial", dialStarted, err)
if remote != "" && err != nil {
if errors.Is(err, transport.ErrIncompatibleVersion) {
return nil, fmt.Errorf("Environment connection rejected: %w: %w", transport.ErrPermanent, transport.ErrIncompatibleVersion)
Expand Down
1 change: 1 addition & 0 deletions apps/daemon/internal/cli/connect_environment.go
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,7 @@ func environmentBootstrap(ctx context.Context, prof auth.Profile, remote string)
}

func runEnvironmentConnect(parent context.Context, rc *runContext, profile string, background bool, remote, environment, credentialFile string) error {
initializeRuntimeObservations(rc)
base, err := environmentBase(remote)
if err != nil {
return err
Expand Down
30 changes: 30 additions & 0 deletions apps/daemon/internal/cli/startup_observation.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
package cli

import (
"context"
"errors"
"log/slog"
"time"

obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log"
)

func observeRuntimeStartup(ctx context.Context, stage string, started time.Time, err error) {
status := "ok"
if err != nil {
status = "error"
if errors.Is(err, context.Canceled) {
status = "cancelled"
}
if errors.Is(err, context.DeadlineExceeded) {
status = "timeout"
}
}
obslog.Info(ctx, "runtime startup stage", "stage", stage,
"duration_ms", float64(time.Since(started))/float64(time.Millisecond), "status", status)
}

func initializeRuntimeObservations(rc *runContext) {
obslog.Init(obslog.Config{Format: "text", Level: slog.LevelInfo, Out: rc.stderr})
obslog.Bg().Info("runtime process starting", "daemon_version", Version)
}
15 changes: 15 additions & 0 deletions deploy/kubernetes/BETA_CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,3 +28,18 @@ The integration adds a typed suspension lifecycle, E2B pause/resume with durable
### Verification

Upstream [CI run 37126282818](https://github.com/MiniMax-AI/OpenAgentCore/actions/runs/37126282818) passed its selected checks at the integrated source commit. Local validation passed the Go sandbox, process configuration and Runtime bootstrap tests, focused deployment/execution/server tests, 188 E2B helper tests, 11 E2B template tests, the generated helper contract check and SQL generation freshness check. Database lifecycle integration and fresh cloud execution remain separate qualification gates. This ledger records source integration; it does not establish production activation or fresh E2B Session/Turn acceptance.

## 2026-10-06 — Runtime readiness observations

| Item | Value |
| --- | --- |
| Fork beta baseline | `8f6ab3cd3272c156249360d2640383e3aee36a97` |
| Feature source commit | `b19ad377d357621c1fc48f24c575327cf24297c5` |
| Feature branch | `codex/runtime-readiness-observations` |
| Schema migrations / DDL / SQL changes | None |
| Runtime wire, helper protocol, native dependency pins | Unchanged |
| Required release assets | Matching Core/Web release and a newly built combined E2B Runtime template |

The release adds authenticated transport registration and confirmed capability observations, daemon startup stage timings, and coalesced capability-driven scheduler hints. [Execution latency](../../docs/getting-started/operations.md#execution-latency) owns the log boundaries and limitations. Existing allocations retain their immutable template generation; new allocations use the new template only after its selection is updated through the administrative deployment API.

Validation passed focused gateway/execution/daemon regressions and race checks, isolated PostgreSQL scheduling and Worker admission/device-isolation fixtures, Runtime contract checks, static checks, the name guard, documentation checks and all translation checks. The source review found no blocking issues. Cloud build readiness and production health remain separate from real-model and suspension acceptance.
4 changes: 3 additions & 1 deletion docs/getting-started/operations.md
Original file line number Diff line number Diff line change
Expand Up @@ -198,6 +198,8 @@ Sandboxes are the isolation boundary ([Runtime and outer isolation](../concepts.

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.
`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. `runtime transport registered` records authenticated WebSocket registration. `runtime capability snapshot observed` records validated declaration changes with device ID and kind counts; transport registration alone does not establish execution readiness. Capability changes with an available Harness send a coalesced scheduling hint, and the Worker rechecks its ordinary ownership, capacity, lifecycle and capability requirements; polling recovers missed hints. 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.

Daemon startup logs record `runtime process starting` with its build version and `runtime startup stage` for each Harness discovery, bootstrap and transport dial attempt. Durations use a local monotonic clock and omit credentials and raw errors. Discovery runs before transport registration. Compare daemon timestamps with Core registration and capability observations to separate startup from periodic observation and scheduling. These daemon records require a Runtime built with this instrumentation; updating Core alone does not update an existing Runtime template. A managed startup receipt confirms process spawn, not connection or capability readiness.
6 changes: 5 additions & 1 deletion docs/zh/getting-started/operations.md
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
---
title: "管理你的安装"
source: docs/getting-started/operations.md
source_hash: 263a48e467c9aeec1f8357f0e11b40c6f1b329dd6162f65eda3a1038acd4f551
source_hash: ee5e9211cc09dbf5a4bfba412c2c97af0aefcb7f9d2d5ffd6c0666729fbe2809
---

安装运维人员负责 Core 主机、存储和可用性。节点主机运行各自的服务;参阅[节点](nodes.md)。设置见[配置参考](../configuration.md)。
Expand Down Expand Up @@ -204,3 +204,7 @@ Web 使用 Core 密钥让管理员登录,检查每个请求来源,并用保
`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 原始错误。

`runtime transport registered` 记录认证后的 WebSocket 注册,`runtime capability snapshot observed` 记录合法能力声明变化及设备 ID、能力数量。连接注册本身不代表执行就绪。包含可用 Harness 的能力变化会发出合并后的调度提示;Worker 仍检查所有权、容量、生命周期和能力要求,轮询负责兜底。

Daemon 启动日志记录 `runtime process starting` 和构建版本,`runtime startup stage` 记录各 Harness 探测、bootstrap 和每次连接尝试的本地单调时钟耗时,不输出凭据或原始错误。探测发生在连接注册之前。结合这些边界与已观察到的持久化连接时间,区分启动、周期观察和调度等待。Daemon 观测要求 Runtime 使用带有这些埋点的构建;仅升级 Core 不会升级现有 Runtime 模板。托管启动回执只确认进程已启动,不确认连接或能力就绪。
7 changes: 3 additions & 4 deletions services/core/internal/execution/worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -298,6 +298,8 @@ func (w *Worker) Run(ctx context.Context) (runErr error) {
rescanOnCompletion = false
case <-w.scheduleWake:
rescanOnCompletion = true
case <-w.dispatcher.Registry.CapabilityHints():
rescanOnCompletion = true
case <-ticker.C:
maintenance = true
}
Expand Down Expand Up @@ -327,10 +329,7 @@ func (w *Worker) Run(ctx context.Context) (runErr error) {
w.observeSchedulerPoll(0, nil)
continue
}
if !maintenance {
schedule.nextEnvironmentScan = time.Time{}
}
work, err := schedule.selectWork(ctx, w, devices, active)
work, err := schedule.selectWork(ctx, w, devices, active, !maintenance)
w.observeSlots(len(active))
if err != nil {
w.observeSchedulerPoll(0, err)
Expand Down
66 changes: 66 additions & 0 deletions services/core/internal/execution/worker_capability_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
package execution

import (
"context"
"encoding/json"
"strings"
"sync/atomic"
"testing"
"time"

"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)

type candidateQueryObserver struct{ scans atomic.Int32 }

func (o *candidateQueryObserver) TraceQueryStart(ctx context.Context, _ *pgx.Conn, data pgx.TraceQueryStartData) context.Context {
if strings.Contains(data.SQL, "-- name: ListEnvironmentInputWork") {
o.scans.Add(1)
}
return ctx
}
func (*candidateQueryObserver) TraceQueryEnd(context.Context, *pgx.Conn, pgx.TraceQueryEndData) {}

func TestSchedulingHintScansBeforeDeadlineAndRestoresPollingDelay(t *testing.T) {
observer := &candidateQueryObserver{}
pool := pgtest.OpenIsolated(t, func(cfg *pgxpool.Config) { cfg.ConnConfig.Tracer = observer })
worker := &Worker{dispatcher: &Dispatcher{Store: store.New(pool)}, concurrency: 1}
schedule := workerSchedule{nextEnvironmentScan: time.Now().Add(time.Hour)}
for _, hinted := range []bool{false, true, false} {
if _, err := schedule.selectWork(t.Context(), worker, nil, map[string]bool{}, hinted); err != nil {
t.Fatal(err)
}
}
if observer.scans.Load() != 1 {
t.Fatalf("candidate scans = %d; hint must bypass the deadline once", observer.scans.Load())
}
}

type phaseSessionReader struct{ sessions.Reader }

func (phaseSessionReader) GetSessionEnvironment(context.Context, string, string) (sessions.Environment, error) {
return sessions.Environment{ID: "environment", Configuration: json.RawMessage(`{"type":"openai_hosted"}`)}, nil
}
func (phaseSessionReader) GetSessionDevice(context.Context, string, string) (sessions.ExecutionDevice, error) {
panic("non-running compute must not reach device eligibility")
}
func TestCapabilityHintDoesNotBypassComputePhase(t *testing.T) {
for _, phase := range []string{"quiescing", "suspending", "suspended", "restoring", "waking"} {
t.Run(phase, func(t *testing.T) {
reader := &strictDeploymentReader{t: t, environmentAllocation: func(context.Context, deployment.AllocationKey) (deployment.Allocation, error) {
return deployment.Allocation{ComputePhase: phase}, nil
}}
worker := &Worker{dispatcher: &Dispatcher{SessionsReader: phaseSessionReader{}, DeploymentReader: reader}}
session := sessions.Session{Configuration: json.RawMessage(`{"environment":{"type":"openai_hosted"}}`)}
ready, err := worker.bindSessionDevice(t.Context(), session, func(string) bool { t.Fatal("phase bypassed"); return true })
if err != nil || ready {
t.Fatal("non-running compute selected", ready, err)
}
})
}
}
5 changes: 4 additions & 1 deletion services/core/internal/execution/worker_schedule.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,10 @@ type scheduledWork struct {
reservationID string
}

func (s *workerSchedule) selectWork(ctx context.Context, w *Worker, devices []string, active map[string]bool) ([]scheduledWork, error) {
func (s *workerSchedule) selectWork(ctx context.Context, w *Worker, devices []string, active map[string]bool, hinted bool) ([]scheduledWork, error) {
if hinted {
s.nextEnvironmentScan = time.Time{}
}
turns, err := w.dispatcher.Store.ListExecutionWork(ctx, s.turnCursor, []string{sessions.TurnQueued}, devices)
if err != nil {
return nil, err
Expand Down
16 changes: 16 additions & 0 deletions services/core/internal/execution/worker_wakeup_test.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
package execution

import (
"context"
"errors"
"sync"
"testing"
)
Expand Down Expand Up @@ -33,3 +35,17 @@ func TestSchedulerWakeCoalescesConcurrentAdmissionsAndKeepsNextHint(t *testing.T
default:
}
}

func TestWorkerCancelledWithoutGatewayPreservesShutdown(t *testing.T) {
worker := &Worker{dispatcher: &Dispatcher{}, lease: heldLease{}, stopped: make(chan struct{}), concurrency: 1}
ctx, cancel := context.WithCancel(t.Context())
cancel()
if err := worker.Run(ctx); !errors.Is(err, context.Canceled) {
t.Fatal("gateway-free admission worker did not stop", err)
}
select {
case <-worker.stopped:
default:
t.Fatal("worker did not publish shutdown")
}
}
Loading
Loading