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
18 changes: 14 additions & 4 deletions apps/daemon/internal/dispatch/assignment.go
Original file line number Diff line number Diff line change
@@ -1,19 +1,25 @@
package dispatch

import (
"bytes"
"context"
"sync"
"time"

"github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto"
"github.com/MiniMax-AI/OpenAgentCore/internal/sandboxbootstrap"
)

// assignmentState is the Router's record of one Session's assignment. Router.mu
// protects it. A released assignment stays recorded, so its frames stay fenced.
type assignmentState struct {
ref proto.AssignmentRef
environmentID string
released bool
// resource and grant are the bind's Link resource and attach grant; the
// resource's Kind is empty when the bind carried none.
resource sandboxbootstrap.Resource
grant []byte
released bool
// work counts the Session's admitted reads, writes, exports and Runtime
// preparations until each has sent its terminal result. A release waits
// for it, and the released assignment admits no more.
Expand Down Expand Up @@ -55,17 +61,21 @@ func (r *Router) admitRunLocked(ref proto.AssignmentRef, state *sessionState) st
func (r *Router) handleAssignmentBind(ctx context.Context, env proto.Envelope) error {
var input proto.AssignmentBindPayload
ref, code := env.Assignment, ""
if env.ID == "" || env.DecodeRequest(&input) != nil || !ref.Valid() {
if env.ID == "" || env.DecodeRequest(&input) != nil || input.Validate() != nil || !ref.Valid() {
code = "invalid_request"
} else {
var resource sandboxbootstrap.Resource
if input.Resource != nil {
resource = *input.Resource
}
r.mu.Lock()
a := r.assignments[ref.SessionID]
switch {
case a == nil:
r.assignments[ref.SessionID] = &assignmentState{ref: ref, environmentID: input.EnvironmentID}
r.assignments[ref.SessionID] = &assignmentState{ref: ref, environmentID: input.EnvironmentID, resource: resource, grant: input.AttachGrant}
case a.ref.AssignmentID == ref.AssignmentID && (ref.Epoch < a.ref.Epoch || ref.Epoch == a.ref.Epoch && a.released):
code = proto.AssignmentStale
case a.ref != ref || a.environmentID != input.EnvironmentID:
case a.ref != ref || a.environmentID != input.EnvironmentID || a.resource != resource || !bytes.Equal(a.grant, input.AttachGrant):
code = proto.AssignmentConflict
}
r.mu.Unlock()
Expand Down
39 changes: 39 additions & 0 deletions apps/daemon/internal/dispatch/assignment_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,12 @@ import (
"testing"
"time"

"github.com/google/uuid"

"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent"
"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/dispatch"
"github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto"
"github.com/MiniMax-AI/OpenAgentCore/internal/sandboxbootstrap"
)

// frameFor returns the last frame of kind correlated with id.
Expand Down Expand Up @@ -168,3 +171,39 @@ func TestUnknownEnvelopeGetsCorrelatedProtocolError(t *testing.T) {
t.Fatalf("protocol_error = %+v %+v", frame, got)
}
}

func TestAssignmentBindCarriesLink(t *testing.T) {
h := newHarness(t)
defer h.router.Shutdown(context.Background())
environment := uuid.NewString()
resource := &sandboxbootstrap.Resource{TenantID: uuid.NewString(), EnvironmentID: environment, Kind: "allocation", ID: uuid.NewString(), Generation: 1}
other := *resource
other.EnvironmentID = uuid.NewString()
bind := func(id string, payload proto.AssignmentBindPayload) proto.AssignmentStatusPayload {
if err := h.router.Handle(t.Context(), scoped(t, "s", proto.TypeAssignmentBind, id, payload)); err != nil {
t.Fatal(err)
}
return waitAssignmentStatus(t, h.sender, id)
}
for id, payload := range map[string]proto.AssignmentBindPayload{
"no grant": {EnvironmentID: environment, Resource: resource},
"no resource": {EnvironmentID: environment, AttachGrant: []byte("grant")},
"other environment": {EnvironmentID: environment, Resource: &other, AttachGrant: []byte("grant")},
} {
if got := bind(id, payload); got.ErrorCode != "invalid_request" {
t.Fatalf("%s: bind = %+v", id, got)
}
}
link := proto.AssignmentBindPayload{EnvironmentID: environment, Resource: resource, AttachGrant: []byte("grant")}
if got := bind("bind", link); got.State != proto.AssignmentBound {
t.Fatalf("bind = %+v", got)
}
if got := bind("again", link); got.State != proto.AssignmentBound {
t.Fatalf("identical bind = %+v", got)
}
changed := link
changed.AttachGrant = []byte("other grant")
if got := bind("changed", changed); got.ErrorCode != proto.AssignmentConflict {
t.Fatalf("bind with another grant = %+v", got)
}
}
4 changes: 2 additions & 2 deletions docs/runtime-protocol.md
Original file line number Diff line number Diff line change
Expand Up @@ -118,9 +118,9 @@ An assignment binds one Session to the Runtime that runs it. `Envelope.assignmen

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.

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`.
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`. When the Runtime is an agent host and the Environment has a live [Link](./sandbox-link-protocol.md) resource, the bind also carries `resource`, that resource as the [bootstrap input](./sandbox-bootstrap.md#launch-input) names it, and `attach_grant`, the base64 grant with which the agent host opens services on that resource generation under this assignment and epoch. The grant is secret. Core sends neither field to any other Runtime. A bind with only one of them, or with a resource of another Environment, fails with `invalid_request`. A repeated bind of the same assignment with the same Environment, resource and grant is `bound` again; any other bind of it fails with `assignment_conflict`. 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`.

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; a release that fails backs off, and the release due longest goes first, so failing releases cannot delay the rest. A quiesced Runtime admits only a release and the matching `environment_resume`, which carries the assignment that quiesced it.
Core records a release and advances the epoch before it sends anything, which withdraws the assignment's attach grant; it then has the relay revoke the Environment's Link resource at its current generation, so the attachments opened under the grant close before the Runtime receives the release. 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; a release that fails backs off, and the release due longest goes first, so failing releases cannot delay the rest. A quiesced Runtime admits only a release and the matching `environment_resume`, which carries the assignment that quiesced it.

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.

Expand Down
6 changes: 3 additions & 3 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: 5f4fc0a408f12938550bb548f6b95c336316be2320f5cbb8c6801a823d88863e
source_hash: 3d5083634593391e5a3627961d54e1dceb7c4589f74539b3e01b677a9efbe04d
---

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

每个 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 不携带分配。

在一条连接上执行 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` 失败。
在一条连接上执行 Session 的第一个操作之前,包括没有 Turn 的 Environment 初始化和文件操作,Core 发送带 Session 的 Environment ID 的 `assignment_bind`,并等待 `assignment_status` `bound`。当 Runtime 是 agent host 且 Environment 有存活的 [Link](./sandbox-link-protocol.md) resource 时,绑定还携带 `resource` 和 `attach_grant`:前者是该 resource,形式与[引导输入](./sandbox-bootstrap.md#launch-input)中的相同;后者是 base64 编码的 grant,agent host 凭它在此分配和 epoch 下打开该 resource generation 上的服务。grant 是机密。Core 不向其他任何 Runtime 发送这两个字段。只带其中一个字段、或带其他 Environment 的 resource 的绑定以 `invalid_request` 失败。以相同的 Environment、resource 和 grant 重复绑定同一分配仍得到 `bound`;该分配的其他绑定以 `assignment_conflict` 失败。Runtime 只在其已绑定的分配下准入 Session frame:较旧的 epoch 或已释放的分配以 `assignment_stale` 失败;其他分配、Session 或 Environment 以 `assignment_conflict` 失败。已启动 Run 的 frame,包括其取消回执,在释放前仍可在启动它的分配下准入。Runtime 已记录回执的重复函数结果或决策只在应用它的分配下得到回答;其他分配以 `assignment_conflict` 失败。

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 的分配。
Core 先记录释放并推进 epoch,再发送任何消息;记录即撤回该分配的 attach grant。随后 Core 让 relay 吊销 Environment 的 Link resource 的当前 generation,使凭该 grant 打开的 attachment 在 Runtime 收到释放之前关闭。删除 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 的分配。

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

Expand Down
29 changes: 27 additions & 2 deletions internal/agentdaemon/proto/assignment.go
Original file line number Diff line number Diff line change
@@ -1,5 +1,12 @@
package proto

import (
"errors"

"github.com/MiniMax-AI/OpenAgentCore/internal/sandboxbootstrap"
"github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink"
)

// Assignment traffic binds a Session to the Runtime that runs it. Envelope.ID
// correlates a request with its assignment_status, and Envelope.Assignment
// names the assignment on both.
Expand Down Expand Up @@ -53,9 +60,27 @@ func (r AssignmentRef) Valid() bool {
}

// AssignmentBindPayload is what the Runtime fences the Session's frames with.
// EnvironmentID is empty for a Session without an Environment.
// EnvironmentID is empty for a Session without an Environment. Resource and
// AttachGrant come together, only to an agent host whose Environment has a
// Link resource: the agent host attaches to Resource with AttachGrant, which
// is secret.
type AssignmentBindPayload struct {
EnvironmentID string `json:"environment_id,omitempty"`
EnvironmentID string `json:"environment_id,omitempty"`
Resource *sandboxbootstrap.Resource `json:"resource,omitempty"`
AttachGrant []byte `json:"attach_grant,omitempty"`
}

// Validate checks that Resource and AttachGrant come together and that
// Resource is a resource of the Environment.
func (p AssignmentBindPayload) Validate() error {
if p.Resource == nil && len(p.AttachGrant) == 0 {
return nil
}
if p.Resource == nil || len(p.AttachGrant) == 0 || len(p.AttachGrant) > sandboxlink.MaxGrantBytes ||
p.Resource.Validate() != nil || p.Resource.EnvironmentID != p.EnvironmentID {
return errors.New("assignment_bind requires a valid resource of its Environment with an attach grant")
}
return nil
}

// AssignmentReleasePayload asks the Runtime to remove the Session's native
Expand Down
5 changes: 3 additions & 2 deletions internal/sandboxbootstrap/bootstrap.go
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,7 @@ func (in Input) Validate() error {
if in.Version != Version || sandboxlink.CheckRelayURL(in.LinkURL) != nil ||
in.Credential == "" || len(in.Credential) > sandboxlink.MaxCredentialBytes ||
strings.ContainsFunc(in.Credential, func(r rune) bool { return r == 0 || unicode.IsSpace(r) }) ||
in.Resource.validate() != nil {
in.Resource.Validate() != nil {
return ErrInvalid
}
raw, err := json.Marshal(in)
Expand All @@ -58,7 +58,8 @@ func (in Input) Validate() error {
return nil
}

func (r Resource) validate() error {
// Validate checks that the resource names one Link resource.
func (r Resource) Validate() error {
for _, id := range []string{r.TenantID, r.EnvironmentID, r.ID} {
if u, err := uuid.Parse(id); err != nil || u == uuid.Nil || u.String() != id {
return ErrInvalid
Expand Down
Loading
Loading