From 8657e6923933c31664ddb4627ab2db7f24dfb8c3 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Thu, 8 Oct 2026 03:29:30 +0000 Subject: [PATCH] Follow Serve liveness for sandbox readiness --- apps/daemon/internal/dispatch/suspend.go | 6 +- apps/daemon/internal/dispatch/suspend_test.go | 32 ++++- docs/runtime-protocol.md | 6 +- docs/sandbox-link-protocol.md | 2 + docs/zh/runtime-protocol.md | 8 +- docs/zh/sandbox-link-protocol.md | 4 +- internal/sandboxlink/relay/relay.go | 10 ++ internal/sandboxlink/relay/relay_test.go | 22 +++ .../db/queries/environment_initialization.sql | 9 +- .../core/internal/db/queries/sandbox_link.sql | 9 ++ .../db/sqlc/environment_initialization.sql.go | 46 +++++-- .../core/internal/db/sqlc/sandbox_link.sql.go | 46 +++++++ .../archive_cancellation_cleanup_test.go | 2 +- .../internal/execution/runtime_compute.go | 5 +- .../execution/runtime_compute_wake.go | 6 +- .../internal/execution/runtime_connections.go | 78 +++++++---- .../execution/runtime_connections_test.go | 24 ---- .../execution/runtime_initialization.go | 32 ++++- .../internal/execution/runtime_lifecycle.go | 18 +-- .../internal/execution/runtime_manager.go | 3 +- .../core/internal/execution/runtime_setup.go | 10 +- .../sandbox_deployment_drain_test.go | 4 +- .../sandbox_deployment_setup_test.go | 12 +- .../sandbox_deployment_switch_test.go | 8 +- .../execution/sandbox_generations_test.go | 2 +- .../sandbox_provider_contract_test.go | 2 +- .../internal/execution/sandbox_reset_test.go | 8 +- .../execution/sandbox_snapshot_budget_test.go | 2 +- services/core/internal/execution/worker.go | 9 +- .../postgres/sessionpg/environment.go | 4 +- .../sessionpg/execution_environment.go | 9 ++ .../persistence/postgres/sessionpg/link.go | 16 +++ services/core/internal/sessions/devices.go | 10 ++ .../core/internal/sessions/environment.go | 5 + .../sessions/execution_environment.go | 11 ++ .../sessions/execution_environment_test.go | 5 + .../internal/sessions/transaction_test.go | 1 + .../000096_session_runtime_assignments.sql | 7 +- .../environment_initialization_test.go | 2 +- ...sted_initialization_failure_public_test.go | 5 +- .../tests/integration/link_authority_test.go | 114 ++++++++-------- .../integration/runtime_connection_test.go | 126 ++++++------------ .../runtime_initialization_peer_test.go | 21 ++- .../runtime_initialization_test.go | 2 +- .../integration/runtime_lifecycle_test.go | 16 ++- .../runtime_node_lifecycle_fixture_test.go | 17 ++- 46 files changed, 517 insertions(+), 279 deletions(-) delete mode 100644 services/core/internal/execution/runtime_connections_test.go diff --git a/apps/daemon/internal/dispatch/suspend.go b/apps/daemon/internal/dispatch/suspend.go index fd14f1327..64b3d44bf 100644 --- a/apps/daemon/internal/dispatch/suspend.go +++ b/apps/daemon/internal/dispatch/suspend.go @@ -90,8 +90,8 @@ func (r *Router) handleQuiesce(ctx context.Context, env proto.Envelope) error { } // handleResume reopens the quiesced Environment that the resume names. A -// rollback of an Environment this connection has not quiesced is accepted -// and changes nothing. +// resume of an Environment this connection has not quiesced, as after a +// restart, is accepted and changes nothing. func (r *Router) handleResume(ctx context.Context, env proto.Envelope) error { var request proto.EnvironmentSuspendPayload if env.ID == "" || len(env.ID) > 128 || env.DecodeRequest(&request) != nil { @@ -99,7 +99,7 @@ func (r *Router) handleResume(ctx context.Context, env proto.Envelope) error { } r.mu.Lock() err := r.resumeLocked(env.Assignment, request) - if errors.Is(err, errNotSuspended) && request.Rollback { + if errors.Is(err, errNotSuspended) && r.suspensions[request.EnvironmentID] == nil { err = nil } r.mu.Unlock() diff --git a/apps/daemon/internal/dispatch/suspend_test.go b/apps/daemon/internal/dispatch/suspend_test.go index f15b11fc0..afe1debca 100644 --- a/apps/daemon/internal/dispatch/suspend_test.go +++ b/apps/daemon/internal/dispatch/suspend_test.go @@ -157,6 +157,34 @@ func TestResumeRequiresExactSuspensionAndAssignment(t *testing.T) { } } +// A restarted Runtime has quiesced nothing and holds no assignment, so Core's +// resume of the Environment it quiesced before the restart succeeds. +func TestResumeOnFreshRouterSucceeds(t *testing.T) { + frames := make(chan proto.Envelope, 1) + r, err := New(Config{Registry: agent.NewRegistry(), Sender: suspendSender(func(_ context.Context, env proto.Envelope) error { frames <- env; return nil }), IdleTimeout: time.Hour}) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = r.Shutdown(context.Background()) }) + env, err := proto.NewEnvelope(proto.TypeEnvironmentResume, "resume", proto.EnvironmentSuspendPayload{EnvironmentID: "env", SuspendID: "attempt"}) + if err != nil { + t.Fatal(err) + } + env.Assignment = suspendRef + if err := r.Handle(t.Context(), env); err != nil { + t.Fatal(err) + } + select { + case frame := <-frames: + var result proto.EnvironmentSuspendResultPayload + if frame.Type != proto.TypeEnvironmentResumed || frame.DecodePayload(&result) != nil || !result.Accepted { + t.Fatalf("resume = %+v %+v", frame, result) + } + case <-time.After(3 * time.Second): + t.Fatal("resume has no result") + } +} + func TestQuiescingOneEnvironmentLeavesAnotherRunning(t *testing.T) { frames := make(chan proto.Envelope, 16) r := suspensionRouter(t, suspendSender(func(_ context.Context, env proto.Envelope) error { frames <- env; return nil })) @@ -217,9 +245,9 @@ func TestQuiescingOneEnvironmentLeavesAnotherRunning(t *testing.T) { code string }{ "foreign": {other, request, proto.AssignmentConflict}, - "unpaused": {other, proto.EnvironmentSuspendPayload{EnvironmentID: "other", SuspendID: "attempt"}, "not_suspended"}, - "rollback": {other, proto.EnvironmentSuspendPayload{EnvironmentID: "other", SuspendID: "attempt", Rollback: true}, ""}, + "unpaused": {other, proto.EnvironmentSuspendPayload{EnvironmentID: "other", SuspendID: "attempt"}, ""}, "obsolete": {suspendRef, proto.EnvironmentSuspendPayload{EnvironmentID: "env", SuspendID: "obsolete"}, "not_suspended"}, + "rollback": {suspendRef, proto.EnvironmentSuspendPayload{EnvironmentID: "env", SuspendID: "obsolete", Rollback: true}, "not_suspended"}, } { if got := suspend(proto.TypeEnvironmentResume, id, test.ref, test.request); got.ErrorCode != test.code || got.Accepted != (test.code == "") { t.Fatalf("%s resume = %+v", id, got) diff --git a/docs/runtime-protocol.md b/docs/runtime-protocol.md index 2ea7cb6ad..8eb0956d2 100644 --- a/docs/runtime-protocol.md +++ b/docs/runtime-protocol.md @@ -110,11 +110,11 @@ 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. 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`. +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. A bind of another assignment, or of a released assignment at a higher epoch, fails with `assignment_conflict`. 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`. +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 that supersedes the epoch takes the owner over once that attachment is closed. 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`. -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, then releases what its Environment owner holds 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. `environment_quiesce` applies to the named Environment only: it fails with `resource_busy` while one of the Environment's Sessions has work in progress; otherwise the Runtime closes the Environment's Executors, has its owners release what they hold and replies `environment_quiesced`. Until the matching `environment_resume`, which carries the assignment that quiesced it, the Runtime admits for that Environment only the releases of its Sessions and answers any other of its frames, including a bind, with `protocol_error` `resource_unavailable`; Sessions of other Environments keep running. +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, then releases what its Environment owner holds 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. `environment_quiesce` applies to the named Environment only: it fails with `resource_busy` while one of the Environment's Sessions has work in progress; otherwise the Runtime closes the Environment's Executors, has its owners release what they hold and replies `environment_quiesced`. Until the matching `environment_resume`, which carries the assignment that quiesced it, the Runtime admits for that Environment only the releases of its Sessions and answers any other of its frames, including a bind, with `protocol_error` `resource_unavailable`; Sessions of other Environments keep running. A resume of an Environment that the connection has not quiesced, as after the Runtime restarts, succeeds and changes nothing; one that names another suspension fails with `not_suspended`. 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. diff --git a/docs/sandbox-link-protocol.md b/docs/sandbox-link-protocol.md index 3d3d61c92..9d791b233 100644 --- a/docs/sandbox-link-protocol.md +++ b/docs/sandbox-link-protocol.md @@ -53,6 +53,8 @@ Each link carries at most 256 concurrent service streams, and each resource has The owner of the relay implements `Authority` from its durable records, and the relay consults it for every Hello, Open and renewal. To revoke, withdraw the authority first, then call `RevokeAttachment` or `RevokeResource` so the relay closes what it holds. +`Serving(ref)` reports whether the relay holds a serve peer of `ref` at `ref.Generation`. Like the rest of the relay state, it is this process's view: the owner uses it to tell whether a resource can carry streams now, and keeps its durable records as the judge of whether the resource exists. + Tests use `sandboxlinktest.NewAuthority`, which holds static credentials and grants, and `sandboxlinktest.StartRelay`, which runs a relay on an `httptest` TLS server and returns its URL and a TLS configuration that trusts it. ## Framing diff --git a/docs/zh/runtime-protocol.md b/docs/zh/runtime-protocol.md index 9bb451de0..66104ed95 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: b823b0a4d564b9191687005326939a7db10803ef99667d088336427b906d06d5 +source_hash: 7adbec57ec29c33da6a85b1204490e54a3bc580b1a37dd871a9f2a8b587ec278 --- 此协议在 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,11 +112,11 @@ 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 上的服务。若 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 的第一个操作之前,包括没有 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` 失败。其他分配的绑定,或以更高 epoch 绑定已释放的分配,以 `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`。 +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。它把每个 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`。 -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,随后释放其 Environment owner 持有的资源,并在要求时删除原生 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 连接时重发所有未确认的释放;失败的释放退避重试,等待最久的释放先发送,因此失败的释放不会拖延其他释放。`environment_quiesce` 只作用于指定的 Environment:该 Environment 的某个 Session 仍有进行中的工作时,它以 `resource_busy` 失败;否则 Runtime 关闭该 Environment 的 Executor,让其 owner 释放所持有的资源,并回复 `environment_quiesced`。在匹配的 `environment_resume`(携带使其 quiesce 的分配)到达之前,Runtime 对该 Environment 只准入其 Session 的释放,并以 `protocol_error` `resource_unavailable` 回答它的其他任何 frame(包括绑定);其他 Environment 的 Session 继续运行。 +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,随后释放其 Environment owner 持有的资源,并在要求时删除原生 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 连接时重发所有未确认的释放;失败的释放退避重试,等待最久的释放先发送,因此失败的释放不会拖延其他释放。`environment_quiesce` 只作用于指定的 Environment:该 Environment 的某个 Session 仍有进行中的工作时,它以 `resource_busy` 失败;否则 Runtime 关闭该 Environment 的 Executor,让其 owner 释放所持有的资源,并回复 `environment_quiesced`。在匹配的 `environment_resume`(携带使其 quiesce 的分配)到达之前,Runtime 对该 Environment 只准入其 Session 的释放,并以 `protocol_error` `resource_unavailable` 回答它的其他任何 frame(包括绑定);其他 Environment 的 Session 继续运行。对该连接未 quiesce 的 Environment 的 resume(例如 Runtime 重启之后)会成功且不改变任何状态;指定其他暂停的 resume 以 `not_suspended` 失败。 Runtime 对无法路由的 Core frame 回复 `protocol_error`,回显请求 ID,并携带其类型和错误码。 diff --git a/docs/zh/sandbox-link-protocol.md b/docs/zh/sandbox-link-protocol.md index 99ca95126..989859245 100644 --- a/docs/zh/sandbox-link-protocol.md +++ b/docs/zh/sandbox-link-protocol.md @@ -1,7 +1,7 @@ --- title: "沙箱 Link 协议" source: docs/sandbox-link-protocol.md -source_hash: f9c047631990775331946f78c7fad835a0663932ac6a2d4bd7c2514405f00b6f +source_hash: 429127c0e7407a565f72d3839739a4530b8336a6a2034a562514862339d2423c --- Link 协议通过 relay 连接沙箱 I/O 的两端。Sandbox I/O 服务运行在沙箱内并为其提供服务,是 serve peer。agent host 上的 Runtime 在沙箱外运行 Harness,并通过该服务使用沙箱,是 attach peer。每个 peer 各自向 relay 认证自己的 link。relay 授权 attach peer 打开的每个服务 stream,将其绑定到该资源当前的 serve peer,然后在两个 stream 之间复制字节而不读取内容。服务帧从不携带凭据或 grant。 @@ -55,6 +55,8 @@ attachment 的生命周期长于其 link。重连后,Runtime 使用相同的 b relay 的 owner 基于其持久记录实现 `Authority`,relay 对每个 Hello、Open 和续期都咨询它。撤销时,先撤回授权,再调用 `RevokeAttachment` 或 `RevokeResource`,让 relay 关闭其持有的对象。 +`Serving(ref)` 报告 relay 是否持有 `ref` 在 `ref.Generation` 的 serve peer。与 relay 的其他状态一样,它是本进程的视图:owner 用它判断资源当前能否承载 stream,而资源是否存在仍由其持久记录裁决。 + 测试使用 `sandboxlinktest.NewAuthority`(持有静态凭据和 grant)和 `sandboxlinktest.StartRelay`(在 `httptest` TLS server 上运行 relay,并返回其 URL 和信任它的 TLS 配置)。 ## 分帧 {#framing} diff --git a/internal/sandboxlink/relay/relay.go b/internal/sandboxlink/relay/relay.go index ec671954d..d08bdfa98 100644 --- a/internal/sandboxlink/relay/relay.go +++ b/internal/sandboxlink/relay/relay.go @@ -153,6 +153,16 @@ func (rl *Relay) RevokeResource(ref sandboxlink.ResourceRef) { rl.releaseLocked(key, r) } +// Serving reports whether the relay holds a serve peer of ref at +// ref.Generation. It is this process's view, held in memory like the rest of +// the relay. +func (rl *Relay) Serving(ref sandboxlink.ResourceRef) bool { + rl.mu.Lock() + defer rl.mu.Unlock() + r := rl.resources[keyOf(ref)] + return r != nil && r.serve != nil && r.serve.hello.Resource.Generation == ref.Generation +} + type resourceKey struct { tenant, environment, id sandboxwire.ID kind sandboxlink.ResourceKind diff --git a/internal/sandboxlink/relay/relay_test.go b/internal/sandboxlink/relay/relay_test.go index 9db74bad5..e26eabbc4 100644 --- a/internal/sandboxlink/relay/relay_test.go +++ b/internal/sandboxlink/relay/relay_test.go @@ -327,6 +327,28 @@ func TestGenerations(t *testing.T) { } } +// A resource is served at the generation of its current serve peer only, and +// no longer once the relay revokes it. +func TestServing(t *testing.T) { + f := newFixture(t) + if f.srv.Relay.Serving(resource(1)) { + t.Fatal("serving before any serve peer") + } + f.serve(1) + if !f.srv.Relay.Serving(resource(1)) || f.srv.Relay.Serving(resource(2)) { + t.Fatal("generation 1's serve peer does not serve generation 1 alone") + } + f.serve(2) + if f.srv.Relay.Serving(resource(1)) || !f.srv.Relay.Serving(resource(2)) { + t.Fatal("generation 2's serve peer does not replace generation 1") + } + f.auth.RemoveServe([]byte("serve credential 2")) + f.srv.Relay.RevokeResource(resource(2)) + if f.srv.Relay.Serving(resource(2)) { + t.Fatal("serving after revocation") + } +} + func TestLeaseExpiry(t *testing.T) { f := newFixture(t) p := f.serve(1) diff --git a/services/core/internal/db/queries/environment_initialization.sql b/services/core/internal/db/queries/environment_initialization.sql index 30d7ca519..10a97c766 100644 --- a/services/core/internal/db/queries/environment_initialization.sql +++ b/services/core/internal/db/queries/environment_initialization.sql @@ -1,7 +1,10 @@ -- name: ListEnvironmentInitializations :many -SELECT e.id, e.session_id, s.tenant_id, s.engine, e.initialization, b.runtime_id, b.assignment_id, b.epoch +SELECT e.id, e.session_id, s.tenant_id, s.engine, e.initialization, b.runtime_id, b.assignment_id, b.epoch, + r.tenant_id AS resource_tenant_id, r.environment_id AS resource_environment_id, r.kind AS resource_kind, + r.id AS resource_id, r.generation AS resource_generation, COALESCE(r.live, false)::boolean AS resource_live FROM environments e JOIN sessions s ON s.id = e.session_id LEFT JOIN session_runtime_assignments b ON b.session_id = s.id AND b.desired_state = 'bound' +LEFT JOIN sandbox_resources r ON r.environment_id = e.id WHERE e.id > $1 AND s.deleted_at IS NULL AND e.status NOT IN ('failed', 'expired') AND e.initialization IN ('pending', 'running') ORDER BY e.id LIMIT 32; @@ -10,6 +13,10 @@ ORDER BY e.id LIMIT 32; UPDATE environments SET initialization = 'running' WHERE id = $1 AND initialization = 'pending' AND status NOT IN ('failed', 'expired'); +-- name: UnclaimEnvironmentInitialization :execrows +UPDATE environments SET initialization = 'pending' +WHERE id = $1 AND initialization = 'running' AND status NOT IN ('failed', 'expired'); + -- name: CompleteEnvironmentInitialization :execrows UPDATE environments SET initialization = 'complete' WHERE id = $1 AND initialization = 'running' AND status NOT IN ('failed', 'expired'); diff --git a/services/core/internal/db/queries/sandbox_link.sql b/services/core/internal/db/queries/sandbox_link.sql index ebf6849de..39ae37b94 100644 --- a/services/core/internal/db/queries/sandbox_link.sql +++ b/services/core/internal/db/queries/sandbox_link.sql @@ -2,6 +2,15 @@ SELECT tenant_id, environment_id, kind, generation, credential_hash FROM sandbox_resources WHERE id = $1 AND live; +-- name: ListLiveSandboxResources :many +-- Every live Link resource. Quiesced compute is between a quiesce and the +-- wake that resumes it. +SELECT r.tenant_id, r.environment_id, r.kind, r.id, r.generation, + COALESCE(a.compute_phase NOT IN ('disabled', 'running'), false)::boolean AS quiesced +FROM sandbox_resources r +LEFT JOIN runtime_allocations a ON r.kind = 'allocation' AND a.id = r.id +WHERE r.live; + -- name: GetAgentHostCredential :one SELECT COALESCE(credential_hash, '')::text AS credential_hash, credential_revision FROM devices WHERE id = $1 AND agent_host AND revoked_at IS NULL; diff --git a/services/core/internal/db/sqlc/environment_initialization.sql.go b/services/core/internal/db/sqlc/environment_initialization.sql.go index 21c56093d..56fd64470 100644 --- a/services/core/internal/db/sqlc/environment_initialization.sql.go +++ b/services/core/internal/db/sqlc/environment_initialization.sql.go @@ -47,23 +47,32 @@ func (q *Queries) FailEnvironmentInitialization(ctx context.Context, id pgtype.U } const listEnvironmentInitializations = `-- name: ListEnvironmentInitializations :many -SELECT e.id, e.session_id, s.tenant_id, s.engine, e.initialization, b.runtime_id, b.assignment_id, b.epoch +SELECT e.id, e.session_id, s.tenant_id, s.engine, e.initialization, b.runtime_id, b.assignment_id, b.epoch, + r.tenant_id AS resource_tenant_id, r.environment_id AS resource_environment_id, r.kind AS resource_kind, + r.id AS resource_id, r.generation AS resource_generation, COALESCE(r.live, false)::boolean AS resource_live FROM environments e JOIN sessions s ON s.id = e.session_id LEFT JOIN session_runtime_assignments b ON b.session_id = s.id AND b.desired_state = 'bound' +LEFT JOIN sandbox_resources r ON r.environment_id = e.id WHERE e.id > $1 AND s.deleted_at IS NULL AND e.status NOT IN ('failed', 'expired') AND e.initialization IN ('pending', 'running') ORDER BY e.id LIMIT 32 ` type ListEnvironmentInitializationsRow struct { - ID pgtype.UUID `json:"id"` - SessionID pgtype.UUID `json:"session_id"` - TenantID pgtype.UUID `json:"tenant_id"` - Engine string `json:"engine"` - Initialization string `json:"initialization"` - RuntimeID pgtype.UUID `json:"runtime_id"` - AssignmentID pgtype.UUID `json:"assignment_id"` - Epoch pgtype.Int8 `json:"epoch"` + ID pgtype.UUID `json:"id"` + SessionID pgtype.UUID `json:"session_id"` + TenantID pgtype.UUID `json:"tenant_id"` + Engine string `json:"engine"` + Initialization string `json:"initialization"` + RuntimeID pgtype.UUID `json:"runtime_id"` + AssignmentID pgtype.UUID `json:"assignment_id"` + Epoch pgtype.Int8 `json:"epoch"` + ResourceTenantID pgtype.UUID `json:"resource_tenant_id"` + ResourceEnvironmentID pgtype.UUID `json:"resource_environment_id"` + ResourceKind pgtype.Text `json:"resource_kind"` + ResourceID pgtype.UUID `json:"resource_id"` + ResourceGeneration pgtype.Int8 `json:"resource_generation"` + ResourceLive bool `json:"resource_live"` } func (q *Queries) ListEnvironmentInitializations(ctx context.Context, id pgtype.UUID) ([]ListEnvironmentInitializationsRow, error) { @@ -84,6 +93,12 @@ func (q *Queries) ListEnvironmentInitializations(ctx context.Context, id pgtype. &i.RuntimeID, &i.AssignmentID, &i.Epoch, + &i.ResourceTenantID, + &i.ResourceEnvironmentID, + &i.ResourceKind, + &i.ResourceID, + &i.ResourceGeneration, + &i.ResourceLive, ); err != nil { return nil, err } @@ -94,3 +109,16 @@ func (q *Queries) ListEnvironmentInitializations(ctx context.Context, id pgtype. } return items, nil } + +const unclaimEnvironmentInitialization = `-- name: UnclaimEnvironmentInitialization :execrows +UPDATE environments SET initialization = 'pending' +WHERE id = $1 AND initialization = 'running' AND status NOT IN ('failed', 'expired') +` + +func (q *Queries) UnclaimEnvironmentInitialization(ctx context.Context, id pgtype.UUID) (int64, error) { + result, err := q.db.Exec(ctx, unclaimEnvironmentInitialization, id) + if err != nil { + return 0, err + } + return result.RowsAffected(), nil +} diff --git a/services/core/internal/db/sqlc/sandbox_link.sql.go b/services/core/internal/db/sqlc/sandbox_link.sql.go index acc0e7e7c..f9c23dfb3 100644 --- a/services/core/internal/db/sqlc/sandbox_link.sql.go +++ b/services/core/internal/db/sqlc/sandbox_link.sql.go @@ -105,6 +105,52 @@ func (q *Queries) GetSandboxServeAuthority(ctx context.Context, id pgtype.UUID) return i, err } +const listLiveSandboxResources = `-- name: ListLiveSandboxResources :many +SELECT r.tenant_id, r.environment_id, r.kind, r.id, r.generation, + COALESCE(a.compute_phase NOT IN ('disabled', 'running'), false)::boolean AS quiesced +FROM sandbox_resources r +LEFT JOIN runtime_allocations a ON r.kind = 'allocation' AND a.id = r.id +WHERE r.live +` + +type ListLiveSandboxResourcesRow struct { + TenantID pgtype.UUID `json:"tenant_id"` + EnvironmentID pgtype.UUID `json:"environment_id"` + Kind string `json:"kind"` + ID pgtype.UUID `json:"id"` + Generation int64 `json:"generation"` + Quiesced bool `json:"quiesced"` +} + +// Every live Link resource. Quiesced compute is between a quiesce and the +// wake that resumes it. +func (q *Queries) ListLiveSandboxResources(ctx context.Context) ([]ListLiveSandboxResourcesRow, error) { + rows, err := q.db.Query(ctx, listLiveSandboxResources) + if err != nil { + return nil, err + } + defer rows.Close() + items := []ListLiveSandboxResourcesRow{} + for rows.Next() { + var i ListLiveSandboxResourcesRow + if err := rows.Scan( + &i.TenantID, + &i.EnvironmentID, + &i.Kind, + &i.ID, + &i.Generation, + &i.Quiesced, + ); err != nil { + return nil, err + } + items = append(items, i) + } + if err := rows.Err(); err != nil { + return nil, err + } + return items, nil +} + const registerAgentHost = `-- name: RegisterAgentHost :one INSERT INTO devices (id, name, credential_hash, agent_host) VALUES ($1, 'agent-host', $2, true) diff --git a/services/core/internal/execution/archive_cancellation_cleanup_test.go b/services/core/internal/execution/archive_cancellation_cleanup_test.go index 8129521b9..440f80ca2 100644 --- a/services/core/internal/execution/archive_cancellation_cleanup_test.go +++ b/services/core/internal/execution/archive_cancellation_cleanup_test.go @@ -173,7 +173,7 @@ func TestArchiveWaitingCleanupReceiptBarrier(t *testing.T) { t.Fatal("Kill bypassed durable cleanup ownership", allocation, err) } }} - lifecycle := &runtimeLifecycle{sessions: sessionReader, sessionExecution: leased.Sessions, deployment: leased.Deployment, deployments: deployments, reader: reader, lease: leased.Lease, registry: registry, links: relay.New(nil), config: RuntimeProvider{InstallationID: installation, Provider: provider}, connections: map[string]*runtimeConnection{}} + lifecycle := &runtimeLifecycle{sessions: sessionReader, sessionExecution: leased.Sessions, deployment: leased.Deployment, deployments: deployments, reader: reader, lease: leased.Lease, registry: registry, links: relay.New(nil), config: RuntimeProvider{InstallationID: installation, Provider: provider}} if checkpoint { lifecycle.config.Provider = waitingCleanupCheckpoint{beforeKill: provider.beforeKill} } diff --git a/services/core/internal/execution/runtime_compute.go b/services/core/internal/execution/runtime_compute.go index 4be7eba50..6f64bc441 100644 --- a/services/core/internal/execution/runtime_compute.go +++ b/services/core/internal/execution/runtime_compute.go @@ -167,7 +167,10 @@ func (r *runtimeLifecycle) idleCompute(ctx context.Context, p sandbox.SandboxPro return err } // Publish disconnected only after receiving the daemon's receipt barrier. - if err := observeRuntimeConnection(ctx, r.sessionExecution, r.connections, owner.TenantID, owner.EnvironmentID, nil, false); err != nil { + r.connections.mu.Lock() + err = observeRuntimeConnection(ctx, r.sessionExecution, r.connections.current, owner.TenantID, owner.EnvironmentID, nil, false) + r.connections.mu.Unlock() + if err != nil { return err } // The Session-locked phase commit checks pending work and wake requests. diff --git a/services/core/internal/execution/runtime_compute_wake.go b/services/core/internal/execution/runtime_compute_wake.go index 635a072e6..dac40bd71 100644 --- a/services/core/internal/execution/runtime_compute_wake.go +++ b/services/core/internal/execution/runtime_compute_wake.go @@ -70,10 +70,8 @@ func (r *runtimeLifecycle) wakeCompute(ctx context.Context, p sandbox.SandboxPro if err != nil { return err } - if err := r.deployment.ClearWake(ctx, next, owner.ComputeActivityAt); err != nil { - return err - } - return r.observeConnection(ctx, next) + // The Worker's pass publishes the Environment connected again. + return r.deployment.ClearWake(ctx, next, owner.ComputeActivityAt) } func (r *runtimeLifecycle) cleanupCompute(ctx context.Context, p sandbox.SandboxProvider, owner deployment.Allocation, state runtimeCompute) error { diff --git a/services/core/internal/execution/runtime_connections.go b/services/core/internal/execution/runtime_connections.go index 10c4bfd1a..6338da9f0 100644 --- a/services/core/internal/execution/runtime_connections.go +++ b/services/core/internal/execution/runtime_connections.go @@ -3,47 +3,78 @@ package execution import ( "context" "errors" + "sync" "github.com/google/uuid" - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" ) -// Access is serialized by the existing lifecycle gate. Durable generations fence -// old observations; this map only remembers the currently observed socket. +// environmentConnections holds the connection generation of each Environment +// with a live Link resource. The Worker's pass and the hosted lifecycle's +// quiesce publish through it; mu orders a whole pass before or after a +// quiesce's publish, so no pass republishes an Environment that quiesced. +type environmentConnections struct { + mu sync.Mutex + current map[string]*runtimeConnection +} + +// Durable generations fence old observations; a runtimeConnection only +// remembers this process's generation of an Environment's connection. type runtimeConnection struct { peer *runtimegateway.Session + tenant string generation string revision int64 connected bool } -func (r *runtimeLifecycle) observeConnection(ctx context.Context, owner deployment.Allocation) error { - bound, err := r.sessions.GetSessionRuntimeDevice(ctx, owner.TenantID, owner.SessionID) - if errors.Is(err, sessions.ErrNotFound) { - return nil - } +// observeSandboxConnections publishes each Environment with a live Link +// resource as connected while the relay holds its serve peer, and as +// disconnected once its resource is gone. It leaves a quiesced Environment to +// the hosted lifecycle, which publishes it disconnected once its Runtime has +// quiesced, until the wake that resumes its compute. +func (w *Worker) observeSandboxConnections(ctx context.Context) error { + c := w.connections + c.mu.Lock() + defer c.mu.Unlock() + resources, err := w.dispatcher.SessionsReader.ListLiveSandboxResources(ctx) if err != nil { return err } - if bound.ID != owner.DeviceID || bound.EnvironmentID != owner.EnvironmentID { - return sessions.ErrDeviceBindingConflict - } - if !owner.CreateSettled || owner.State != "running" { - return nil + live := make(map[string]bool, len(resources)) + for _, resource := range resources { + live[resource.Resource.EnvironmentID] = true + if resource.Quiesced { + continue + } + serving := w.dispatcher.Links.Serving(resource.Resource.Ref()) + if err := observeRuntimeConnection(ctx, w.dispatcher.sessionExecution, c.current, resource.Resource.TenantID, resource.Resource.EnvironmentID, nil, serving); err != nil && !environmentGone(err) { + return err + } } - peer, err := authorizedRuntimePeer(ctx, r.sessions, r.registry, owner.DeviceID) - connected := err == nil - if err != nil && !errors.Is(err, sessions.ErrNotFound) && !errors.Is(err, runtimegateway.ErrSessionClosed) && !errors.Is(err, runtimegateway.ErrDeviceNotRegistered) { - return err + for environment, current := range c.current { + if live[environment] { + continue + } + if err := observeRuntimeConnection(ctx, w.dispatcher.sessionExecution, c.current, current.tenant, environment, nil, false); err != nil && !environmentGone(err) { + return err + } + delete(c.current, environment) } - return observeRuntimeConnection(ctx, r.sessionExecution, r.connections, owner.TenantID, owner.EnvironmentID, peer, connected) + return nil +} + +// environmentGone reports a connection observation of an Environment that is +// gone, failed or expired. +func environmentGone(err error) bool { + return errors.Is(err, sessions.ErrNotFound) || errors.Is(err, sessions.ErrInvalidInput) } -// Each Environment has one observer: the hosted lifecycle or the Worker loop for -// enrolled user compute. Both publish the same durable generation/revision rules. +// Each Environment has one observer: the Worker's pass over Link resources, +// with the hosted lifecycle's quiesce, or the Worker loop for enrolled user +// compute. Both publish the same durable generation/revision rules. func observeRuntimeConnection(ctx context.Context, operations *sessions.ExecutionOperations, connections map[string]*runtimeConnection, tenant, environment string, peer *runtimegateway.Session, connected bool) error { current := connections[environment] if connected && (current == nil || current.peer != peer) { @@ -51,7 +82,7 @@ func observeRuntimeConnection(ctx context.Context, operations *sessions.Executio if err := operations.ReplaceEnvironmentConnection(ctx, tenant, environment, generation); err != nil { return err } - current = &runtimeConnection{peer: peer, generation: generation} + current = &runtimeConnection{peer: peer, tenant: tenant, generation: generation} connections[environment] = current } if current == nil || current.connected == connected { @@ -78,10 +109,7 @@ func (w *Worker) observeEnrolledRuntimes(ctx context.Context) error { if err != nil && !errors.Is(err, sessions.ErrNotFound) && !errors.Is(err, runtimegateway.ErrSessionClosed) && !errors.Is(err, runtimegateway.ErrDeviceNotRegistered) { return err } - if err := observeRuntimeConnection(ctx, w.dispatcher.sessionExecution, w.enrolledConnections, bound.TenantID, bound.EnvironmentID, peer, connected); err != nil { - if errors.Is(err, sessions.ErrNotFound) || errors.Is(err, sessions.ErrInvalidInput) { - continue - } + if err := observeRuntimeConnection(ctx, w.dispatcher.sessionExecution, w.enrolledConnections, bound.TenantID, bound.EnvironmentID, peer, connected); err != nil && !environmentGone(err) { return err } } diff --git a/services/core/internal/execution/runtime_connections_test.go b/services/core/internal/execution/runtime_connections_test.go deleted file mode 100644 index c5b3d0624..000000000 --- a/services/core/internal/execution/runtime_connections_test.go +++ /dev/null @@ -1,24 +0,0 @@ -package execution - -import ( - "testing" - - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment" -) - -func TestRuntimeCleanupDropsOnlyOwnedEnvironmentConnection(t *testing.T) { - first := deployment.Allocation{ID: "allocation-one", EnvironmentID: "environment-one"} - second := deployment.Allocation{ID: "allocation-two", EnvironmentID: "environment-two"} - retained := &runtimeConnection{} - r := &runtimeLifecycle{ - connections: map[string]*runtimeConnection{first.EnvironmentID: {}, second.EnvironmentID: retained}, - } - r.clearRuntimeState(first) - if len(r.connections) != 1 || r.connections[second.EnvironmentID] != retained { - t.Fatal("cleanup retained the retired connection or discarded another Runtime's state") - } - r.clearRuntimeState(second) - if len(r.connections) != 0 { - t.Fatal("cleanup retained owned connection or initialization") - } -} diff --git a/services/core/internal/execution/runtime_initialization.go b/services/core/internal/execution/runtime_initialization.go index e619198dd..33ef2657d 100644 --- a/services/core/internal/execution/runtime_initialization.go +++ b/services/core/internal/execution/runtime_initialization.go @@ -4,6 +4,7 @@ import ( "context" "encoding/json" "errors" + "fmt" "sync" "time" @@ -56,6 +57,9 @@ func (w *Worker) runEnvironmentInitializations(ctx context.Context) error { if len(active) >= w.executionConcurrency() { continue } + if owner.Resource.Kind != "" && (!owner.ResourceLive || !w.dispatcher.Links.Serving(owner.Resource.Ref())) { + continue + } peer, err := w.dispatcher.authorizedPeer(ctx, owner.DeviceID) if err != nil { if errors.Is(err, sessions.ErrNotFound) || errors.Is(err, runtimegateway.ErrSessionClosed) || errors.Is(err, runtimegateway.ErrDeviceNotRegistered) { @@ -98,15 +102,25 @@ func (w *Worker) initializeEnvironment(ctx context.Context, owner sessions.Envir if err == nil { err = w.dispatcher.sessionExecution.CompleteEnvironmentInitialization(operation, owner) } - if err != nil { - log.Warn(ctx, "Environment preparation failed", "environment_id", owner.EnvironmentID, "session_id", owner.SessionID) - // A later scan settles an unrecorded failure; it never retries the setup. - record, stop := context.WithTimeout(context.WithoutCancel(ctx), 5*time.Second) - defer stop() - _ = w.dispatcher.sessionExecution.FailEnvironmentInitialization(record, owner, failure) + if err == nil { + return } + record, stop := context.WithTimeout(context.WithoutCancel(ctx), 5*time.Second) + defer stop() + // Nothing took effect, so a later scan starts the preparation again. + if errors.Is(err, errBeforeEffect) && w.dispatcher.sessionExecution.UnclaimEnvironmentInitialization(record, owner) == nil { + return + } + log.Warn(ctx, "Environment preparation failed", "environment_id", owner.EnvironmentID, "session_id", owner.SessionID) + // A later scan settles an unrecorded failure; it never retries the setup. + _ = w.dispatcher.sessionExecution.FailEnvironmentInitialization(record, owner, failure) } +// errBeforeEffect ends a preparation before any of its steps took effect: its +// bind failed because the Runtime is gone or the Environment has no live Link +// resource, or the Runtime rejected its first step with resource_unavailable. +var errBeforeEffect = errors.New("environment preparation ended before any effect") + func (w *Worker) prepareEnvironment(ctx context.Context, owner sessions.EnvironmentInitialization, failure *sessions.ProvisioningFailure) error { environment, err := w.dispatcher.SessionsReader.GetEnvironment(ctx, owner.TenantID, owner.EnvironmentID) if err != nil { @@ -123,6 +137,9 @@ func (w *Worker) prepareEnvironment(ctx context.Context, owner sessions.Environm return err } peer, err := w.dispatcher.assignedPeer(ctx, sessions.ExecutionDevice{ID: owner.DeviceID, Assignment: owner.Assignment, SessionEnvironmentID: owner.EnvironmentID}) + if errors.Is(err, runtimegateway.ErrNoLinkResource) || errors.Is(err, runtimegateway.ErrSessionClosed) || errors.Is(err, runtimegateway.ErrDeviceNotRegistered) { + return fmt.Errorf("%w: %w", errBeforeEffect, err) + } if err != nil { return err } @@ -155,6 +172,9 @@ func (w *Worker) prepareEnvironment(ctx context.Context, owner sessions.Environm if err != nil { var confirmed *runtimeStepFailure if errors.As(err, &confirmed) { + if index == 0 && confirmed.unavailable { + return fmt.Errorf("%w: %w", errBeforeEffect, err) + } candidate.ExitCode = confirmed.exitCode *failure = candidate } diff --git a/services/core/internal/execution/runtime_lifecycle.go b/services/core/internal/execution/runtime_lifecycle.go index 5129d3e24..2c77cd305 100644 --- a/services/core/internal/execution/runtime_lifecycle.go +++ b/services/core/internal/execution/runtime_lifecycle.go @@ -59,11 +59,11 @@ type runtimeLifecycle struct { reconcileCancel context.CancelFunc cursor string pendingCursor string - connections map[string]*runtimeConnection + connections *environmentConnections wakeHints chan struct{} } -func newRuntimeManager(owner Owner, deployments *deployment.Service, deploymentReader deployment.Reader, sessionReader sessions.Reader, registry *runtimegateway.Registry, links *relay.Relay, config *RuntimeProvider) (*runtimeManager, error) { +func newRuntimeManager(owner Owner, deployments *deployment.Service, deploymentReader deployment.Reader, sessionReader sessions.Reader, registry *runtimegateway.Registry, links *relay.Relay, connections *environmentConnections, config *RuntimeProvider) (*runtimeManager, error) { if config == nil { return nil, nil } @@ -72,7 +72,7 @@ func newRuntimeManager(owner Owner, deployments *deployment.Service, deploymentR return nil, sandbox.ErrInvalid } ctx, stop := context.WithCancel(context.Background()) - return &runtimeManager{sessions: sessionReader, sessionExecution: owner.Sessions, deployment: owner.Deployment, deploymentService: deployments, deploymentReader: deploymentReader, lease: owner.Lease, registry: registry, links: links, setupInstallationID: config.InstallationID, loadDeployment: config.loadDeployment, prepareDeployment: config.prepareDeployment, publishUnconfigured: config.PublishUnconfigured, setupGate: make(chan struct{}, 1), mutationGate: make(chan struct{}, 1), ctx: ctx, cancel: stop, nodes: make(map[string]*runtimeNode), failed: make(chan error, 1), inventory: make(chan struct{}, 1)}, nil + return &runtimeManager{sessions: sessionReader, sessionExecution: owner.Sessions, deployment: owner.Deployment, deploymentService: deployments, deploymentReader: deploymentReader, lease: owner.Lease, registry: registry, links: links, connections: connections, setupInstallationID: config.InstallationID, loadDeployment: config.loadDeployment, prepareDeployment: config.prepareDeployment, publishUnconfigured: config.PublishUnconfigured, setupGate: make(chan struct{}, 1), mutationGate: make(chan struct{}, 1), ctx: ctx, cancel: stop, nodes: make(map[string]*runtimeNode), failed: make(chan error, 1), inventory: make(chan struct{}, 1)}, nil } func validatedRuntimeProvider(config *RuntimeProvider, registry *runtimegateway.Registry) (RuntimeProvider, error) { @@ -312,11 +312,6 @@ func (r *runtimeLifecycle) observe(ctx context.Context, owner deployment.Allocat return nil } } - r.clearRuntimeState(owner) - } else if owner.ComputePhase == "disabled" || owner.ComputePhase == "running" { - if err := r.observeConnection(ctx, owner); err != nil { - return err - } } if owner.ComputePhase != "disabled" { return r.observeCompute(ctx, owner) @@ -388,7 +383,7 @@ func (r *runtimeLifecycle) observe(ctx context.Context, owner deployment.Allocat if r.config.Suspension != nil && environment.Initialization == "complete" { return r.enableCompute(ctx, owner) } - if peer, err := r.registry.LookupDevice(owner.DeviceID); err != nil || peer.IsClosed() { + if !r.links.Serving(serveResource(owner).Ref()) { return nil } renewed, err := provider.Renew(ctx, runtimeReference(owner)) @@ -419,11 +414,6 @@ func serveResource(owner deployment.Allocation) sandboxbootstrap.Resource { return sandboxbootstrap.Resource{TenantID: owner.TenantID, EnvironmentID: owner.EnvironmentID, Kind: "allocation", ID: owner.ID, Generation: owner.ServeGeneration} } -// Environment identity owns connectivity; preparation has an independent owner. -func (r *runtimeLifecycle) clearRuntimeState(owner deployment.Allocation) { - delete(r.connections, owner.EnvironmentID) -} - func runtimeReference(owner deployment.Allocation) sandbox.Reference { return sandbox.Reference{TenantID: owner.TenantID, EnvironmentID: owner.EnvironmentID, AllocationID: owner.ID} } diff --git a/services/core/internal/execution/runtime_manager.go b/services/core/internal/execution/runtime_manager.go index f80503a97..0e5914747 100644 --- a/services/core/internal/execution/runtime_manager.go +++ b/services/core/internal/execution/runtime_manager.go @@ -26,6 +26,7 @@ type runtimeManager struct { lease Ownership registry *runtimegateway.Registry links *relay.Relay + connections *environmentConnections config RuntimeProvider setupInstallationID string loadDeployment func(context.Context) (*RuntimeProvider, error) @@ -93,7 +94,7 @@ func (m *runtimeManager) node(id string) (*runtimeNode, error) { deployment: m.deployment, deployments: m.deploymentService, reader: m.deploymentReader, lease: m.lease, registry: m.registry, links: m.links, config: m.config, nodeID: id, gate: make(chan struct{}, 1), ctx: ctx, stop: stop, - connections: make(map[string]*runtimeConnection), wakeHints: make(chan struct{}, 1), + connections: m.connections, wakeHints: make(chan struct{}, 1), }} m.nodes[id] = n } diff --git a/services/core/internal/execution/runtime_setup.go b/services/core/internal/execution/runtime_setup.go index f02994f70..a1003e955 100644 --- a/services/core/internal/execution/runtime_setup.go +++ b/services/core/internal/execution/runtime_setup.go @@ -19,7 +19,13 @@ type runtimeSetupOperation struct { Index int } -type runtimeStepFailure struct{ exitCode int } +// runtimeStepFailure is a step the Runtime rejected or failed. unavailable +// marks a rejection with resource_unavailable, which ends a step before any +// effect. +type runtimeStepFailure struct { + exitCode int + unavailable bool +} func (*runtimeStepFailure) Error() string { return "environment initialization operation failed" } @@ -92,7 +98,7 @@ func runRuntimeSetup(ctx context.Context, peer runtimePreparer, ref proto.Assign return nil } if err == nil && (result.Outcome == "rejected" || result.Outcome == "failed") { - return &runtimeStepFailure{exitCode: result.ExitCode} + return &runtimeStepFailure{exitCode: result.ExitCode, unavailable: result.Outcome == "rejected" && result.ErrorCode == "resource_unavailable"} } return errors.New("environment initialization operation unconfirmed") } diff --git a/services/core/internal/execution/sandbox_deployment_drain_test.go b/services/core/internal/execution/sandbox_deployment_drain_test.go index 5a48bd2a4..ed9723ce0 100644 --- a/services/core/internal/execution/sandbox_deployment_drain_test.go +++ b/services/core/internal/execution/sandbox_deployment_drain_test.go @@ -94,7 +94,7 @@ func testLifecycleCancellationPreservesLease(t *testing.T, mode string) { defer hub.Close() id := uuid.NewString() configuration := &RuntimeProvider{InstallationID: id, ProviderKind: "docker", Mode: "nodes", Generation: 1, CoreURL: "https://core.example/api/v1", SandboxLink: "wss://core.example/api/v1/sandbox-link", BackendFingerprint: strings.Repeat("a", 64), Provider: hub.Proxy(uuid.NewString(), docker.Operations(), 1)} - m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return configuration, nil }, unusedPreparation(t))) + m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return configuration, nil }, unusedPreparation(t))) if err != nil { t.Fatal(err) } @@ -232,7 +232,7 @@ func TestSandboxDeploymentDrainFailureCannotReactivate(t *testing.T) { defer hub.Close() id := uuid.NewString() configuration := &RuntimeProvider{InstallationID: id, ProviderKind: "docker", Mode: "nodes", Generation: 1, CoreURL: "https://core.example/api/v1", SandboxLink: "wss://core.example/api/v1/sandbox-link", BackendFingerprint: strings.Repeat("a", 64), Provider: hub.Proxy(uuid.NewString(), docker.Operations(), 1)} - m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return configuration, nil }, unusedPreparation(t))) + m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return configuration, nil }, unusedPreparation(t))) if err != nil { t.Fatal(err) } diff --git a/services/core/internal/execution/sandbox_deployment_setup_test.go b/services/core/internal/execution/sandbox_deployment_setup_test.go index 37f2a88a4..f0a0c08ea 100644 --- a/services/core/internal/execution/sandbox_deployment_setup_test.go +++ b/services/core/internal/execution/sandbox_deployment_setup_test.go @@ -26,7 +26,7 @@ func TestDeferredSandboxDeploymentLoadsOnceBeforeNodeCreation(t *testing.T) { var selected atomic.Bool var loads atomic.Int32 configuration := &RuntimeProvider{InstallationID: id, ProviderKind: "docker", Mode: "nodes", CoreURL: "https://core.example/api/v1", SandboxLink: "wss://core.example/api/v1/sandbox-link", BackendFingerprint: strings.Repeat("a", 64), Provider: hub.Proxy(uuid.NewString(), docker.Operations(), 1)} - m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { loads.Add(1) if !selected.Load() { return nil, nil @@ -71,7 +71,7 @@ func TestDeferredSandboxDeploymentLoadsOnceBeforeNodeCreation(t *testing.T) { func TestDeferredSandboxDeploymentShutdownCancelsLoad(t *testing.T) { entered := make(chan struct{}) - m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(uuid.NewString(), func(ctx context.Context) (*RuntimeProvider, error) { + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, NewDeferredRuntimeProvider(uuid.NewString(), func(ctx context.Context) (*RuntimeProvider, error) { close(entered) <-ctx.Done() return nil, ctx.Err() @@ -96,7 +96,7 @@ func TestDeferredSandboxProviderFailureKeepsRecoveryAvailable(t *testing.T) { available := false loadErr := ErrExecutionUnavailable configuration := &RuntimeProvider{InstallationID: id, ProviderKind: "docker", Mode: "nodes", CoreURL: "https://core.example/api/v1", SandboxLink: "wss://core.example/api/v1/sandbox-link", BackendFingerprint: strings.Repeat("a", 64), Provider: hub.Proxy(uuid.NewString(), docker.Operations(), 1)} - m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { if !available { return nil, loadErr } @@ -128,7 +128,7 @@ func TestRejectedSandboxCandidatePreservesActiveGeneration(t *testing.T) { id := uuid.NewString() config := &RuntimeProvider{InstallationID: id, ProviderKind: "docker", Mode: "nodes", Generation: 1, CoreURL: "https://core.example/api/v1", SandboxLink: "wss://core.example/api/v1/sandbox-link", BackendFingerprint: strings.Repeat("a", 64), Provider: hub.Proxy(uuid.NewString(), docker.Operations(), 1)} rejected := errors.New("candidate provider unavailable") - m, err := newRuntimeManager(Owner{Lease: heldLease{}}, unitDeploymentService(t), nil, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return config, nil }, + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, unitDeploymentService(t), nil, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return config, nil }, func(context.Context, deployment.Setup) (PreparedRuntimeDeployment, error) { return PreparedRuntimeDeployment{}, rejected })) @@ -160,7 +160,7 @@ func TestRejectedSandboxCandidatePreservesActiveGeneration(t *testing.T) { func TestSandboxCandidateValidationDoesNotHoldManagerLock(t *testing.T) { id := uuid.NewString() entered, release := make(chan struct{}), make(chan struct{}) - m, err := newRuntimeManager(Owner{Lease: heldLease{}}, unitDeploymentService(t), nil, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil }, + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, unitDeploymentService(t), nil, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil }, func(context.Context, deployment.Setup) (PreparedRuntimeDeployment, error) { close(entered) <-release @@ -193,7 +193,7 @@ func TestCommittedSandboxCandidatePublishesAfterShutdown(t *testing.T) { hub := node.NewHub(node.HubOptions{}) defer hub.Close() id := uuid.NewString() - m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil }, unusedPreparation(t))) + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil }, unusedPreparation(t))) if err != nil { t.Fatal(err) } diff --git a/services/core/internal/execution/sandbox_deployment_switch_test.go b/services/core/internal/execution/sandbox_deployment_switch_test.go index 3dde18a09..767463d81 100644 --- a/services/core/internal/execution/sandbox_deployment_switch_test.go +++ b/services/core/internal/execution/sandbox_deployment_switch_test.go @@ -20,7 +20,7 @@ func TestSandboxManagerSwitchDrainsBeforeDirectActivation(t *testing.T) { defer hub.Close() id := uuid.NewString() config := &RuntimeProvider{InstallationID: id, ProviderKind: "docker", Mode: "nodes", Generation: 1, CoreURL: "https://core.example/api/v1", SandboxLink: "wss://core.example/api/v1/sandbox-link", BackendFingerprint: strings.Repeat("a", 64), Provider: hub.Proxy(uuid.NewString(), docker.Operations(), 1)} - m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return config, nil }, unusedPreparation(t))) + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return config, nil }, unusedPreparation(t))) if err != nil { t.Fatal(err) } @@ -81,7 +81,7 @@ func TestSandboxManagerSwitchDrainsBeforeDirectActivation(t *testing.T) { func TestSandboxManagerFailedActivationStaysPaused(t *testing.T) { id := uuid.NewString() - m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, errors.New("provider unavailable") }, unusedPreparation(t))) + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, errors.New("provider unavailable") }, unusedPreparation(t))) if err != nil { t.Fatal(err) } @@ -102,7 +102,7 @@ func TestSandboxManagerCancelledSwitchCannotResumeBeforeDrain(t *testing.T) { defer hub.Close() id := uuid.NewString() config := &RuntimeProvider{InstallationID: id, ProviderKind: "docker", Mode: "nodes", Generation: 1, CoreURL: "https://core.example/api/v1", SandboxLink: "wss://core.example/api/v1/sandbox-link", BackendFingerprint: strings.Repeat("a", 64), Provider: hub.Proxy(uuid.NewString(), docker.Operations(), 1)} - m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return config, nil }, unusedPreparation(t))) + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return config, nil }, unusedPreparation(t))) if err != nil { t.Fatal(err) } @@ -156,7 +156,7 @@ func TestSandboxActivationCannotBypassOutstandingDrain(t *testing.T) { hub := node.NewHub(node.HubOptions{}) defer hub.Close() id := uuid.NewString() - m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil }, func(_ context.Context, setup deployment.Setup) (PreparedRuntimeDeployment, error) { return PreparedRuntimeDeployment{Config: &RuntimeProvider{InstallationID: setup.InstallationID, ProviderKind: setup.Provider, Mode: setup.Mode, CoreURL: "https://core.example/api/v1", SandboxLink: "wss://core.example/api/v1/sandbox-link", BackendFingerprint: setup.BackendFingerprint, Provider: hub.Proxy(uuid.NewString(), docker.Operations(), 1)}}, nil diff --git a/services/core/internal/execution/sandbox_generations_test.go b/services/core/internal/execution/sandbox_generations_test.go index 1405e8465..bdfae1e0b 100644 --- a/services/core/internal/execution/sandbox_generations_test.go +++ b/services/core/internal/execution/sandbox_generations_test.go @@ -51,7 +51,7 @@ func TestE2BReplacementVerifiesTwiceAndNeverPublishesFailedCommit(t *testing.T) }, FenceCredential: func(context.Context) (func(), error) { fenced++; return func() { released++ }, nil }, Publish: func(*RuntimeProvider) { published++ }}, nil }) - m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(nil), config) + m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, config) if err != nil { t.Fatal(err) } diff --git a/services/core/internal/execution/sandbox_provider_contract_test.go b/services/core/internal/execution/sandbox_provider_contract_test.go index ee3328381..f86a7b4bd 100644 --- a/services/core/internal/execution/sandbox_provider_contract_test.go +++ b/services/core/internal/execution/sandbox_provider_contract_test.go @@ -19,7 +19,7 @@ func TestSandboxProviderRegistrationDoesNotRequireAnExecutionVendorBranch(t *tes id := uuid.NewString() config := &RuntimeProvider{InstallationID: id, ProviderKind: "contract-fixture", Mode: mode, CoreURL: "https://core.example/api/v1", SandboxLink: "wss://core.example/api/v1/sandbox-link", BackendFingerprint: strings.Repeat("a", 64), Provider: &lifecycleOnlySandbox{}} - m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return config, nil }, unusedPreparation(t))) + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return config, nil }, unusedPreparation(t))) if err != nil { t.Fatal(err) } diff --git a/services/core/internal/execution/sandbox_reset_test.go b/services/core/internal/execution/sandbox_reset_test.go index aab7ac473..7deeee337 100644 --- a/services/core/internal/execution/sandbox_reset_test.go +++ b/services/core/internal/execution/sandbox_reset_test.go @@ -91,7 +91,7 @@ func TestSandboxResetPageTimeoutRecoversCommittedOwner(t *testing.T) { } return &RuntimeProvider{InstallationID: id, ProviderKind: setup.Provider, Mode: setup.Mode, Generation: setup.Generation, CoreURL: "https://core.example/api/v1", SandboxLink: "wss://core.example/api/v1/sandbox-link", BackendFingerprint: setup.BackendFingerprint, Provider: hub.Proxy(uuid.NewString(), docker.Operations(), 1)}, nil }, unusedPreparation(t)) - m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(nil), config) + m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, config) if err != nil { t.Fatal(err) } @@ -212,7 +212,7 @@ func TestSandboxResetPublishesCommittedGenerationWithoutReading(t *testing.T) { } published = append(published, generation) } - m, err := newRuntimeManager(owner, deployments, adapter, nil, runtimegateway.NewRegistry(), relay.New(nil), config) + m, err := newRuntimeManager(owner, deployments, adapter, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, config) if err != nil { t.Fatal(err) } @@ -247,7 +247,7 @@ func TestCommittedResetViewStopsOwnerWithoutLease(t *testing.T) { return deployment.Snapshot{Record: deployment.Record{InstallationID: id, Generation: 1}}, nil }} deployments, operations := deploymentOperations(t, &strictDeploymentStorage{t: t}, reader, &strictExecutionStorage{t: t}) - m, err := newRuntimeManager(Owner{Lease: lostLease{}, Deployment: operations}, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil }, unusedPreparation(t))) + m, err := newRuntimeManager(Owner{Lease: lostLease{}, Deployment: operations}, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil }, unusedPreparation(t))) if err != nil { t.Fatal(err) } @@ -269,7 +269,7 @@ func TestCommittedResetViewStopsOwnerWithoutLease(t *testing.T) { func TestSandboxResetChangesReturnViewReadAfterCommit(t *testing.T) { owner, deployments, reader := resetManager(t) id := initializeE2BDeployment(t, owner) - m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil }, unusedPreparation(t))) + m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil }, unusedPreparation(t))) if err != nil { t.Fatal(err) } diff --git a/services/core/internal/execution/sandbox_snapshot_budget_test.go b/services/core/internal/execution/sandbox_snapshot_budget_test.go index 7c42780b0..065622a0d 100644 --- a/services/core/internal/execution/sandbox_snapshot_budget_test.go +++ b/services/core/internal/execution/sandbox_snapshot_budget_test.go @@ -75,7 +75,7 @@ func TestSandboxResetSnapshotFitsPageBudget(t *testing.T) { } return &RuntimeProvider{InstallationID: id, ProviderKind: setup.Provider, Generation: setup.Generation, Mode: setup.Mode, CoreURL: "https://core.example/api/v1", SandboxLink: "wss://core.example/api/v1/sandbox-link", BackendFingerprint: setup.BackendFingerprint, Provider: hub.Proxy(uuid.NewString(), docker.Operations(), 1)}, nil }, unusedPreparation(t)) - m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(nil), configuration) + m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(nil), nil, configuration) if err != nil { t.Fatal(err) } diff --git a/services/core/internal/execution/worker.go b/services/core/internal/execution/worker.go index 0ecd39e57..bcc90e47e 100644 --- a/services/core/internal/execution/worker.go +++ b/services/core/internal/execution/worker.go @@ -26,6 +26,7 @@ type Worker struct { stopped chan struct{} stopOnce sync.Once runtimes *runtimeManager + connections *environmentConnections enrolledConnections map[string]*runtimeConnection } @@ -76,8 +77,8 @@ func StartWorker(ctx context.Context, dispatcher *Dispatcher, owner Owner) (_ *W return nil, errors.New("execution worker requires the Link relay") } owned.notifications = &executionNotifications{} - worker := &Worker{concurrency: dispatcher.MaxConcurrentExecutions, dispatcher: owned, lease: owner.Lease, directoryReads: make(chan directoryReadRequest), fileWrites: make(chan fileWriteRequest), stopped: make(chan struct{}), scheduleWake: make(chan struct{}, 1), enrolledConnections: make(map[string]*runtimeConnection)} - worker.runtimes, err = newRuntimeManager(owner, owned.Deployment, owned.DeploymentReader, owned.SessionsReader, owned.Registry, owned.Links, owned.ManagedRuntimes) + worker := &Worker{concurrency: dispatcher.MaxConcurrentExecutions, dispatcher: owned, lease: owner.Lease, directoryReads: make(chan directoryReadRequest), fileWrites: make(chan fileWriteRequest), stopped: make(chan struct{}), scheduleWake: make(chan struct{}, 1), connections: &environmentConnections{current: make(map[string]*runtimeConnection)}, enrolledConnections: make(map[string]*runtimeConnection)} + worker.runtimes, err = newRuntimeManager(owner, owned.Deployment, owned.DeploymentReader, owned.SessionsReader, owned.Registry, owned.Links, worker.connections, owned.ManagedRuntimes) if err != nil { return nil, err } @@ -295,6 +296,10 @@ func (w *Worker) Run(ctx context.Context) (runErr error) { w.observeSchedulerPoll(0, err) return err } + if err := w.observeSandboxConnections(ctx); err != nil { + w.observeSchedulerPoll(0, err) + return err + } if err := w.observeEnrolledRuntimes(ctx); err != nil { w.observeSchedulerPoll(0, err) return err diff --git a/services/core/internal/persistence/postgres/sessionpg/environment.go b/services/core/internal/persistence/postgres/sessionpg/environment.go index 25e56c911..6be26b7ef 100644 --- a/services/core/internal/persistence/postgres/sessionpg/environment.go +++ b/services/core/internal/persistence/postgres/sessionpg/environment.go @@ -64,7 +64,9 @@ func (s *Store) ListEnvironmentInitializations(ctx context.Context, after string result = append(result, sessions.EnvironmentInitialization{ EnvironmentID: optionalID(row.ID), SessionID: optionalID(row.SessionID), TenantID: optionalID(row.TenantID), DeviceID: optionalID(row.RuntimeID), State: row.Initialization, Engine: row.Engine, - Assignment: assignmentRef(row.SessionID, row.AssignmentID, row.Epoch.Int64), + Assignment: assignmentRef(row.SessionID, row.AssignmentID, row.Epoch.Int64), + Resource: linkResource(row.ResourceTenantID, row.ResourceEnvironmentID, row.ResourceKind, row.ResourceID, row.ResourceGeneration), + ResourceLive: row.ResourceLive, }) } return result, nil diff --git a/services/core/internal/persistence/postgres/sessionpg/execution_environment.go b/services/core/internal/persistence/postgres/sessionpg/execution_environment.go index be9ee9816..a62dba4f8 100644 --- a/services/core/internal/persistence/postgres/sessionpg/execution_environment.go +++ b/services/core/internal/persistence/postgres/sessionpg/execution_environment.go @@ -138,6 +138,15 @@ func (t *environmentTx) ClaimInitialization(ctx context.Context, environment str return count == 1, err } +func (t *environmentTx) UnclaimInitialization(ctx context.Context, environment string) (bool, error) { + id, err := parseID(environment) + if err != nil { + return false, err + } + count, err := t.q.UnclaimEnvironmentInitialization(ctx, id) + return count == 1, err +} + func (t *environmentTx) CompleteInitialization(ctx context.Context, environment string) (bool, error) { id, err := parseID(environment) if err != nil { diff --git a/services/core/internal/persistence/postgres/sessionpg/link.go b/services/core/internal/persistence/postgres/sessionpg/link.go index 3219ec32a..266af9b8c 100644 --- a/services/core/internal/persistence/postgres/sessionpg/link.go +++ b/services/core/internal/persistence/postgres/sessionpg/link.go @@ -12,6 +12,7 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/db/sqlc" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgunit" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimedevice" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" ) // attachGrantPurpose is the credential key purpose attach grants are signed @@ -38,6 +39,21 @@ func (s *Store) GetServeAuthority(ctx context.Context, id string) (runtimedevice }, true, nil } +func (s *Store) ListLiveSandboxResources(ctx context.Context) ([]sessions.SandboxResource, error) { + rows, err := s.units.Queries().ListLiveSandboxResources(ctx) + if err != nil { + return nil, err + } + result := make([]sessions.SandboxResource, 0, len(rows)) + for _, row := range rows { + result = append(result, sessions.SandboxResource{ + Resource: linkResource(row.TenantID, row.EnvironmentID, pgtype.Text{String: row.Kind, Valid: true}, row.ID, pgtype.Int8{Int64: row.Generation, Valid: true}), + Quiesced: row.Quiesced, + }) + } + return result, nil +} + // GetAgentHostCredential reads a live agent host's credential. A malformed or // unknown device, or one that is not an agent host, has none. func (s *Store) GetAgentHostCredential(ctx context.Context, runtime string) (runtimedevice.AgentHost, bool, error) { diff --git a/services/core/internal/sessions/devices.go b/services/core/internal/sessions/devices.go index d7fd40cf9..dc8ba7d40 100644 --- a/services/core/internal/sessions/devices.go +++ b/services/core/internal/sessions/devices.go @@ -6,6 +6,7 @@ import ( "fmt" "strings" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxbootstrap" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimedevice" ) @@ -44,6 +45,13 @@ type EnrolledRuntimeBinding struct { DeviceID, TenantID, EnvironmentID, SessionID string } +// SandboxResource is a live Link resource. Quiesced compute is between a +// quiesce and the wake that resumes it. +type SandboxResource struct { + Resource sandboxbootstrap.Resource + Quiesced bool +} + // EnrollmentAuthority is what an executor credential authorizes when it // enrolls a Runtime: the key and the Environment's workspace directory. type EnrollmentAuthority struct { @@ -72,6 +80,8 @@ type DeviceReader interface { // ListEnrolledRuntimeBindings lists the enrolled user-managed Runtimes of // live Environments. ListEnrolledRuntimeBindings(ctx context.Context) ([]EnrolledRuntimeBinding, error) + // ListLiveSandboxResources lists the live Link resources. + ListLiveSandboxResources(ctx context.Context) ([]SandboxResource, error) // GetSessionExecutionBinding reads the Runtime device that executes the // Session's Turns, with the native session that continues its history, // once its Environment preparation completed; before that, and without an diff --git a/services/core/internal/sessions/environment.go b/services/core/internal/sessions/environment.go index ee02013a3..91406df7e 100644 --- a/services/core/internal/sessions/environment.go +++ b/services/core/internal/sessions/environment.go @@ -10,6 +10,7 @@ import ( "github.com/google/uuid" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxbootstrap" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/environmentconfig" ) @@ -77,6 +78,10 @@ type EnvironmentInitialization struct { EnvironmentID, SessionID, TenantID, DeviceID, State, Engine string // Assignment is the Session's bound assignment to DeviceID. Assignment proto.AssignmentRef + // Resource is the Environment's Link resource, or the zero Resource when + // it has none; ResourceLive reports whether it may be served. + Resource sandboxbootstrap.Resource + ResourceLive bool } // EnvironmentInputActivity is the reservation-owned override before a newer Turn exists. diff --git a/services/core/internal/sessions/execution_environment.go b/services/core/internal/sessions/execution_environment.go index 3392c0ac9..3256015f6 100644 --- a/services/core/internal/sessions/execution_environment.go +++ b/services/core/internal/sessions/execution_environment.go @@ -81,6 +81,9 @@ type InitializationTx interface { // ClaimInitialization starts the pending preparation of the live // Environment and reports whether it did. ClaimInitialization(ctx context.Context, environment string) (bool, error) + // UnclaimInitialization returns the running preparation of the live + // Environment to pending and reports whether it did. + UnclaimInitialization(ctx context.Context, environment string) (bool, error) // CompleteInitialization completes the running preparation of the live // Environment and reports whether it did. CompleteInitialization(ctx context.Context, environment string) (bool, error) @@ -199,6 +202,14 @@ func (o *ExecutionOperations) ClaimEnvironmentInitialization(ctx context.Context }) } +// UnclaimEnvironmentInitialization returns the claimed preparation to +// pending, for a preparation that ended before any step took effect. +func (o *ExecutionOperations) UnclaimEnvironmentInitialization(ctx context.Context, owner EnvironmentInitialization) error { + return o.advanceInitialization(ctx, owner, func(ctx context.Context, tx InitializationTx, environment Environment) (bool, error) { + return tx.UnclaimInitialization(ctx, environment.ID) + }) +} + // CompleteEnvironmentInitialization completes the claimed preparation. func (o *ExecutionOperations) CompleteEnvironmentInitialization(ctx context.Context, owner EnvironmentInitialization) error { return o.advanceInitialization(ctx, owner, func(ctx context.Context, tx InitializationTx, environment Environment) (bool, error) { diff --git a/services/core/internal/sessions/execution_environment_test.go b/services/core/internal/sessions/execution_environment_test.go index fb1e66c5d..d62ecad03 100644 --- a/services/core/internal/sessions/execution_environment_test.go +++ b/services/core/internal/sessions/execution_environment_test.go @@ -21,6 +21,11 @@ func (f *fakeTx) ClaimInitialization(_ context.Context, environment string) (boo return f.claimInitialization() } +func (f *fakeTx) UnclaimInitialization(_ context.Context, environment string) (bool, error) { + f.record("UnclaimInitialization", f.unclaimInitialization != nil, environment) + return f.unclaimInitialization() +} + func (f *fakeTx) CompleteInitialization(_ context.Context, environment string) (bool, error) { f.record("CompleteInitialization", f.completeInitialization != nil, environment) return f.completeInitialization() diff --git a/services/core/internal/sessions/transaction_test.go b/services/core/internal/sessions/transaction_test.go index 318013ab3..0fba5f086 100644 --- a/services/core/internal/sessions/transaction_test.go +++ b/services/core/internal/sessions/transaction_test.go @@ -41,6 +41,7 @@ type fakeTx struct { insertEnvironmentDevice func() error loadSessionDevice func() (ExecutionDevice, bool, error) claimInitialization func() (bool, error) + unclaimInitialization func() (bool, error) completeInitialization func() (bool, error) failInitialization func() error loadConnection func() (EnvironmentConnection, bool, error) diff --git a/services/core/migrations/000096_session_runtime_assignments.sql b/services/core/migrations/000096_session_runtime_assignments.sql index dd9243aac..a7e5c6bc6 100644 --- a/services/core/migrations/000096_session_runtime_assignments.sql +++ b/services/core/migrations/000096_session_runtime_assignments.sql @@ -1,8 +1,9 @@ -- +goose Up -- A Session's binding to its Runtime becomes a fenced assignment. Core --- advances epoch with each change of desired_state; applied_epoch is the --- latest released epoch the Runtime acknowledged, or that no Runtime can act --- on because its Runtime lost authority. +-- advances epoch when it releases the assignment and when a later release +-- adds home removal; applied_epoch is the latest released epoch the Runtime +-- acknowledged, or that no Runtime can act on because its Runtime lost +-- authority. ALTER TABLE session_devices RENAME TO session_runtime_assignments; ALTER TABLE session_runtime_assignments RENAME COLUMN device_id TO runtime_id; ALTER INDEX session_devices_pkey RENAME TO session_runtime_assignments_pkey; diff --git a/services/core/tests/integration/environment_initialization_test.go b/services/core/tests/integration/environment_initialization_test.go index d21de7390..994add168 100644 --- a/services/core/tests/integration/environment_initialization_test.go +++ b/services/core/tests/integration/environment_initialization_test.go @@ -94,7 +94,7 @@ func TestUserManagedPreparationUsesAuthenticatedRuntimeWithoutAllocation(t *test var mu sync.Mutex var actions []string peer := &initializationPeer{unavailable: outcome == "unavailable"} - peer.setRuntimeGateway(t, "ws"+strings.TrimPrefix(server.URL, "http"), registry) + peer.setRuntimeGateway(t, "ws"+strings.TrimPrefix(server.URL, "http"), registry, nil) peer.apply = func(request proto.RuntimePreparePayload, data []byte) proto.RuntimePrepareResultPayload { if request.EnvironmentID != environment.ID || request.SessionID != session.ID { t.Error("wrong authorization binding") diff --git a/services/core/tests/integration/hosted_initialization_failure_public_test.go b/services/core/tests/integration/hosted_initialization_failure_public_test.go index 4364a0e8b..04ec8f9ca 100644 --- a/services/core/tests/integration/hosted_initialization_failure_public_test.go +++ b/services/core/tests/integration/hosted_initialization_failure_public_test.go @@ -18,6 +18,7 @@ import ( v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/sandboxlinktest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/credentialcrypto" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/environmentconfig" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimedevice" @@ -64,8 +65,8 @@ type hostedFailureProvider struct { steps []string } -func (p *hostedFailureProvider) setRuntimeGateway(t *testing.T, endpoint string, registry *runtimegateway.Registry) { - p.initializationPeer.setRuntimeGateway(t, endpoint, registry) +func (p *hostedFailureProvider) setRuntimeGateway(t *testing.T, endpoint string, registry *runtimegateway.Registry, link *sandboxlinktest.Server) { + p.initializationPeer.setRuntimeGateway(t, endpoint, registry, link) p.apply = p.prepare } func (p *hostedFailureProvider) Create(ctx context.Context, b sandbox.Bootstrap) (sandbox.Info, error) { diff --git a/services/core/tests/integration/link_authority_test.go b/services/core/tests/integration/link_authority_test.go index df0f83bc7..5a31fc3b0 100644 --- a/services/core/tests/integration/link_authority_test.go +++ b/services/core/tests/integration/link_authority_test.go @@ -132,7 +132,8 @@ func (l *linkHarness) exec(query string, args ...any) { // 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. +// bind, so bind reopens the initialization, which has nothing to install. The +// Worker claims it only while the Environment's resource is Serving. func (l *linkHarness) bind() (proto.AssignmentRef, proto.AssignmentBindPayload) { t := l.t t.Helper() @@ -333,6 +334,7 @@ func TestLinkAuthorityAgentHost(t *testing.T) { // carries no grant and that its device can neither attach nor be marked. func TestLinkAuthorityGuest(t *testing.T) { l := newLinkHarness(t, true) + within(t, startLinkServe(t, l.relay, l.serve, l.resource.Ref()).connected) if _, payload := l.bind(); payload.Resource != nil || payload.AttachGrant != nil { t.Fatalf("guest bind = %+v", payload) } @@ -550,61 +552,61 @@ func TestRegisteredAgentHostAuthenticates(t *testing.T) { } // 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. +// agent host and runs the Environment's initialization once its Link resource +// is Serving. The bind carries the resource and an attach grant. A bind that +// fails before any effect, here because the agent host's connection closes, +// leaves the initialization unclaimed, and a later pass completes it. 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") - }) + 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} + device, err := sessionService(t, s).CreateDevice(t.Context(), tenant, "sandbox", runtimedevice.HashCredential(uuid.NewString())) + if err != nil { + t.Fatal(err) + } + serve := []byte(uuid.NewString()) + insertAllocation(t, s, resource, device.ID, serve) + 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) }) + link := startLinkRoute(t, s) + runWorker(t, startWorker(t, t.Context(), s, &execution.Dispatcher{Registry: registry, Links: link.Relay})) + within(t, startLinkServe(t, link, serve, resource.Ref()).connected) + + dropped := &initializationPeer{apply: completedInitialization, binds: make(chan proto.AssignmentBindPayload, 1), closeOnBind: true} + dropped.setRuntimeGateway(t, endpoint, registry, nil) + if err := dropped.connect(sandbox.Bootstrap{DeviceID: host, Credential: credential}); err != nil { + t.Fatal(err) + } + within(t, dropped.binds) + awaitInitialization(t, s, tenant, session.Environment.ID, "pending") + + peer := &initializationPeer{apply: completedInitialization, binds: make(chan proto.AssignmentBindPayload, 1)} + peer.setRuntimeGateway(t, endpoint, registry, nil) + if err := peer.connect(sandbox.Bootstrap{DeviceID: host, Credential: credential}); err != nil { + t.Fatal(err) + } + 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_connection_test.go b/services/core/tests/integration/runtime_connection_test.go index 8b2ec41ce..e965938da 100644 --- a/services/core/tests/integration/runtime_connection_test.go +++ b/services/core/tests/integration/runtime_connection_test.go @@ -3,88 +3,61 @@ package integration import ( "context" "errors" - "net/http" - "net/http/httptest" - "net/url" "testing" - "time" "github.com/google/uuid" - "github.com/gorilla/websocket" - "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/sandboxlinktest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution" - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtime" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" ) -func TestManagedRuntimeConnectionTracksAuthenticatedSocket(t *testing.T) { +// TestManagedRuntimeConnectionFollowsServe checks that a hosted Environment is +// connected while the relay holds its allocation's serve peer, and that a +// restarted Worker publishes it connected again. +func TestManagedRuntimeConnectionFollowsServe(t *testing.T) { s, _ := newManagedTestStore(t) key := webDeployment(t, s, "e2b") tenant, session, environment := managedSession(t, s) - server := httptest.NewUnstartedServer(nil) - wsURL := "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)), wsURL) - if err != nil { - t.Fatal(err) - } - server.Config.Handler = handler - server.Start() - t.Cleanup(func() { server.Close(); runtime.CloseConnections(registry) }) + link := sandboxlinktest.StartRelay(t, runtimegateway.NewLinkAuthority(sessionAdapter(s))) p := &lifecycleProvider{resources: map[string]sandbox.Info{}} - start := func() *execution.Worker { return startWebWorker(t, s, registry, key, p, nil) } - stop := func(w *execution.Worker) { - ctx, cancel := context.WithCancel(context.Background()) - cancel() - _ = w.Run(ctx) + start := func() (*execution.Worker, func()) { + w, err := startNextWorker(t.Context(), s, &execution.Dispatcher{Registry: runtimegateway.NewRegistry(), Links: link.Relay, ManagedRuntimes: webRuntimes(t, s, key, p, nil)}) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithCancel(t.Context()) + done := make(chan error, 1) + go func() { done <- w.Run(ctx) }() + stop := func() { + cancel() + if err := <-done; err != nil && !errors.Is(err, context.Canceled) { + t.Error(err) + } + } + return w, stop } - w := start() - t.Cleanup(func() { stop(w) }) + w, stop := start() owner, err := w.ProvisionEnvironment(t.Context(), tenant, environment.ID, key) if err != nil { t.Fatal(err) } - assertStatus := func(want string) { - t.Helper() - for deadline := time.Now().Add(3 * time.Second); time.Now().Before(deadline); { - if err := w.ReconcileManagedRuntimes(t.Context()); err != nil { - t.Fatal(err) - } - got, err := sessionAdapter(s).GetEnvironment(t.Context(), tenant, environment.ID) - if err != nil { - t.Fatal(err) - } - if got.Status == want { - return - } - time.Sleep(10 * time.Millisecond) - } - t.Fatal("Environment did not reach", want) + if got, err := sessionAdapter(s).GetEnvironment(t.Context(), tenant, environment.ID); err != nil || got.Status != "pending" { + t.Fatal("compute existence connected the Environment", got.Status, err) } - dial := func(token string) (*websocket.Conn, error) { - u, _ := url.Parse(wsURL) - u.RawQuery = url.Values{"device_id": {owner.DeviceID}, "version": {proto.Version}}.Encode() - conn, response, err := websocket.DefaultDialer.Dial(u.String(), http.Header{"Authorization": {"Bearer " + token}}) - if response != nil && response.Body != nil { - response.Body.Close() - } - return conn, err - } - assertStatus("pending") // Compute existence alone is insufficient. - if conn, err := dial(uuid.NewString()); err == nil { - conn.Close() - t.Fatal("unrelated credential connected") + p.mu.Lock() + io := p.serve + p.mu.Unlock() + serve := func() *linkServe { + serve := startLinkServe(t, link, []byte(io.Credential), io.Resource.Ref()) + within(t, serve.connected) + return serve } - assertStatus("pending") - conn, err := dial(p.credential) - if err != nil { - t.Fatal("authorized connection failed") - } - defer conn.Close() - assertStatus("connected") + served := serve() + awaitEnvironmentConnectionState(t, t.Context(), s, tenant, environment.ID, "connected") got, err := sessionAdapter(s).GetSession(t.Context(), tenant, session.ID) if err != nil || got.LastTurn != nil || got.EnvironmentInputActivity != nil { t.Fatal("connection fabricated native execution", err) @@ -92,29 +65,16 @@ func TestManagedRuntimeConnectionTracksAuthenticatedSocket(t *testing.T) { if _, err := sessionAdapter(s).GetEnvironment(t.Context(), uuid.NewString(), environment.ID); !errors.Is(err, sessions.ErrNotFound) { t.Fatal("foreign Environment access", err) } - p.unavailable = true - conn.Close() - assertStatus("disconnected") // Provider outage cannot conceal socket loss. - p.unavailable = false - conn, err = dial(p.credential) - if err != nil { - t.Fatal("authorized reconnection failed") - } - defer conn.Close() - assertStatus("connected") - stop(w) - w = start() - assertStatus("connected") + served.stop() + awaitEnvironmentConnectionState(t, t.Context(), s, tenant, environment.ID, "disconnected") + serve() + awaitEnvironmentConnectionState(t, t.Context(), s, tenant, environment.ID, "connected") + stop() + _, stop = start() + t.Cleanup(stop) + awaitEnvironmentConnectionState(t, t.Context(), s, tenant, environment.ID, "connected") retained, err := deploymentStore(s).EnvironmentAllocation(t.Context(), deployment.AllocationKey{TenantID: tenant, EnvironmentID: environment.ID}) - if err != nil || retained.ID != owner.ID || retained.DeviceID != owner.DeviceID || p.creates != 1 { - t.Fatal("restart replaced Runtime identity", err) - } - if err := sessionService(t, s).DeleteSession(t.Context(), sessions.DeleteSessionCommand{TenantID: tenant, SessionID: session.ID}); err != nil { - t.Fatal(err) - } - reconcileManagedState(t, w, s, tenant, environment.ID, "released") - if conn, err := dial(p.credential); err == nil { - conn.Close() - t.Fatal("released Runtime reconnected") + if err != nil || retained.ID != owner.ID || p.creates != 1 { + t.Fatal("restart replaced the allocation", err) } } diff --git a/services/core/tests/integration/runtime_initialization_peer_test.go b/services/core/tests/integration/runtime_initialization_peer_test.go index a4ed781da..fd5facab1 100644 --- a/services/core/tests/integration/runtime_initialization_peer_test.go +++ b/services/core/tests/integration/runtime_initialization_peer_test.go @@ -12,17 +12,20 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto/prototest" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/sandboxlinktest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" "github.com/gorilla/websocket" ) // initializationPeer exercises the real authenticated gateway and chunk receipts. -// The provider fixture bootstraps its socket; all initialization runs on that peer. +// The provider fixture bootstraps its socket, and Serves the bootstrap's Link +// resource at link; all initialization runs on that peer. type initializationPeer struct { t *testing.T endpoint string registry *runtimegateway.Registry + link *sandboxlinktest.Server apply func(proto.RuntimePreparePayload, []byte) proto.RuntimePrepareResultPayload writes atomic.Int32 commandCalls atomic.Int32 @@ -30,16 +33,24 @@ type initializationPeer struct { unavailable bool bootstrap sandbox.Bootstrap binds chan proto.AssignmentBindPayload // when not nil, receives each bind's payload + closeOnBind bool // close the socket at a bind instead of replying } -func (p *initializationPeer) setRuntimeGateway(t *testing.T, endpoint string, registry *runtimegateway.Registry) { - p.t, p.endpoint, p.registry = t, endpoint, registry +func (p *initializationPeer) setRuntimeGateway(t *testing.T, endpoint string, registry *runtimegateway.Registry, link *sandboxlinktest.Server) { + p.t, p.endpoint, p.registry, p.link = t, endpoint, registry, link } func (p *initializationPeer) connect(b sandbox.Bootstrap) error { p.bootstrap = b if p.deferred { return nil } + if p.link != nil && b.SandboxIO.Credential != "" { + select { + case <-startLinkServe(p.t, p.link, []byte(b.SandboxIO.Credential), b.SandboxIO.Resource.Ref()).connected: + case <-time.After(linkWait): + return context.DeadlineExceeded + } + } c, _, err := websocket.DefaultDialer.Dial(p.endpoint+"?device_id="+b.DeviceID+"&version="+proto.Version, http.Header{"Authorization": {"Bearer " + b.Credential}}) if err != nil { return err @@ -62,6 +73,10 @@ func (p *initializationPeer) connect(b sandbox.Bootstrap) error { if p.binds != nil && env.DecodePayload(&bind) == nil { p.binds <- bind } + if p.closeOnBind { + _ = c.Close() + return + } if c.WriteJSON(reply) != nil { return } diff --git a/services/core/tests/integration/runtime_initialization_test.go b/services/core/tests/integration/runtime_initialization_test.go index 8acc75925..51fc6b646 100644 --- a/services/core/tests/integration/runtime_initialization_test.go +++ b/services/core/tests/integration/runtime_initialization_test.go @@ -106,7 +106,7 @@ func TestEnvironmentInitializationCompletionUnknownAndRestart(t *testing.T) { t.Fatal("missing socket consumed initialization") } p.deferred = false - if err := p.connect(sandbox.Bootstrap{DeviceID: owner.DeviceID, Credential: p.credential}); err != nil { + if err := p.connect(p.bootstrap); err != nil { t.Fatal(err) } want := "complete" diff --git a/services/core/tests/integration/runtime_lifecycle_test.go b/services/core/tests/integration/runtime_lifecycle_test.go index 770504f8f..a9dab0db4 100644 --- a/services/core/tests/integration/runtime_lifecycle_test.go +++ b/services/core/tests/integration/runtime_lifecycle_test.go @@ -13,6 +13,7 @@ import ( "github.com/google/uuid" "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxbootstrap" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/sandboxlinktest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimedevice" @@ -85,16 +86,21 @@ func managedWorker(t *testing.T, s *Store, key string, p sandbox.SandboxProvider // managedWorkerMode is managedWorker that also runs the Worker when run is set. func managedWorkerMode(t *testing.T, s *Store, key string, p sandbox.SandboxProvider, run bool) (*execution.Worker, func()) { t.Helper() - registry := runtimegateway.NewRegistry() + dispatcher := &execution.Dispatcher{Registry: runtimegateway.NewRegistry(), ManagedRuntimes: webRuntimes(t, s, key, p, nil)} if peer, ok := p.(interface { - setRuntimeGateway(*testing.T, string, *runtimegateway.Registry) + setRuntimeGateway(*testing.T, string, *runtimegateway.Registry, *sandboxlinktest.Server) }); ok { - handler := runtimegateway.NewHandler(runtimegateway.HandlerConfig{Authenticator: runtimegateway.NewAuthenticator(sessionAdapter(s)), Registry: registry}) + handler := runtimegateway.NewHandler(runtimegateway.HandlerConfig{Authenticator: runtimegateway.NewAuthenticator(sessionAdapter(s)), Registry: dispatcher.Registry}) server := httptest.NewServer(http.HandlerFunc(handler.WS)) t.Cleanup(server.Close) - peer.setRuntimeGateway(t, "ws"+strings.TrimPrefix(server.URL, "http"), registry) + link := sandboxlinktest.StartRelay(t, runtimegateway.NewLinkAuthority(sessionAdapter(s))) + peer.setRuntimeGateway(t, "ws"+strings.TrimPrefix(server.URL, "http"), dispatcher.Registry, link) + dispatcher.Links = link.Relay + } + w, err := startNextWorker(t.Context(), s, dispatcher) + if err != nil { + t.Fatal(err) } - w := startWebWorker(t, s, registry, key, p, nil) if run { ctx, cancel := context.WithCancel(t.Context()) done := make(chan error, 1) diff --git a/services/core/tests/integration/runtime_node_lifecycle_fixture_test.go b/services/core/tests/integration/runtime_node_lifecycle_fixture_test.go index 80769a8ce..4dd787771 100644 --- a/services/core/tests/integration/runtime_node_lifecycle_fixture_test.go +++ b/services/core/tests/integration/runtime_node_lifecycle_fixture_test.go @@ -14,6 +14,7 @@ import ( "time" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/sandboxlinktest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/credentialcrypto" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/environmentconfig" @@ -26,8 +27,10 @@ import ( "github.com/jackc/pgx/v5/pgxpool" ) +// nodeIsolationProvider Serves each allocation it creates at link. type nodeIsolationProvider struct { *fakeCheckpointProvider + link *sandboxlinktest.Server blockMu sync.Mutex blocked map[string]bool mode string @@ -55,6 +58,13 @@ func (p *nodeIsolationProvider) Create(ctx context.Context, b sandbox.Bootstrap) if err == nil { err = p.connect(ctx, b) } + if err == nil { + select { + case <-startLinkServe(p.preparation.t, p.link, []byte(b.SandboxIO.Credential), b.SandboxIO.Resource.Ref()).connected: + case <-ctx.Done(): + err = ctx.Err() + } + } return info, err } func (p *nodeIsolationProvider) GetInfo(ctx context.Context, r sandbox.Reference) (sandbox.Info, error) { @@ -98,7 +108,7 @@ func newNodeIsolationFixture(t *testing.T, mode string) *nodeIsolationFixture { s := NewWithCredentialCipher(pool, cipher) registry := runtimegateway.NewRegistry() cp := &fakeCheckpointProvider{lifecycleProvider: lifecycleProvider{resources: map[string]sandbox.Info{}}, computes: map[string]sandbox.ComputeState{}, snapshots: map[string]sandbox.SnapshotIdentity{}, bootstraps: map[string]sandbox.Bootstrap{}, peers: map[string]*websocket.Conn{}, registry: registry} - p := &nodeIsolationProvider{fakeCheckpointProvider: cp, blocked: map[string]bool{}, mode: mode, entered: make(chan struct{})} + p := &nodeIsolationProvider{fakeCheckpointProvider: cp, link: sandboxlinktest.StartRelay(t, runtimegateway.NewLinkAuthority(sessionAdapter(s))), blocked: map[string]bool{}, mode: mode, entered: make(chan struct{})} preparationContext, cancelPreparation := context.WithCancel(t.Context()) t.Cleanup(cancelPreparation) cp.preparation = &initializationPeer{t: t, apply: func(request proto.RuntimePreparePayload, data []byte) proto.RuntimePrepareResultPayload { @@ -126,7 +136,10 @@ func newNodeIsolationFixture(t *testing.T, mode string) *nodeIsolationFixture { f := &nodeIsolationFixture{initializationCancel: cancelPreparation, t: t, store: s, nodes: deploymentService(t, s), pool: pool, provider: p, key: webDeployment(t, s, "microsandbox"), nodeA: uuid.NewString(), nodeB: uuid.NewString()} // Keep restored compute awake throughout the isolation assertions. // The suspension setup explicitly dates its activity two minutes in the past. - f.worker = startWebWorker(t, s, registry, f.key, p, &execution.RuntimeSuspensionPolicy{IdleTimeout: time.Minute, Retention: time.Hour}) + f.worker, err = startNextWorker(t.Context(), s, &execution.Dispatcher{Registry: registry, Links: p.link.Relay, ManagedRuntimes: webRuntimes(t, s, f.key, p, &execution.RuntimeSuspensionPolicy{IdleTimeout: time.Minute, Retention: time.Hour})}) + if err != nil { + t.Fatal(err) + } t.Cleanup(f.stop) f.epoch = fixtureOwnerEpoch(t, s) f.enroll(f.nodeA)