diff --git a/docs/runtime-protocol.md b/docs/runtime-protocol.md index 397c0e1c3..6ed8f348e 100644 --- a/docs/runtime-protocol.md +++ b/docs/runtime-protocol.md @@ -110,7 +110,7 @@ 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` and `function_result`; 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`. 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 and epoch with the same Environment, resource and grant is `bound` again; one with anything else fails with `assignment_conflict`. A bind of the bound assignment at a higher epoch supersedes the earlier epoch, with any resource or grant: the Runtime fences and cleans up the earlier epoch's work as a release does, keeping the home, then binds the new epoch and replies `bound`. A Session's Environment never changes, so such a bind that names another Environment fails with `assignment_conflict` before anything is fenced. Unfinished cleanup replies `failed` with `cleanup_unconfirmed`, and a retry at the same epoch repeats it; a release or a later bind in the meantime fails it with `assignment_stale`. 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. When the Environment has no live resource, Core sends the agent host no bind, and the operation that needed the bind fails. 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 and epoch with the same Environment, resource and grant is `bound` again; one with anything else fails with `assignment_conflict`. A bind of the bound assignment at a higher epoch supersedes the earlier epoch, with any resource or grant: the Runtime fences and cleans up the earlier epoch's work as a release does, keeping the home, then binds the new epoch and replies `bound`. A Session's Environment never changes, so such a bind that names another Environment fails with `assignment_conflict` before anything is fenced. Unfinished cleanup replies `failed` with `cleanup_unconfirmed`, and a retry at the same epoch repeats it; a release or a later bind in the meantime fails it with `assignment_stale`. 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`. A Session's first bind resolves its Environment owner, which holds the Environment's resources and performs every effect on them; it does not change afterwards. The owner checks each `execution_prepare` configuration against the Environment, including the read-only profile, and fills the installed capabilities before an Executor starts. It applies `runtime_prepare`, lists directories for `workspace_read`, writes files for `workspace_write` and exports outputs for `workspace_export`. The Runtime's dispatcher keeps admission, transfer framing and fencing, and never substitutes another implementation. A self-hosted Runtime's owner is its bound local workspace, which outlives each assignment; the Runtime rejects the bind of any other Session with `assignment_conflict`. An agent host's owner works in the Session's sandbox through File and Process on a [Link](./sandbox-link-protocol.md) attachment of its own, opened under the bind's grant on first use. It lasts from the Session's first bind until its home is removed and outlives the Session's Executors and connections. Quiescing its Environment, releasing the assignment or superseding its epoch closes its attachment; a bind under a later assignment or epoch takes the owner over once that attachment is closed, and a bind it cannot take over, including one of another Environment, fails with `assignment_conflict`. It runs each setup step as the Process operation whose ID is the `runtime_prepare` envelope ID. A `workspace_write` or `runtime_prepare` that cannot reach the sandbox before any effect ends `rejected` with `resource_unavailable`. It fails a Plugin whose MCP server declares literal `http_headers`, or is a stdio server with `env_vars`, before staging any of it. A file mutation or setup step whose outcome it cannot observe quarantines the owner until the home is removed: it sends no further mutation, and every later `workspace_write` and `runtime_prepare` ends `unknown`. A Session without an owner supports none of these operations, and the Runtime rejects each with its typed code: `unsupported_read_preparation` for a read-only preparation, `invalid_configuration` for an Executor configuration with a `local_environment`, `runtime_preparation_unsupported` for `runtime_prepare`, `write_unsupported` for `workspace_write`, and `read_unsupported` for `workspace_read` and `workspace_export`. diff --git a/docs/zh/runtime-protocol.md b/docs/zh/runtime-protocol.md index d5b39a551..4a030e550 100644 --- a/docs/zh/runtime-protocol.md +++ b/docs/zh/runtime-protocol.md @@ -1,7 +1,7 @@ --- title: "Core–Runtime 协议" source: docs/runtime-protocol.md -source_hash: 9ec9470b3e5cd72b1a81f83fc739a61afe2e2abd424aa646dfa366c067be1148 +source_hash: 103462f4f02379eeb587294e3124e366ede0743cc174d335f9b5355e37d84c62 --- 此协议在 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)。 @@ -112,7 +112,7 @@ Usage frame 和最终 usage snapshot 都携带当前执行的累计测量,替 每个 Session frame 都携带分配:`execution_prepare`、`execution_start` 和 `execution_release`;`prompt_cancel`、`prompt_steer` 和 `function_result`;`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`。当 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` 失败。以相同的 epoch、Environment、resource 和 grant 重复绑定同一分配仍得到 `bound`;其他字段不同的同 epoch 绑定以 `assignment_conflict` 失败。以更高 epoch 绑定已绑定的分配会取代较早的 epoch,resource 和 grant 均可不同:Runtime 像释放那样 fence 并清理较早 epoch 的工作,但保留 home,然后绑定新 epoch 并回复 `bound`。Session 的 Environment 从不改变,因此这样的绑定若指定其他 Environment,会在任何 fence 之前以 `assignment_conflict` 失败。清理未完成时回复 `failed` 和 `cleanup_unconfirmed`,以同一 epoch 重试会重复清理;期间到达的释放或更晚的绑定使其以 `assignment_stale` 失败。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 上的服务。若 Environment 没有存活的 resource,Core 不向 agent host 发送绑定,需要该绑定的操作失败。grant 是机密。Core 不向其他任何 Runtime 发送这两个字段。只带其中一个字段、或带其他 Environment 的 resource 的绑定以 `invalid_request` 失败。以相同的 epoch、Environment、resource 和 grant 重复绑定同一分配仍得到 `bound`;其他字段不同的同 epoch 绑定以 `assignment_conflict` 失败。以更高 epoch 绑定已绑定的分配会取代较早的 epoch,resource 和 grant 均可不同:Runtime 像释放那样 fence 并清理较早 epoch 的工作,但保留 home,然后绑定新 epoch 并回复 `bound`。Session 的 Environment 从不改变,因此这样的绑定若指定其他 Environment,会在任何 fence 之前以 `assignment_conflict` 失败。清理未完成时回复 `failed` 和 `cleanup_unconfirmed`,以同一 epoch 重试会重复清理;期间到达的释放或更晚的绑定使其以 `assignment_stale` 失败。Runtime 只在其已绑定的分配下准入 Session frame:较旧的 epoch 或已释放的分配以 `assignment_stale` 失败;其他分配、Session 或 Environment 以 `assignment_conflict` 失败。已启动 Run 的 frame,包括其取消回执,在释放前仍可在启动它的分配下准入。Runtime 已记录回执的重复函数结果或决策只在应用它的分配下得到回答;其他分配以 `assignment_conflict` 失败。 Session 的第一次绑定确定其 Environment owner,此后不再改变;owner 持有 Environment 的资源,并执行对这些资源的每个作用。owner 根据 Environment 检查每个 `execution_prepare` 配置(包括只读 profile),并在 Executor 启动前填入已安装的能力。它应用 `runtime_prepare`,为 `workspace_read` 列举目录,为 `workspace_write` 写入文件,为 `workspace_export` 导出输出。Runtime 的 dispatcher 保留准入、传输分帧和 fencing,从不替换为其他实现。self-hosted Runtime 的 owner 是其绑定的本地工作区,该工作区比每个分配存续得更久;Runtime 以 `assignment_conflict` 拒绝任何其他 Session 的绑定。agent host 的 owner 通过自己的一个 [Link](./sandbox-link-protocol.md) attachment,用 File 和 Process 在 Session 的沙箱中工作;该 attachment 在首次使用时凭绑定的 grant 打开。owner 从 Session 的第一次绑定存续到其 home 被删除,比 Session 的 Executor 和连接存续得更久。其 Environment 被 quiesce、分配被释放或其 epoch 被取代时关闭其 attachment;之后的分配或 epoch 下的绑定在该 attachment 关闭后接管 owner,它无法接管的绑定(包括其他 Environment 的绑定)以 `assignment_conflict` 失败。它把每个 setup 步骤作为 Process 操作运行,操作 ID 即 `runtime_prepare` 的 envelope ID。在产生任何作用前无法连到沙箱的 `workspace_write` 或 `runtime_prepare` 以 `rejected` 和 `resource_unavailable` 结束。若 Plugin 的 MCP server 声明了字面量 `http_headers`,或是带 `env_vars` 的 stdio server,owner 会在暂存其任何内容之前使其失败。无法观察到结果的文件变更或 setup 步骤会隔离 owner,直到 home 被删除:它不再发送任何变更,之后每个 `workspace_write` 和 `runtime_prepare` 都以 `unknown` 结束。没有 owner 的 Session 不支持上述任何操作,Runtime 以各自的类型化错误码拒绝:只读 preparation 为 `unsupported_read_preparation`;带 `local_environment` 的 Executor 配置为 `invalid_configuration`;`runtime_prepare` 为 `runtime_preparation_unsupported`;`workspace_write` 为 `write_unsupported`;`workspace_read` 和 `workspace_export` 为 `read_unsupported`。 diff --git a/services/core/internal/db/queries/devices.sql b/services/core/internal/db/queries/devices.sql index 47d2b5ca4..78d99e1d0 100644 --- a/services/core/internal/db/queries/devices.sql +++ b/services/core/internal/db/queries/devices.sql @@ -40,26 +40,28 @@ WHERE session_runtime_assignments.runtime_id = EXCLUDED.runtime_id AND session_r RETURNING runtime_id; -- name: GetSessionDevice :one -SELECT d.id, d.name, d.environment_id, b.assignment_id, b.epoch FROM session_runtime_assignments b +-- The Session's bound Runtime: a device of its tenant or the deployment's +-- agent host. environment_id is the device's own Environment and +-- session_environment_id the Session's, which the assignment binds. +SELECT d.id, d.name, d.environment_id, e.id AS session_environment_id, b.assignment_id, b.epoch FROM session_runtime_assignments b JOIN sessions s ON s.id = b.session_id -JOIN devices d ON d.id = b.runtime_id AND d.tenant_id = s.tenant_id +JOIN devices d ON d.id = b.runtime_id AND (d.tenant_id = s.tenant_id OR (d.agent_host AND d.tenant_id IS NULL)) +LEFT JOIN environments e ON e.session_id = s.id WHERE s.tenant_id = $1 AND s.id = $2 AND d.revoked_at IS NULL AND b.desired_state = 'bound' AND EXISTS (SELECT 1 FROM runtime_device_authority a WHERE a.id = d.id) -AND (d.environment_id IS NULL OR EXISTS ( - SELECT 1 FROM environments e WHERE e.id = d.environment_id AND e.session_id = s.id -)); +AND (d.environment_id IS NULL OR d.environment_id = e.id); -- name: GetSessionExecutionBinding :one -SELECT d.id, d.name, b.native_session_id, d.environment_id, b.assignment_id, b.epoch, +-- The bound Runtime as GetSessionDevice reads it, with the native session. +SELECT d.id, d.name, b.native_session_id, d.environment_id, e.id AS session_environment_id, b.assignment_id, b.epoch, EXISTS (SELECT 1 FROM turns t WHERE t.session_id = s.id AND t.started_at IS NOT NULL) AS has_started_turn FROM session_runtime_assignments b JOIN sessions s ON s.id = b.session_id -JOIN devices d ON d.id = b.runtime_id AND d.tenant_id = s.tenant_id +JOIN devices d ON d.id = b.runtime_id AND (d.tenant_id = s.tenant_id OR (d.agent_host AND d.tenant_id IS NULL)) +LEFT JOIN environments e ON e.session_id = s.id WHERE s.tenant_id = $1 AND s.id = $2 AND d.revoked_at IS NULL AND b.desired_state = 'bound' AND EXISTS (SELECT 1 FROM runtime_device_authority a WHERE a.id = d.id) -AND (d.environment_id IS NULL OR EXISTS ( - SELECT 1 FROM environments e WHERE e.id = d.environment_id AND e.session_id = s.id -)); +AND (d.environment_id IS NULL OR d.environment_id = e.id); -- name: RememberNativeSession :execrows UPDATE session_runtime_assignments SET native_session_id = $2 WHERE session_id = $1; diff --git a/services/core/internal/db/sqlc/devices.sql.go b/services/core/internal/db/sqlc/devices.sql.go index 8dce11cf0..dcdd2e9b0 100644 --- a/services/core/internal/db/sqlc/devices.sql.go +++ b/services/core/internal/db/sqlc/devices.sql.go @@ -129,14 +129,13 @@ func (q *Queries) GetDeviceCredential(ctx context.Context, id pgtype.UUID) (GetD } const getSessionDevice = `-- name: GetSessionDevice :one -SELECT d.id, d.name, d.environment_id, b.assignment_id, b.epoch FROM session_runtime_assignments b +SELECT d.id, d.name, d.environment_id, e.id AS session_environment_id, b.assignment_id, b.epoch FROM session_runtime_assignments b JOIN sessions s ON s.id = b.session_id -JOIN devices d ON d.id = b.runtime_id AND d.tenant_id = s.tenant_id +JOIN devices d ON d.id = b.runtime_id AND (d.tenant_id = s.tenant_id OR (d.agent_host AND d.tenant_id IS NULL)) +LEFT JOIN environments e ON e.session_id = s.id WHERE s.tenant_id = $1 AND s.id = $2 AND d.revoked_at IS NULL AND b.desired_state = 'bound' AND EXISTS (SELECT 1 FROM runtime_device_authority a WHERE a.id = d.id) -AND (d.environment_id IS NULL OR EXISTS ( - SELECT 1 FROM environments e WHERE e.id = d.environment_id AND e.session_id = s.id -)) +AND (d.environment_id IS NULL OR d.environment_id = e.id) ` type GetSessionDeviceParams struct { @@ -145,13 +144,17 @@ type GetSessionDeviceParams struct { } type GetSessionDeviceRow struct { - ID pgtype.UUID `json:"id"` - Name string `json:"name"` - EnvironmentID pgtype.UUID `json:"environment_id"` - AssignmentID pgtype.UUID `json:"assignment_id"` - Epoch int64 `json:"epoch"` + ID pgtype.UUID `json:"id"` + Name string `json:"name"` + EnvironmentID pgtype.UUID `json:"environment_id"` + SessionEnvironmentID pgtype.UUID `json:"session_environment_id"` + AssignmentID pgtype.UUID `json:"assignment_id"` + Epoch int64 `json:"epoch"` } +// The Session's bound Runtime: a device of its tenant or the deployment's +// agent host. environment_id is the device's own Environment and +// session_environment_id the Session's, which the assignment binds. func (q *Queries) GetSessionDevice(ctx context.Context, arg GetSessionDeviceParams) (GetSessionDeviceRow, error) { row := q.db.QueryRow(ctx, getSessionDevice, arg.TenantID, arg.ID) var i GetSessionDeviceRow @@ -159,6 +162,7 @@ func (q *Queries) GetSessionDevice(ctx context.Context, arg GetSessionDevicePara &i.ID, &i.Name, &i.EnvironmentID, + &i.SessionEnvironmentID, &i.AssignmentID, &i.Epoch, ) @@ -166,16 +170,15 @@ func (q *Queries) GetSessionDevice(ctx context.Context, arg GetSessionDevicePara } const getSessionExecutionBinding = `-- name: GetSessionExecutionBinding :one -SELECT d.id, d.name, b.native_session_id, d.environment_id, b.assignment_id, b.epoch, +SELECT d.id, d.name, b.native_session_id, d.environment_id, e.id AS session_environment_id, b.assignment_id, b.epoch, EXISTS (SELECT 1 FROM turns t WHERE t.session_id = s.id AND t.started_at IS NOT NULL) AS has_started_turn FROM session_runtime_assignments b JOIN sessions s ON s.id = b.session_id -JOIN devices d ON d.id = b.runtime_id AND d.tenant_id = s.tenant_id +JOIN devices d ON d.id = b.runtime_id AND (d.tenant_id = s.tenant_id OR (d.agent_host AND d.tenant_id IS NULL)) +LEFT JOIN environments e ON e.session_id = s.id WHERE s.tenant_id = $1 AND s.id = $2 AND d.revoked_at IS NULL AND b.desired_state = 'bound' AND EXISTS (SELECT 1 FROM runtime_device_authority a WHERE a.id = d.id) -AND (d.environment_id IS NULL OR EXISTS ( - SELECT 1 FROM environments e WHERE e.id = d.environment_id AND e.session_id = s.id -)) +AND (d.environment_id IS NULL OR d.environment_id = e.id) ` type GetSessionExecutionBindingParams struct { @@ -184,15 +187,17 @@ type GetSessionExecutionBindingParams struct { } type GetSessionExecutionBindingRow struct { - ID pgtype.UUID `json:"id"` - Name string `json:"name"` - NativeSessionID string `json:"native_session_id"` - EnvironmentID pgtype.UUID `json:"environment_id"` - AssignmentID pgtype.UUID `json:"assignment_id"` - Epoch int64 `json:"epoch"` - HasStartedTurn bool `json:"has_started_turn"` + ID pgtype.UUID `json:"id"` + Name string `json:"name"` + NativeSessionID string `json:"native_session_id"` + EnvironmentID pgtype.UUID `json:"environment_id"` + SessionEnvironmentID pgtype.UUID `json:"session_environment_id"` + AssignmentID pgtype.UUID `json:"assignment_id"` + Epoch int64 `json:"epoch"` + HasStartedTurn bool `json:"has_started_turn"` } +// The bound Runtime as GetSessionDevice reads it, with the native session. func (q *Queries) GetSessionExecutionBinding(ctx context.Context, arg GetSessionExecutionBindingParams) (GetSessionExecutionBindingRow, error) { row := q.db.QueryRow(ctx, getSessionExecutionBinding, arg.TenantID, arg.ID) var i GetSessionExecutionBindingRow @@ -201,6 +206,7 @@ func (q *Queries) GetSessionExecutionBinding(ctx context.Context, arg GetSession &i.Name, &i.NativeSessionID, &i.EnvironmentID, + &i.SessionEnvironmentID, &i.AssignmentID, &i.Epoch, &i.HasStartedTurn, diff --git a/services/core/internal/execution/device_authority.go b/services/core/internal/execution/device_authority.go index 5600e383c..4ea9eaca3 100644 --- a/services/core/internal/execution/device_authority.go +++ b/services/core/internal/execution/device_authority.go @@ -43,7 +43,7 @@ func assignedRuntimePeer(ctx context.Context, devices sessions.DeviceReader, reg if err != nil { return nil, err } - if err := peer.Bind(ctx, bound.Assignment, bound.EnvironmentID); err != nil { + if err := peer.Bind(ctx, bound.Assignment, bound.SessionEnvironmentID); err != nil { return nil, err } return peer, nil diff --git a/services/core/internal/execution/dispatcher.go b/services/core/internal/execution/dispatcher.go index 8dfda86ff..512a57ac5 100644 --- a/services/core/internal/execution/dispatcher.go +++ b/services/core/internal/execution/dispatcher.go @@ -108,7 +108,7 @@ func (d *Dispatcher) Run(ctx context.Context, tenantID, sessionID, turnID string return sessions.Turn{}, err } req.DisableExecutionEnvironment = true - if err := peer.Bind(ctx, bound.Device.Assignment, bound.Device.EnvironmentID); err != nil { + if err := peer.Bind(ctx, bound.Device.Assignment, bound.Device.SessionEnvironmentID); err != nil { return sessions.Turn{}, err } prepared, err := d.prepareTurnExecutor(ctx, peer, bound.Device.Assignment, tenantID, sessionID, turnID, req, sessions.TurnQueued) diff --git a/services/core/internal/execution/prepared_dispatch.go b/services/core/internal/execution/prepared_dispatch.go index 79943b272..6faa49132 100644 --- a/services/core/internal/execution/prepared_dispatch.go +++ b/services/core/internal/execution/prepared_dispatch.go @@ -73,7 +73,7 @@ func (d *Dispatcher) RunEnvironmentInput(ctx context.Context, lease Ownership, t if err := d.configurePreparedEnvironment(session, environment, bound.Device, &req); err != nil { return run, err } - if err := peer.Bind(owner, bound.Device.Assignment, bound.Device.EnvironmentID); err != nil { + if err := peer.Bind(owner, bound.Device.Assignment, bound.Device.SessionEnvironmentID); err != nil { return run, err } prepared, err := newPreparedStart(peer, bound.Device.Assignment) diff --git a/services/core/internal/execution/runtime_initialization.go b/services/core/internal/execution/runtime_initialization.go index 85dc81fbb..e619198dd 100644 --- a/services/core/internal/execution/runtime_initialization.go +++ b/services/core/internal/execution/runtime_initialization.go @@ -122,7 +122,7 @@ func (w *Worker) prepareEnvironment(ctx context.Context, owner sessions.Environm if err != nil { return err } - peer, err := w.dispatcher.assignedPeer(ctx, sessions.ExecutionDevice{ID: owner.DeviceID, EnvironmentID: owner.EnvironmentID, Assignment: owner.Assignment}) + peer, err := w.dispatcher.assignedPeer(ctx, sessions.ExecutionDevice{ID: owner.DeviceID, Assignment: owner.Assignment, SessionEnvironmentID: owner.EnvironmentID}) if err != nil { return err } diff --git a/services/core/internal/persistence/postgres/sessionpg/devices.go b/services/core/internal/persistence/postgres/sessionpg/devices.go index f14e32ae4..c5e23051a 100644 --- a/services/core/internal/persistence/postgres/sessionpg/devices.go +++ b/services/core/internal/persistence/postgres/sessionpg/devices.go @@ -45,7 +45,7 @@ func (s *Store) GetSessionExecutionBinding(ctx context.Context, tenant, session } return sessions.ExecutionBinding{ Device: sessions.ExecutionDevice{ID: uuid.UUID(row.ID.Bytes).String(), Name: row.Name, EnvironmentID: optionalID(row.EnvironmentID), - Assignment: assignmentRef(lookup.ID, row.AssignmentID, row.Epoch)}, + Assignment: assignmentRef(lookup.ID, row.AssignmentID, row.Epoch), SessionEnvironmentID: optionalID(row.SessionEnvironmentID)}, NativeSessionID: row.NativeSessionID, HasStartedTurn: row.HasStartedTurn, }, nil @@ -101,7 +101,7 @@ func loadSessionDevice(ctx context.Context, q *sqlc.Queries, tenant, session pgt return sessions.ExecutionDevice{}, false, err } return sessions.ExecutionDevice{ID: uuid.UUID(row.ID.Bytes).String(), Name: row.Name, EnvironmentID: optionalID(row.EnvironmentID), - Assignment: assignmentRef(session, row.AssignmentID, row.Epoch)}, true, nil + Assignment: assignmentRef(session, row.AssignmentID, row.Epoch), SessionEnvironmentID: optionalID(row.SessionEnvironmentID)}, true, nil } // assignmentRef is the reference of a Session's assignment; an absent diff --git a/services/core/internal/runtimegateway/assignment.go b/services/core/internal/runtimegateway/assignment.go index 963c720cb..0ce761ca0 100644 --- a/services/core/internal/runtimegateway/assignment.go +++ b/services/core/internal/runtimegateway/assignment.go @@ -15,7 +15,8 @@ import ( // assignment_bind and waits for bound, once per connection and reference. // Every Session operation on the connection follows its Bind. A bind of an // Environment with a live Link resource to an agent host carries the -// resource and its attach grant. +// resource and its attach grant; without one, it sends nothing and fails with +// ErrNoLinkResource. func (s *Session) Bind(ctx context.Context, ref proto.AssignmentRef, environmentID string) error { s.assignmentMu.Lock() bound := s.assignments[ref.SessionID] == ref diff --git a/services/core/internal/runtimegateway/link.go b/services/core/internal/runtimegateway/link.go index 37481b580..3c2a9210d 100644 --- a/services/core/internal/runtimegateway/link.go +++ b/services/core/internal/runtimegateway/link.go @@ -6,6 +6,7 @@ import ( "crypto/sha256" "crypto/subtle" "encoding/binary" + "errors" "fmt" "net/netip" "time" @@ -160,18 +161,26 @@ func (l *LinkAuthority) authorize(ctx context.Context, peer sandboxlink.AttachPe return sandboxlink.Identity{Resource: resource, SessionID: sandboxwire.ID(session), AssignmentID: id, AssignmentEpoch: epoch}, assignment, nil } +// ErrNoLinkResource fails an agent host's bind to an Environment without a live +// Link resource: the agent host reaches an Environment only through one. +var ErrNoLinkResource = errors.New("agentdaemon gateway: the Environment has no live Link resource") + // bindLink returns the Link resource and attach grant that runtimeID's bind of // ref in environmentID carries: none unless runtimeID is the agent host that -// holds ref bound and environmentID has a live Link resource. +// holds ref bound, and then environmentID's live Link resource, without which +// it is ErrNoLinkResource. func (l *LinkAuthority) bindLink(ctx context.Context, runtimeID string, ref proto.AssignmentRef, environmentID string) (*sandboxbootstrap.Resource, []byte, error) { if environmentID == "" { return nil, nil, nil } assignment, found, err := l.store.GetLinkAssignment(ctx, ref.AssignmentID) - if err != nil || !found || !assignment.AgentHost || !assignment.Bound || assignment.Resource.EnvironmentID != environmentID || + if err != nil || !found || !assignment.AgentHost || !assignment.Bound || assignment.RuntimeID != runtimeID || assignment.SessionID != ref.SessionID || assignment.Epoch != ref.Epoch { return nil, nil, err } + if assignment.Resource.EnvironmentID != environmentID { + return nil, nil, ErrNoLinkResource + } id, err := uuid.Parse(ref.AssignmentID) if err != nil { return nil, nil, err diff --git a/services/core/internal/sessions/execution_environment.go b/services/core/internal/sessions/execution_environment.go index 8d277d460..3392c0ac9 100644 --- a/services/core/internal/sessions/execution_environment.go +++ b/services/core/internal/sessions/execution_environment.go @@ -140,14 +140,14 @@ func checkInitializationOwner(owner EnvironmentInitialization, environment Envir return nil } -// checkInitializationDevice confirms that the device the initialization was -// listed with is still the Session's device for this Environment. A Session -// without a device is ErrNotFound and another binding ErrTurnConflict. -func checkInitializationDevice(owner EnvironmentInitialization, environment Environment, bound ExecutionDevice, found bool) error { +// checkInitializationDevice confirms that the Runtime the initialization was +// listed with is still the one the Session is bound to. A Session without a +// bound Runtime is ErrNotFound and another binding ErrTurnConflict. +func checkInitializationDevice(owner EnvironmentInitialization, bound ExecutionDevice, found bool) error { if !found { return ErrNotFound } - if bound.ID != owner.DeviceID || bound.EnvironmentID != environment.ID { + if bound.ID != owner.DeviceID { return ErrTurnConflict } return nil @@ -180,7 +180,7 @@ func (o *ExecutionOperations) advanceInitialization(ctx context.Context, owner E if err != nil { return err } - if err := checkInitializationDevice(owner, environment, bound, found); err != nil { + if err := checkInitializationDevice(owner, bound, found); err != nil { return err } applied, err := transition(ctx, tx, environment) diff --git a/services/core/internal/sessions/execution_environment_test.go b/services/core/internal/sessions/execution_environment_test.go index adb4ab649..fb1e66c5d 100644 --- a/services/core/internal/sessions/execution_environment_test.go +++ b/services/core/internal/sessions/execution_environment_test.go @@ -156,19 +156,17 @@ func TestInitializationOwnerAndDeviceRules(t *testing.T) { t.Fatalf("owner %s: %v", test.name, err) } } - environment := Environment{ID: "environment"} for _, test := range []struct { name string bound ExecutionDevice found bool want error }{ - {"listed device", ExecutionDevice{ID: "device", EnvironmentID: "environment"}, true, nil}, + {"listed device", ExecutionDevice{ID: "device"}, true, nil}, {"no device", ExecutionDevice{}, false, ErrNotFound}, - {"another device", ExecutionDevice{ID: "other", EnvironmentID: "environment"}, true, ErrTurnConflict}, - {"another Environment's device", ExecutionDevice{ID: "device", EnvironmentID: "other"}, true, ErrTurnConflict}, + {"another device", ExecutionDevice{ID: "other"}, true, ErrTurnConflict}, } { - if err := checkInitializationDevice(owner, environment, test.bound, test.found); !errors.Is(err, test.want) { + if err := checkInitializationDevice(owner, test.bound, test.found); !errors.Is(err, test.want) { t.Fatalf("device %s: %v", test.name, err) } } diff --git a/services/core/internal/sessions/turn.go b/services/core/internal/sessions/turn.go index c4323dacb..6d9e76e77 100644 --- a/services/core/internal/sessions/turn.go +++ b/services/core/internal/sessions/turn.go @@ -84,12 +84,16 @@ type ExecutionEvent struct { type ExecutionWork struct{ TenantID, SessionID, TurnID, Status string } // ExecutionDevice contains safe identity only, never a device credential. -// Assignment is the Session's bound assignment to the device. +// EnvironmentID is the device's own Environment, empty for a device of none. +// Assignment is the Session's bound assignment to the device, and +// SessionEnvironmentID the Session's Environment, which that assignment binds; +// it is empty for a Session without one. type ExecutionDevice struct { - ID string - Name string - EnvironmentID string - Assignment proto.AssignmentRef + ID string + Name string + EnvironmentID string + Assignment proto.AssignmentRef + SessionEnvironmentID string } // ExecutionBinding identifies the Runtime and native history selected for one API Session. diff --git a/services/core/tests/integration/link_authority_test.go b/services/core/tests/integration/link_authority_test.go index 4a514dad6..c25e1c117 100644 --- a/services/core/tests/integration/link_authority_test.go +++ b/services/core/tests/integration/link_authority_test.go @@ -27,8 +27,10 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/api" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/environmentconfig" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/processconfig" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtime" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimedevice" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" @@ -48,7 +50,8 @@ type linkHarness struct { } // newLinkHarness binds the Session to an unmarked operator device, or, for a -// guest, to the allocation's own device, as an in-sandbox daemon is bound. +// guest, to the allocation's own device, as an in-sandbox daemon is bound, and +// runs a Worker with the Link route's relay. func newLinkHarness(t *testing.T, guest bool) *linkHarness { t.Helper() h := newDispatchHarnessForSession(t, []byte(`{"agent":{"model":"test-model"},"environment":{"type":"openai_hosted","network":{"access":"disabled"}}}`), guest) @@ -62,11 +65,29 @@ func newLinkHarness(t *testing.T, guest bool) *linkHarness { } l := &linkHarness{dispatchHarness: h, relay: startLinkRoute(t, h.s), serve: []byte(uuid.NewString()), resource: sandboxbootstrap.Resource{TenantID: h.tenant, EnvironmentID: h.session.Environment.ID, Kind: "allocation", ID: uuid.NewString(), Generation: 1}} - if _, err := h.s.pool.Exec(t.Context(), `INSERT INTO runtime_allocations(id,environment_id,device_id,provider_key,state,create_settled,deployment_generation,serve_credential_hash) - VALUES($1,$2,$3,$4,'running',true,(SELECT generation FROM runtime_deployment),$5)`, l.resource.ID, l.resource.EnvironmentID, device, uuid.NewString(), serveHash(l.serve)); err != nil { + insertAllocation(t, h.s, l.resource, device, l.serve) + dispatcher := *h.d + dispatcher.Links = l.relay.Relay + runWorker(t, startWorker(t, t.Context(), h.s, &dispatcher)) + return l +} + +// insertAllocation inserts a running allocation of device that is resource +// and Serves with the credential. +func insertAllocation(t *testing.T, s *Store, resource sandboxbootstrap.Resource, device string, serve []byte) { + t.Helper() + if _, err := s.pool.Exec(t.Context(), `INSERT INTO runtime_allocations(id,environment_id,device_id,provider_key,state,create_settled,deployment_generation,serve_credential_hash) + VALUES($1,$2,$3,$4,'running',true,(SELECT generation FROM runtime_deployment),$5)`, resource.ID, resource.EnvironmentID, device, uuid.NewString(), serveHash(serve)); err != nil { t.Fatal(err) } - return l +} + +// runWorker runs w until the test ends. +func runWorker(t *testing.T, w *execution.Worker) { + ctx, cancel := context.WithCancel(t.Context()) + done := make(chan error, 1) + go func() { done <- w.Run(ctx) }() + t.Cleanup(func() { cancel(); <-done }) } // startLinkRoute serves Core's Link Authority at the Link route of the API @@ -107,23 +128,14 @@ func (l *linkHarness) exec(query string, args ...any) { } } -// bind binds the Session's current assignment to h.device through Core's -// gateway and returns it with the payload the Runtime received. +// bind has the Worker bind the Session's current assignment to h.device and +// returns it with the payload the Runtime received. Until Sessions are placed +// on an agent host, the Environment's initialization is its one production +// bind, so bind reopens the initialization, which has nothing to install. func (l *linkHarness) bind() (proto.AssignmentRef, proto.AssignmentBindPayload) { t := l.t t.Helper() - ref := proto.AssignmentRef{SessionID: l.session.ID} - var epoch int64 - if err := l.s.pool.QueryRow(t.Context(), "SELECT assignment_id::text, epoch FROM session_runtime_assignments WHERE session_id=$1", l.session.ID).Scan(&ref.AssignmentID, &epoch); err != nil { - t.Fatal(err) - } - ref.Epoch = uint64(epoch) - peer, err := l.registry.LookupDevice(l.device.ID) - if err != nil { - t.Fatal(err) - } - bound := make(chan error, 1) - go func() { bound <- peer.Bind(t.Context(), ref, l.session.Environment.ID) }() + l.exec("UPDATE environments SET initialization = 'pending' WHERE id = $1", l.session.Environment.ID) frame := l.read(proto.TypeAssignmentBind) var payload proto.AssignmentBindPayload if err := frame.DecodePayload(&payload); err != nil { @@ -131,15 +143,13 @@ func (l *linkHarness) bind() (proto.AssignmentRef, proto.AssignmentBindPayload) } reply, _ := assignmentReply(frame) l.writeMu.Lock() - err = l.conn.WriteJSON(reply) + err := l.conn.WriteJSON(reply) l.writeMu.Unlock() if err != nil { t.Fatal(err) } - if err := within(t, bound); err != nil { - t.Fatal(err) - } - return ref, payload + awaitInitialization(t, l.s, l.tenant, l.session.Environment.ID, "complete") + return frame.Assignment, payload } func (l *linkHarness) attach(runtime string, credential []byte) (*sandboxlink.AttachLink, error) { @@ -356,13 +366,6 @@ func TestLinkAuthorityReleaseRevokesBeforeSend(t *testing.T) { if _, err := l.open(link, sandboxlink.ServiceFile, ref, payload.AttachGrant); linkCode(err) != sandboxlink.ResourceNotFound { t.Fatal("a released assignment opened a service", err) } - dispatcher := *l.d - dispatcher.Links = l.relay.Relay - worker := startWorker(t, t.Context(), l.s, &dispatcher) - ctx, cancel := context.WithCancel(t.Context()) - done := make(chan error, 1) - go func() { done <- worker.Run(ctx) }() - t.Cleanup(func() { cancel(); <-done }) if release := l.read(proto.TypeAssignmentRelease); release.Assignment.AssignmentID != ref.AssignmentID || release.Assignment.Epoch != ref.Epoch+1 { t.Fatal("release of", release.Assignment) } @@ -537,3 +540,63 @@ func TestRegisteredAgentHostAuthenticates(t *testing.T) { t.Fatal("registration took over a tenant device") } } + +// TestInitializationBindsAgentHost binds a hosted Session to the registered +// agent host and runs the Environment's initialization. The bind carries the +// Environment's live Link resource and an attach grant; without a live +// resource, initialization fails before any bind is sent. +func TestInitializationBindsAgentHost(t *testing.T) { + for name, live := range map[string]bool{"live resource": true, "no live resource": false} { + t.Run(name, func(t *testing.T) { + s, _ := newManagedTestStore(t) + tenant, host, credential := uuid.NewString(), uuid.NewString(), uuid.NewString() + if err := sessionAdapter(s).RegisterAgentHost(t.Context(), host, runtimedevice.HashCredential(credential)); err != nil { + t.Fatal(err) + } + session, err := s.CreateSession(t.Context(), tenant, WithFixtureModelProvider(sessions.CreateSession{Creator: FixtureCreator(), Engine: "codex", IdempotencyKey: uuid.NewString(), + Configuration: json.RawMessage(`{"agent":{"model":"test-model"},"environment":{"type":"openai_hosted"}}`), + InitialFiles: []environmentconfig.InitialFile{{Type: "inline", Path: "/workspace/input", Data: []byte("frozen")}}})) + if err != nil { + t.Fatal(err) + } + // Placement does not choose the agent host yet. + if _, err := s.pool.Exec(t.Context(), "INSERT INTO session_runtime_assignments(session_id, runtime_id) VALUES($1, $2)", session.ID, host); err != nil { + t.Fatal(err) + } + resource := sandboxbootstrap.Resource{TenantID: tenant, EnvironmentID: session.Environment.ID, Kind: "allocation", ID: uuid.NewString(), Generation: 1} + if live { + device, err := sessionService(t, s).CreateDevice(t.Context(), tenant, "sandbox", runtimedevice.HashCredential(uuid.NewString())) + if err != nil { + t.Fatal(err) + } + insertAllocation(t, s, resource, device.ID, []byte(uuid.NewString())) + } + server := httptest.NewUnstartedServer(nil) + endpoint := "ws://" + server.Listener.Addr().String() + "/api/v1/agent-daemon/ws" + handler, registry, err := runtime.NewGateway(sessionAdapter(s), sessionService(t, s), sessionAdapter(s), runtimegateway.NewLinkAuthority(sessionAdapter(s)), endpoint) + if err != nil { + t.Fatal(err) + } + server.Config.Handler = handler + server.Start() + t.Cleanup(func() { server.Close(); runtime.CloseConnections(registry) }) + runWorker(t, startWorker(t, t.Context(), s, &execution.Dispatcher{Registry: registry})) + peer := &initializationPeer{apply: completedInitialization, binds: make(chan proto.AssignmentBindPayload, 1)} + peer.setRuntimeGateway(t, endpoint, registry) + if err := peer.connect(sandbox.Bootstrap{DeviceID: host, Credential: credential}); err != nil { + t.Fatal(err) + } + if !live { + awaitInitialization(t, s, tenant, session.Environment.ID, "failed") + if len(peer.binds) != 0 { + t.Fatal("bound without a live resource", <-peer.binds) + } + return + } + if bind := within(t, peer.binds); bind.EnvironmentID != session.Environment.ID || bind.Resource == nil || *bind.Resource != resource || len(bind.AttachGrant) == 0 { + t.Fatalf("agent host bind = %+v", bind) + } + awaitInitialization(t, s, tenant, session.Environment.ID, "complete") + }) + } +} diff --git a/services/core/tests/integration/runtime_initialization_peer_test.go b/services/core/tests/integration/runtime_initialization_peer_test.go index 35e6163dd..a4ed781da 100644 --- a/services/core/tests/integration/runtime_initialization_peer_test.go +++ b/services/core/tests/integration/runtime_initialization_peer_test.go @@ -29,6 +29,7 @@ type initializationPeer struct { deferred bool unavailable bool bootstrap sandbox.Bootstrap + binds chan proto.AssignmentBindPayload // when not nil, receives each bind's payload } func (p *initializationPeer) setRuntimeGateway(t *testing.T, endpoint string, registry *runtimegateway.Registry) { @@ -57,6 +58,10 @@ func (p *initializationPeer) connect(b sandbox.Bootstrap) error { return } if reply, ok := assignmentReply(env); ok { + var bind proto.AssignmentBindPayload + if p.binds != nil && env.DecodePayload(&bind) == nil { + p.binds <- bind + } if c.WriteJSON(reply) != nil { return }