diff --git a/apps/daemon/internal/dispatch/assignment.go b/apps/daemon/internal/dispatch/assignment.go index 209fc4be6..ae5b82b41 100644 --- a/apps/daemon/internal/dispatch/assignment.go +++ b/apps/daemon/internal/dispatch/assignment.go @@ -1,11 +1,13 @@ package dispatch import ( + "bytes" "context" "sync" "time" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxbootstrap" ) // assignmentState is the Router's record of one Session's assignment. Router.mu @@ -13,7 +15,11 @@ import ( type assignmentState struct { ref proto.AssignmentRef environmentID string - released bool + // resource and grant are the bind's Link resource and attach grant; the + // resource's Kind is empty when the bind carried none. + resource sandboxbootstrap.Resource + grant []byte + released bool // work counts the Session's admitted reads, writes, exports and Runtime // preparations until each has sent its terminal result. A release waits // for it, and the released assignment admits no more. @@ -55,17 +61,21 @@ func (r *Router) admitRunLocked(ref proto.AssignmentRef, state *sessionState) st func (r *Router) handleAssignmentBind(ctx context.Context, env proto.Envelope) error { var input proto.AssignmentBindPayload ref, code := env.Assignment, "" - if env.ID == "" || env.DecodeRequest(&input) != nil || !ref.Valid() { + if env.ID == "" || env.DecodeRequest(&input) != nil || input.Validate() != nil || !ref.Valid() { code = "invalid_request" } else { + var resource sandboxbootstrap.Resource + if input.Resource != nil { + resource = *input.Resource + } r.mu.Lock() a := r.assignments[ref.SessionID] switch { case a == nil: - r.assignments[ref.SessionID] = &assignmentState{ref: ref, environmentID: input.EnvironmentID} + r.assignments[ref.SessionID] = &assignmentState{ref: ref, environmentID: input.EnvironmentID, resource: resource, grant: input.AttachGrant} case a.ref.AssignmentID == ref.AssignmentID && (ref.Epoch < a.ref.Epoch || ref.Epoch == a.ref.Epoch && a.released): code = proto.AssignmentStale - case a.ref != ref || a.environmentID != input.EnvironmentID: + case a.ref != ref || a.environmentID != input.EnvironmentID || a.resource != resource || !bytes.Equal(a.grant, input.AttachGrant): code = proto.AssignmentConflict } r.mu.Unlock() diff --git a/apps/daemon/internal/dispatch/assignment_test.go b/apps/daemon/internal/dispatch/assignment_test.go index fbef907e8..c6d585e6a 100644 --- a/apps/daemon/internal/dispatch/assignment_test.go +++ b/apps/daemon/internal/dispatch/assignment_test.go @@ -6,9 +6,12 @@ import ( "testing" "time" + "github.com/google/uuid" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/dispatch" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxbootstrap" ) // frameFor returns the last frame of kind correlated with id. @@ -168,3 +171,39 @@ func TestUnknownEnvelopeGetsCorrelatedProtocolError(t *testing.T) { t.Fatalf("protocol_error = %+v %+v", frame, got) } } + +func TestAssignmentBindCarriesLink(t *testing.T) { + h := newHarness(t) + defer h.router.Shutdown(context.Background()) + environment := uuid.NewString() + resource := &sandboxbootstrap.Resource{TenantID: uuid.NewString(), EnvironmentID: environment, Kind: "allocation", ID: uuid.NewString(), Generation: 1} + other := *resource + other.EnvironmentID = uuid.NewString() + bind := func(id string, payload proto.AssignmentBindPayload) proto.AssignmentStatusPayload { + if err := h.router.Handle(t.Context(), scoped(t, "s", proto.TypeAssignmentBind, id, payload)); err != nil { + t.Fatal(err) + } + return waitAssignmentStatus(t, h.sender, id) + } + for id, payload := range map[string]proto.AssignmentBindPayload{ + "no grant": {EnvironmentID: environment, Resource: resource}, + "no resource": {EnvironmentID: environment, AttachGrant: []byte("grant")}, + "other environment": {EnvironmentID: environment, Resource: &other, AttachGrant: []byte("grant")}, + } { + if got := bind(id, payload); got.ErrorCode != "invalid_request" { + t.Fatalf("%s: bind = %+v", id, got) + } + } + link := proto.AssignmentBindPayload{EnvironmentID: environment, Resource: resource, AttachGrant: []byte("grant")} + if got := bind("bind", link); got.State != proto.AssignmentBound { + t.Fatalf("bind = %+v", got) + } + if got := bind("again", link); got.State != proto.AssignmentBound { + t.Fatalf("identical bind = %+v", got) + } + changed := link + changed.AttachGrant = []byte("other grant") + if got := bind("changed", changed); got.ErrorCode != proto.AssignmentConflict { + t.Fatalf("bind with another grant = %+v", got) + } +} diff --git a/docs/runtime-protocol.md b/docs/runtime-protocol.md index 97fef0b42..2ae65bd40 100644 --- a/docs/runtime-protocol.md +++ b/docs/runtime-protocol.md @@ -118,9 +118,9 @@ An assignment binds one Session to the Runtime that runs it. `Envelope.assignmen Every Session frame carries the assignment: `execution_prepare`, `execution_start` and `execution_release`; `prompt_cancel`, `prompt_steer`, `function_result`, `permission_decision` and `prompt_for_user_choice_decision`; every frame of `runtime_prepare`, `workspace_read`, `workspace_write` and `workspace_export`; and `environment_quiesce` and `environment_resume`. A reply echoes its request's assignment, and a Run's frames carry the assignment that started it; Core rejects a reply or Run frame that names another. Heartbeats carry none. -Before a Session's first operation on a connection, including Environment initialization and file work without a Turn, Core sends `assignment_bind` with the Session's Environment ID and waits for `assignment_status` `bound`. A repeated bind of the same assignment is `bound` again. The Runtime admits a Session frame only under the assignment it bound: an older epoch, or a released one, fails with `assignment_stale`; another assignment, Session or Environment fails with `assignment_conflict`. A started Run's frames, including its cancellation receipt, stay admissible under the assignment that started it until the release. A repeated function result or decision whose receipt the Runtime already recorded is answered only under the assignment that applied it; another fails with `assignment_conflict`. +Before a Session's first operation on a connection, including Environment initialization and file work without a Turn, Core sends `assignment_bind` with the Session's Environment ID and waits for `assignment_status` `bound`. When the Runtime is an agent host and the Environment has a live [Link](./sandbox-link-protocol.md) resource, the bind also carries `resource`, that resource as the [bootstrap input](./sandbox-bootstrap.md#launch-input) names it, and `attach_grant`, the base64 grant with which the agent host opens services on that resource generation under this assignment and epoch. The grant is secret. Core sends neither field to any other Runtime. A bind with only one of them, or with a resource of another Environment, fails with `invalid_request`. A repeated bind of the same assignment with the same Environment, resource and grant is `bound` again; any other bind of it fails with `assignment_conflict`. The Runtime admits a Session frame only under the assignment it bound: an older epoch, or a released one, fails with `assignment_stale`; another assignment, Session or Environment fails with `assignment_conflict`. A started Run's frames, including its cancellation receipt, stay admissible under the assignment that started it until the release. A repeated function result or decision whose receipt the Runtime already recorded is answered only under the assignment that applied it; another fails with `assignment_conflict`. -Core records a release and advances the epoch before it sends anything. Deleting a Session releases its assignment with `remove_home: true`; releasing its Environment sends `false`. A deletion never revokes a shared Runtime credential. `assignment_release` fences the assignment at once. The Runtime then stops the Session's work: a transfer still receiving its body, or committed but not yet applied, ends with `assignment_stale`; it releases read-only preparations and waits until every workspace read, write, export and Runtime preparation has sent its result. It closes the Session's Executors and, when asked, removes the native home; only then does it reply `released` or `home_removed`. Unfinished cleanup replies `failed` with `cleanup_unconfirmed`, and a retry at the same epoch repeats it. A Runtime that declares `home_removal` unsupported answers `remove_home: true` with `unsupported_operation`, and Core asks it only to release. Core records the release as applied from a matching `released` or `home_removed`, or at once when no Runtime is left to act on it: a release to a Runtime without authority is settled when recorded, and revoking a Runtime settles its releases. Core resends every unacknowledged release to a Runtime when it connects; a release that fails backs off, and the release due longest goes first, so failing releases cannot delay the rest. A quiesced Runtime admits only a release and the matching `environment_resume`, which carries the assignment that quiesced it. +Core records a release and advances the epoch before it sends anything, which withdraws the assignment's attach grant; it then has the relay revoke the Environment's Link resource at its current generation, so the attachments opened under the grant close before the Runtime receives the release. Deleting a Session releases its assignment with `remove_home: true`; releasing its Environment sends `false`. A deletion never revokes a shared Runtime credential. `assignment_release` fences the assignment at once. The Runtime then stops the Session's work: a transfer still receiving its body, or committed but not yet applied, ends with `assignment_stale`; it releases read-only preparations and waits until every workspace read, write, export and Runtime preparation has sent its result. It closes the Session's Executors and, when asked, removes the native home; only then does it reply `released` or `home_removed`. Unfinished cleanup replies `failed` with `cleanup_unconfirmed`, and a retry at the same epoch repeats it. A Runtime that declares `home_removal` unsupported answers `remove_home: true` with `unsupported_operation`, and Core asks it only to release. Core records the release as applied from a matching `released` or `home_removed`, or at once when no Runtime is left to act on it: a release to a Runtime without authority is settled when recorded, and revoking a Runtime settles its releases. Core resends every unacknowledged release to a Runtime when it connects; a release that fails backs off, and the release due longest goes first, so failing releases cannot delay the rest. A quiesced Runtime admits only a release and the matching `environment_resume`, which carries the assignment that quiesced it. The Runtime answers a Core frame it cannot route with `protocol_error`, which echoes the request's ID and carries its type and an error code. diff --git a/docs/zh/runtime-protocol.md b/docs/zh/runtime-protocol.md index 4b9131b00..0eb136233 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: 5f4fc0a408f12938550bb548f6b95c336316be2320f5cbb8c6801a823d88863e +source_hash: 3d5083634593391e5a3627961d54e1dceb7c4589f74539b3e01b677a9efbe04d --- 此协议在 Runtime daemon 获取机器凭据后连接 Core 与 daemon,定义 daemon 连接上消息的含义和顺序。wire 类型、限制和验证器仅在 [`internal/agentdaemon/proto`](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/internal/agentdaemon/proto) 中定义一次;Core 的 [gateway](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/services/core/internal/runtimegateway) 与参考 Runtime 的 [dispatcher](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/apps/daemon/internal/dispatch) 都使用它们,因此无需同步第二套 payload schema。签发凭据和打开连接的 HTTP 路由见[机器连接 API](../../contracts/agents-api/zh/machine-api.md)。 @@ -120,9 +120,9 @@ Usage frame 和最终 usage snapshot 都携带当前执行的累计测量,替 每个 Session frame 都携带分配:`execution_prepare`、`execution_start` 和 `execution_release`;`prompt_cancel`、`prompt_steer`、`function_result`、`permission_decision` 和 `prompt_for_user_choice_decision`;`runtime_prepare`、`workspace_read`、`workspace_write` 和 `workspace_export` 的每个 frame;以及 `environment_quiesce` 和 `environment_resume`。回复回显请求的分配,Run 的 frame 携带启动它的分配;Core 拒绝指明其他分配的回复或 Run frame。heartbeat 不携带分配。 -在一条连接上执行 Session 的第一个操作之前,包括没有 Turn 的 Environment 初始化和文件操作,Core 发送带 Session 的 Environment ID 的 `assignment_bind`,并等待 `assignment_status` `bound`。重复绑定同一分配仍得到 `bound`。Runtime 只在其已绑定的分配下准入 Session frame:较旧的 epoch 或已释放的分配以 `assignment_stale` 失败;其他分配、Session 或 Environment 以 `assignment_conflict` 失败。已启动 Run 的 frame,包括其取消回执,在释放前仍可在启动它的分配下准入。Runtime 已记录回执的重复函数结果或决策只在应用它的分配下得到回答;其他分配以 `assignment_conflict` 失败。 +在一条连接上执行 Session 的第一个操作之前,包括没有 Turn 的 Environment 初始化和文件操作,Core 发送带 Session 的 Environment ID 的 `assignment_bind`,并等待 `assignment_status` `bound`。当 Runtime 是 agent host 且 Environment 有存活的 [Link](./sandbox-link-protocol.md) resource 时,绑定还携带 `resource` 和 `attach_grant`:前者是该 resource,形式与[引导输入](./sandbox-bootstrap.md#launch-input)中的相同;后者是 base64 编码的 grant,agent host 凭它在此分配和 epoch 下打开该 resource generation 上的服务。grant 是机密。Core 不向其他任何 Runtime 发送这两个字段。只带其中一个字段、或带其他 Environment 的 resource 的绑定以 `invalid_request` 失败。以相同的 Environment、resource 和 grant 重复绑定同一分配仍得到 `bound`;该分配的其他绑定以 `assignment_conflict` 失败。Runtime 只在其已绑定的分配下准入 Session frame:较旧的 epoch 或已释放的分配以 `assignment_stale` 失败;其他分配、Session 或 Environment 以 `assignment_conflict` 失败。已启动 Run 的 frame,包括其取消回执,在释放前仍可在启动它的分配下准入。Runtime 已记录回执的重复函数结果或决策只在应用它的分配下得到回答;其他分配以 `assignment_conflict` 失败。 -Core 先记录释放并推进 epoch,再发送任何消息。删除 Session 以 `remove_home: true` 释放其分配;释放其 Environment 发送 `false`。删除从不吊销共享的 Runtime 凭据。`assignment_release` 立即约束该分配。随后 Runtime 停止 Session 的工作:仍在接收内容、或已提交但尚未应用的传输以 `assignment_stale` 结束;它释放只读准备,并等待每个 workspace 读取、写入、导出和 Runtime 准备发送结果。它关闭 Session 的 Executor,并在要求时删除原生 home;此后才回复 `released` 或 `home_removed`。未完成的清理回复 `failed` 和 `cleanup_unconfirmed`,同一 epoch 的重试会重复清理。声明 `home_removal` 不支持的 Runtime 以 `unsupported_operation` 回答 `remove_home: true`,Core 只要求它释放。Core 根据匹配的 `released` 或 `home_removed` 记录释放已应用;没有 Runtime 能处理该释放时立即记录:发给无授权 Runtime 的释放在记录时即结清,吊销 Runtime 会结清它的释放。Core 在 Runtime 连接时重发所有未确认的释放;失败的释放退避重试,等待最久的释放先发送,因此失败的释放不会拖延其他释放。已 quiesce 的 Runtime 只准入释放和匹配的 `environment_resume`,后者携带使其 quiesce 的分配。 +Core 先记录释放并推进 epoch,再发送任何消息;记录即撤回该分配的 attach grant。随后 Core 让 relay 吊销 Environment 的 Link resource 的当前 generation,使凭该 grant 打开的 attachment 在 Runtime 收到释放之前关闭。删除 Session 以 `remove_home: true` 释放其分配;释放其 Environment 发送 `false`。删除从不吊销共享的 Runtime 凭据。`assignment_release` 立即约束该分配。随后 Runtime 停止 Session 的工作:仍在接收内容、或已提交但尚未应用的传输以 `assignment_stale` 结束;它释放只读准备,并等待每个 workspace 读取、写入、导出和 Runtime 准备发送结果。它关闭 Session 的 Executor,并在要求时删除原生 home;此后才回复 `released` 或 `home_removed`。未完成的清理回复 `failed` 和 `cleanup_unconfirmed`,同一 epoch 的重试会重复清理。声明 `home_removal` 不支持的 Runtime 以 `unsupported_operation` 回答 `remove_home: true`,Core 只要求它释放。Core 根据匹配的 `released` 或 `home_removed` 记录释放已应用;没有 Runtime 能处理该释放时立即记录:发给无授权 Runtime 的释放在记录时即结清,吊销 Runtime 会结清它的释放。Core 在 Runtime 连接时重发所有未确认的释放;失败的释放退避重试,等待最久的释放先发送,因此失败的释放不会拖延其他释放。已 quiesce 的 Runtime 只准入释放和匹配的 `environment_resume`,后者携带使其 quiesce 的分配。 Runtime 对无法路由的 Core frame 回复 `protocol_error`,回显请求 ID,并携带其类型和错误码。 diff --git a/internal/agentdaemon/proto/assignment.go b/internal/agentdaemon/proto/assignment.go index 84470cda0..1dc69fd1d 100644 --- a/internal/agentdaemon/proto/assignment.go +++ b/internal/agentdaemon/proto/assignment.go @@ -1,5 +1,12 @@ package proto +import ( + "errors" + + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxbootstrap" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink" +) + // Assignment traffic binds a Session to the Runtime that runs it. Envelope.ID // correlates a request with its assignment_status, and Envelope.Assignment // names the assignment on both. @@ -53,9 +60,27 @@ func (r AssignmentRef) Valid() bool { } // AssignmentBindPayload is what the Runtime fences the Session's frames with. -// EnvironmentID is empty for a Session without an Environment. +// EnvironmentID is empty for a Session without an Environment. Resource and +// AttachGrant come together, only to an agent host whose Environment has a +// Link resource: the agent host attaches to Resource with AttachGrant, which +// is secret. type AssignmentBindPayload struct { - EnvironmentID string `json:"environment_id,omitempty"` + EnvironmentID string `json:"environment_id,omitempty"` + Resource *sandboxbootstrap.Resource `json:"resource,omitempty"` + AttachGrant []byte `json:"attach_grant,omitempty"` +} + +// Validate checks that Resource and AttachGrant come together and that +// Resource is a resource of the Environment. +func (p AssignmentBindPayload) Validate() error { + if p.Resource == nil && len(p.AttachGrant) == 0 { + return nil + } + if p.Resource == nil || len(p.AttachGrant) == 0 || len(p.AttachGrant) > sandboxlink.MaxGrantBytes || + p.Resource.Validate() != nil || p.Resource.EnvironmentID != p.EnvironmentID { + return errors.New("assignment_bind requires a valid resource of its Environment with an attach grant") + } + return nil } // AssignmentReleasePayload asks the Runtime to remove the Session's native diff --git a/internal/sandboxbootstrap/bootstrap.go b/internal/sandboxbootstrap/bootstrap.go index 7509532d4..96210497e 100644 --- a/internal/sandboxbootstrap/bootstrap.go +++ b/internal/sandboxbootstrap/bootstrap.go @@ -48,7 +48,7 @@ func (in Input) Validate() error { if in.Version != Version || sandboxlink.CheckRelayURL(in.LinkURL) != nil || in.Credential == "" || len(in.Credential) > sandboxlink.MaxCredentialBytes || strings.ContainsFunc(in.Credential, func(r rune) bool { return r == 0 || unicode.IsSpace(r) }) || - in.Resource.validate() != nil { + in.Resource.Validate() != nil { return ErrInvalid } raw, err := json.Marshal(in) @@ -58,7 +58,8 @@ func (in Input) Validate() error { return nil } -func (r Resource) validate() error { +// Validate checks that the resource names one Link resource. +func (r Resource) Validate() error { for _, id := range []string{r.TenantID, r.EnvironmentID, r.ID} { if u, err := uuid.Parse(id); err != nil || u == uuid.Nil || u.String() != id { return ErrInvalid diff --git a/scripts/build-core.sh b/scripts/build-core.sh index 3feabc800..a0baf5f16 100755 --- a/scripts/build-core.sh +++ b/scripts/build-core.sh @@ -26,7 +26,7 @@ tar -C "$repo_root" -cf - \ go.mod go.sum \ contracts/agents-api/v1 \ contracts/agents-api/openapi.go contracts/agents-api/openapi.yaml contracts/agents-api/core.openapi.yaml contracts/agents-api/runtime.openapi.yaml \ - internal/agentdaemon/proto \ + internal/agentdaemon/proto internal/sandboxlink internal/sandboxwire internal/sandboxbootstrap internal/sandboxfs \ internal/runtimefs internal/runtimebootstrap internal/agentnetwork internal/agentbundle internal/agentcapabilities internal/agentplugin internal/agentskill internal/harnessconfig internal/modelprovider internal/providerassets internal/obs/log services/core \ | tar -C "$build_context" -xf - diff --git a/services/core/IMPLEMENTATION.md b/services/core/IMPLEMENTATION.md index c6d817b57..486d1814c 100644 --- a/services/core/IMPLEMENTATION.md +++ b/services/core/IMPLEMENTATION.md @@ -92,6 +92,8 @@ The Worker scans pending inputs with the same scheduling slots, Session locks, d `services/core/internal/runtimegateway` is Core's daemon connection implementation; its persistence interfaces use `services/core/internal/runtimedevice`, and the frames and validators live in the shared `internal/agentdaemon/proto`. It is a single-process registry: connectivity comes from the live registry, never a persisted online flag, and `last_seen_at` is diagnostic only. Session-to-device bindings are tenant-scoped and immutable. Revocation denies new connections and binding reads at once, and an open connection closes at its next heartbeat. +`runtimegateway.LinkAuthority` is Core's [Link](../../docs/sandbox-link-protocol.md) `Authority`. `cmd/server` builds it over `sessionpg.Store` and gives it to one `relay.New`, which mounts no route yet; the Worker revokes through that relay as `execution.Dispatcher.Links`. Every Hello, Open and renewal rereads the database. A Serve credential authenticates only its own resource while the `sandbox_resources` view marks it live, at that resource's current generation: an allocation in `creating` or `running` by its `serve_credential_hash`, or a `sandbox_enrollments` row by its executor key while the key would still authenticate for the Environment. Rotating the key advances the generation of each of its enrollments in the same transaction, so Opens and renewals for the old generation are refused and its attachments end within one lease. A key without an Environment restriction may hold several enrollments, and its holder is trusted for every Environment the key authenticates for: it may Serve any of them, replacing that enrollment's serve peer. Only a device marked `agent_host` Attaches, and its `credential_revision` is the peer's `Revision`. An attach grant is the assignment ID, epoch and resource generation followed by their keyed digest under the credential key, so Core stores none. It opens a service only while its assignment is the Session's bound assignment at that epoch, held by the peer's Runtime, and its generation is current. File gets the `world` export, Network gets every destination while the Session's network access is enabled and none otherwise, and each lease lasts one minute. Cleanup of an allocation and a release first commit, which withdraws their authority, then revoke the resource at the relay, and only then destroy the compute or send the release. + A dedicated self-hosted device is bound to exactly one Environment's Session and is excluded from general device selection, even within the tenant. Enrollment creates or recovers the device and binding atomically under the Session lock; the frozen workspace and capability directories come from the Session configuration and must match the local binding. Core rechecks the persisted Environment and device binding for preparation and active reads; capability discovery never selects or authorizes a device for this placement. Connection observations use the execution lease and the Session lock. A separate `environment_connections` row holds the current generation and revision, and `environments.status` commits together with its Session Environment event. The producer serializes replacements and numbers socket observations within each generation; duplicate or older revisions and superseded generations are inert, and a replacement retires the previous connected observation before registering the new one. Registration alone creates no `connected` event. Event payloads carry only public Environment identity, type, status and nullable error, never configuration, credentials, registration IDs or revisions, and have no Turn association. `connected` and `disconnected` are distinct from native readiness: never cast resource `expired` into this vocabulary or emit `ready` for a self-hosted connection. On restart the Worker reconciles old observations before admitting new ones, and a failed or stale observation never establishes a connection. diff --git a/services/core/cmd/server/main.go b/services/core/cmd/server/main.go index f734986f2..31773c625 100644 --- a/services/core/cmd/server/main.go +++ b/services/core/cmd/server/main.go @@ -32,6 +32,7 @@ import ( "time" "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/relay" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/agents" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/api" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/coremetrics" @@ -273,6 +274,7 @@ func run() error { } var daemonHandler http.Handler var registry *runtimegateway.Registry + var linkRelay *relay.Relay var executorURL string var nativeInstaller *api.NativeInstaller if public != "" { @@ -280,11 +282,14 @@ func run() error { if err != nil { return err } - daemonHandler, registry, err = runtime.NewGateway(sessionStore, sessionService, sessionStore, executorURL) + links := runtimegateway.NewLinkAuthority(sessionStore) + daemonHandler, registry, err = runtime.NewGateway(sessionStore, sessionService, sessionStore, links, executorURL) if err != nil { return err } defer runtime.CloseConnections(registry) + linkRelay = relay.New(links) + defer linkRelay.Close() var catalog *nativeinstaller.Catalog if directory := os.Getenv("OAC_NATIVE_INSTALLER_DIR"); directory != "" { catalog, err = nativeinstaller.Load(directory, buildRevision) @@ -302,6 +307,7 @@ func run() error { Credentials: vaultService, Observer: modelConfigurationStore, Deployment: deploymentService, DeploymentReader: deploymentStore, Sessions: sessionService, SessionsReader: sessionStore, + Links: linkRelay, ManagedRuntimes: managed, MaxConcurrentExecutions: concurrency} lease, err := pgunit.AcquireLease(ctx, pool) if err != nil { diff --git a/services/core/internal/db/queries/devices.sql b/services/core/internal/db/queries/devices.sql index fc0f5f362..47d2b5ca4 100644 --- a/services/core/internal/db/queries/devices.sql +++ b/services/core/internal/db/queries/devices.sql @@ -80,9 +80,16 @@ SET desired_state = 'released', epoch = b.epoch + 1, remove_home = b.remove_home WHERE b.session_id = sqlc.arg(session_id) AND (b.desired_state = 'bound' OR (sqlc.arg(remove_home)::boolean AND NOT b.remove_home)); -- name: ListPendingAssignmentReleases :many -SELECT session_id, runtime_id, assignment_id, epoch, remove_home FROM session_runtime_assignments -WHERE desired_state = 'released' AND applied_epoch < epoch AND runtime_id = ANY(sqlc.arg(runtime_ids)::uuid[]) -ORDER BY runtime_id, session_id; +-- Each release carries its Session's Link resource, live or not, which the +-- release revokes before it is sent. +SELECT b.session_id, b.runtime_id, b.assignment_id, b.epoch, b.remove_home, + 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 +FROM session_runtime_assignments b +LEFT JOIN environments e ON e.session_id = b.session_id +LEFT JOIN sandbox_resources r ON r.environment_id = e.id +WHERE b.desired_state = 'released' AND b.applied_epoch < b.epoch AND b.runtime_id = ANY(sqlc.arg(runtime_ids)::uuid[]) +ORDER BY b.runtime_id, b.session_id; -- name: AcknowledgeAssignmentRelease :execrows UPDATE session_runtime_assignments SET applied_epoch = epoch diff --git a/services/core/internal/db/queries/environment_executor_credentials.sql b/services/core/internal/db/queries/environment_executor_credentials.sql index 241875d79..a8b4e484f 100644 --- a/services/core/internal/db/queries/environment_executor_credentials.sql +++ b/services/core/internal/db/queries/environment_executor_credentials.sql @@ -30,11 +30,19 @@ WHERE c.key_id = sqlc.arg(key_id) AND c.tenant_id = sqlc.arg(tenant_id) AND p.organization_id = sqlc.arg(organization_id) AND p.project_id = sqlc.arg(project_id); -- name: RotateExecutorCredential :one -UPDATE environment_executor_credentials -SET token_sha256 = sqlc.arg(token_sha256), issued_at = clock_timestamp(), revoked_at = NULL -WHERE key_id = sqlc.arg(key_id) AND tenant_id = sqlc.arg(tenant_id) - AND subject_kind = sqlc.arg(subject_kind) AND subject_id = sqlc.arg(subject_id) -RETURNING key_id, environment_id; +-- Rotation advances the generation of the key's enrollments, so the Link +-- authority refuses what the old secret served. +WITH rotated AS ( + UPDATE environment_executor_credentials + SET token_sha256 = sqlc.arg(token_sha256), issued_at = clock_timestamp(), revoked_at = NULL + WHERE key_id = sqlc.arg(key_id) AND tenant_id = sqlc.arg(tenant_id) + AND subject_kind = sqlc.arg(subject_kind) AND subject_id = sqlc.arg(subject_id) + RETURNING key_id, environment_id +), advanced AS ( + UPDATE sandbox_enrollments n SET generation = n.generation + 1 + FROM rotated r WHERE n.executor_key_id = r.key_id +) +SELECT key_id, environment_id FROM rotated; -- name: RevokeExecutorCredential :execrows UPDATE environment_executor_credentials SET revoked_at = COALESCE(revoked_at, clock_timestamp()) diff --git a/services/core/internal/db/queries/sandbox_link.sql b/services/core/internal/db/queries/sandbox_link.sql new file mode 100644 index 000000000..3126f6c34 --- /dev/null +++ b/services/core/internal/db/queries/sandbox_link.sql @@ -0,0 +1,22 @@ +-- name: GetSandboxServeAuthority :one +SELECT tenant_id, environment_id, kind, generation, credential_hash FROM sandbox_resources +WHERE id = $1 AND 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; + +-- name: GetLinkAssignment :one +-- The assignment with its Runtime's Attach authority, the live Link resource +-- of its Session's Environment and that Environment's network access. +SELECT b.session_id, b.runtime_id, b.epoch, b.desired_state = 'bound' AS bound, + (d.agent_host AND d.revoked_at IS NULL)::boolean AS agent_host, d.credential_revision, + 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(s.configuration->'environment'->'network'->>'access', 'enabled') = 'enabled')::boolean AS network_enabled +FROM session_runtime_assignments b +JOIN sessions s ON s.id = b.session_id +JOIN devices d ON d.id = b.runtime_id +LEFT JOIN environments e ON e.session_id = b.session_id +LEFT JOIN sandbox_resources r ON r.environment_id = e.id AND r.live +WHERE b.assignment_id = $1; diff --git a/services/core/internal/db/sqlc/devices.sql.go b/services/core/internal/db/sqlc/devices.sql.go index 274b848d3..8dce11cf0 100644 --- a/services/core/internal/db/sqlc/devices.sql.go +++ b/services/core/internal/db/sqlc/devices.sql.go @@ -209,19 +209,31 @@ func (q *Queries) GetSessionExecutionBinding(ctx context.Context, arg GetSession } const listPendingAssignmentReleases = `-- name: ListPendingAssignmentReleases :many -SELECT session_id, runtime_id, assignment_id, epoch, remove_home FROM session_runtime_assignments -WHERE desired_state = 'released' AND applied_epoch < epoch AND runtime_id = ANY($1::uuid[]) -ORDER BY runtime_id, session_id +SELECT b.session_id, b.runtime_id, b.assignment_id, b.epoch, b.remove_home, + 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 +FROM session_runtime_assignments b +LEFT JOIN environments e ON e.session_id = b.session_id +LEFT JOIN sandbox_resources r ON r.environment_id = e.id +WHERE b.desired_state = 'released' AND b.applied_epoch < b.epoch AND b.runtime_id = ANY($1::uuid[]) +ORDER BY b.runtime_id, b.session_id ` type ListPendingAssignmentReleasesRow struct { - SessionID pgtype.UUID `json:"session_id"` - RuntimeID pgtype.UUID `json:"runtime_id"` - AssignmentID pgtype.UUID `json:"assignment_id"` - Epoch int64 `json:"epoch"` - RemoveHome bool `json:"remove_home"` + SessionID pgtype.UUID `json:"session_id"` + RuntimeID pgtype.UUID `json:"runtime_id"` + AssignmentID pgtype.UUID `json:"assignment_id"` + Epoch int64 `json:"epoch"` + RemoveHome bool `json:"remove_home"` + 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"` } +// Each release carries its Session's Link resource, live or not, which the +// release revokes before it is sent. func (q *Queries) ListPendingAssignmentReleases(ctx context.Context, runtimeIds []pgtype.UUID) ([]ListPendingAssignmentReleasesRow, error) { rows, err := q.db.Query(ctx, listPendingAssignmentReleases, runtimeIds) if err != nil { @@ -237,6 +249,11 @@ func (q *Queries) ListPendingAssignmentReleases(ctx context.Context, runtimeIds &i.AssignmentID, &i.Epoch, &i.RemoveHome, + &i.ResourceTenantID, + &i.ResourceEnvironmentID, + &i.ResourceKind, + &i.ResourceID, + &i.ResourceGeneration, ); err != nil { return nil, err } diff --git a/services/core/internal/db/sqlc/environment_executor_credentials.sql.go b/services/core/internal/db/sqlc/environment_executor_credentials.sql.go index 54db31689..d759bf02c 100644 --- a/services/core/internal/db/sqlc/environment_executor_credentials.sql.go +++ b/services/core/internal/db/sqlc/environment_executor_credentials.sql.go @@ -248,11 +248,17 @@ func (q *Queries) RevokeExecutorCredential(ctx context.Context, arg RevokeExecut } const rotateExecutorCredential = `-- name: RotateExecutorCredential :one -UPDATE environment_executor_credentials -SET token_sha256 = $1, issued_at = clock_timestamp(), revoked_at = NULL -WHERE key_id = $2 AND tenant_id = $3 - AND subject_kind = $4 AND subject_id = $5 -RETURNING key_id, environment_id +WITH rotated AS ( + UPDATE environment_executor_credentials + SET token_sha256 = $1, issued_at = clock_timestamp(), revoked_at = NULL + WHERE key_id = $2 AND tenant_id = $3 + AND subject_kind = $4 AND subject_id = $5 + RETURNING key_id, environment_id +), advanced AS ( + UPDATE sandbox_enrollments n SET generation = n.generation + 1 + FROM rotated r WHERE n.executor_key_id = r.key_id +) +SELECT key_id, environment_id FROM rotated ` type RotateExecutorCredentialParams struct { @@ -268,6 +274,8 @@ type RotateExecutorCredentialRow struct { EnvironmentID pgtype.UUID `json:"environment_id"` } +// Rotation advances the generation of the key's enrollments, so the Link +// authority refuses what the old secret served. func (q *Queries) RotateExecutorCredential(ctx context.Context, arg RotateExecutorCredentialParams) (RotateExecutorCredentialRow, error) { row := q.db.QueryRow(ctx, rotateExecutorCredential, arg.TokenSha256, diff --git a/services/core/internal/db/sqlc/models.go b/services/core/internal/db/sqlc/models.go index 14083ca76..3220a0e0d 100644 --- a/services/core/internal/db/sqlc/models.go +++ b/services/core/internal/db/sqlc/models.go @@ -81,6 +81,8 @@ type Device struct { EnvironmentID pgtype.UUID `json:"environment_id"` ExecutorKeyID pgtype.UUID `json:"executor_key_id"` ArchiveCancelTurnID pgtype.UUID `json:"archive_cancel_turn_id"` + AgentHost bool `json:"agent_host"` + CredentialRevision int64 `json:"credential_revision"` } type Environment struct { @@ -250,6 +252,8 @@ type RuntimeAllocation struct { ObservationError string `json:"observation_error"` ComputePhaseChangedAt pgtype.Timestamptz `json:"compute_phase_changed_at"` DeploymentGeneration pgtype.Int8 `json:"deployment_generation"` + ServeCredentialHash pgtype.Text `json:"serve_credential_hash"` + ServeGeneration int64 `json:"serve_generation"` } type RuntimeDeployment struct { @@ -366,6 +370,23 @@ type RuntimePlacement struct { DeploymentGeneration pgtype.Int8 `json:"deployment_generation"` } +type SandboxEnrollment struct { + ID pgtype.UUID `json:"id"` + EnvironmentID pgtype.UUID `json:"environment_id"` + ExecutorKeyID pgtype.UUID `json:"executor_key_id"` + Generation int64 `json:"generation"` +} + +type SandboxResource 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"` + CredentialHash pgtype.Text `json:"credential_hash"` + Live pgtype.Bool `json:"live"` +} + type Session struct { ID pgtype.UUID `json:"id"` TenantID pgtype.UUID `json:"tenant_id"` diff --git a/services/core/internal/db/sqlc/runtime_allocations.sql.go b/services/core/internal/db/sqlc/runtime_allocations.sql.go index 0caae4420..459a93258 100644 --- a/services/core/internal/db/sqlc/runtime_allocations.sql.go +++ b/services/core/internal/db/sqlc/runtime_allocations.sql.go @@ -13,7 +13,7 @@ import ( const createRuntimeAllocation = `-- name: CreateRuntimeAllocation :one INSERT INTO runtime_allocations (id, environment_id, device_id, provider_key, node_id, deployment_generation) -VALUES ($1, $2, $3, $4, $5, $6) RETURNING id, environment_id, device_id, provider_key, state, create_settled, created_at, kept_at, released_at, compute_phase, compute_revision, compute_state, compute_activity_at, compute_wake_requested, compute_retained_until, node_id, observation_error, compute_phase_changed_at, deployment_generation +VALUES ($1, $2, $3, $4, $5, $6) RETURNING id, environment_id, device_id, provider_key, state, create_settled, created_at, kept_at, released_at, compute_phase, compute_revision, compute_state, compute_activity_at, compute_wake_requested, compute_retained_until, node_id, observation_error, compute_phase_changed_at, deployment_generation, serve_credential_hash, serve_generation ` type CreateRuntimeAllocationParams struct { @@ -55,12 +55,14 @@ func (q *Queries) CreateRuntimeAllocation(ctx context.Context, arg CreateRuntime &i.ObservationError, &i.ComputePhaseChangedAt, &i.DeploymentGeneration, + &i.ServeCredentialHash, + &i.ServeGeneration, ) return i, err } const getRuntimeAllocation = `-- name: GetRuntimeAllocation :one -SELECT a.id, a.environment_id, a.device_id, a.provider_key, a.state, a.create_settled, a.created_at, a.kept_at, a.released_at, a.compute_phase, a.compute_revision, a.compute_state, a.compute_activity_at, a.compute_wake_requested, a.compute_retained_until, a.node_id, a.observation_error, a.compute_phase_changed_at, a.deployment_generation, e.session_id, s.tenant_id, s.deleted_at, (CASE WHEN a.compute_phase NOT IN ('disabled', 'running') THEN a.compute_retained_until IS NOT NULL AND a.compute_retained_until <= clock_timestamp() ELSE a.node_id IS NULL AND (SELECT mode FROM runtime_deployment) <> 'direct' AND a.kept_at <= clock_timestamp() - interval '1 hour' END)::boolean AS expired +SELECT a.id, a.environment_id, a.device_id, a.provider_key, a.state, a.create_settled, a.created_at, a.kept_at, a.released_at, a.compute_phase, a.compute_revision, a.compute_state, a.compute_activity_at, a.compute_wake_requested, a.compute_retained_until, a.node_id, a.observation_error, a.compute_phase_changed_at, a.deployment_generation, a.serve_credential_hash, a.serve_generation, e.session_id, s.tenant_id, s.deleted_at, (CASE WHEN a.compute_phase NOT IN ('disabled', 'running') THEN a.compute_retained_until IS NOT NULL AND a.compute_retained_until <= clock_timestamp() ELSE a.node_id IS NULL AND (SELECT mode FROM runtime_deployment) <> 'direct' AND a.kept_at <= clock_timestamp() - interval '1 hour' END)::boolean AS expired FROM runtime_allocations a JOIN environments e ON e.id = a.environment_id JOIN sessions s ON s.id = e.session_id @@ -103,6 +105,8 @@ func (q *Queries) GetRuntimeAllocation(ctx context.Context, arg GetRuntimeAlloca &i.RuntimeAllocation.ObservationError, &i.RuntimeAllocation.ComputePhaseChangedAt, &i.RuntimeAllocation.DeploymentGeneration, + &i.RuntimeAllocation.ServeCredentialHash, + &i.RuntimeAllocation.ServeGeneration, &i.SessionID, &i.TenantID, &i.DeletedAt, @@ -115,7 +119,7 @@ const keepRuntimeAllocation = `-- name: KeepRuntimeAllocation :one UPDATE runtime_allocations SET kept_at = clock_timestamp() WHERE id = $1 AND state = 'running' AND (node_id IS NOT NULL OR (SELECT mode FROM runtime_deployment) = 'direct' OR kept_at > clock_timestamp() - interval '1 hour') -RETURNING id, environment_id, device_id, provider_key, state, create_settled, created_at, kept_at, released_at, compute_phase, compute_revision, compute_state, compute_activity_at, compute_wake_requested, compute_retained_until, node_id, observation_error, compute_phase_changed_at, deployment_generation +RETURNING id, environment_id, device_id, provider_key, state, create_settled, created_at, kept_at, released_at, compute_phase, compute_revision, compute_state, compute_activity_at, compute_wake_requested, compute_retained_until, node_id, observation_error, compute_phase_changed_at, deployment_generation, serve_credential_hash, serve_generation ` func (q *Queries) KeepRuntimeAllocation(ctx context.Context, id pgtype.UUID) (RuntimeAllocation, error) { @@ -141,12 +145,14 @@ func (q *Queries) KeepRuntimeAllocation(ctx context.Context, id pgtype.UUID) (Ru &i.ObservationError, &i.ComputePhaseChangedAt, &i.DeploymentGeneration, + &i.ServeCredentialHash, + &i.ServeGeneration, ) return i, err } const listRuntimeAllocations = `-- name: ListRuntimeAllocations :many -SELECT a.id, a.environment_id, a.device_id, a.provider_key, a.state, a.create_settled, a.created_at, a.kept_at, a.released_at, a.compute_phase, a.compute_revision, a.compute_state, a.compute_activity_at, a.compute_wake_requested, a.compute_retained_until, a.node_id, a.observation_error, a.compute_phase_changed_at, a.deployment_generation, e.session_id, s.tenant_id, s.deleted_at, (CASE WHEN a.compute_phase NOT IN ('disabled', 'running') THEN a.compute_retained_until IS NOT NULL AND a.compute_retained_until <= clock_timestamp() ELSE a.node_id IS NULL AND (SELECT mode FROM runtime_deployment) <> 'direct' AND a.kept_at <= clock_timestamp() - interval '1 hour' END)::boolean AS expired +SELECT a.id, a.environment_id, a.device_id, a.provider_key, a.state, a.create_settled, a.created_at, a.kept_at, a.released_at, a.compute_phase, a.compute_revision, a.compute_state, a.compute_activity_at, a.compute_wake_requested, a.compute_retained_until, a.node_id, a.observation_error, a.compute_phase_changed_at, a.deployment_generation, a.serve_credential_hash, a.serve_generation, e.session_id, s.tenant_id, s.deleted_at, (CASE WHEN a.compute_phase NOT IN ('disabled', 'running') THEN a.compute_retained_until IS NOT NULL AND a.compute_retained_until <= clock_timestamp() ELSE a.node_id IS NULL AND (SELECT mode FROM runtime_deployment) <> 'direct' AND a.kept_at <= clock_timestamp() - interval '1 hour' END)::boolean AS expired FROM runtime_allocations a JOIN environments e ON e.id = a.environment_id JOIN sessions s ON s.id = e.session_id @@ -191,6 +197,8 @@ func (q *Queries) ListRuntimeAllocations(ctx context.Context, id pgtype.UUID) ([ &i.RuntimeAllocation.ObservationError, &i.RuntimeAllocation.ComputePhaseChangedAt, &i.RuntimeAllocation.DeploymentGeneration, + &i.RuntimeAllocation.ServeCredentialHash, + &i.RuntimeAllocation.ServeGeneration, &i.SessionID, &i.TenantID, &i.DeletedAt, @@ -253,7 +261,7 @@ const observeRuntimeRunning = `-- name: ObserveRuntimeRunning :one UPDATE runtime_allocations SET state = 'running', create_settled = true WHERE id = $1 AND state IN ('creating', 'running') AND (node_id IS NOT NULL OR (SELECT mode FROM runtime_deployment) = 'direct' OR kept_at > clock_timestamp() - interval '1 hour') -RETURNING id, environment_id, device_id, provider_key, state, create_settled, created_at, kept_at, released_at, compute_phase, compute_revision, compute_state, compute_activity_at, compute_wake_requested, compute_retained_until, node_id, observation_error, compute_phase_changed_at, deployment_generation +RETURNING id, environment_id, device_id, provider_key, state, create_settled, created_at, kept_at, released_at, compute_phase, compute_revision, compute_state, compute_activity_at, compute_wake_requested, compute_retained_until, node_id, observation_error, compute_phase_changed_at, deployment_generation, serve_credential_hash, serve_generation ` func (q *Queries) ObserveRuntimeRunning(ctx context.Context, id pgtype.UUID) (RuntimeAllocation, error) { @@ -279,13 +287,15 @@ func (q *Queries) ObserveRuntimeRunning(ctx context.Context, id pgtype.UUID) (Ru &i.ObservationError, &i.ComputePhaseChangedAt, &i.DeploymentGeneration, + &i.ServeCredentialHash, + &i.ServeGeneration, ) return i, err } const releaseRuntimeAllocation = `-- name: ReleaseRuntimeAllocation :one UPDATE runtime_allocations SET state = 'released', released_at = clock_timestamp() -WHERE id = $1 AND state = 'cleanup_pending' AND create_settled RETURNING id, environment_id, device_id, provider_key, state, create_settled, created_at, kept_at, released_at, compute_phase, compute_revision, compute_state, compute_activity_at, compute_wake_requested, compute_retained_until, node_id, observation_error, compute_phase_changed_at, deployment_generation +WHERE id = $1 AND state = 'cleanup_pending' AND create_settled RETURNING id, environment_id, device_id, provider_key, state, create_settled, created_at, kept_at, released_at, compute_phase, compute_revision, compute_state, compute_activity_at, compute_wake_requested, compute_retained_until, node_id, observation_error, compute_phase_changed_at, deployment_generation, serve_credential_hash, serve_generation ` func (q *Queries) ReleaseRuntimeAllocation(ctx context.Context, id pgtype.UUID) (RuntimeAllocation, error) { @@ -311,13 +321,15 @@ func (q *Queries) ReleaseRuntimeAllocation(ctx context.Context, id pgtype.UUID) &i.ObservationError, &i.ComputePhaseChangedAt, &i.DeploymentGeneration, + &i.ServeCredentialHash, + &i.ServeGeneration, ) return i, err } const requestRuntimeCleanup = `-- name: RequestRuntimeCleanup :one UPDATE runtime_allocations SET state = 'cleanup_pending' -WHERE id = $1 AND state <> 'released' RETURNING id, environment_id, device_id, provider_key, state, create_settled, created_at, kept_at, released_at, compute_phase, compute_revision, compute_state, compute_activity_at, compute_wake_requested, compute_retained_until, node_id, observation_error, compute_phase_changed_at, deployment_generation +WHERE id = $1 AND state <> 'released' RETURNING id, environment_id, device_id, provider_key, state, create_settled, created_at, kept_at, released_at, compute_phase, compute_revision, compute_state, compute_activity_at, compute_wake_requested, compute_retained_until, node_id, observation_error, compute_phase_changed_at, deployment_generation, serve_credential_hash, serve_generation ` func (q *Queries) RequestRuntimeCleanup(ctx context.Context, id pgtype.UUID) (RuntimeAllocation, error) { @@ -343,13 +355,15 @@ func (q *Queries) RequestRuntimeCleanup(ctx context.Context, id pgtype.UUID) (Ru &i.ObservationError, &i.ComputePhaseChangedAt, &i.DeploymentGeneration, + &i.ServeCredentialHash, + &i.ServeGeneration, ) return i, err } const settleRuntimeCreation = `-- name: SettleRuntimeCreation :one UPDATE runtime_allocations SET create_settled = true -WHERE id = $1 AND state <> 'released' RETURNING id, environment_id, device_id, provider_key, state, create_settled, created_at, kept_at, released_at, compute_phase, compute_revision, compute_state, compute_activity_at, compute_wake_requested, compute_retained_until, node_id, observation_error, compute_phase_changed_at, deployment_generation +WHERE id = $1 AND state <> 'released' RETURNING id, environment_id, device_id, provider_key, state, create_settled, created_at, kept_at, released_at, compute_phase, compute_revision, compute_state, compute_activity_at, compute_wake_requested, compute_retained_until, node_id, observation_error, compute_phase_changed_at, deployment_generation, serve_credential_hash, serve_generation ` func (q *Queries) SettleRuntimeCreation(ctx context.Context, id pgtype.UUID) (RuntimeAllocation, error) { @@ -375,6 +389,8 @@ func (q *Queries) SettleRuntimeCreation(ctx context.Context, id pgtype.UUID) (Ru &i.ObservationError, &i.ComputePhaseChangedAt, &i.DeploymentGeneration, + &i.ServeCredentialHash, + &i.ServeGeneration, ) return i, err } diff --git a/services/core/internal/db/sqlc/runtime_lifecycle_nodes.sql.go b/services/core/internal/db/sqlc/runtime_lifecycle_nodes.sql.go index 6036aa233..e352b438e 100644 --- a/services/core/internal/db/sqlc/runtime_lifecycle_nodes.sql.go +++ b/services/core/internal/db/sqlc/runtime_lifecycle_nodes.sql.go @@ -50,7 +50,7 @@ func (q *Queries) GetRuntimeLifecyclePlacement(ctx context.Context, arg GetRunti } const listRuntimeAllocationsForNode = `-- name: ListRuntimeAllocationsForNode :many -SELECT a.id, a.environment_id, a.device_id, a.provider_key, a.state, a.create_settled, a.created_at, a.kept_at, a.released_at, a.compute_phase, a.compute_revision, a.compute_state, a.compute_activity_at, a.compute_wake_requested, a.compute_retained_until, a.node_id, a.observation_error, a.compute_phase_changed_at, a.deployment_generation, e.session_id, s.tenant_id, s.deleted_at, (CASE WHEN a.compute_phase NOT IN ('disabled', 'running') THEN a.compute_retained_until IS NOT NULL AND a.compute_retained_until <= clock_timestamp() ELSE a.node_id IS NULL AND (SELECT mode FROM runtime_deployment) <> 'direct' AND a.kept_at <= clock_timestamp() - interval '1 hour' END)::boolean AS expired +SELECT a.id, a.environment_id, a.device_id, a.provider_key, a.state, a.create_settled, a.created_at, a.kept_at, a.released_at, a.compute_phase, a.compute_revision, a.compute_state, a.compute_activity_at, a.compute_wake_requested, a.compute_retained_until, a.node_id, a.observation_error, a.compute_phase_changed_at, a.deployment_generation, a.serve_credential_hash, a.serve_generation, e.session_id, s.tenant_id, s.deleted_at, (CASE WHEN a.compute_phase NOT IN ('disabled', 'running') THEN a.compute_retained_until IS NOT NULL AND a.compute_retained_until <= clock_timestamp() ELSE a.node_id IS NULL AND (SELECT mode FROM runtime_deployment) <> 'direct' AND a.kept_at <= clock_timestamp() - interval '1 hour' END)::boolean AS expired FROM runtime_allocations a JOIN environments e ON e.id=a.environment_id JOIN sessions s ON s.id=e.session_id @@ -101,6 +101,8 @@ func (q *Queries) ListRuntimeAllocationsForNode(ctx context.Context, arg ListRun &i.RuntimeAllocation.ObservationError, &i.RuntimeAllocation.ComputePhaseChangedAt, &i.RuntimeAllocation.DeploymentGeneration, + &i.RuntimeAllocation.ServeCredentialHash, + &i.RuntimeAllocation.ServeGeneration, &i.SessionID, &i.TenantID, &i.DeletedAt, diff --git a/services/core/internal/db/sqlc/runtime_suspension.sql.go b/services/core/internal/db/sqlc/runtime_suspension.sql.go index 5c4867839..003c30eae 100644 --- a/services/core/internal/db/sqlc/runtime_suspension.sql.go +++ b/services/core/internal/db/sqlc/runtime_suspension.sql.go @@ -139,7 +139,7 @@ WHERE runtime_allocations.id = $4 AND compute_revision = $5 AND state = 'running' AND EXISTS (SELECT 1 FROM environments e WHERE e.id = runtime_allocations.environment_id AND e.initialization = 'complete') AND ((compute_phase IN ('disabled','running') AND (node_id IS NOT NULL OR (SELECT mode FROM runtime_deployment) = 'direct' OR kept_at > clock_timestamp() - interval '1 hour')) OR (compute_phase NOT IN ('disabled','running') AND compute_retained_until > clock_timestamp())) -RETURNING id, environment_id, device_id, provider_key, state, create_settled, created_at, kept_at, released_at, compute_phase, compute_revision, compute_state, compute_activity_at, compute_wake_requested, compute_retained_until, node_id, observation_error, compute_phase_changed_at, deployment_generation +RETURNING id, environment_id, device_id, provider_key, state, create_settled, created_at, kept_at, released_at, compute_phase, compute_revision, compute_state, compute_activity_at, compute_wake_requested, compute_retained_until, node_id, observation_error, compute_phase_changed_at, deployment_generation, serve_credential_hash, serve_generation ` type SetRuntimeComputeParams struct { @@ -179,6 +179,8 @@ func (q *Queries) SetRuntimeCompute(ctx context.Context, arg SetRuntimeComputePa &i.ObservationError, &i.ComputePhaseChangedAt, &i.DeploymentGeneration, + &i.ServeCredentialHash, + &i.ServeGeneration, ) return i, err } diff --git a/services/core/internal/db/sqlc/sandbox_link.sql.go b/services/core/internal/db/sqlc/sandbox_link.sql.go new file mode 100644 index 000000000..6c3767167 --- /dev/null +++ b/services/core/internal/db/sqlc/sandbox_link.sql.go @@ -0,0 +1,106 @@ +// Code generated by sqlc. DO NOT EDIT. +// versions: +// sqlc v1.29.0 +// source: sandbox_link.sql + +package sqlc + +import ( + "context" + + "github.com/jackc/pgx/v5/pgtype" +) + +const getAgentHostCredential = `-- 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 +` + +type GetAgentHostCredentialRow struct { + CredentialHash string `json:"credential_hash"` + CredentialRevision int64 `json:"credential_revision"` +} + +func (q *Queries) GetAgentHostCredential(ctx context.Context, id pgtype.UUID) (GetAgentHostCredentialRow, error) { + row := q.db.QueryRow(ctx, getAgentHostCredential, id) + var i GetAgentHostCredentialRow + err := row.Scan(&i.CredentialHash, &i.CredentialRevision) + return i, err +} + +const getLinkAssignment = `-- name: GetLinkAssignment :one +SELECT b.session_id, b.runtime_id, b.epoch, b.desired_state = 'bound' AS bound, + (d.agent_host AND d.revoked_at IS NULL)::boolean AS agent_host, d.credential_revision, + 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(s.configuration->'environment'->'network'->>'access', 'enabled') = 'enabled')::boolean AS network_enabled +FROM session_runtime_assignments b +JOIN sessions s ON s.id = b.session_id +JOIN devices d ON d.id = b.runtime_id +LEFT JOIN environments e ON e.session_id = b.session_id +LEFT JOIN sandbox_resources r ON r.environment_id = e.id AND r.live +WHERE b.assignment_id = $1 +` + +type GetLinkAssignmentRow struct { + SessionID pgtype.UUID `json:"session_id"` + RuntimeID pgtype.UUID `json:"runtime_id"` + Epoch int64 `json:"epoch"` + Bound bool `json:"bound"` + AgentHost bool `json:"agent_host"` + CredentialRevision int64 `json:"credential_revision"` + 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"` + NetworkEnabled bool `json:"network_enabled"` +} + +// The assignment with its Runtime's Attach authority, the live Link resource +// of its Session's Environment and that Environment's network access. +func (q *Queries) GetLinkAssignment(ctx context.Context, assignmentID pgtype.UUID) (GetLinkAssignmentRow, error) { + row := q.db.QueryRow(ctx, getLinkAssignment, assignmentID) + var i GetLinkAssignmentRow + err := row.Scan( + &i.SessionID, + &i.RuntimeID, + &i.Epoch, + &i.Bound, + &i.AgentHost, + &i.CredentialRevision, + &i.ResourceTenantID, + &i.ResourceEnvironmentID, + &i.ResourceKind, + &i.ResourceID, + &i.ResourceGeneration, + &i.NetworkEnabled, + ) + return i, err +} + +const getSandboxServeAuthority = `-- name: GetSandboxServeAuthority :one +SELECT tenant_id, environment_id, kind, generation, credential_hash FROM sandbox_resources +WHERE id = $1 AND live +` + +type GetSandboxServeAuthorityRow struct { + TenantID pgtype.UUID `json:"tenant_id"` + EnvironmentID pgtype.UUID `json:"environment_id"` + Kind string `json:"kind"` + Generation int64 `json:"generation"` + CredentialHash pgtype.Text `json:"credential_hash"` +} + +func (q *Queries) GetSandboxServeAuthority(ctx context.Context, id pgtype.UUID) (GetSandboxServeAuthorityRow, error) { + row := q.db.QueryRow(ctx, getSandboxServeAuthority, id) + var i GetSandboxServeAuthorityRow + err := row.Scan( + &i.TenantID, + &i.EnvironmentID, + &i.Kind, + &i.Generation, + &i.CredentialHash, + ) + return i, err +} diff --git a/services/core/internal/deployment/allocation.go b/services/core/internal/deployment/allocation.go index c825dda4f..a0ae9a37b 100644 --- a/services/core/internal/deployment/allocation.go +++ b/services/core/internal/deployment/allocation.go @@ -22,6 +22,8 @@ type Allocation struct { ID, EnvironmentID string SessionID, TenantID string DeviceID string + // ServeGeneration is the generation of the allocation's Link resource. + ServeGeneration uint64 // ProviderKey is the installation the allocation was provisioned for. ProviderKey string State string diff --git a/services/core/internal/execution/archive_cancellation_cleanup_test.go b/services/core/internal/execution/archive_cancellation_cleanup_test.go index a0cde3cd0..5c48792e7 100644 --- a/services/core/internal/execution/archive_cancellation_cleanup_test.go +++ b/services/core/internal/execution/archive_cancellation_cleanup_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/relay" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/adminaudit" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/identity" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtime" @@ -118,7 +119,7 @@ func TestArchiveWaitingCleanupReceiptBarrier(t *testing.T) { server := httptest.NewUnstartedServer(nil) wsURL := "ws://" + server.Listener.Addr().String() + "/api/v1/agent-daemon/ws" credentials, heartbeat := testSessions(t, pool, testCredentialCipher(t)) - handler, liveRegistry, err := runtime.NewGateway(credentials, heartbeat, credentials, wsURL) + handler, liveRegistry, err := runtime.NewGateway(credentials, heartbeat, credentials, runtimegateway.NewLinkAuthority(credentials), wsURL) if err != nil { t.Fatal(err) } @@ -171,7 +172,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, 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}, connections: map[string]*runtimeConnection{}} if checkpoint { lifecycle.config.Provider = waitingCleanupCheckpoint{beforeKill: provider.beforeKill} } diff --git a/services/core/internal/execution/assignment_releases.go b/services/core/internal/execution/assignment_releases.go index f75f6d1ab..0698e59a4 100644 --- a/services/core/internal/execution/assignment_releases.go +++ b/services/core/internal/execution/assignment_releases.go @@ -103,10 +103,15 @@ func (r releaseRetries) keep(releases []sessions.AssignmentRelease) { } // releaseAssignment sends one release and records its acknowledgement. Home -// removal is requested only from a Runtime that declares it. +// removal is requested only from a Runtime that declares it. The committed +// release already denies the assignment's grant; the relay closes the +// attachments it opened before the Runtime receives the release. func (w *Worker) releaseAssignment(ctx context.Context, release sessions.AssignmentRelease) bool { ctx, cancel := context.WithTimeout(ctx, 2*time.Minute) defer cancel() + if release.Resource.Kind != "" { + w.dispatcher.Links.RevokeResource(release.Resource.Ref()) + } peer, err := w.dispatcher.authorizedPeer(ctx, release.RuntimeID) if err != nil { return false diff --git a/services/core/internal/execution/dispatcher.go b/services/core/internal/execution/dispatcher.go index 3107832d3..b43d495a3 100644 --- a/services/core/internal/execution/dispatcher.go +++ b/services/core/internal/execution/dispatcher.go @@ -10,6 +10,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/relay" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/modelconfiguration" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" @@ -49,6 +50,10 @@ type Dispatcher struct { Sessions *sessions.Service // SessionsReader serves the plain Session reads. It is required. SessionsReader sessions.Reader + // Links is the Link relay. The Worker revokes a resource there after its + // durable Serve authority or an assignment's attach authority ends, and + // before it destroys the resource or sends the release. It is required. + Links *relay.Relay // ManagedRuntimes is optional internal provisioning; it does not admit hosted API requests. ManagedRuntimes *RuntimeProvider // MaxConcurrentExecutions bounds work admitted by this Core execution owner. diff --git a/services/core/internal/execution/owner_test.go b/services/core/internal/execution/owner_test.go index 57d676ee9..fd5f35f7d 100644 --- a/services/core/internal/execution/owner_test.go +++ b/services/core/internal/execution/owner_test.go @@ -7,6 +7,7 @@ import ( "testing" "time" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/relay" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/modelconfiguration" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgunit" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" @@ -134,7 +135,7 @@ func TestStartWorkerFailureClosesLeaseOnce(t *testing.T) { lease.inner = owner.Lease id := uuid.NewString() service, reader := unusedSessions(t) - dispatcher := &Dispatcher{Registry: runtimegateway.NewRegistry(), Credentials: &recordingCredentials{}, Observer: unusedObserver{t}, Deployment: deployments, DeploymentReader: deploymentReader, Sessions: service, SessionsReader: reader, ManagedRuntimes: NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil })} + dispatcher := &Dispatcher{Registry: runtimegateway.NewRegistry(), Credentials: &recordingCredentials{}, Observer: unusedObserver{t}, Deployment: deployments, DeploymentReader: deploymentReader, Sessions: service, SessionsReader: reader, Links: relay.New(nil), ManagedRuntimes: NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil })} _, err := StartWorker(canceled, dispatcher, Owner{Lease: lease, Deployment: owner.Deployment, Sessions: owner.Sessions}) if ping := owner.Lease.CheckOwnership(t.Context()); !errors.Is(ping, pgunit.ErrLeaseClosed) { t.Error("failed start kept the database lease", ping) @@ -174,6 +175,7 @@ func TestStartWorkerChecksDeploymentAfterItsDependencies(t *testing.T) { {"missing Session reader", Dispatcher{Credentials: credentials, Observer: observer, Deployment: deployments, DeploymentReader: deploymentReader, Sessions: service}, bound, "execution worker requires the Session reader"}, {"missing Session operations", Dispatcher{Credentials: credentials, Observer: observer, Deployment: deployments, DeploymentReader: deploymentReader, Sessions: service, SessionsReader: reader}, Owner{}, "execution requires the Session execution operations"}, {"missing deployment", Dispatcher{Credentials: credentials, Observer: observer, Deployment: deployments, DeploymentReader: deploymentReader, Sessions: service, SessionsReader: reader}, bound, "execution worker requires the deployment execution operations"}, + {"missing Link relay", Dispatcher{Credentials: credentials, Observer: observer, Deployment: deployments, DeploymentReader: deploymentReader, Sessions: service, SessionsReader: reader}, Owner{Sessions: owner.Sessions, Deployment: owner.Deployment}, "execution worker requires the Link relay"}, } { t.Run(test.name, func(t *testing.T) { lease := &closeCountingLease{t: t} @@ -193,7 +195,7 @@ func TestWorkerRunClosesLeaseAfterDrain(t *testing.T) { id := uuid.NewString() // The Worker's first reconciliation scans the Session work. reader, service := testSessions(t, pool, nil) - dispatcher := &Dispatcher{Registry: runtimegateway.NewRegistry(), Credentials: &recordingCredentials{}, Observer: unusedObserver{t}, Deployment: deployments, DeploymentReader: deploymentReader, Sessions: service, SessionsReader: reader, ManagedRuntimes: NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil })} + dispatcher := &Dispatcher{Registry: runtimegateway.NewRegistry(), Credentials: &recordingCredentials{}, Observer: unusedObserver{t}, Deployment: deployments, DeploymentReader: deploymentReader, Sessions: service, SessionsReader: reader, Links: relay.New(nil), ManagedRuntimes: NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil })} worker, err := StartWorker(t.Context(), dispatcher, Owner{Lease: lease, Deployment: owner.Deployment, Sessions: owner.Sessions}) if err != nil { t.Fatal(err) diff --git a/services/core/internal/execution/runtime_lifecycle.go b/services/core/internal/execution/runtime_lifecycle.go index d6ca1b27c..040853289 100644 --- a/services/core/internal/execution/runtime_lifecycle.go +++ b/services/core/internal/execution/runtime_lifecycle.go @@ -12,6 +12,8 @@ import ( "github.com/google/uuid" "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxbootstrap" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/relay" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimedevice" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" @@ -50,6 +52,7 @@ type runtimeLifecycle struct { reader deployment.Reader lease Ownership registry *runtimegateway.Registry + links *relay.Relay config RuntimeProvider nodeID string gate chan struct{} @@ -63,7 +66,7 @@ type runtimeLifecycle struct { wakeHints chan struct{} } -func newRuntimeManager(owner Owner, deployments *deployment.Service, deploymentReader deployment.Reader, sessionReader sessions.Reader, registry *runtimegateway.Registry, config *RuntimeProvider) (*runtimeManager, error) { +func newRuntimeManager(owner Owner, deployments *deployment.Service, deploymentReader deployment.Reader, sessionReader sessions.Reader, registry *runtimegateway.Registry, links *relay.Relay, config *RuntimeProvider) (*runtimeManager, error) { if config == nil { return nil, nil } @@ -81,7 +84,7 @@ func newRuntimeManager(owner Owner, deployments *deployment.Service, deploymentR } } 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, config: copied, 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, config: copied, 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) { @@ -251,7 +254,7 @@ func (r *runtimeLifecycle) provision(ctx context.Context, tenant, environment, p // Explicit invalid/foreign bootstrap cannot become an authorized Runtime. // Other failures may hide a successful Create; retain recovery for those. if errors.Is(err, sandbox.ErrInvalid) || errors.Is(err, sandbox.ErrOwnership) { - if _, cleanupErr := r.deployment.RequestCleanup(ctx, owner); cleanupErr != nil { + if _, cleanupErr := r.requestCleanup(ctx, owner); cleanupErr != nil { return owner, cleanupErr } } @@ -320,7 +323,7 @@ func (r *runtimeLifecycle) observe(ctx context.Context, owner deployment.Allocat } if owner.SessionDeleted || owner.Expired || owner.State == "cleanup_pending" { var err error - owner, err = r.deployment.RequestCleanup(ctx, owner) + owner, err = r.requestCleanup(ctx, owner) if err != nil { return err } @@ -371,7 +374,7 @@ func (r *runtimeLifecycle) observe(ctx context.Context, owner deployment.Allocat } } if owner.SessionDeleted || owner.Expired || owner.State == "cleanup_pending" { - owner, err = r.deployment.RequestCleanup(ctx, owner) + owner, err = r.requestCleanup(ctx, owner) if err != nil { return err } @@ -423,6 +426,19 @@ func (r *runtimeLifecycle) observe(ctx context.Context, owner deployment.Allocat return err } +// requestCleanup records the allocation's cleanup, which withdraws its Serve +// authority, then revokes its Link resource at the relay. The caller destroys +// the compute afterwards. +func (r *runtimeLifecycle) requestCleanup(ctx context.Context, owner deployment.Allocation) (deployment.Allocation, error) { + pending, err := r.deployment.RequestCleanup(ctx, owner) + if err != nil { + return pending, err + } + r.links.RevokeResource(sandboxbootstrap.Resource{TenantID: pending.TenantID, EnvironmentID: pending.EnvironmentID, + Kind: "allocation", ID: pending.ID, Generation: pending.ServeGeneration}.Ref()) + return pending, nil +} + // Environment identity owns connectivity; preparation has an independent owner. func (r *runtimeLifecycle) clearRuntimeState(owner deployment.Allocation) { delete(r.connections, owner.EnvironmentID) diff --git a/services/core/internal/execution/runtime_manager.go b/services/core/internal/execution/runtime_manager.go index 24f0ec8cb..8f09bdc50 100644 --- a/services/core/internal/execution/runtime_manager.go +++ b/services/core/internal/execution/runtime_manager.go @@ -7,6 +7,7 @@ import ( "sync" "time" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/relay" "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" @@ -24,6 +25,7 @@ type runtimeManager struct { deploymentReader deployment.Reader lease Ownership registry *runtimegateway.Registry + links *relay.Relay config RuntimeProvider setupInstallationID string loadDeployment func(context.Context) (*RuntimeProvider, error) @@ -89,7 +91,7 @@ func (m *runtimeManager) node(id string) (*runtimeNode, error) { n = &runtimeNode{lifecycle: &runtimeLifecycle{ sessions: m.sessions, sessionExecution: m.sessionExecution, deployment: m.deployment, deployments: m.deploymentService, reader: m.deploymentReader, - lease: m.lease, registry: m.registry, config: m.config, nodeID: id, + 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), }} diff --git a/services/core/internal/execution/sandbox_deployment_drain_test.go b/services/core/internal/execution/sandbox_deployment_drain_test.go index 3104220af..0dd9235d2 100644 --- a/services/core/internal/execution/sandbox_deployment_drain_test.go +++ b/services/core/internal/execution/sandbox_deployment_drain_test.go @@ -10,6 +10,7 @@ import ( "testing" "time" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/relay" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgunit" @@ -93,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", BackendFingerprint: strings.Repeat("a", 64), Provider: hub.Proxy(uuid.NewString(), "docker", docker.Operations(), 1)} - m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return configuration, nil })) + m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return configuration, nil })) if err != nil { t.Fatal(err) } @@ -231,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", BackendFingerprint: strings.Repeat("a", 64), Provider: hub.Proxy(uuid.NewString(), "docker", docker.Operations(), 1)} - m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return configuration, nil })) + m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return configuration, nil })) 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 35271c554..5cf03759c 100644 --- a/services/core/internal/execution/sandbox_deployment_setup_test.go +++ b/services/core/internal/execution/sandbox_deployment_setup_test.go @@ -9,6 +9,7 @@ import ( "testing" "time" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/relay" "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/sandbox" @@ -25,7 +26,7 @@ func TestDeferredSandboxDeploymentLoadsOnceBeforeNodeCreation(t *testing.T) { var selected atomic.Bool var loads atomic.Int32 configuration := &RuntimeProvider{InstallationID: id, ProviderKind: "docker", CoreURL: "https://core.example/api/v1", BackendFingerprint: strings.Repeat("a", 64), Provider: hub.Proxy(uuid.NewString(), "docker", docker.Operations(), 1)} - m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { loads.Add(1) if !selected.Load() { return nil, nil @@ -70,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(), NewDeferredRuntimeProvider(uuid.NewString(), func(ctx context.Context) (*RuntimeProvider, error) { + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(uuid.NewString(), func(ctx context.Context) (*RuntimeProvider, error) { close(entered) <-ctx.Done() return nil, ctx.Err() @@ -95,7 +96,7 @@ func TestDeferredSandboxProviderFailureKeepsRecoveryAvailable(t *testing.T) { available := false loadErr := ErrExecutionUnavailable configuration := &RuntimeProvider{InstallationID: id, ProviderKind: "docker", CoreURL: "https://core.example/api/v1", BackendFingerprint: strings.Repeat("a", 64), Provider: hub.Proxy(uuid.NewString(), "docker", docker.Operations(), 1)} - m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { if !available { return nil, loadErr } @@ -127,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", BackendFingerprint: strings.Repeat("a", 64), Provider: hub.Proxy(uuid.NewString(), "docker", docker.Operations(), 1)} rejected := errors.New("candidate provider unavailable") - m, err := newRuntimeManager(Owner{Lease: heldLease{}}, unitDeploymentService(t), nil, nil, runtimegateway.NewRegistry(), 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), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return config, nil }, func(context.Context, deployment.Setup) (PreparedRuntimeDeployment, error) { return PreparedRuntimeDeployment{}, rejected })) @@ -159,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(), 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), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil }, func(context.Context, deployment.Setup) (PreparedRuntimeDeployment, error) { close(entered) <-release @@ -192,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(), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil })) + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil })) 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 2d31d3d32..a0d63173c 100644 --- a/services/core/internal/execution/sandbox_deployment_switch_test.go +++ b/services/core/internal/execution/sandbox_deployment_switch_test.go @@ -7,6 +7,7 @@ import ( "testing" "time" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/relay" "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/sandbox/docker" @@ -19,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", BackendFingerprint: strings.Repeat("a", 64), Provider: hub.Proxy(uuid.NewString(), "docker", docker.Operations(), 1)} - m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return config, nil })) + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return config, nil })) if err != nil { t.Fatal(err) } @@ -80,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(), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, errors.New("provider unavailable") })) + 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") })) if err != nil { t.Fatal(err) } @@ -101,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", BackendFingerprint: strings.Repeat("a", 64), Provider: hub.Proxy(uuid.NewString(), "docker", docker.Operations(), 1)} - m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return config, nil })) + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return config, nil })) if err != nil { t.Fatal(err) } @@ -155,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(), NewDeferredRuntimeProvider(id, + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(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", BackendFingerprint: setup.BackendFingerprint, Provider: hub.Proxy(uuid.NewString(), "docker", 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 32fc48961..c8927d983 100644 --- a/services/core/internal/execution/sandbox_generations_test.go +++ b/services/core/internal/execution/sandbox_generations_test.go @@ -5,6 +5,7 @@ import ( "errors" "testing" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/relay" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/adminaudit" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" @@ -50,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(), config) + m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(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 531b3fbbc..33de667df 100644 --- a/services/core/internal/execution/sandbox_provider_contract_test.go +++ b/services/core/internal/execution/sandbox_provider_contract_test.go @@ -5,6 +5,7 @@ import ( "strings" "testing" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/relay" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" "github.com/google/uuid" @@ -18,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", BackendFingerprint: strings.Repeat("a", 64), Provider: &lifecycleOnlySandbox{}} - m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return config, nil })) + m, err := newRuntimeManager(Owner{Lease: heldLease{}}, nil, nil, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return config, nil })) 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 9ec0bb050..f4a174314 100644 --- a/services/core/internal/execution/sandbox_reset_test.go +++ b/services/core/internal/execution/sandbox_reset_test.go @@ -8,6 +8,7 @@ import ( "testing" "time" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/relay" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/adminaudit" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/credentialcrypto" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment" @@ -102,7 +103,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", BackendFingerprint: setup.BackendFingerprint, Provider: hub.Proxy(uuid.NewString(), "docker", docker.Operations(), 1)}, nil }) - m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), config) + m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(nil), config) if err != nil { t.Fatal(err) } @@ -223,7 +224,7 @@ func TestSandboxResetPublishesCommittedGenerationWithoutReading(t *testing.T) { } published = append(published, generation) } - m, err := newRuntimeManager(owner, deployments, adapter, nil, runtimegateway.NewRegistry(), config) + m, err := newRuntimeManager(owner, deployments, adapter, nil, runtimegateway.NewRegistry(), relay.New(nil), config) if err != nil { t.Fatal(err) } @@ -258,7 +259,7 @@ func TestCommittedResetViewStopsOwnerWithoutLease(t *testing.T) { return deployment.Snapshot{Record: deployment.Record{InstallationID: id, WebManaged: true, 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(), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil })) + 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 })) if err != nil { t.Fatal(err) } @@ -280,7 +281,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(), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil })) + m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(nil), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil })) 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 467596e9a..bcdaf4d5b 100644 --- a/services/core/internal/execution/sandbox_snapshot_budget_test.go +++ b/services/core/internal/execution/sandbox_snapshot_budget_test.go @@ -8,6 +8,7 @@ import ( "testing" "time" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/relay" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/adminaudit" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" @@ -74,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", BackendFingerprint: setup.BackendFingerprint, Provider: hub.Proxy(uuid.NewString(), "docker", docker.Operations(), 1)}, nil }) - m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), configuration) + m, err := newRuntimeManager(owner, deployments, reader, nil, runtimegateway.NewRegistry(), relay.New(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 0766df958..e3119423d 100644 --- a/services/core/internal/execution/worker.go +++ b/services/core/internal/execution/worker.go @@ -74,9 +74,12 @@ func StartWorker(ctx context.Context, dispatcher *Dispatcher, owner Owner) (_ *W if owner.Deployment == nil { return nil, errors.New("execution worker requires the deployment execution operations") } + if dispatcher.Links == nil { + 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.ManagedRuntimes) + worker.runtimes, err = newRuntimeManager(owner, owned.Deployment, owned.DeploymentReader, owned.SessionsReader, owned.Registry, owned.Links, owned.ManagedRuntimes) if err != nil { return nil, err } diff --git a/services/core/internal/persistence/postgres/deploymentpg/allocations.go b/services/core/internal/persistence/postgres/deploymentpg/allocations.go index 4d24c722f..2c9550d7e 100644 --- a/services/core/internal/persistence/postgres/deploymentpg/allocations.go +++ b/services/core/internal/persistence/postgres/deploymentpg/allocations.go @@ -24,7 +24,7 @@ func allocation(row sqlc.RuntimeAllocation, session, tenant pgtype.UUID, deleted ComputePhase: row.ComputePhase, ComputeRevision: row.ComputeRevision, ComputeState: row.ComputeState, ComputeActivityAt: row.ComputeActivityAt.Time, ComputeWakeRequested: row.ComputeWakeRequested, ComputeRetainedUntil: timestamp(row.ComputeRetainedUntil), ID: uuidString(row.ID), EnvironmentID: uuidString(row.EnvironmentID), SessionID: uuidString(session), TenantID: uuidString(tenant), - DeviceID: uuidString(row.DeviceID), ProviderKey: uuidString(row.ProviderKey), State: row.State, CreateSettled: row.CreateSettled, + DeviceID: uuidString(row.DeviceID), ProviderKey: uuidString(row.ProviderKey), ServeGeneration: uint64(row.ServeGeneration), State: row.State, CreateSettled: row.CreateSettled, SessionDeleted: deleted.Valid, Expired: expired, CreatedAt: row.CreatedAt.Time, KeptAt: row.KeptAt.Time, } } diff --git a/services/core/internal/persistence/postgres/sessionpg/execution_environment.go b/services/core/internal/persistence/postgres/sessionpg/execution_environment.go index 821171713..be9ee9816 100644 --- a/services/core/internal/persistence/postgres/sessionpg/execution_environment.go +++ b/services/core/internal/persistence/postgres/sessionpg/execution_environment.go @@ -95,7 +95,8 @@ func (e *Execution) ListAssignmentReleases(ctx context.Context, runtimes []strin } releases := make([]sessions.AssignmentRelease, 0, len(rows)) for _, row := range rows { - releases = append(releases, sessions.AssignmentRelease{RuntimeID: optionalID(row.RuntimeID), Assignment: assignmentRef(row.SessionID, row.AssignmentID, row.Epoch), RemoveHome: row.RemoveHome}) + releases = append(releases, sessions.AssignmentRelease{RuntimeID: optionalID(row.RuntimeID), Assignment: assignmentRef(row.SessionID, row.AssignmentID, row.Epoch), RemoveHome: row.RemoveHome, + Resource: linkResource(row.ResourceTenantID, row.ResourceEnvironmentID, row.ResourceKind, row.ResourceID, row.ResourceGeneration)}) } return releases, nil } diff --git a/services/core/internal/persistence/postgres/sessionpg/link.go b/services/core/internal/persistence/postgres/sessionpg/link.go new file mode 100644 index 000000000..043fa74ae --- /dev/null +++ b/services/core/internal/persistence/postgres/sessionpg/link.go @@ -0,0 +1,94 @@ +package sessionpg + +import ( + "context" + "errors" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgtype" + + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxbootstrap" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/credentialcrypto" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgunit" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimedevice" +) + +// attachGrantPurpose is the credential key purpose attach grants are signed +// under. +const attachGrantPurpose = "sandbox-link-attach-grant" + +// GetServeAuthority reads the live Link resource with the ID. A malformed or +// unknown resource has none. +func (s *Store) GetServeAuthority(ctx context.Context, id string) (runtimedevice.ServeAuthority, bool, error) { + resource, err := pgunit.ParseID(id) + if err != nil { + return runtimedevice.ServeAuthority{}, false, nil + } + row, err := s.units.Queries().GetSandboxServeAuthority(ctx, resource) + if errors.Is(err, pgx.ErrNoRows) { + return runtimedevice.ServeAuthority{}, false, nil + } + if err != nil { + return runtimedevice.ServeAuthority{}, false, err + } + return runtimedevice.ServeAuthority{ + Resource: linkResource(row.TenantID, row.EnvironmentID, pgtype.Text{String: row.Kind, Valid: true}, resource, pgtype.Int8{Int64: row.Generation, Valid: true}), + CredentialHash: row.CredentialHash.String, + }, true, 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) { + id, err := pgunit.ParseID(runtime) + if err != nil { + return runtimedevice.AgentHost{}, false, nil + } + row, err := s.units.Queries().GetAgentHostCredential(ctx, id) + if errors.Is(err, pgx.ErrNoRows) { + return runtimedevice.AgentHost{}, false, nil + } + if err != nil { + return runtimedevice.AgentHost{}, false, err + } + return runtimedevice.AgentHost{CredentialHash: row.CredentialHash, Revision: uint64(row.CredentialRevision)}, true, nil +} + +// GetLinkAssignment reads an assignment as the Link authority sees it. A +// malformed or unknown assignment ID has none. +func (s *Store) GetLinkAssignment(ctx context.Context, assignment string) (runtimedevice.LinkAssignment, bool, error) { + id, err := pgunit.ParseID(assignment) + if err != nil { + return runtimedevice.LinkAssignment{}, false, nil + } + row, err := s.units.Queries().GetLinkAssignment(ctx, id) + if errors.Is(err, pgx.ErrNoRows) { + return runtimedevice.LinkAssignment{}, false, nil + } + if err != nil { + return runtimedevice.LinkAssignment{}, false, err + } + return runtimedevice.LinkAssignment{ + SessionID: optionalID(row.SessionID), RuntimeID: optionalID(row.RuntimeID), Epoch: uint64(row.Epoch), Bound: row.Bound, + AgentHost: row.AgentHost, Revision: uint64(row.CredentialRevision), NetworkEnabled: row.NetworkEnabled, + Resource: linkResource(row.ResourceTenantID, row.ResourceEnvironmentID, row.ResourceKind, row.ResourceID, row.ResourceGeneration), + }, true, nil +} + +// SignAttachGrant returns the keyed digest of an attach grant's payload. A +// service without the credential key cannot sign or verify a grant. +func (s *Store) SignAttachGrant(_ context.Context, payload string) (string, error) { + if s.cipher == nil { + return "", credentialcrypto.ErrUnavailable + } + return s.cipher.Fingerprint(attachGrantPurpose, payload) +} + +// linkResource returns a resource of the sandbox_resources view, or the zero +// Resource when the row has none. +func linkResource(tenant, environment pgtype.UUID, kind pgtype.Text, id pgtype.UUID, generation pgtype.Int8) sandboxbootstrap.Resource { + if !kind.Valid { + return sandboxbootstrap.Resource{} + } + return sandboxbootstrap.Resource{TenantID: optionalID(tenant), EnvironmentID: optionalID(environment), Kind: kind.String, ID: optionalID(id), Generation: uint64(generation.Int64)} +} diff --git a/services/core/internal/runtime/gateway.go b/services/core/internal/runtime/gateway.go index 56a4a06e1..f7e1c609f 100644 --- a/services/core/internal/runtime/gateway.go +++ b/services/core/internal/runtime/gateway.go @@ -13,18 +13,18 @@ import ( // NewGateway serves the V1 daemon executor transport for both managed and // user-managed Runtime. It authenticates devices with credentials, records -// their heartbeats with heartbeat and drains an archived Session's cancellation -// receipts through cancellations. Its credentials never grant public Session -// API access. -func NewGateway(credentials runtimegateway.RuntimeStore, heartbeat runtimegateway.HeartbeatTouch, cancellations runtimegateway.ArchivedCancellationStore, publicWSURL string) (http.Handler, *runtimegateway.Registry, error) { +// their heartbeats with heartbeat, drains an archived Session's cancellation +// receipts through cancellations and gives agent hosts' binds their Link +// fields from links. Its credentials never grant public Session API access. +func NewGateway(credentials runtimegateway.RuntimeStore, heartbeat runtimegateway.HeartbeatTouch, cancellations runtimegateway.ArchivedCancellationStore, links *runtimegateway.LinkAuthority, publicWSURL string) (http.Handler, *runtimegateway.Registry, error) { u, err := url.Parse(publicWSURL) - if err != nil || credentials == nil || heartbeat == nil || cancellations == nil || (u.Scheme != "ws" && u.Scheme != "wss") || u.Hostname() == "" || u.User != nil || u.RawQuery != "" || u.Fragment != "" || u.Path != "/api/v1/agent-daemon/ws" { + if err != nil || credentials == nil || heartbeat == nil || cancellations == nil || links == nil || (u.Scheme != "ws" && u.Scheme != "wss") || u.Hostname() == "" || u.User != nil || u.RawQuery != "" || u.Fragment != "" || u.Path != "/api/v1/agent-daemon/ws" { return nil, nil, errors.New("daemon URL must be an absolute ws(s) URL ending in /api/v1/agent-daemon/ws") } registry := runtimegateway.NewRegistry() h := runtimegateway.NewHandler(runtimegateway.HandlerConfig{ Authenticator: runtimegateway.NewAuthenticator(credentials), Registry: registry, - Heartbeat: heartbeat, ArchivedCancellations: cancellations, PublicWSURL: publicWSURL, + Heartbeat: heartbeat, ArchivedCancellations: cancellations, Links: links, PublicWSURL: publicWSURL, }) r := chi.NewRouter() r.Route("/api/v1", func(r chi.Router) { runtimegateway.RegisterRoutes(r, h) }) diff --git a/services/core/internal/runtimedevice/link.go b/services/core/internal/runtimedevice/link.go new file mode 100644 index 000000000..e7a8a66c3 --- /dev/null +++ b/services/core/internal/runtimedevice/link.go @@ -0,0 +1,34 @@ +package runtimedevice + +import "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxbootstrap" + +// ServeAuthority is a live Link resource and the SHA-256 hex digest of the +// credential that serves it. +type ServeAuthority struct { + Resource sandboxbootstrap.Resource + CredentialHash string +} + +// LinkAssignment is a Session's assignment as the Link authority reads it. +type LinkAssignment struct { + SessionID string + RuntimeID string + Epoch uint64 + Bound bool + // AgentHost reports that the Runtime is a live agent host, and Revision + // is its credential revision. + AgentHost bool + Revision uint64 + // Resource is the live Link resource of the Session's Environment; its + // Kind is empty when there is none. + Resource sandboxbootstrap.Resource + // NetworkEnabled reports that the Environment's network access is + // enabled. + NetworkEnabled bool +} + +// AgentHost is a live agent host's credential digest and revision. +type AgentHost struct { + CredentialHash string + Revision uint64 +} diff --git a/services/core/internal/runtimegateway/assignment.go b/services/core/internal/runtimegateway/assignment.go index 701e922a8..963c720cb 100644 --- a/services/core/internal/runtimegateway/assignment.go +++ b/services/core/internal/runtimegateway/assignment.go @@ -13,7 +13,9 @@ import ( // Bind binds the Session's assignment to this connection's Runtime: it sends // assignment_bind and waits for bound, once per connection and reference. -// Every Session operation on the connection follows its Bind. +// Every Session operation on the connection follows its Bind. A bind of an +// Environment with a live Link resource to an agent host carries the +// resource and its attach grant. func (s *Session) Bind(ctx context.Context, ref proto.AssignmentRef, environmentID string) error { s.assignmentMu.Lock() bound := s.assignments[ref.SessionID] == ref @@ -23,7 +25,14 @@ func (s *Session) Bind(ctx context.Context, ref proto.AssignmentRef, environment } ctx, cancel := context.WithTimeout(ctx, 10*time.Second) defer cancel() - status, err := s.exchangeAssignment(ctx, proto.TypeAssignmentBind, ref, proto.AssignmentBindPayload{EnvironmentID: environmentID}) + payload := proto.AssignmentBindPayload{EnvironmentID: environmentID} + if s.links != nil { + var err error + if payload.Resource, payload.AttachGrant, err = s.links.bindLink(ctx, s.DeviceID, ref, environmentID); err != nil { + return err + } + } + status, err := s.exchangeAssignment(ctx, proto.TypeAssignmentBind, ref, payload) if err != nil { return err } diff --git a/services/core/internal/runtimegateway/handler.go b/services/core/internal/runtimegateway/handler.go index 291715528..2745aae2d 100644 --- a/services/core/internal/runtimegateway/handler.go +++ b/services/core/internal/runtimegateway/handler.go @@ -40,6 +40,10 @@ type HandlerConfig struct { // cancellation still owes a connection's delivery. nil drains nothing. ArchivedCancellations ArchivedCancellationStore + // Links gives an agent host's binds their Link resource and attach grant. + // nil sends binds without them. + Links *LinkAuthority + // PublicWSURL is the wss://... URL returned in the bootstrap // response so deployments behind a TLS terminator can advertise // the externally-reachable URL. @@ -130,6 +134,7 @@ func (h *Handler) WS(w http.ResponseWriter, r *http.Request) { sess := NewSession(conn, auth.DeviceID, auth.WorkspaceID, version, h.cfg.Registry, h.cfg.Log) sess.heartbeat = h.cfg.Heartbeat sess.archivedCancellations = h.cfg.ArchivedCancellations + sess.links = h.cfg.Links sess.credentialHash = runtimedevice.HashCredential(token) h.cfg.Log("agentdaemon gateway: ws upgrade ok, registering device_id=%s waiters=%d", auth.DeviceID, len(h.cfg.Registry.PendingWaiters(auth.DeviceID))) diff --git a/services/core/internal/runtimegateway/link.go b/services/core/internal/runtimegateway/link.go new file mode 100644 index 000000000..08f502ea6 --- /dev/null +++ b/services/core/internal/runtimegateway/link.go @@ -0,0 +1,207 @@ +package runtimegateway + +import ( + "context" + "crypto/hmac" + "crypto/sha256" + "crypto/subtle" + "encoding/binary" + "encoding/hex" + "fmt" + "net/netip" + "time" + + "github.com/google/uuid" + + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxbootstrap" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimedevice" +) + +// LinkStore reads the Link authority Core keeps and signs attach grants with +// the credential key. +type LinkStore interface { + // GetServeAuthority reads the live Link resource with the ID. + GetServeAuthority(ctx context.Context, id string) (runtimedevice.ServeAuthority, bool, error) + // GetAgentHostCredential reads a live agent host's credential. + GetAgentHostCredential(ctx context.Context, runtimeID string) (runtimedevice.AgentHost, bool, error) + // GetLinkAssignment reads an assignment by its ID. + GetLinkAssignment(ctx context.Context, assignmentID string) (runtimedevice.LinkAssignment, bool, error) + // SignAttachGrant returns the keyed hex digest of a grant's payload. + SignAttachGrant(ctx context.Context, payload string) (string, error) +} + +// linkLease bounds how long an attachment outlives a withdrawal that the +// relay did not see. +const linkLease = time.Minute + +// LinkAuthority is Core's sandboxlink.Authority. A Serve credential serves +// only its live resource. Only an agent host attaches, and it opens a service +// only with the grant its current bound assignment carries, for the current +// generation of its Session's resource. Every authorization rereads that +// state. +type LinkAuthority struct { + store LinkStore +} + +func NewLinkAuthority(store LinkStore) *LinkAuthority { + return &LinkAuthority{store: store} +} + +var _ sandboxlink.Authority = (*LinkAuthority)(nil) + +func (l *LinkAuthority) AuthenticateServe(ctx context.Context, hello sandboxlink.ServeHello) (sandboxlink.ServePeer, error) { + authority, found, err := l.store.GetServeAuthority(ctx, uuid.UUID(hello.Resource.ID).String()) + if err != nil { + return sandboxlink.ServePeer{}, err + } + digest := sha256.Sum256(hello.Credential) + if !found || subtle.ConstantTimeCompare([]byte(hex.EncodeToString(digest[:])), []byte(authority.CredentialHash)) != 1 { + return sandboxlink.ServePeer{}, sandboxlink.Fail(sandboxlink.AuthenticationFailed) + } + current := authority.Resource.Ref() + switch { + case current.SameResource(hello.Resource) && hello.Resource.Generation < current.Generation: + return sandboxlink.ServePeer{}, sandboxlink.Fail(sandboxlink.StaleGeneration) + case current != hello.Resource: + return sandboxlink.ServePeer{}, sandboxlink.Fail(sandboxlink.PermissionDenied) + } + return sandboxlink.ServePeer{PeerID: current.ID, Resource: current}, nil +} + +func (l *LinkAuthority) AuthenticateAttach(ctx context.Context, hello sandboxlink.AttachHello) (sandboxlink.AttachPeer, error) { + host, found, err := l.store.GetAgentHostCredential(ctx, uuid.UUID(hello.RuntimeID).String()) + if err != nil { + return sandboxlink.AttachPeer{}, err + } + presented := runtimedevice.HashCredential(string(hello.Credential)) + if !found || subtle.ConstantTimeCompare([]byte(presented), []byte(host.CredentialHash)) != 1 { + return sandboxlink.AttachPeer{}, sandboxlink.Fail(sandboxlink.AuthenticationFailed) + } + return sandboxlink.AttachPeer{RuntimeID: hello.RuntimeID, Revision: host.Revision}, nil +} + +// AuthorizeOpen grants File the world export and Network every address +// while the Environment's network access is enabled, and nothing otherwise. +func (l *LinkAuthority) AuthorizeOpen(ctx context.Context, peer sandboxlink.AttachPeer, open sandboxlink.Open) (sandboxlink.Authorization, error) { + identity, assignment, err := l.authorize(ctx, peer, open.AttachGrant) + if err != nil { + return sandboxlink.Authorization{}, err + } + identity.AttachmentID = open.AttachmentID + if identity != open.Identity() { + return sandboxlink.Authorization{}, sandboxlink.Fail(sandboxlink.PermissionDenied) + } + auth := sandboxlink.Authorization{Identity: identity, Service: open.Service, LeaseExpiresAt: time.Now().Add(linkLease)} + switch { + case open.Service == sandboxlink.ServiceFile: + auth.Exports = []sandboxlink.ExportGrant{{ID: sandboxfs.WorldExport}} + case open.Service == sandboxlink.ServiceNetwork && assignment.NetworkEnabled: + auth.Egress = allEgress + } + return auth, nil +} + +func (l *LinkAuthority) Renew(ctx context.Context, peer sandboxlink.AttachPeer, renew sandboxlink.RenewAttachment) (sandboxlink.Authorization, error) { + identity, _, err := l.authorize(ctx, peer, renew.AttachGrant) + if err != nil { + return sandboxlink.Authorization{}, err + } + identity.AttachmentID = renew.AttachmentID + return sandboxlink.Authorization{Identity: identity, LeaseExpiresAt: time.Now().Add(linkLease)}, nil +} + +var allEgress = []sandboxlink.EgressRule{ + {Prefix: netip.MustParsePrefix("0.0.0.0/0"), PortFirst: 1, PortLast: 65535}, + {Prefix: netip.MustParsePrefix("::/0"), PortFirst: 1, PortLast: 65535}, +} + +// authorize returns the attachment identity, without its attachment ID, that +// grant authorizes for peer: the grant's assignment must be the current bound +// assignment of peer's Runtime, and its resource generation the current one. +func (l *LinkAuthority) authorize(ctx context.Context, peer sandboxlink.AttachPeer, grant []byte) (sandboxlink.Identity, runtimedevice.LinkAssignment, error) { + if len(grant) != grantBytes { + return sandboxlink.Identity{}, runtimedevice.LinkAssignment{}, sandboxlink.Fail(sandboxlink.PermissionDenied) + } + id := sandboxwire.ID(grant[:16]) + epoch, generation := binary.BigEndian.Uint64(grant[16:24]), binary.BigEndian.Uint64(grant[24:32]) + assignment, found, err := l.store.GetLinkAssignment(ctx, uuid.UUID(id).String()) + switch { + case err != nil: + return sandboxlink.Identity{}, assignment, err + case !found || assignment.RuntimeID != uuid.UUID(peer.RuntimeID).String(): + return sandboxlink.Identity{}, assignment, sandboxlink.Fail(sandboxlink.PermissionDenied) + case !assignment.AgentHost || assignment.Revision != peer.Revision: + return sandboxlink.Identity{}, assignment, sandboxlink.Fail(sandboxlink.AuthenticationFailed) + case assignment.Resource.Kind == "": + return sandboxlink.Identity{}, assignment, sandboxlink.Fail(sandboxlink.ResourceNotFound) + } + want, err := l.grant(ctx, id, epoch, assignment.Resource.Kind, assignment.Resource.ID, generation) + if err != nil { + return sandboxlink.Identity{}, assignment, err + } + resource := assignment.Resource.Ref() + switch { + case !hmac.Equal(want, grant): + return sandboxlink.Identity{}, assignment, sandboxlink.Fail(sandboxlink.PermissionDenied) + case epoch != assignment.Epoch || !assignment.Bound: + return sandboxlink.Identity{}, assignment, sandboxlink.Fail(sandboxlink.StaleAssignment) + case generation < resource.Generation: + return sandboxlink.Identity{}, assignment, sandboxlink.Fail(sandboxlink.StaleGeneration) + case generation != resource.Generation: + return sandboxlink.Identity{}, assignment, sandboxlink.Fail(sandboxlink.PermissionDenied) + } + session, err := uuid.Parse(assignment.SessionID) + if err != nil { + return sandboxlink.Identity{}, assignment, err + } + return sandboxlink.Identity{Resource: resource, SessionID: sandboxwire.ID(session), AssignmentID: id, AssignmentEpoch: epoch}, assignment, nil +} + +// bindLink returns the Link resource and attach grant that runtimeID's bind of +// ref in environmentID carries: none unless runtimeID is the agent host that +// holds ref bound and environmentID has a live Link resource. +func (l *LinkAuthority) bindLink(ctx context.Context, runtimeID string, ref proto.AssignmentRef, environmentID string) (*sandboxbootstrap.Resource, []byte, error) { + if environmentID == "" { + return nil, nil, nil + } + assignment, found, err := l.store.GetLinkAssignment(ctx, ref.AssignmentID) + if err != nil || !found || !assignment.AgentHost || !assignment.Bound || assignment.Resource.EnvironmentID != environmentID || + assignment.RuntimeID != runtimeID || assignment.SessionID != ref.SessionID || assignment.Epoch != ref.Epoch { + return nil, nil, err + } + id, err := uuid.Parse(ref.AssignmentID) + if err != nil { + return nil, nil, err + } + resource := assignment.Resource + grant, err := l.grant(ctx, sandboxwire.ID(id), ref.Epoch, resource.Kind, resource.ID, resource.Generation) + if err != nil { + return nil, nil, err + } + return &resource, grant, nil +} + +// grantBytes is the size of an attach grant: the assignment ID, epoch and +// resource generation, then the hex digest that binds them to the resource. +const grantBytes = 16 + 8 + 8 + 2*sha256.Size + +// grant returns the attach grant of an assignment epoch on a resource +// generation. +func (l *LinkAuthority) grant(ctx context.Context, assignment sandboxwire.ID, epoch uint64, kind, resource string, generation uint64) ([]byte, error) { + header := make([]byte, 0, grantBytes) + header = append(header, assignment[:]...) + header = binary.BigEndian.AppendUint64(header, epoch) + header = binary.BigEndian.AppendUint64(header, generation) + digest, err := l.store.SignAttachGrant(ctx, fmt.Sprintf("%x/%s/%s", header, kind, resource)) + if err != nil { + return nil, err + } + if len(digest) != 2*sha256.Size { + return nil, fmt.Errorf("attach grant digest of %d bytes", len(digest)) + } + return append(header, digest...), nil +} diff --git a/services/core/internal/runtimegateway/session.go b/services/core/internal/runtimegateway/session.go index 1067da96d..d13143632 100644 --- a/services/core/internal/runtimegateway/session.go +++ b/services/core/internal/runtimegateway/session.go @@ -87,6 +87,8 @@ type Session struct { // archivedCancellations reads the receipt an archived Session's // cancellation owes this connection's delivery. archivedCancellations ArchivedCancellationStore + // links supplies the Link fields of this connection's binds. + links *LinkAuthority hbMu sync.Mutex lastSeenAt time.Time diff --git a/services/core/internal/sessions/execution_environment.go b/services/core/internal/sessions/execution_environment.go index 107b8ed5b..8d277d460 100644 --- a/services/core/internal/sessions/execution_environment.go +++ b/services/core/internal/sessions/execution_environment.go @@ -7,6 +7,7 @@ import ( "github.com/google/uuid" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxbootstrap" ) // EnvironmentKey names one Environment of a tenant. @@ -20,11 +21,13 @@ type EnvironmentConnection struct { } // AssignmentRelease is a released assignment whose Runtime has not -// acknowledged the release. +// acknowledged the release. Resource is the Link resource of the Session's +// Environment, live or not; its Kind is empty when there is none. type AssignmentRelease struct { RuntimeID string Assignment proto.AssignmentRef RemoveHome bool + Resource sandboxbootstrap.Resource } // EnvironmentExecution is the lease-bound storage of the Environment diff --git a/services/core/internal/sessions/executor_credentials.go b/services/core/internal/sessions/executor_credentials.go index dcde3026d..7f857ea1d 100644 --- a/services/core/internal/sessions/executor_credentials.go +++ b/services/core/internal/sessions/executor_credentials.go @@ -127,8 +127,8 @@ type ExecutorCredentialTx interface { IssueExecutorCredential(ctx context.Context, grant ExecutorCredentialGrant) (IssuedExecutorCredential, error) // RotateExecutorCredential replaces the secret of the subject's key with // the one whose SHA-256 digest is given, restores the key if it was - // revoked, and returns its key and restriction. An unknown key is - // ErrNotFound. + // revoked, advances the generation of the key's Link enrollments, and + // returns its key and restriction. An unknown key is ErrNotFound. RotateExecutorCredential(ctx context.Context, subject identity.Subject, key, digest string) (IssuedExecutorCredential, error) // RevokeExecutorCredential revokes the subject's key; revoking it again // changes nothing. An unknown key is ErrNotFound. diff --git a/services/core/migrations/000093_sandbox_link_authority.sql b/services/core/migrations/000093_sandbox_link_authority.sql new file mode 100644 index 000000000..ff220f5a8 --- /dev/null +++ b/services/core/migrations/000093_sandbox_link_authority.sql @@ -0,0 +1,53 @@ +-- +goose Up +-- Link authority. An allocation's Serve credential serves only its own +-- resource at serve_generation. A self_hosted enrollment's resource is served +-- with its executor credential while that key is authorized for the +-- Environment. Only a marked agent host may Attach; credential_revision fences +-- links that a rotated credential authenticated. +ALTER TABLE runtime_allocations + ADD COLUMN serve_credential_hash text CHECK (serve_credential_hash ~ '^[0-9a-f]{64}$'), + ADD COLUMN serve_generation bigint NOT NULL DEFAULT 1 CHECK (serve_generation > 0); + +CREATE TABLE sandbox_enrollments ( + id uuid PRIMARY KEY, + environment_id uuid NOT NULL UNIQUE REFERENCES environments(id), + executor_key_id uuid NOT NULL REFERENCES environment_executor_credentials(key_id), + generation bigint NOT NULL DEFAULT 1 CHECK (generation > 0) +); + +ALTER TABLE devices + ADD COLUMN agent_host boolean NOT NULL DEFAULT false, + ADD COLUMN credential_revision bigint NOT NULL DEFAULT 1 CHECK (credential_revision > 0), + ADD CONSTRAINT devices_agent_host CHECK (NOT agent_host OR (environment_id IS NULL AND executor_key_id IS NULL)); + +-- Every Link resource with its Serve credential hash. A resource is live +-- while that credential may Serve it. +CREATE VIEW sandbox_resources AS +SELECT s.tenant_id, a.environment_id, 'allocation'::text AS kind, a.id, a.serve_generation AS generation, + a.serve_credential_hash AS credential_hash, + a.state IN ('creating', 'running') AND s.deleted_at IS NULL AS live +FROM runtime_allocations a +JOIN environments e ON e.id = a.environment_id +JOIN sessions s ON s.id = e.session_id +WHERE a.serve_credential_hash IS NOT NULL +UNION ALL +SELECT s.tenant_id, n.environment_id, 'enrollment'::text, n.id, n.generation, c.token_sha256, + c.revoked_at IS NULL AND c.tenant_id = s.tenant_id AND c.subject_kind = s.creator_kind AND c.subject_id = s.creator_id + AND (c.environment_id IS NULL OR c.environment_id = e.id) AND s.deleted_at IS NULL + AND e.status NOT IN ('failed', 'expired') AND s.configuration->'environment'->>'type' = 'self_hosted' + AND EXISTS (SELECT 1 FROM execution_project_scopes p WHERE p.tenant_id = c.tenant_id) +FROM sandbox_enrollments n +JOIN environments e ON e.id = n.environment_id +JOIN sessions s ON s.id = e.session_id +JOIN environment_executor_credentials c ON c.key_id = n.executor_key_id; + +-- +goose Down +DROP VIEW sandbox_resources; +ALTER TABLE devices + DROP CONSTRAINT devices_agent_host, + DROP COLUMN credential_revision, + DROP COLUMN agent_host; +DROP TABLE sandbox_enrollments; +ALTER TABLE runtime_allocations + DROP COLUMN serve_generation, + DROP COLUMN serve_credential_hash; diff --git a/services/core/tests/integration/archive_cancellation_test.go b/services/core/tests/integration/archive_cancellation_test.go index 4430f0eee..90addbd8e 100644 --- a/services/core/tests/integration/archive_cancellation_test.go +++ b/services/core/tests/integration/archive_cancellation_test.go @@ -85,7 +85,7 @@ func TestArchiveWaitingCancellationReceipts(t *testing.T) { } 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), wsURL) + handler, registry, err := runtime.NewGateway(sessionAdapter(s), sessionService(t, s), sessionAdapter(s), runtimegateway.NewLinkAuthority(sessionAdapter(s)), wsURL) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/devices_test.go b/services/core/tests/integration/devices_test.go index 270e48ae7..6461f34e1 100644 --- a/services/core/tests/integration/devices_test.go +++ b/services/core/tests/integration/devices_test.go @@ -124,7 +124,7 @@ func TestStandaloneGatewayUsesExecutionCredentials(t *testing.T) { _, foreignSecret := registerTestDevice(t, s, uuid.NewString()) 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), wsURL) + handler, registry, err := runtime.NewGateway(sessionAdapter(s), sessionService(t, s), sessionAdapter(s), runtimegateway.NewLinkAuthority(sessionAdapter(s)), wsURL) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/dispatch_test.go b/services/core/tests/integration/dispatch_test.go index 37ad16047..1d2ac99b6 100644 --- a/services/core/tests/integration/dispatch_test.go +++ b/services/core/tests/integration/dispatch_test.go @@ -88,7 +88,7 @@ func newDispatchHarnessForSession(t *testing.T, configuration []byte, local bool } server := httptest.NewUnstartedServer(nil) wsURL := "ws://" + server.Listener.Addr().String() + "/api/v1/agent-daemon/ws" - server.Config.Handler, h.registry, err = runtime.NewGateway(sessionAdapter(s), sessionService(t, s), sessionAdapter(s), wsURL) + server.Config.Handler, h.registry, err = runtime.NewGateway(sessionAdapter(s), sessionService(t, s), sessionAdapter(s), runtimegateway.NewLinkAuthority(sessionAdapter(s)), wsURL) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/executor_principals_migration_test.go b/services/core/tests/integration/executor_principals_migration_test.go index 808e69a2a..bc6769d4e 100644 --- a/services/core/tests/integration/executor_principals_migration_test.go +++ b/services/core/tests/integration/executor_principals_migration_test.go @@ -125,6 +125,9 @@ func TestExecutorPrincipalMigrationRetiresUnknownAuthority(t *testing.T) { if _, err := provider.DownTo(ctx, 25); err == nil || !strings.Contains(err.Error(), "Cannot remove durable executor principal identities") { t.Fatal("downgrade lost principal keys", err) } + if _, err := provider.Up(ctx); err != nil { + t.Fatal(err) + } if _, err := credentials.RotateExecutorCredential(ctx, p, key.KeyID); err != nil { t.Fatal("failed downgrade damaged identity", err) } diff --git a/services/core/tests/integration/link_authority_test.go b/services/core/tests/integration/link_authority_test.go new file mode 100644 index 000000000..18f084613 --- /dev/null +++ b/services/core/tests/integration/link_authority_test.go @@ -0,0 +1,448 @@ +package integration + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "strings" + "testing" + "time" + + "github.com/google/uuid" + + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxbootstrap" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/sandboxlinktest" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimedevice" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" +) + +const linkWait = 10 * time.Second + +// linkHarness is a hosted Session bound to h.device whose Environment has an +// allocation with a Serve credential, and Core's Link Authority behind a +// relay. +type linkHarness struct { + *dispatchHarness + relay *sandboxlinktest.Server + resource sandboxbootstrap.Resource + serve []byte +} + +// newLinkHarness binds the Session to an unmarked operator device, or, for a +// guest, to the allocation's own device, as an in-sandbox daemon is bound. +func newLinkHarness(t *testing.T, guest bool) *linkHarness { + t.Helper() + h := newDispatchHarnessForSession(t, []byte(`{"agent":{"model":"test-model"},"environment":{"type":"openai_hosted","network":{"access":"disabled"}}}`), guest) + device := h.device.ID + if !guest { + allocated, err := sessionService(t, h.s).CreateDevice(t.Context(), h.tenant, "sandbox", runtimedevice.HashCredential(uuid.NewString())) + if err != nil { + t.Fatal(err) + } + device = allocated.ID + } + l := &linkHarness{dispatchHarness: h, relay: sandboxlinktest.StartRelay(t, runtimegateway.NewLinkAuthority(sessionAdapter(h.s))), serve: []byte(uuid.NewString()), + resource: sandboxbootstrap.Resource{TenantID: h.tenant, EnvironmentID: h.session.Environment.ID, Kind: "allocation", ID: uuid.NewString(), Generation: 1}} + if _, err := h.s.pool.Exec(t.Context(), `INSERT INTO runtime_allocations(id,environment_id,device_id,provider_key,state,create_settled,deployment_generation,serve_credential_hash) + VALUES($1,$2,$3,$4,'running',true,(SELECT generation FROM runtime_deployment),$5)`, l.resource.ID, l.resource.EnvironmentID, device, uuid.NewString(), serveHash(l.serve)); err != nil { + t.Fatal(err) + } + return l +} + +func serveHash(credential []byte) string { + digest := sha256.Sum256(credential) + return hex.EncodeToString(digest[:]) +} + +func (l *linkHarness) exec(query string, args ...any) { + l.t.Helper() + if _, err := l.s.pool.Exec(l.t.Context(), query, args...); err != nil { + l.t.Fatal(err) + } +} + +// bind binds the Session's current assignment to h.device through Core's +// gateway and returns it with the payload the Runtime received. +func (l *linkHarness) bind() (proto.AssignmentRef, proto.AssignmentBindPayload) { + t := l.t + t.Helper() + ref := proto.AssignmentRef{SessionID: l.session.ID} + var epoch int64 + if err := l.s.pool.QueryRow(t.Context(), "SELECT assignment_id::text, epoch FROM session_runtime_assignments WHERE session_id=$1", l.session.ID).Scan(&ref.AssignmentID, &epoch); err != nil { + t.Fatal(err) + } + ref.Epoch = uint64(epoch) + peer, err := l.registry.LookupDevice(l.device.ID) + if err != nil { + t.Fatal(err) + } + bound := make(chan error, 1) + go func() { bound <- peer.Bind(t.Context(), ref, l.session.Environment.ID) }() + frame := l.read(proto.TypeAssignmentBind) + var payload proto.AssignmentBindPayload + if err := frame.DecodePayload(&payload); err != nil { + t.Fatal(err) + } + reply, _ := assignmentReply(frame) + l.writeMu.Lock() + err = l.conn.WriteJSON(reply) + l.writeMu.Unlock() + if err != nil { + t.Fatal(err) + } + if err := within(t, bound); err != nil { + t.Fatal(err) + } + return ref, payload +} + +func (l *linkHarness) attach(runtime string, credential []byte) (*sandboxlink.AttachLink, error) { + return attachLink(l.t, l.relay, runtime, credential) +} + +func attachLink(t *testing.T, srv *sandboxlinktest.Server, runtime string, credential []byte) (*sandboxlink.AttachLink, error) { + ctx, cancel := context.WithTimeout(t.Context(), linkWait) + defer cancel() + link, err := sandboxlink.DialAttach(ctx, sandboxlink.AttachConfig{URL: srv.URL, TLS: srv.TLS, RuntimeID: wireID(runtime), Credential: credential}) + if err == nil { + t.Cleanup(func() { link.Close() }) + } + return link, err +} + +// open opens service on a new attachment of the harness's resource under ref. +func (l *linkHarness) open(link *sandboxlink.AttachLink, service sandboxlink.Service, ref proto.AssignmentRef, grant []byte) (sandboxwire.ID, error) { + ctx, cancel := context.WithTimeout(l.t.Context(), linkWait) + defer cancel() + attachment := sandboxwire.NewID() + stream, _, err := link.OpenService(ctx, sandboxlink.Open{Service: service, Version: 1, Resource: l.resource.Ref(), AttachmentID: attachment, + SessionID: wireID(ref.SessionID), AssignmentID: wireID(ref.AssignmentID), AssignmentEpoch: ref.Epoch, AttachGrant: grant}) + if err == nil { + l.t.Cleanup(func() { stream.Reset() }) + } + return attachment, err +} + +func (l *linkHarness) renew(link *sandboxlink.AttachLink, attachment sandboxwire.ID, grant []byte) error { + ctx, cancel := context.WithTimeout(l.t.Context(), linkWait) + defer cancel() + _, err := link.Renew(ctx, sandboxlink.RenewAttachment{AttachmentID: attachment, AttachGrant: grant}) + return err +} + +// linkServe is a serve peer of File and Network whose handlers hold their +// streams until the attachment closes. +type linkServe struct { + connected chan struct{} + binds chan sandboxlink.Bind + done chan struct{} + err error // Serve's result once done is closed + cancel context.CancelFunc +} + +func startLinkServe(t *testing.T, srv *sandboxlinktest.Server, credential []byte, resource sandboxlink.ResourceRef) *linkServe { + p := &linkServe{connected: make(chan struct{}, 16), binds: make(chan sandboxlink.Bind, 16), done: make(chan struct{})} + hold := func(ctx context.Context, b sandboxlink.Bind, _ uint64, _ sandboxlink.Stream) { + p.binds <- b + <-ctx.Done() + } + var ctx context.Context + ctx, p.cancel = context.WithCancel(context.Background()) + go func() { + defer close(p.done) + p.err = sandboxlink.Serve(ctx, sandboxlink.ServeConfig{URL: srv.URL, TLS: srv.TLS, Credential: credential, Resource: resource, ServerInstanceID: sandboxwire.NewID(), + Services: []sandboxlink.ServiceHandler{{Service: sandboxlink.ServiceFile, Version: 1, Serve: hold}, {Service: sandboxlink.ServiceNetwork, Version: 1, Serve: hold}}, + OnConnected: func() { p.connected <- struct{}{} }, MinBackoff: 10 * time.Millisecond, MaxBackoff: 50 * time.Millisecond}) + }() + t.Cleanup(p.stop) + return p +} + +func (p *linkServe) stop() { + p.cancel() + <-p.done +} + +// refused returns the code of the failure that ended Serve. +func (p *linkServe) refused(t *testing.T) sandboxlink.Code { + t.Helper() + within(t, p.done) + return linkCode(p.err) +} + +func linkCode(err error) sandboxlink.Code { + var failure *sandboxlink.Error + if !errors.As(err, &failure) { + return 0 + } + return failure.Code +} + +func wireID(id string) sandboxwire.ID { return sandboxwire.ID(uuid.MustParse(id)) } + +func within[T any](t *testing.T, ch <-chan T) T { + t.Helper() + select { + case v := <-ch: + return v + case <-time.After(linkWait): + t.Fatal("timed out") + panic("unreachable") + } +} + +// TestLinkAuthorityAgentHost serves an allocation and opens services on it +// from a marked agent host, and checks that each part of the grant's +// authority is current at every Open and renewal. +func TestLinkAuthorityAgentHost(t *testing.T) { + l := newLinkHarness(t, false) + p := startLinkServe(t, l.relay, l.serve, l.resource.Ref()) + within(t, p.connected) + + otherID, otherEnvironment := l.resource, l.resource + otherID.ID, otherEnvironment.EnvironmentID = uuid.NewString(), uuid.NewString() + for _, test := range []struct { + resource sandboxbootstrap.Resource + want sandboxlink.Code + }{{otherID, sandboxlink.AuthenticationFailed}, {otherEnvironment, sandboxlink.PermissionDenied}} { + if got := startLinkServe(t, l.relay, l.serve, test.resource.Ref()).refused(t); got != test.want { + t.Fatalf("Serve of another resource refused with %v, want %v", got, test.want) + } + } + if _, err := l.attach(l.device.ID, []byte(l.credential)); linkCode(err) != sandboxlink.AuthenticationFailed { + t.Fatal("an unmarked device attached", err) + } + + l.exec("UPDATE devices SET agent_host = true WHERE id = $1", l.device.ID) + link, err := l.attach(l.device.ID, []byte(l.credential)) + if err != nil { + t.Fatal(err) + } + if _, err := l.attach(l.device.ID, l.serve); linkCode(err) != sandboxlink.AuthenticationFailed { + t.Fatal("a Serve credential attached", err) + } + ref, payload := l.bind() + if payload.Resource == nil || *payload.Resource != l.resource || len(payload.AttachGrant) == 0 { + t.Fatalf("agent host bind = %+v", payload) + } + grant := payload.AttachGrant + file, err := l.open(link, sandboxlink.ServiceFile, ref, grant) + if err != nil { + t.Fatal(err) + } + if b := within(t, p.binds); len(b.Exports) != 1 || b.Exports[0] != (sandboxlink.ExportGrant{ID: sandboxfs.WorldExport}) { + t.Fatalf("File bind exports %+v", b.Exports) + } + if _, err := l.open(link, sandboxlink.ServiceNetwork, ref, grant); err != nil { + t.Fatal(err) + } + if b := within(t, p.binds); len(b.Egress) != 0 { + t.Fatalf("Network bind egress %+v with network access disabled", b.Egress) + } + if err := l.renew(link, file, grant); err != nil { + t.Fatal(err) + } + + wrong := bytes.Clone(grant) + wrong[len(wrong)-1] ^= 1 + if _, err := l.open(link, sandboxlink.ServiceFile, ref, wrong); linkCode(err) != sandboxlink.PermissionDenied { + t.Fatal("wrong grant", err) + } + l.exec("UPDATE devices SET credential_revision = credential_revision + 1 WHERE id = $1", l.device.ID) + if _, err := l.open(link, sandboxlink.ServiceFile, ref, grant); linkCode(err) != sandboxlink.AuthenticationFailed { + t.Fatal("stale credential revision", err) + } + if link, err = l.attach(l.device.ID, []byte(l.credential)); err != nil { + t.Fatal(err) + } + if _, err := l.open(link, sandboxlink.ServiceFile, ref, grant); err != nil { + t.Fatal("current credential revision", err) + } + l.exec("UPDATE session_runtime_assignments SET epoch = epoch + 1 WHERE session_id = $1", l.session.ID) + if _, err := l.open(link, sandboxlink.ServiceFile, ref, grant); linkCode(err) != sandboxlink.StaleAssignment { + t.Fatal("stale epoch", err) + } + next, payload := l.bind() + if _, err := l.open(link, sandboxlink.ServiceFile, next, payload.AttachGrant); err != nil { + t.Fatal("current epoch", err) + } + l.exec("UPDATE runtime_allocations SET serve_generation = 2 WHERE id = $1", l.resource.ID) + if _, err := l.open(link, sandboxlink.ServiceFile, next, payload.AttachGrant); linkCode(err) != sandboxlink.StaleGeneration { + t.Fatal("stale generation", err) + } +} + +// TestLinkAuthorityGuest checks that an in-sandbox daemon's assignment +// carries no grant and that its device can neither attach nor be marked. +func TestLinkAuthorityGuest(t *testing.T) { + l := newLinkHarness(t, true) + if _, payload := l.bind(); payload.Resource != nil || payload.AttachGrant != nil { + t.Fatalf("guest bind = %+v", payload) + } + if _, err := l.attach(l.device.ID, []byte(l.credential)); linkCode(err) != sandboxlink.AuthenticationFailed { + t.Fatal("a guest attached", err) + } + if _, err := l.s.pool.Exec(t.Context(), "UPDATE devices SET agent_host = true WHERE id = $1", l.device.ID); err == nil || !strings.Contains(err.Error(), "devices_agent_host") { + t.Fatal("a guest device was marked as an agent host", err) + } +} + +// TestLinkAuthorityReleaseRevokesBeforeSend checks that a released +// assignment's grant opens nothing from the commit on, and that the Worker +// has the relay close its attachments before it sends the release. +func TestLinkAuthorityReleaseRevokesBeforeSend(t *testing.T) { + l := newLinkHarness(t, false) + p := startLinkServe(t, l.relay, l.serve, l.resource.Ref()) + within(t, p.connected) + l.exec("UPDATE devices SET agent_host = true WHERE id = $1", l.device.ID) + ref, payload := l.bind() + link, err := l.attach(l.device.ID, []byte(l.credential)) + if err != nil { + t.Fatal(err) + } + file, err := l.open(link, sandboxlink.ServiceFile, ref, payload.AttachGrant) + if err != nil { + t.Fatal(err) + } + if err := sessionService(t, l.s).DeleteSession(t.Context(), sessions.DeleteSessionCommand{TenantID: l.tenant, SessionID: l.session.ID}); err != nil { + t.Fatal(err) + } + if _, err := l.open(link, sandboxlink.ServiceFile, ref, payload.AttachGrant); linkCode(err) != sandboxlink.ResourceNotFound { + t.Fatal("a released assignment opened a service", err) + } + dispatcher := *l.d + dispatcher.Links = l.relay.Relay + worker := startWorker(t, t.Context(), l.s, &dispatcher) + ctx, cancel := context.WithCancel(t.Context()) + done := make(chan error, 1) + go func() { done <- worker.Run(ctx) }() + t.Cleanup(func() { cancel(); <-done }) + if release := l.read(proto.TypeAssignmentRelease); release.Assignment.AssignmentID != ref.AssignmentID || release.Assignment.Epoch != ref.Epoch+1 { + t.Fatal("release of", release.Assignment) + } + // The relay no longer holds the attachment, so the revocation preceded the release. + if err := l.renew(link, file, payload.AttachGrant); linkCode(err) != sandboxlink.LeaseExpired { + t.Fatal("the release reached the Runtime before the relay closed its attachment", err) + } +} + +// TestLinkAuthorityEnrollment serves a self_hosted enrollment with its +// executor key until the key is revoked. +func TestLinkAuthorityEnrollment(t *testing.T) { + s, _ := NewModelTestStore(t) + principal := FixtureExecutorPrincipal(t, s, uuid.NewString()) + session, err := s.CreateSession(t.Context(), principal.TenantID, sessions.CreateSession{ + Creator: principal.Subject(), Engine: "codex", IdempotencyKey: uuid.NewString(), + Configuration: json.RawMessage(`{"agent":{"model":"fixture"},"environment":{"type":"self_hosted","workspace_directory":"/workspace"}}`), + }) + if err != nil { + t.Fatal(err) + } + key, err := sessionService(t, s).IssueExecutorCredential(t.Context(), principal, uuid.NewString(), session.Environment.ID) + if err != nil { + t.Fatal(err) + } + resource := sandboxbootstrap.Resource{TenantID: principal.TenantID, EnvironmentID: session.Environment.ID, Kind: "enrollment", ID: uuid.NewString(), Generation: 1} + if _, err := s.pool.Exec(t.Context(), "INSERT INTO sandbox_enrollments(id, environment_id, executor_key_id) VALUES($1, $2, $3)", resource.ID, resource.EnvironmentID, key.KeyID); err != nil { + t.Fatal(err) + } + srv := sandboxlinktest.StartRelay(t, runtimegateway.NewLinkAuthority(sessionAdapter(s))) + served := startLinkServe(t, srv, []byte(key.Token), resource.Ref()) + within(t, served.connected) + served.stop() + if err := sessionService(t, s).RevokeExecutorCredential(t.Context(), principal, key.KeyID); err != nil { + t.Fatal(err) + } + if got := startLinkServe(t, srv, []byte(key.Token), resource.Ref()).refused(t); got != sandboxlink.AuthenticationFailed { + t.Fatal("a revoked executor key served", got) + } +} + +// TestLinkAuthorityEnrollmentRotation checks that rotating an executor key +// advances its enrollment's generation: the serve peer of the old secret stays +// connected but no Open or renewal for the old generation is authorized, and +// the new secret serves the new generation. +func TestLinkAuthorityEnrollmentRotation(t *testing.T) { + l := newLinkHarness(t, false) + // The Session's Link resource becomes an enrollment of an executor key. + l.exec("UPDATE runtime_allocations SET serve_credential_hash = NULL WHERE id = $1", l.resource.ID) + l.exec(`UPDATE sessions SET configuration = jsonb_set(configuration, '{environment,type}', '"self_hosted"') WHERE id = $1`, l.session.ID) + principal := FixtureExecutorPrincipal(t, l.s, l.tenant) + key, err := sessionService(t, l.s).IssueExecutorCredential(t.Context(), principal, uuid.NewString(), l.session.Environment.ID) + if err != nil { + t.Fatal(err) + } + l.resource.Kind, l.resource.ID = "enrollment", uuid.NewString() + l.exec("INSERT INTO sandbox_enrollments(id, environment_id, executor_key_id) VALUES($1, $2, $3)", l.resource.ID, l.resource.EnvironmentID, key.KeyID) + p := startLinkServe(t, l.relay, []byte(key.Token), l.resource.Ref()) + within(t, p.connected) + l.exec("UPDATE devices SET agent_host = true WHERE id = $1", l.device.ID) + ref, payload := l.bind() + if payload.Resource == nil || *payload.Resource != l.resource { + t.Fatalf("agent host bind = %+v", payload) + } + link, err := l.attach(l.device.ID, []byte(l.credential)) + if err != nil { + t.Fatal(err) + } + file, err := l.open(link, sandboxlink.ServiceFile, ref, payload.AttachGrant) + if err != nil { + t.Fatal(err) + } + within(t, p.binds) + + rotated, err := sessionService(t, l.s).RotateExecutorCredential(t.Context(), principal, key.KeyID) + if err != nil { + t.Fatal(err) + } + if _, err := l.open(link, sandboxlink.ServiceFile, ref, payload.AttachGrant); linkCode(err) != sandboxlink.StaleGeneration { + t.Fatal("an Open reached the old secret's serve peer", err) + } + if err := l.renew(link, file, payload.AttachGrant); linkCode(err) != sandboxlink.StaleGeneration { + t.Fatal("an attachment to the old secret's serve peer renewed", err) + } + next := l.resource + next.Generation++ + within(t, startLinkServe(t, l.relay, []byte(rotated.Token), next.Ref()).connected) +} + +// TestLinkAuthorityDestroyedAllocation checks that the Worker's cleanup of an +// allocation withdraws its Serve authority and then disconnects its serve +// peer, which cannot serve again. +func TestLinkAuthorityDestroyedAllocation(t *testing.T) { + s, _ := newManagedTestStore(t) + tenant, session, environment := managedSession(t, s) + key := uuid.NewString() + srv := sandboxlinktest.StartRelay(t, runtimegateway.NewLinkAuthority(sessionAdapter(s))) + w := startWorker(t, t.Context(), s, &execution.Dispatcher{Registry: runtimegateway.NewRegistry(), Links: srv.Relay, ManagedRuntimes: &execution.RuntimeProvider{ + CoreURL: "http://core.invalid/api/v1", InstallationID: key, BackendFingerprint: strings.Repeat("a", 64), Provider: &lifecycleProvider{resources: map[string]sandbox.Info{}}}}) + t.Cleanup(func() { ctx, cancel := context.WithCancel(context.Background()); cancel(); _ = w.Run(ctx) }) + owner, err := w.ProvisionEnvironment(t.Context(), tenant, environment.ID, key) + if err != nil { + t.Fatal(err) + } + credential := []byte(uuid.NewString()) + if _, err := s.pool.Exec(t.Context(), "UPDATE runtime_allocations SET serve_credential_hash = $2 WHERE id = $1", owner.ID, serveHash(credential)); err != nil { + t.Fatal(err) + } + p := startLinkServe(t, srv, credential, sandboxbootstrap.Resource{TenantID: tenant, EnvironmentID: environment.ID, Kind: "allocation", ID: owner.ID, Generation: 1}.Ref()) + within(t, p.connected) + 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 got := p.refused(t); got != sandboxlink.AuthenticationFailed { + t.Fatal("a destroyed allocation's Serve credential served", got) + } +} diff --git a/services/core/tests/integration/runtime_connection_test.go b/services/core/tests/integration/runtime_connection_test.go index 19947a340..4d3108f1d 100644 --- a/services/core/tests/integration/runtime_connection_test.go +++ b/services/core/tests/integration/runtime_connection_test.go @@ -16,6 +16,7 @@ import ( "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" ) @@ -25,7 +26,7 @@ func TestManagedRuntimeConnectionTracksAuthenticatedSocket(t *testing.T) { 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), wsURL) + handler, registry, err := runtime.NewGateway(sessionAdapter(s), sessionService(t, s), sessionAdapter(s), runtimegateway.NewLinkAuthority(sessionAdapter(s)), wsURL) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/runtime_enrollment_connection_test.go b/services/core/tests/integration/runtime_enrollment_connection_test.go index 791fe6dd9..effbdf6a9 100644 --- a/services/core/tests/integration/runtime_enrollment_connection_test.go +++ b/services/core/tests/integration/runtime_enrollment_connection_test.go @@ -16,6 +16,7 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtime" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeenrollment" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" "github.com/google/uuid" "github.com/gorilla/websocket" @@ -46,7 +47,7 @@ func TestEnrolledDaemonConnectionRevocationAndRestart(t *testing.T) { } 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), wsURL) + handler, registry, err := runtime.NewGateway(sessionAdapter(s), sessionService(t, s), sessionAdapter(s), runtimegateway.NewLinkAuthority(sessionAdapter(s)), wsURL) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/worker_fixture_test.go b/services/core/tests/integration/worker_fixture_test.go index 872c77345..9a1abac84 100644 --- a/services/core/tests/integration/worker_fixture_test.go +++ b/services/core/tests/integration/worker_fixture_test.go @@ -5,10 +5,12 @@ import ( "errors" "testing" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/relay" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/modelconfigurationpg" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgunit" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/sessionpg" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" ) @@ -73,6 +75,9 @@ func startOwnedWorkerErr(ctx context.Context, s *Store, dispatcher *execution.Di owned.DeploymentReader = deploymentStore(s) owned.Sessions = service owned.SessionsReader = sessionAdapter(s) + if owned.Links == nil { + owned.Links = relay.New(runtimegateway.NewLinkAuthority(sessionAdapter(s))) + } return execution.StartWorker(ctx, &owned, owner) }