Skip to content

Commit 4fa91ee

Browse files
committed
Document release fencing and echo assignments in integration fakes
Cover the release fence on an unfinished workspace write, make the integration Runtime fakes echo the assignment of each request and Run, and document echoed assignments, release fencing, settlement without authority and release backoff.
1 parent 0b75b12 commit 4fa91ee

6 files changed

Lines changed: 69 additions & 16 deletions

File tree

‎apps/daemon/internal/dispatch/workspace_write_test.go‎

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -168,3 +168,41 @@ func TestLocalUploadReportsDestinationConflictsAndReleasesOwner(t *testing.T) {
168168
})
169169
}
170170
}
171+
172+
func TestReleaseFencesUnfinishedWorkspaceWrite(t *testing.T) {
173+
r, sender, request, workspace := localWriterRouter(t)
174+
id := uuid.NewString()
175+
if err := r.Handle(t.Context(), mustEnv(t, proto.TypeWorkspaceWrite, id, request)); err != nil {
176+
t.Fatal(err)
177+
}
178+
waitWorkspaceWrite(t, sender, id, "ready")
179+
if err := r.Handle(t.Context(), mustEnv(t, proto.TypeWorkspaceWrite, id, proto.WorkspaceWritePayload{Step: "chunk", Data: []byte("abc")})); err != nil {
180+
t.Fatal(err)
181+
}
182+
waitWorkspaceWrite(t, sender, id, "received")
183+
release(t, r, preparationSessionID, "release", 2, false)
184+
if got := waitWorkspaceWrite(t, sender, id, "rejected"); got.ErrorCode != proto.AssignmentStale {
185+
t.Fatal("release did not fence the write", got)
186+
}
187+
if got := waitAssignmentStatus(t, sender, "release"); got.State != proto.AssignmentReleased {
188+
t.Fatal(got)
189+
}
190+
// The release replies only after the write's result.
191+
for _, frame := range sender.snapshot() {
192+
if frame.Type == proto.TypeAssignmentStatus && frame.ID == "release" {
193+
t.Fatal("release replied before the write settled")
194+
}
195+
if frame.Type == proto.TypeWorkspaceWriteResult && frame.ID == id && frame.Assignment == ref(preparationSessionID) {
196+
var result proto.WorkspaceWriteResultPayload
197+
if frame.DecodePayload(&result) == nil && result.Outcome == "rejected" {
198+
break
199+
}
200+
}
201+
}
202+
if err := r.Handle(t.Context(), mustEnv(t, proto.TypeWorkspaceWrite, id, proto.WorkspaceWritePayload{Step: "commit"})); err != nil {
203+
t.Fatal(err)
204+
}
205+
if _, err := os.Stat(filepath.Join(workspace, "file")); !os.IsNotExist(err) {
206+
t.Fatal("a released assignment's write applied", err)
207+
}
208+
}

‎docs/runtime-protocol.md‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -115,11 +115,11 @@ Usage frames and the final usage snapshot each carry the cumulative measurement
115115

116116
An assignment binds one Session to the Runtime that runs it. `Envelope.assignment` names it as `session_id`, `assignment_id` and `epoch`, and is the only place a frame carries it. Core advances the epoch whenever it changes the assignment's desired state, so a lower epoch is stale.
117117

118-
Every Session frame carries the assignment: `execution_prepare`, `execution_start` and `execution_release`; `prompt_cancel`, `prompt_steer`, `function_result`, `permission_decision` and `prompt_for_user_choice_decision`; every frame of `runtime_prepare`, `workspace_read`, `workspace_write` and `workspace_export`; and `environment_quiesce` and `environment_resume`. A reply echoes its request's assignment. Heartbeats carry none.
118+
Every Session frame carries the assignment: `execution_prepare`, `execution_start` and `execution_release`; `prompt_cancel`, `prompt_steer`, `function_result`, `permission_decision` and `prompt_for_user_choice_decision`; every frame of `runtime_prepare`, `workspace_read`, `workspace_write` and `workspace_export`; and `environment_quiesce` and `environment_resume`. A reply echoes its request's assignment, and a Run's frames carry the assignment that started it; Core rejects a reply or Run frame that names another. Heartbeats carry none.
119119

120-
Before a Session's first operation on a connection, including Environment initialization and file work without a Turn, Core sends `assignment_bind` with the Session's Environment ID and waits for `assignment_status` `bound`. A repeated bind of the same assignment is `bound` again. The Runtime admits a Session frame only under the assignment it bound: an older epoch, or a released one, fails with `assignment_stale`; another assignment, Session or Environment fails with `assignment_conflict`. A started Run's frames, including its cancellation receipt, stay admissible under the assignment that started it until the release.
120+
Before a Session's first operation on a connection, including Environment initialization and file work without a Turn, Core sends `assignment_bind` with the Session's Environment ID and waits for `assignment_status` `bound`. A repeated bind of the same assignment is `bound` again. The Runtime admits a Session frame only under the assignment it bound: an older epoch, or a released one, fails with `assignment_stale`; another assignment, Session or Environment fails with `assignment_conflict`. A started Run's frames, including its cancellation receipt, stay admissible under the assignment that started it until the release. A repeated function result or decision whose receipt the Runtime already recorded is answered only under the assignment that applied it; another fails with `assignment_conflict`.
121121

122-
Core records a release and advances the epoch before it sends anything. Deleting a Session releases its assignment with `remove_home: true`; releasing its Environment sends `false`. A deletion never revokes a shared Runtime credential. `assignment_release` fences the assignment at once. The Runtime then stops the Session's work, closes its Executors and, when asked, removes the native home; only then does it reply `released` or `home_removed`. Unfinished cleanup replies `failed` with `cleanup_unconfirmed`, and a retry at the same epoch repeats it. A Runtime that declares `home_removal` unsupported answers `remove_home: true` with `unsupported_operation`, and Core asks it only to release. Core records the release as applied only from a matching `released` or `home_removed`, and resends every unacknowledged release to a Runtime when it connects. A quiesced Runtime admits only a release and the matching `environment_resume`, which carries the assignment that quiesced it.
122+
Core records a release and advances the epoch before it sends anything. Deleting a Session releases its assignment with `remove_home: true`; releasing its Environment sends `false`. A deletion never revokes a shared Runtime credential. `assignment_release` fences the assignment at once. The Runtime then stops the Session's work: a transfer still receiving its body, or committed but not yet applied, ends with `assignment_stale`; it releases read-only preparations and waits until every workspace read, write, export and Runtime preparation has sent its result. It closes the Session's Executors and, when asked, removes the native home; only then does it reply `released` or `home_removed`. Unfinished cleanup replies `failed` with `cleanup_unconfirmed`, and a retry at the same epoch repeats it. A Runtime that declares `home_removal` unsupported answers `remove_home: true` with `unsupported_operation`, and Core asks it only to release. Core records the release as applied from a matching `released` or `home_removed`, or at once when no Runtime is left to act on it: a release to a Runtime without authority is settled when recorded, and revoking a Runtime settles its releases. Core resends every unacknowledged release to a Runtime when it connects and backs off a release that fails. A quiesced Runtime admits only a release and the matching `environment_resume`, which carries the assignment that quiesced it.
123123

124124
The Runtime answers a Core frame it cannot route with `protocol_error`, which echoes the request's ID and carries its type and an error code.
125125

‎docs/zh/runtime-protocol.md‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
---
22
title: "Core–Runtime 协议"
33
source: docs/runtime-protocol.md
4-
source_hash: ca1f05d7f95e6a91f57e44919da83d05e94b146809d727a4e87ad38c94a9d6fd
4+
source_hash: 84a2ec5b9308a5e1a118b05c1ea1c92f8741871785524f3878321286b7e2d25b
55
---
66

77
此协议在 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)。
@@ -117,11 +117,11 @@ Usage frame 和最终 usage snapshot 都携带当前执行的累计测量,替
117117

118118
分配(assignment)把一个 Session 绑定到运行它的 Runtime。`Envelope.assignment` 以 `session_id`、`assignment_id` 和 `epoch` 命名它,这是 frame 携带分配的唯一位置。Core 每次改变分配的期望状态时推进 epoch,因此较低的 epoch 是陈旧的。
119119

120-
每个 Session frame 都携带分配:`execution_prepare`、`execution_start` 和 `execution_release`;`prompt_cancel`、`prompt_steer`、`function_result`、`permission_decision` 和 `prompt_for_user_choice_decision`;`runtime_prepare`、`workspace_read`、`workspace_write` 和 `workspace_export` 的每个 frame;以及 `environment_quiesce` 和 `environment_resume`。回复回显请求的分配。heartbeat 不携带分配。
120+
每个 Session frame 都携带分配:`execution_prepare`、`execution_start` 和 `execution_release`;`prompt_cancel`、`prompt_steer`、`function_result`、`permission_decision` 和 `prompt_for_user_choice_decision`;`runtime_prepare`、`workspace_read`、`workspace_write` 和 `workspace_export` 的每个 frame;以及 `environment_quiesce` 和 `environment_resume`。回复回显请求的分配,Run 的 frame 携带启动它的分配;Core 拒绝指明其他分配的回复或 Run frame。heartbeat 不携带分配。
121121

122-
在一条连接上执行 Session 的第一个操作之前,包括没有 Turn 的 Environment 初始化和文件操作,Core 发送带 Session 的 Environment ID 的 `assignment_bind`,并等待 `assignment_status` `bound`。重复绑定同一分配仍得到 `bound`。Runtime 只在其已绑定的分配下准入 Session frame:较旧的 epoch 或已释放的分配以 `assignment_stale` 失败;其他分配、Session 或 Environment 以 `assignment_conflict` 失败。已启动 Run 的 frame,包括其取消回执,在释放前仍可在启动它的分配下准入。
122+
在一条连接上执行 Session 的第一个操作之前,包括没有 Turn 的 Environment 初始化和文件操作,Core 发送带 Session 的 Environment ID 的 `assignment_bind`,并等待 `assignment_status` `bound`。重复绑定同一分配仍得到 `bound`。Runtime 只在其已绑定的分配下准入 Session frame:较旧的 epoch 或已释放的分配以 `assignment_stale` 失败;其他分配、Session 或 Environment 以 `assignment_conflict` 失败。已启动 Run 的 frame,包括其取消回执,在释放前仍可在启动它的分配下准入。Runtime 已记录回执的重复函数结果或决策只在应用它的分配下得到回答;其他分配以 `assignment_conflict` 失败。
123123

124-
Core 先记录释放并推进 epoch,再发送任何消息。删除 Session 以 `remove_home: true` 释放其分配;释放其 Environment 发送 `false`。删除从不吊销共享的 Runtime 凭据。`assignment_release` 立即约束该分配。随后 Runtime 停止 Session 的工作,关闭其 Executor,并在要求时删除原生 home;此后才回复 `released` 或 `home_removed`。未完成的清理回复 `failed` 和 `cleanup_unconfirmed`,同一 epoch 的重试会重复清理。声明 `home_removal` 不支持的 Runtime 以 `unsupported_operation` 回答 `remove_home: true`,Core 只要求它释放。Core 只根据匹配的 `released` 或 `home_removed` 记录释放已应用,并在 Runtime 连接时重发所有未确认的释放。已 quiesce 的 Runtime 只准入释放和匹配的 `environment_resume`,后者携带使其 quiesce 的分配。
124+
Core 先记录释放并推进 epoch,再发送任何消息。删除 Session 以 `remove_home: true` 释放其分配;释放其 Environment 发送 `false`。删除从不吊销共享的 Runtime 凭据。`assignment_release` 立即约束该分配。随后 Runtime 停止 Session 的工作:仍在接收内容、或已提交但尚未应用的传输以 `assignment_stale` 结束;它释放只读准备,并等待每个 workspace 读取、写入、导出和 Runtime 准备发送结果。它关闭 Session 的 Executor,并在要求时删除原生 home;此后才回复 `released` 或 `home_removed`。未完成的清理回复 `failed` 和 `cleanup_unconfirmed`,同一 epoch 的重试会重复清理。声明 `home_removal` 不支持的 Runtime 以 `unsupported_operation` 回答 `remove_home: true`,Core 只要求它释放。Core 根据匹配的 `released` 或 `home_removed` 记录释放已应用;没有 Runtime 能处理该释放时立即记录:发给无授权 Runtime 的释放在记录时即结清,吊销 Runtime 会结清它的释放。Core 在 Runtime 连接时重发所有未确认的释放,并对失败的释放退避重试。已 quiesce 的 Runtime 只准入释放和匹配的 `environment_resume`,后者携带使其 quiesce 的分配。
125125

126126
Runtime 对无法路由的 Core frame 回复 `protocol_error`,回显请求 ID,并携带其类型和错误码。
127127

‎services/core/internal/runtimegateway/exchange_test.go‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,10 @@ func TestFramesMustNameTheRequestAssignment(t *testing.T) {
1515
s := NewSession(newFakeConn(), "device", "tenant", "test", NewRegistry(), nil)
1616
defer s.Close("test")
1717
read := make(chan error, 1)
18-
go func() { _, err := s.ReadWorkspaceFile(t.Context(), testAssignment, workspaceReadRequest()); read <- err }()
18+
go func() {
19+
_, err := s.ReadWorkspaceFile(t.Context(), testAssignment, workspaceReadRequest())
20+
read <- err
21+
}()
1922
request := <-s.sendCh
2023
reply, _ := request.Reply(proto.TypeWorkspaceReadResult, proto.WorkspaceReadResultPayload{Outcome: "completed", CloseAcknowledged: true})
2124
reply.Assignment = ref

‎services/core/tests/integration/dispatch_test.go‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,8 @@ import (
2727

2828
type dispatchHarness struct {
2929
writeMu sync.Mutex
30-
assignment proto.AssignmentRef // the latest assignment Core named; writeMu guards it
30+
assignments map[string]proto.AssignmentRef // by frame and Run ID; writeMu guards it
31+
assignment proto.AssignmentRef // the latest assignment Core named
3132
admissions map[string]fixtureAdmission
3233
t *testing.T
3334
s *Store
@@ -146,7 +147,9 @@ func (h *dispatchHarness) write(run, kind string, payload any) {
146147
if err != nil {
147148
h.t.Fatal(err)
148149
}
149-
if kind != proto.TypeHeartbeat {
150+
if ref, ok := h.assignments[run]; ok {
151+
env.Assignment = ref
152+
} else if kind != proto.TypeHeartbeat {
150153
env.Assignment = h.assignment
151154
}
152155
if err = h.conn.WriteJSON(env); err != nil {

‎services/core/tests/integration/executor_fixture_test.go‎

Lines changed: 15 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -26,13 +26,22 @@ func assignmentReply(env proto.Envelope) (proto.Envelope, bool) {
2626
return reply, err == nil
2727
}
2828

29-
// observe remembers the assignment env names. The fixture Runtime writes its
30-
// frames under it, as a Runtime echoes its request's assignment.
29+
// observe remembers the assignment env names for its ID and started Run. The
30+
// fixture Runtime writes its frames under it, as a Runtime echoes its
31+
// request's assignment.
3132
func (h *dispatchHarness) observe(env proto.Envelope) {
32-
if env.Assignment.Valid() {
33-
h.writeMu.Lock()
34-
h.assignment = env.Assignment
35-
h.writeMu.Unlock()
33+
if !env.Assignment.Valid() {
34+
return
35+
}
36+
h.writeMu.Lock()
37+
defer h.writeMu.Unlock()
38+
if h.assignments == nil {
39+
h.assignments = make(map[string]proto.AssignmentRef)
40+
}
41+
h.assignment, h.assignments[env.ID] = env.Assignment, env.Assignment
42+
var start proto.ExecutionStartPayload
43+
if env.Type == proto.TypeExecutionStart && env.DecodePayload(&start) == nil && start.RunID != "" {
44+
h.assignments[start.RunID] = env.Assignment
3645
}
3746
}
3847

0 commit comments

Comments
 (0)