From b19ad377d357621c1fc48f24c575327cf24297c5 Mon Sep 17 00:00:00 2001 From: sam Date: Tue, 6 Oct 2026 21:07:24 +0800 Subject: [PATCH 1/3] Observe runtime readiness and accelerate capability-driven scheduling --- apps/daemon/internal/cli/agent_discovery.go | 7 + apps/daemon/internal/cli/connect.go | 5 + .../internal/cli/connect_environment.go | 1 + .../internal/cli/startup_observation.go | 30 +++ docs/getting-started/operations.md | 4 +- docs/zh/getting-started/operations.md | 6 +- services/core/internal/execution/worker.go | 7 +- .../execution/worker_capability_test.go | 66 ++++++ .../internal/execution/worker_schedule.go | 5 +- .../runtimegateway/capability_hint_test.go | 192 ++++++++++++++++++ .../core/internal/runtimegateway/handler.go | 3 + .../core/internal/runtimegateway/registry.go | 45 +++- .../core/internal/runtimegateway/session.go | 25 ++- 13 files changed, 380 insertions(+), 16 deletions(-) create mode 100644 apps/daemon/internal/cli/startup_observation.go create mode 100644 services/core/internal/execution/worker_capability_test.go create mode 100644 services/core/internal/runtimegateway/capability_hint_test.go diff --git a/apps/daemon/internal/cli/agent_discovery.go b/apps/daemon/internal/cli/agent_discovery.go index d29b85acc..d0c987806 100644 --- a/apps/daemon/internal/cli/agent_discovery.go +++ b/apps/daemon/internal/cli/agent_discovery.go @@ -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} @@ -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 } diff --git a/apps/daemon/internal/cli/connect.go b/apps/daemon/internal/cli/connect.go index 7fb6e82f8..b3eba97fa 100644 --- a/apps/daemon/internal/cli/connect.go +++ b/apps/daemon/internal/cli/connect.go @@ -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 @@ -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 @@ -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) } @@ -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, @@ -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) diff --git a/apps/daemon/internal/cli/connect_environment.go b/apps/daemon/internal/cli/connect_environment.go index ec8ab6592..55147716f 100644 --- a/apps/daemon/internal/cli/connect_environment.go +++ b/apps/daemon/internal/cli/connect_environment.go @@ -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 diff --git a/apps/daemon/internal/cli/startup_observation.go b/apps/daemon/internal/cli/startup_observation.go new file mode 100644 index 000000000..6844f1282 --- /dev/null +++ b/apps/daemon/internal/cli/startup_observation.go @@ -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) +} diff --git a/docs/getting-started/operations.md b/docs/getting-started/operations.md index 8b210c993..cc3331593 100644 --- a/docs/getting-started/operations.md +++ b/docs/getting-started/operations.md @@ -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. diff --git a/docs/zh/getting-started/operations.md b/docs/zh/getting-started/operations.md index 4d4e07649..ecf961525 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: 263a48e467c9aeec1f8357f0e11b40c6f1b329dd6162f65eda3a1038acd4f551 +source_hash: ee5e9211cc09dbf5a4bfba412c2c97af0aefcb7f9d2d5ffd6c0666729fbe2809 --- 安装运维人员负责 Core 主机、存储和可用性。节点主机运行各自的服务;参阅[节点](nodes.md)。设置见[配置参考](../configuration.md)。 @@ -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 模板。托管启动回执只确认进程已启动,不确认连接或能力就绪。 diff --git a/services/core/internal/execution/worker.go b/services/core/internal/execution/worker.go index 4f2113be4..4e9a901dd 100644 --- a/services/core/internal/execution/worker.go +++ b/services/core/internal/execution/worker.go @@ -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 } @@ -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) diff --git a/services/core/internal/execution/worker_capability_test.go b/services/core/internal/execution/worker_capability_test.go new file mode 100644 index 000000000..6b6c5d2ff --- /dev/null +++ b/services/core/internal/execution/worker_capability_test.go @@ -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) + } + }) + } +} diff --git a/services/core/internal/execution/worker_schedule.go b/services/core/internal/execution/worker_schedule.go index 2a901f5f2..17ef451a0 100644 --- a/services/core/internal/execution/worker_schedule.go +++ b/services/core/internal/execution/worker_schedule.go @@ -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 diff --git a/services/core/internal/runtimegateway/capability_hint_test.go b/services/core/internal/runtimegateway/capability_hint_test.go new file mode 100644 index 000000000..55d93c9ec --- /dev/null +++ b/services/core/internal/runtimegateway/capability_hint_test.go @@ -0,0 +1,192 @@ +package runtimegateway + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto/prototest" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimedevice" + "log/slog" + "strings" + "testing" + "time" +) + +func expectCapabilityHint(t *testing.T, reg *Registry, want bool) { + t.Helper() + select { + case <-reg.CapabilityHints(): + if !want { + t.Fatal("unexpected capability hint") + } + default: + if want { + t.Fatal("missing capability hint") + } + } +} + +func TestCapabilityHintsRequireCurrentAvailableSnapshot(t *testing.T) { + reg := NewRegistry() + old := NewSession(newFakeConn(), "device", "tenant", "version", reg, nil) + reg.Register(old) + expectCapabilityHint(t, reg, false) // Transport alone is insufficient. + publishCapabilityHeartbeat(t, old, []runtimedevice.SupportedAgentKind{{Kind: "fake", Available: false}}) + expectCapabilityHint(t, reg, false) + kinds := []runtimedevice.SupportedAgentKind{{Kind: "fake", Available: true}} + publishCapabilityHeartbeat(t, old, kinds) + expectCapabilityHint(t, reg, true) + publishCapabilityHeartbeat(t, old, kinds) + expectCapabilityHint(t, reg, false) // Ordinary heartbeats do not rescan. + newer := NewSession(newFakeConn(), "device", "tenant", "version", reg, nil) + reg.Register(newer) + kinds[0].Version = "changed" + publishCapabilityHeartbeat(t, old, kinds) + expectCapabilityHint(t, reg, false) // Superseded socket cannot accelerate work. + publishCapabilityHeartbeat(t, newer, kinds) + expectCapabilityHint(t, reg, true) +} + +func TestCapabilityHintsCoalesceAndIgnoreDisconnectedPeers(t *testing.T) { + reg := NewRegistry() + s := NewSession(newFakeConn(), "device", "tenant", "version", reg, nil) + reg.Register(s) + kinds := []runtimedevice.SupportedAgentKind{{Kind: "fake", Available: true}} + publishCapabilityHeartbeat(t, s, kinds) + kinds[0].Capabilities.Streaming = true + publishCapabilityHeartbeat(t, s, kinds) + expectCapabilityHint(t, reg, true) + expectCapabilityHint(t, reg, false) + reg.Deregister(s) + kinds[0].Capabilities.Steering = true + publishCapabilityHeartbeat(t, s, kinds) + expectCapabilityHint(t, reg, false) +} + +func TestInvalidHeartbeatCannotWakeScheduler(t *testing.T) { + reg := NewRegistry() + s := NewSession(newFakeConn(), "device", "tenant", "version", reg, nil) + reg.Register(s) + s.handleHeartbeat(proto.Envelope{Type: proto.TypeHeartbeat, Payload: json.RawMessage(`{"ts":1,"active_requests":0,"supported_agent_kinds":[{"kind":"fake","available":true,"capabilities":{"streaming":"invented"}}]}`)}) + expectCapabilityHint(t, reg, false) + if !s.IsClosed() { + t.Fatal("malformed capabilities did not close transport") + } +} + +func publishCapabilityHeartbeat(t *testing.T, s *Session, kinds []runtimedevice.SupportedAgentKind) { + t.Helper() + advertised := make([]proto.SupportedAgentKind, len(kinds)) + for i, k := range kinds { + advertised[i] = proto.SupportedAgentKind{Kind: k.Kind, Available: k.Available, Version: k.Version, + Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{ + Streaming: proto.CapabilityFromBool(k.Capabilities.Streaming), + Steering: proto.CapabilityFromBool(k.Capabilities.Steering), + })} + } + env, err := proto.NewEnvelope(proto.TypeHeartbeat, "", proto.HeartbeatPayload{SupportedAgentKinds: advertised}) + if err != nil { + t.Fatal(err) + } + s.handleHeartbeat(env) +} + +func TestCapabilityObservationsExcludeRejectedAndStalePeers(t *testing.T) { + var logs bytes.Buffer + previous := slog.Default() + slog.SetDefault(slog.New(slog.NewJSONHandler(&logs, nil))) + t.Cleanup(func() { slog.SetDefault(previous) }) + reg := NewRegistry() + old := NewSession(newFakeConn(), "device", "tenant", "version", reg, nil) + reg.Register(old) + kinds := []runtimedevice.SupportedAgentKind{{Kind: "fake", Available: true}} + publishCapabilityHeartbeat(t, old, kinds) + if !strings.Contains(logs.String(), "runtime capability snapshot observed") { + t.Fatal("valid declaration was not observed") + } + logs.Reset() + old.handleHeartbeat(proto.Envelope{Type: proto.TypeHeartbeat, Payload: json.RawMessage(`{"supported_agent_kinds":[{"kind":"fake","available":true,"capabilities":{"streaming":"invented"}}]}`)}) + if logs.Len() != 0 { + t.Fatal("invalid declaration logged as an accepted snapshot") + } + newer := NewSession(newFakeConn(), "device", "tenant", "version", reg, nil) + reg.Register(newer) + publishCapabilityHeartbeat(t, old, kinds) + if logs.Len() != 0 { + t.Fatal("superseded connection logged as an accepted snapshot") + } + publishCapabilityHeartbeat(t, newer, kinds) + if !strings.Contains(logs.String(), "runtime capability snapshot observed") { + t.Fatal("replacement declaration was not observed") + } + logs.Reset() + reg.Deregister(newer) + kinds[0].Version = "changed" + publishCapabilityHeartbeat(t, newer, kinds) + if logs.Len() != 0 { + t.Fatal("disconnected connection logged as an accepted snapshot") + } +} + +type capabilityAuthorityStore struct { + HeartbeatTouch + deleted bool + err error +} + +func (a *capabilityAuthorityStore) TouchAgentDaemonHeartbeat(context.Context, runtimedevice.Heartbeat) (runtimedevice.HeartbeatStatus, error) { + return runtimedevice.HeartbeatStatus{Deleted: a.deleted}, a.err +} +func TestCapabilityConfirmationRetriesWithoutObservingRejectedAuthority(t *testing.T) { + for _, mode := range []string{"error", "deleted", "draining"} { + t.Run(mode, func(t *testing.T) { + var logs bytes.Buffer + previous := slog.Default() + slog.SetDefault(slog.New(slog.NewJSONHandler(&logs, nil))) + t.Cleanup(func() { slog.SetDefault(previous) }) + reg := NewRegistry() + s := NewSession(newFakeConn(), "device", "tenant", "version", reg, nil) + reg.Register(s) + t.Cleanup(func() { s.Close("test finished") }) + authority := &capabilityAuthorityStore{deleted: mode != "error"} + if mode == "error" { + authority.err = errors.New("private credential must not be observed") + } + s.heartbeat = authority + if mode == "draining" { + s.credentialHash = "original-hash" + s.archivedCancellations = &receiptStore{receipt: runtimedevice.ArchivedCancellationReceipt{RunID: "run", Deadline: time.Now().Add(time.Minute)}} + release, err := s.TrackExecutionDelivery("run") + if err != nil { + t.Fatal(err) + } + t.Cleanup(release) + } + kinds := []runtimedevice.SupportedAgentKind{{Kind: "fake", Available: true}} + publishCapabilityHeartbeat(t, s, kinds) + expectCapabilityHint(t, reg, false) + if logs.Len() != 0 { + t.Fatal("unconfirmed authority logged as accepted capability") + } + if mode == "draining" && s.IsClosed() { + t.Fatal("existing receipt drain was interrupted") + } + if mode == "error" { + authority.err = nil + publishCapabilityHeartbeat(t, s, kinds) + expectCapabilityHint(t, reg, true) + if !strings.Contains(logs.String(), "runtime capability snapshot observed") { + t.Fatal("confirmation recovery lost observation") + } + logs.Reset() + publishCapabilityHeartbeat(t, s, kinds) + expectCapabilityHint(t, reg, false) + if logs.Len() != 0 { + t.Fatal("unchanged confirmed heartbeat repeated observation") + } + } + }) + } +} diff --git a/services/core/internal/runtimegateway/handler.go b/services/core/internal/runtimegateway/handler.go index 85a62d55e..e7e92bf18 100644 --- a/services/core/internal/runtimegateway/handler.go +++ b/services/core/internal/runtimegateway/handler.go @@ -11,6 +11,7 @@ import ( "github.com/gorilla/websocket" "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/runtimedevice" ) @@ -176,6 +177,8 @@ func (h *Handler) WS(w http.ResponseWriter, r *http.Request) { prev.Close("preempted by newer connection from same device_id") } h.cfg.Log("agentdaemon gateway: device_id=%s registered in registry, starting session", auth.DeviceID) + obslog.Info(r.Context(), "runtime transport registered", "device_id", auth.DeviceID, "protocol_version", version, + "allocation_id", auth.RuntimeAllocationID, "node_id", auth.RuntimeNodeID) sess.Start() } diff --git a/services/core/internal/runtimegateway/registry.go b/services/core/internal/runtimegateway/registry.go index 74de0139d..2bb496d21 100644 --- a/services/core/internal/runtimegateway/registry.go +++ b/services/core/internal/runtimegateway/registry.go @@ -12,6 +12,9 @@ import ( "errors" "sync" "time" + + obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimedevice" ) // ErrDeviceNotRegistered is returned by Registry lookups when a caller @@ -49,18 +52,20 @@ type Registry struct { // is inserted into byDevice. Buffered(1) so a Register that // happens between WaitForDevice registering and selecting on the // chan still wakes the waiter. - waiters map[string][]chan *Session + waiters map[string][]chan *Session + capabilityHints chan struct{} } // NewRegistry returns an empty registry. The zero value would also // work but the constructor avoids accidental nil-map panics. func NewRegistry() *Registry { return &Registry{ - byDevice: map[string]*Session{}, - byRun: map[string]*Session{}, - byPerm: map[string]*Session{}, - byAsk: map[string]*Session{}, - waiters: map[string][]chan *Session{}, + byDevice: map[string]*Session{}, + byRun: map[string]*Session{}, + byPerm: map[string]*Session{}, + byAsk: map[string]*Session{}, + waiters: map[string][]chan *Session{}, + capabilityHints: make(chan struct{}, 1), } } @@ -340,3 +345,31 @@ func (r *Registry) removeWaiter(deviceID string, ch chan *Session) { r.waiters[deviceID] = filtered } } + +// CapabilityHints coalesces capability changes for the single execution Worker. +// A hint grants no execution authority; the Worker rechecks its normal gates. +func (r *Registry) CapabilityHints() <-chan struct{} { return r.capabilityHints } + +func (r *Registry) observeCapabilitySnapshot(sess *Session, kinds []runtimedevice.SupportedAgentKind) bool { + r.mu.RLock() + defer r.mu.RUnlock() + if r.byDevice[sess.DeviceID] != sess || sess.IsClosed() { + return false + } + available := 0 + for _, kind := range kinds { + if kind.Available { + available++ + } + } + obslog.Info(context.Background(), "runtime capability snapshot observed", "device_id", sess.DeviceID, + "kind_count", len(kinds), "available_kind_count", available) + if available == 0 { + return true + } + select { + case r.capabilityHints <- struct{}{}: + default: + } + return true +} diff --git a/services/core/internal/runtimegateway/session.go b/services/core/internal/runtimegateway/session.go index cb115fedb..a1b3f2389 100644 --- a/services/core/internal/runtimegateway/session.go +++ b/services/core/internal/runtimegateway/session.go @@ -6,6 +6,7 @@ import ( "encoding/json" "errors" "fmt" + "reflect" "strings" "sync" "time" @@ -100,9 +101,10 @@ type Session struct { // supportedKinds is the latest daemon-advertised agent_kind snapshot, // updated from heartbeat frames and read by the connector before // dispatching prompt_request so unsupported engines fail on the server. - kindsMu sync.RWMutex - kindsSeen bool - supportedKinds []runtimedevice.SupportedAgentKind + kindsMu sync.RWMutex + kindsSeen bool + supportedKinds []runtimedevice.SupportedAgentKind + capabilityObservationPending bool // Subscribers keyed by runID. The read loop only sends on these // channels; Unsubscribe is the only place that closes them. @@ -233,11 +235,25 @@ func (s *Session) setSupportedAgentKinds(kinds []runtimedevice.SupportedAgentKin copyKinds := make([]runtimedevice.SupportedAgentKind, len(kinds)) copy(copyKinds, kinds) s.kindsMu.Lock() + changed := !s.kindsSeen || !reflect.DeepEqual(s.supportedKinds, copyKinds) s.kindsSeen = true s.supportedKinds = copyKinds + if changed { + s.capabilityObservationPending = true + } s.kindsMu.Unlock() } +// A failed authorization or persistence check retains the diagnostic change +// for the next confirmed heartbeat without changing capability publication. +func (s *Session) observeConfirmedCapabilities(kinds []runtimedevice.SupportedAgentKind) { + s.kindsMu.Lock() + defer s.kindsMu.Unlock() + if s.capabilityObservationPending && s.reg.observeCapabilitySnapshot(s, kinds) { + s.capabilityObservationPending = false + } +} + // Close closes the transport and subscriptions with ErrSessionClosed, then // releases connection ownership. It establishes no execution outcome. Idempotent. func (s *Session) Close(reason string) { @@ -464,6 +480,7 @@ func (s *Session) handleHeartbeat(env proto.Envelope) { kinds := deviceKindsFromHeartbeat(p) s.setSupportedAgentKinds(kinds) if s.heartbeat == nil { + s.observeConfirmedCapabilities(kinds) return } @@ -492,7 +509,9 @@ func (s *Session) handleHeartbeat(env proto.Envelope) { // actual admin action; the daemon only sees it's no longer // the current owner. s.CloseWithCode(CloseRuntimeDeleted, "runtime retired") + return } + s.observeConfirmedCapabilities(kinds) } func deviceKindsFromHeartbeat(p proto.HeartbeatPayload) []runtimedevice.SupportedAgentKind { From f2d860ffc30372fd7c3f197f8302a377e93cce4a Mon Sep 17 00:00:00 2001 From: sam Date: Tue, 6 Oct 2026 21:08:16 +0800 Subject: [PATCH 2/3] Record beta runtime readiness release requirements --- deploy/kubernetes/BETA_CHANGELOG.md | 15 +++++++++++++++ 1 file changed, 15 insertions(+) diff --git a/deploy/kubernetes/BETA_CHANGELOG.md b/deploy/kubernetes/BETA_CHANGELOG.md index 598ccafac..6d6ad31cf 100644 --- a/deploy/kubernetes/BETA_CHANGELOG.md +++ b/deploy/kubernetes/BETA_CHANGELOG.md @@ -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. From d321b7b7e28254b58ce48ef925630540c21c7620 Mon Sep 17 00:00:00 2001 From: sam Date: Tue, 6 Oct 2026 21:16:23 +0800 Subject: [PATCH 3/3] Preserve gateway-free shutdown and verify capability-driven scans --- .../internal/execution/worker_wakeup_test.go | 16 ++++++++++++++++ .../core/internal/runtimegateway/registry.go | 7 ++++++- .../store/environment_worker_scan_test.go | 8 ++++---- .../core/internal/store/worker_wakeup_test.go | 6 ++++++ 4 files changed, 32 insertions(+), 5 deletions(-) diff --git a/services/core/internal/execution/worker_wakeup_test.go b/services/core/internal/execution/worker_wakeup_test.go index 8a0a48f15..a4ade2ef5 100644 --- a/services/core/internal/execution/worker_wakeup_test.go +++ b/services/core/internal/execution/worker_wakeup_test.go @@ -1,6 +1,8 @@ package execution import ( + "context" + "errors" "sync" "testing" ) @@ -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") + } +} diff --git a/services/core/internal/runtimegateway/registry.go b/services/core/internal/runtimegateway/registry.go index 2bb496d21..4c5b37096 100644 --- a/services/core/internal/runtimegateway/registry.go +++ b/services/core/internal/runtimegateway/registry.go @@ -348,7 +348,12 @@ func (r *Registry) removeWaiter(deviceID string, ch chan *Session) { // CapabilityHints coalesces capability changes for the single execution Worker. // A hint grants no execution authority; the Worker rechecks its normal gates. -func (r *Registry) CapabilityHints() <-chan struct{} { return r.capabilityHints } +func (r *Registry) CapabilityHints() <-chan struct{} { + if r == nil { + return nil + } + return r.capabilityHints +} func (r *Registry) observeCapabilitySnapshot(sess *Session, kinds []runtimedevice.SupportedAgentKind) bool { r.mu.RLock() diff --git a/services/core/internal/store/environment_worker_scan_test.go b/services/core/internal/store/environment_worker_scan_test.go index 64fe1191f..214c81b82 100644 --- a/services/core/internal/store/environment_worker_scan_test.go +++ b/services/core/internal/store/environment_worker_scan_test.go @@ -8,7 +8,7 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" ) -func TestWorkerEnvironmentRetriesNewlyReadyAtNextScan(t *testing.T) { +func TestWorkerEnvironmentReadinessHintBypassesNextScan(t *testing.T) { h := newDispatchHarness(t) enableWorkerEnvironment(t, h) pending := workerEnvironmentReservation(t, h) @@ -20,7 +20,6 @@ func TestWorkerEnvironmentRetriesNewlyReadyAtNextScan(t *testing.T) { h.session = publicSession(t, h, "scan-barrier") receipt := h.message("barrier", "ordinary work") - scanned := time.Now() _, stop := startEnvironmentExpiryWorker(t, h.db, h.d) // Dispatch starts only after selectWork has examined the pending input on // the same pass, while its exact Runtime is still incapable of preparation. @@ -28,13 +27,14 @@ func TestWorkerEnvironmentRetriesNewlyReadyAtNextScan(t *testing.T) { if barrier.ID != receipt.TurnID { t.Fatal("unexpected scan barrier") } + readyAt := time.Now() awaitFixtureCapabilities(t, runtime, workerEnvironmentCapabilities()) h.write(barrier.ID, proto.TypeDone, proto.DonePayload{Content: "complete"}) waitTurn(t, h, barrier.ID, sessions.TurnCompleted) prepare := nextWorkerFrame(t, frames, proto.TypeExecutionPrepare) - if elapsed := time.Since(scanned); elapsed < 750*time.Millisecond || elapsed > 3*time.Second { - t.Fatal("readiness retry must use the next scan, without an empty scan interval", elapsed) + if elapsed := time.Since(readyAt); elapsed >= 750*time.Millisecond { + t.Fatal("confirmed readiness waited for the periodic candidate scan", elapsed) } if workerRuntimeForPreparation(t, h, prepare) != runtime { t.Fatal("readiness retry moved Runtime ownership") diff --git a/services/core/internal/store/worker_wakeup_test.go b/services/core/internal/store/worker_wakeup_test.go index 1b3813e19..2775b0d54 100644 --- a/services/core/internal/store/worker_wakeup_test.go +++ b/services/core/internal/store/worker_wakeup_test.go @@ -84,6 +84,12 @@ func TestWorkerSchedulerCommittedAdmissionWakesBeforeMaintenance(t *testing.T) { if err != nil { t.Fatal(err) } + // Fixture connection/capability hints precede this admission operation. + // Remove them so a rejected admission is tested independently. + select { + case <-h.registry.CapabilityHints(): + default: + } trace.armed.Store(true) started = true go func() { done <- worker.Run(ctx) }()