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
35 changes: 27 additions & 8 deletions apps/daemon/internal/dispatch/runtime_preparation_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -246,6 +246,7 @@ func TestRuntimePreparationUploadBlocksWorkspaceWriteAndSuspension(t *testing.T)

func TestRuntimePreparationCancellationKeepsOwnershipUntilApplyStops(t *testing.T) {
r, sender, environment, session := capabilitiesTestRouter(t)
sender.frames = make(chan proto.Envelope)
ctx, cancel := context.WithCancel(context.Background())
id := uuid.NewString()
request := proto.RuntimePreparePayload{BudgetMS: 300000, Step: "begin", Action: "finalize", EnvironmentID: environment, SessionID: session, Sources: &agentcapabilities.Input{}}
Expand All @@ -263,7 +264,10 @@ func TestRuntimePreparationCancellationKeepsOwnershipUntilApplyStops(t *testing.
<-ctx.Done()
close(interrupted)
<-release
return os.WriteFile(retained, []byte("retained"), 0400)
if err := os.WriteFile(retained, []byte("retained"), 0400); err != nil {
return err
}
return context.Canceled
})
<-started
wait, stop := context.WithTimeout(context.Background(), 20*time.Millisecond)
Expand All @@ -280,13 +284,22 @@ func TestRuntimePreparationCancellationKeepsOwnershipUntilApplyStops(t *testing.
t.Fatal("cancel released unsettled capability ownership")
}
close(release)
wait, stop = context.WithTimeout(context.Background(), 20*time.Millisecond)
err = r.Shutdown(wait)
stop()
if !errors.Is(err, context.DeadlineExceeded) {
t.Fatalf("shutdown abandoned pending result delivery: %v", err)
}
capabilitiesReceipt(t, sender, id, "unknown")
shutdownCapabilitiesRouter(t, r)
capabilitiesReceipt(t, sender, id, "completed")
if _, err := os.Stat(retained); err != nil {
t.Fatal("shutdown deleted installation result")
}
if r.runtimePreparation != nil {
t.Fatal("confirmed completion retained capacity")
if r.runtimePreparation != owner || !owner.uncertain {
t.Fatal("shutdown erased unknown outcome")
}
if err := r.Handle(t.Context(), capabilityEnvelope(t, uuid.NewString(), request)); !errors.Is(err, ErrRouterClosed) {
t.Fatalf("closed router admitted a successor: %v", err)
}
}

Expand Down Expand Up @@ -341,10 +354,16 @@ func TestRuntimePreparationResultCategoriesAndUnknownOwnership(t *testing.T) {
if !owned {
t.Fatal("unknown mutation released its ownership")
}
wait, stop := context.WithTimeout(context.Background(), time.Second)
defer stop()
if err := r.Shutdown(wait); err == nil {
t.Fatal("shutdown claimed uncertain mutation settled")
next := uuid.NewString()
if err := r.Handle(t.Context(), capabilityEnvelope(t, next, request)); err != nil {
t.Fatal(err)
}
if got := capabilitiesReceipt(t, sender, next, "rejected"); got.ErrorCode != "runtime_preparation_capacity" {
t.Fatalf("unknown operation lost its admission fence: %+v", got)
}
shutdownCapabilitiesRouter(t, r)
if r.runtimePreparation != owner || !owner.uncertain {
t.Fatal("shutdown erased unknown outcome")
}
}

Expand Down
6 changes: 3 additions & 3 deletions apps/daemon/internal/dispatch/shutdown.go
Original file line number Diff line number Diff line change
Expand Up @@ -69,9 +69,9 @@ func (r *Router) runShutdownAttempt(attempt *shutdownAttempt, victims []sessionC

r.shutdownWG.Wait()
r.mu.Lock()
if r.runtimePreparation != nil && r.runtimePreparation.uncertain {
attempt.err = errors.Join(attempt.err, errors.New("dispatch: capability preparation remains uncertain"))
}
// Runtime preparation joins only after ApplyRuntimePreparation stops local
// mutations and its receipt send finishes. An unknown result still fences
// this closed Router, but is not outstanding cleanup. Core owns no-replay.
if r.workspaceWrite != nil && r.workspaceWrite.uncertain {
attempt.err = errors.Join(attempt.err, errors.New("dispatch: local workspace write remains uncertain"))
}
Expand Down
2 changes: 1 addition & 1 deletion docs/runtime-protocol.md
Original file line number Diff line number Diff line change
Expand Up @@ -124,7 +124,7 @@ A preparation reserves a per-Turn admission, not a new Executor. It carries an e

Preparation and start run outside the receive loop and router lock. An admission expires five minutes after it is granted, and retries do not extend that deadline; expiry does not remove the Runtime's obligation to settle cleanup. The Runtime bounds active preparation and execution separately from idle retained resources and counts closing or uncertain resources until their cleanup succeeds. A definite `execution_prepare` rejection with `preparation_capacity` leaves the queued Turn unclaimed for the Worker to retry, including when cleanup holds the capacity; any other error or uncertain delivery authorizes no replay. The Runtime retains at most 64 admission records, and an old handle never consumes a replacement's admission. These records are connection-local, not durable input replay.

Idle expiry of an Executor is a Runtime resource policy, separate from Core's active-Turn concurrency. On shutdown the Runtime closes active and idle Executors, keeps any target whose close failed and allows a later serialized retry. An ordinary disconnection closes the failed transport and keeps the exact router until shutdown succeeds; a wait timeout or failed cleanup never authorizes reconnection, and process shutdown keeps waiting rather than discarding owned native resources. Workspace operations keep their binding and settlement rules across Turn boundaries and Executor closure.
Idle expiry of an Executor is a Runtime resource policy, separate from Core's active-Turn concurrency. On shutdown the Runtime closes active and idle Executors, keeps any target whose close failed and allows a later serialized retry. An ordinary disconnection closes the failed transport and keeps the exact router until shutdown succeeds; a wait timeout or failed cleanup never authorizes reconnection, and process shutdown keeps waiting rather than discarding owned native resources. For Runtime initialization, shutdown joins the synchronous apply operation and its receipt delivery before releasing the connection; an `unknown` result remains unknown and blocks further preparation on that Router, but does not prevent reconnection after local mutations have stopped. Core retains the failed initialization and never replays it. Workspace operations keep their binding and settlement rules across Turn boundaries and Executor closure.

Suspension closes admission and drains admitted work and receipts while retaining settled idle Executors before acknowledging `environment_quiesced`. Active, preparing, invalid or otherwise unsettled owners prevent suspension. Idle expiry is stopped and its callbacks fenced throughout suspension; exact authenticated `environment_resume` restores a fresh idle interval on the same owners. Only this planned reconnect preserves the Router and native owners across sockets; ordinary disconnect, shutdown and cancellation of the connection lifecycle still require confirmed cleanup. A drain failure keeps admission closed until shutdown. Preparation after resume retains the same native identity and immutable-configuration checks, including credential conflicts; it never silently replaces a conflicting resident Executor. The Sandbox Provider remains responsible for the filesystem flush and compute-stop guarantees of its checkpoint implementation; Runtime quiescence alone proves neither those guarantees nor restoration of native RAM or external network connections.

Expand Down
4 changes: 2 additions & 2 deletions docs/zh/runtime-protocol.md
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
---
title: "Core–Runtime 协议"
source: docs/runtime-protocol.md
source_hash: 902f53e5471de13ecac0e914e4d23d7650955a4c50ec1aa01517c13aef7946f7
source_hash: 8349aca995e66669b7941b6aad5f28e6ae3259516d923a0494938a61fb3e711e
---

此协议在 Runtime daemon 获取机器凭据后连接 Core 与 daemon,定义 daemon 连接上消息的含义和顺序。wire 类型、限制和验证器仅在 [`internal/agentdaemon/proto`](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/internal/agentdaemon/proto) 中定义一次;Core 的 [gateway](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/services/core/internal/runtimegateway) 与参考 Runtime 的 [dispatcher](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/apps/daemon/internal/dispatch) 都使用它们,因此无需同步第二套 payload schema。签发凭据和打开连接的 HTTP 路由见[机器连接 API](../../contracts/agents-api/zh/machine-api.md)。
Expand Down Expand Up @@ -126,7 +126,7 @@ preparation 预约每个 Turn 的准入,而不是新 Executor。它携带明

preparation 和 start 在 receive loop 与 router lock 之外运行。admission 在授予五分钟后到期,重试不延长截止时间;到期不解除 Runtime 完成清理结算的义务。Runtime 分别限制活动 preparation、execution 和保留的空闲资源,关闭中或不确定资源持续计入限制,直到清理成功。明确的 `execution_prepare` 拒绝若为 `preparation_capacity`,会让排队 Turn 保持未领取,供 Worker 重试,包括清理占用容量的情况;其他错误或不确定交付都不授权重放。Runtime 最多保留 64 条 admission 记录,旧 handle 不会消耗替代项的 admission。这些记录仅属于连接,不是持久化输入重放。

Executor 空闲到期属于 Runtime 资源策略,与 Core 的活动 Turn 并发限制独立。关闭时 Runtime 关闭活动和空闲 Executor,保留关闭失败的目标,并允许稍后串行重试。普通断连会关闭失败的 transport 并保留原 router,直到 shutdown 成功;等待超时或清理失败不授权重连,进程 shutdown 继续等待,不丢弃自己拥有的原生资源。工作区操作在跨 Turn 和 Executor 关闭后仍保留绑定与结算规则。
Executor 空闲到期属于 Runtime 资源策略,与 Core 的活动 Turn 并发限制独立。关闭时 Runtime 关闭活动和空闲 Executor,保留关闭失败的目标,并允许稍后串行重试。普通断连会关闭失败的 transport 并保留原 router,直到 shutdown 成功;等待超时或清理失败不授权重连,进程 shutdown 继续等待,不丢弃自己拥有的原生资源。对于 Runtime 初始化,shutdown 会等待同步 apply 操作和回执发送结束后才释放连接;`unknown` 结果仍是不确定结果,并阻止该 Router 上的后续准备,但本地变更停止后不再阻止重连。Core 保留失败的初始化状态,绝不重放初始化。工作区操作在跨 Turn 和 Executor 关闭后仍保留绑定与结算规则。


挂起先关闭准入、排空已接纳工作与回执,同时保留已结算的空闲 Executor,随后才确认 `environment_quiesced`。活动、准备中、失效或其他未结算的所有者会阻止挂起。整个挂起期间停止空闲到期计时并隔离其回调;精确匹配且已认证的 `environment_resume` 为同一批所有者恢复一个新的空闲计时间隔。只有这种计划性重连跨 socket 保留 Router 与原生所有者;普通断连、shutdown 和连接生命周期取消仍要求确认清理完成。排空失败后,准入保持关闭直到 shutdown。恢复后的准备仍检查相同的原生身份与不可变配置,包括凭据冲突;绝不静默替换配置冲突的驻留 Executor。Sandbox Provider 仍负责其 checkpoint 实现的文件系统刷新和计算停止保证;Runtime 静止本身既不证明这些保证,也不证明原生 RAM 或外部网络连接已恢复。
Expand Down
18 changes: 16 additions & 2 deletions services/core/tests/integration/runtime_initialization_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -121,8 +121,22 @@ func TestEnvironmentInitializationCompletionUnknownAndRestart(t *testing.T) {
if int(p.writes.Load()) != expectedSteps {
t.Fatal("completed preparation replayed")
}
} else if p.writes.Load() != 1 {
t.Fatal("unknown operation replayed", p.writes.Load())
} else {
if p.writes.Load() != 1 {
t.Fatal("unknown operation replayed", p.writes.Load())
}
stop()
_, _ = managedWorkerMode(t, s, key, p, true)
if err := p.connect(sandbox.Bootstrap{DeviceID: owner.DeviceID, Credential: p.credential}); err != nil {
t.Fatal(err)
}
time.Sleep(350 * time.Millisecond)
if p.writes.Load() != 1 || initializationState(t, s, tenant, env.ID) != "failed" {
t.Fatal("reconnect or Worker restart replayed unknown initialization")
}
if _, err := sessionAdapter(s).GetSessionExecutionBinding(t.Context(), tenant, session.ID); !errors.Is(err, sessions.ErrNotFound) {
t.Fatal("unknown initialization admitted execution", err)
}
}
p.mu.Lock()
kills := p.kills
Expand Down
Loading