diff --git a/docs/sandbox-link-protocol.md b/docs/sandbox-link-protocol.md index 874594390..3d3d61c92 100644 --- a/docs/sandbox-link-protocol.md +++ b/docs/sandbox-link-protocol.md @@ -41,7 +41,15 @@ An attachment outlives its link. After reconnecting, the Runtime opens a stream ## Run a relay -`relay.New` takes an `Authority` and returns a `*relay.Relay`, which is an `http.Handler`. The relay endpoint is served behind the installation's HTTPS ingress, which terminates TLS, so the handler accepts the upgrade on the ingress's plain HTTP hop; peers enforce TLS when they dial. Each link carries at most 256 concurrent service streams. +`relay.New` takes an `Authority` and returns a `*relay.Relay`, which is an `http.Handler`. The relay endpoint is served behind the installation's HTTPS ingress, which terminates TLS, so the handler accepts the upgrade on the ingress's plain HTTP hop; peers enforce TLS when they dial. + +Each link carries at most 256 concurrent service streams, and each resource has at most 2 serve Hellos being decided. The relay also holds at most 4096 resources, 16384 attachments and 4096 attach links. A Hello or an Open takes the slots it needs before the Authority is consulted. When one is not free, it is refused with `LimitExceeded` and leaves nothing behind: + +- A serve Hello takes one of its resource's Hello slots until it is decided. +- A resource takes a slot at the serve Hello that first names it and keeps it while the relay holds anything for it: its serve peer, a Hello being decided, an attachment or an unwritten `AttachmentClosed` event. A reconnect of a held resource needs no new resource slot. A resource has at most one serve link, so the resource limit also bounds serve links and, with the Hello slots, the serve Hellos being decided. +- An Open takes a stream slot of its attach link. An opened stream keeps it until the stream ends; a refused Open gives it back once it is decided, even when the peer reset the stream first. +- An attachment takes a slot at the Open that creates it and keeps it until it is closed and its `AttachmentClosed` events are written or discarded. +- An attach link takes a slot at its Hello and keeps it until it ends and every stream and renewal it carried has finished. The owner of the relay implements `Authority` from its durable records, and the relay consults it for every Hello, Open and renewal. To revoke, withdraw the authority first, then call `RevokeAttachment` or `RevokeResource` so the relay closes what it holds. @@ -194,7 +202,7 @@ AttachmentClosed 1. The peer dials the relay's URL: `wss://`, or `ws://` only when the host is `localhost` or a loopback address. The URL carries no user, query or fragment, and never a credential. `sandboxlink.CheckRelayURL` applies this rule for both peers and the [bootstrap input](./sandbox-bootstrap.md) and refuses any other URL with `sandboxlink.ErrRelayURL`. Over `wss://`, TLS authenticates the relay. 2. The peer starts yamux and opens the control stream. 3. It sends a Hello as the first request of the control stream, with request ID 1: `ServeHello` from the Sandbox I/O service, `AttachHello` from a Runtime. Credentials travel only in the Hello. -4. The relay authenticates the peer with its Authority and answers `HelloAccepted`, or a failure after which the link ends. A Hello of another version is answered `VersionMismatch` without reading past its version. When a revocation lands while the Authority decides a serve Hello, the relay asks again, so a withdrawn credential never installs a serve peer. +4. The relay takes the [slot](#run-a-relay) the Hello needs, authenticates the peer with its Authority and answers `HelloAccepted`, or a failure after which the link ends. A Hello of another version is answered `VersionMismatch` without reading past its version. When a revocation lands while the Authority decides a serve Hello, the relay asks again, so a withdrawn credential never installs a serve peer. Later control requests continue the Hello's request IDs. The relay ends an attach link whose request ID does not increase with `ProtocolViolation`. `Open` and `Bind` are each the only request on their stream and use request ID 1. @@ -206,10 +214,11 @@ The relay and the serve peer bound each handshake step, the WebSocket upgrade, t The attach peer opens a stream and sends `Open`. The relay then: -1. Calls `Authority.AuthorizeOpen`. The Authority checks the grant, the Runtime, the current assignment and its epoch, the resource generation, the permitted service and access, and that the resource's serve authority is current. It returns the binding identity, the service, the lease, the exports for `ServiceFile` and the egress rules for `ServiceNetwork`. When a revocation lands while the Authority decides, the relay asks again. -2. Checks, in order: each link's stream limit (`LimitExceeded`); that the newest generation the relay has seen for the resource is not newer than the Open's (`StaleGeneration`); that a serve peer of the Open's generation is connected and offers the service (`ServiceUnavailable`) at the Open's version (`VersionMismatch`); that a nonzero `ExpectedServerInstanceID` equals the serve peer's (`InstanceChanged`); that the lease lies in the future (`LeaseExpired`); and that an attachment the relay already holds under this `AttachmentID` has the identical identity and Runtime (`AttachmentConflict`). -3. Opens a stream to the serve peer and sends `Bind`. The serve peer answers `Bound`, or a failure: `ServiceUnavailable` for a service it does not serve, `VersionMismatch`, `InstanceChanged` when `ExpectedServerInstanceID` is not its own, `LeaseExpired` for a recently closed attachment, or `ProtocolViolation`. The relay passes a failure on to the attach peer. When the `Bind` began to be sent but no answer arrives, the relay answers `ServiceUnavailable` with `EffectPossible`. -4. Answers `Opened` and splices the two streams. +1. Takes a stream [slot](#run-a-relay) of the link, and an attachment slot when it holds no attachment under the Open's `AttachmentID`, or answers `LimitExceeded` when one is not free. +2. Calls `Authority.AuthorizeOpen`. The Authority checks the grant, the Runtime, the current assignment and its epoch, the resource generation, the permitted service and access, and that the resource's serve authority is current. It returns the binding identity, the service, the lease, the exports for `ServiceFile` and the egress rules for `ServiceNetwork`. When a revocation lands while the Authority decides, the relay asks again. +3. Checks, in order: the serve link's stream limit (`LimitExceeded`); that the newest generation the relay holds for the resource is not newer than the Open's (`StaleGeneration`); that a serve peer of the Open's generation is connected and offers the service (`ServiceUnavailable`) at the Open's version (`VersionMismatch`); that a nonzero `ExpectedServerInstanceID` equals the serve peer's (`InstanceChanged`); that the lease lies in the future and that an attachment the relay held when the Open arrived has not closed since (`LeaseExpired`); and that an attachment the relay already holds under this `AttachmentID` has the identical identity and Runtime (`AttachmentConflict`). +4. Opens a stream to the serve peer and sends `Bind`. The serve peer answers `Bound`, or a failure: `ServiceUnavailable` for a service it does not serve, `VersionMismatch`, `InstanceChanged` when `ExpectedServerInstanceID` is not its own, `LeaseExpired` for a recently closed attachment, or `ProtocolViolation`. The relay passes a failure on to the attach peer. When the `Bind` began to be sent but no answer arrives, the relay answers `ServiceUnavailable` with `EffectPossible`. +5. Answers `Opened` and splices the two streams. An attachment's binding identity is its `AttachmentID`, `Resource`, `SessionID`, `AssignmentID` and `AssignmentEpoch`. Reopening an attachment, on the same link or a later one, requires the identical identity from the same Runtime and a current authorization. @@ -231,7 +240,7 @@ The relay keeps links, attachments and leases in memory. The Authority stays the A method returns a `*sandboxlink.Error` for a typed refusal; any other error is answered `ServiceUnavailable`. Credential revision and allowed access come from the Authority, never from what a peer asserts. -A recreated resource has a higher generation. A serve peer of the same or a higher generation replaces the resource's current serve peer, and a higher generation also closes every attachment of an older generation with `CloseStaleGeneration`. A serve peer or an Open of a generation older than the newest the relay has seen is refused with `StaleGeneration`. +A recreated resource has a higher generation. A serve peer of the same or a higher generation replaces the resource's current serve peer, and a higher generation also closes every attachment of an older generation with `CloseStaleGeneration`. A serve peer or an Open of a generation older than the newest the relay holds for the resource is refused with `StaleGeneration`. The relay keeps a resource's generation only while it [holds the resource](#run-a-relay); after that the Authority alone refuses older generations. ## Leases, closing and revocation @@ -240,7 +249,7 @@ A recreated resource has a higher generation. A serve peer of the same or a high - `CloseAttachment` closes the caller's attachment with `CloseRequested`. Closing an unknown attachment succeeds, and closing another Runtime's returns `PermissionDenied`. - `Relay.RevokeAttachment` closes one attachment. `Relay.RevokeResource` closes every attachment of a resource generation and older, writes the serve peer their `AttachmentClosed` events and then disconnects it. Both close with `CloseRevoked`. -Closing an attachment resets all its streams and sends `AttachmentClosed` to its attach peer, except after `CloseRequested`, and to the serve peer. The relay holds each serve peer's unwritten events in a set with no size limit and removes an event only once it is written, so a serve peer that is disconnected, or whose link drops before the event is written, receives it when it reconnects with the same generation. An Open that a close interrupts fails with `AttachmentConflict`, `LeaseExpired`, `PermissionDenied` or `StaleGeneration`, matching the reason, with `EffectPossible` when its `Bind` may have reached the serve peer. +Closing an attachment resets all its streams and sends `AttachmentClosed` to its attach peer, except after `CloseRequested`, and to the serve peer. The relay holds each peer's unwritten events in a set and removes an event once it is written, so a serve peer that is disconnected, or whose link drops before the event is written, receives it when it reconnects with the same generation. Events no peer can read any more are discarded: an attach peer's when its link ends, and a serve peer's when a newer generation connects or `RevokeResource` revokes its generation. An Open that a close interrupts fails with `AttachmentConflict`, `LeaseExpired`, `PermissionDenied` or `StaleGeneration`, matching the reason, with `EffectPossible` when its `Bind` may have reached the serve peer. Losing a link resets the streams it carries and keeps its attachments until their leases expire or they are closed. @@ -269,7 +278,7 @@ The relay copies each direction through a 32 KiB buffer and holds at most one 25 | 8 | `InstanceChanged` | The serve peer's `ServerInstanceID` is not `ExpectedServerInstanceID` | | 9 | `LeaseExpired` | The lease has passed, or the relay no longer holds the attachment | | 10 | `AttachmentConflict` | The `AttachmentID` is held with another identity or Runtime, or was closed during the Open | -| 11 | `LimitExceeded` | A link's stream limit or its limit of renewals being decided is reached | +| 11 | `LimitExceeded` | A Hello or an Open needs a [slot](#run-a-relay) that is not free, a serve link has its limit of streams, or an attach link has its limit of renewals being decided | | 12 | `ProtocolViolation` | A message is malformed, not allowed where it arrived, or carries a request ID that does not increase | `ServiceUnavailable` and `LimitExceeded` are transient: the same request may succeed later, and `Code.Retryable` reports them. Every other code is final: repeating the request with the same credential, attachment and generation fails again. @@ -278,4 +287,4 @@ An answer that is malformed, or that carries another request ID or operation tha ## Verification -`go test ./internal/sandboxlink/...` covers the golden frames, decode rejection, and the relay's authorization, generation, lease, revocation, renewal bound and reconnect behavior, including orderly end and abort propagation. `go test -run '^$' -fuzz FuzzDecode ./internal/sandboxlink` fuzzes the decoder. +`go test ./internal/sandboxlink/...` covers the golden frames, decode rejection, and the relay's authorization, generation, lease, revocation, capacity and renewal bounds, and reconnect behavior, including orderly end and abort propagation. `go test -run '^$' -fuzz FuzzDecode ./internal/sandboxlink` fuzzes the decoder. diff --git a/docs/zh/sandbox-link-protocol.md b/docs/zh/sandbox-link-protocol.md index 4056fec05..99ca95126 100644 --- a/docs/zh/sandbox-link-protocol.md +++ b/docs/zh/sandbox-link-protocol.md @@ -1,7 +1,7 @@ --- title: "沙箱 Link 协议" source: docs/sandbox-link-protocol.md -source_hash: ea2004916489420736d6bc73b31f1e26b409747406c20bf646729a2fe0478c41 +source_hash: f9c047631990775331946f78c7fad835a0663932ac6a2d4bd7c2514405f00b6f --- Link 协议通过 relay 连接沙箱 I/O 的两端。Sandbox I/O 服务运行在沙箱内并为其提供服务,是 serve peer。agent host 上的 Runtime 在沙箱外运行 Harness,并通过该服务使用沙箱,是 attach peer。每个 peer 各自向 relay 认证自己的 link。relay 授权 attach peer 打开的每个服务 stream,将其绑定到该资源当前的 serve peer,然后在两个 stream 之间复制字节而不读取内容。服务帧从不携带凭据或 grant。 @@ -43,7 +43,15 @@ attachment 的生命周期长于其 link。重连后,Runtime 使用相同的 b ## 运行 relay {#run-a-relay} -`relay.New` 接收 `Authority`,返回 `*relay.Relay`,它是一个 `http.Handler`。relay endpoint 位于安装实例的 HTTPS ingress 之后,由 ingress 终止 TLS,因此 handler 在 ingress 的明文 HTTP 一跳上接受 upgrade;peer 在拨号时强制 TLS。每条 link 最多承载 256 个并发服务 stream。 +`relay.New` 接收 `Authority`,返回 `*relay.Relay`,它是一个 `http.Handler`。relay endpoint 位于安装实例的 HTTPS ingress 之后,由 ingress 终止 TLS,因此 handler 在 ingress 的明文 HTTP 一跳上接受 upgrade;peer 在拨号时强制 TLS。 + +每条 link 最多承载 256 个并发服务 stream,每个资源最多有 2 个裁决中的 serve Hello。relay 还最多持有 4096 个资源、16384 个 attachment 和 4096 条 attach link。Hello 或 Open 在咨询 Authority 之前占用它所需的名额。某个名额没有空闲时,它会以 `LimitExceeded` 被拒绝,且不留下任何状态: + +- serve Hello 在裁决结束前占用其资源的一个 Hello 名额。 +- 资源在首次指明它的 serve Hello 处占用名额,并在 relay 仍为它持有任何东西时保留该名额:其 serve peer、裁决中的 Hello、attachment 或未写出的 `AttachmentClosed` 事件。已持有资源的重连不需要新的资源名额。一个资源最多有一条 serve link,因此资源上限也限制了 serve link 的数量,并与 Hello 名额一起限制了裁决中的 serve Hello 数量。 +- Open 占用其 attach link 的一个 stream 名额。已打开的 stream 保留该名额直到 stream 结束;被拒绝的 Open 在裁决结束后归还名额,即使 peer 已先重置该 stream。 +- attachment 在创建它的 Open 处占用名额,并保留到它被关闭且其 `AttachmentClosed` 事件已写出或丢弃为止。 +- attach link 在其 Hello 处占用名额,并保留到 link 结束且它承载过的每个 stream 和续期都已结束为止。 relay 的 owner 基于其持久记录实现 `Authority`,relay 对每个 Hello、Open 和续期都咨询它。撤销时,先撤回授权,再调用 `RevokeAttachment` 或 `RevokeResource`,让 relay 关闭其持有的对象。 @@ -196,7 +204,7 @@ AttachmentClosed 1. peer 拨号 relay 的 URL:使用 `wss://`;仅当主机为 `localhost` 或 loopback 地址时可用 `ws://`。URL 不含 user、查询或片段,也绝不含凭据。`sandboxlink.CheckRelayURL` 对两个 peer 和[引导输入](./sandbox-bootstrap.md)应用此规则,并以 `sandboxlink.ErrRelayURL` 拒绝其他 URL。使用 `wss://` 时,由 TLS 认证 relay。 2. peer 启动 yamux 并打开 control stream。 3. 它以请求 ID 1 发送 Hello,作为 control stream 的第一个请求:Sandbox I/O 服务发送 `ServeHello`,Runtime 发送 `AttachHello`。凭据只在 Hello 中传输。 -4. relay 通过其 Authority 认证 peer,并回复 `HelloAccepted`,或回复失败,随后结束 link。对其他版本的 Hello,relay 不读取版本之后的内容,直接回复 `VersionMismatch`。如果 Authority 裁决 serve Hello 期间发生撤销,relay 会再次询问,因此已撤回的凭据绝不会建立 serve peer。 +4. relay 占用该 Hello 所需的[名额](#run-a-relay),通过其 Authority 认证 peer,并回复 `HelloAccepted`,或回复失败,随后结束 link。对其他版本的 Hello,relay 不读取版本之后的内容,直接回复 `VersionMismatch`。如果 Authority 裁决 serve Hello 期间发生撤销,relay 会再次询问,因此已撤回的凭据绝不会建立 serve peer。 后续控制请求延续 Hello 的请求 ID。attach link 的请求 ID 未递增时,relay 以 `ProtocolViolation` 结束该 link。`Open` 和 `Bind` 各自是其 stream 上唯一的请求,使用请求 ID 1。 @@ -208,10 +216,11 @@ relay 和 serve peer 以 `sandboxlink.HandshakeTimeout`(10 秒)限制每个 attach peer 打开一个 stream 并发送 `Open`。relay 随后: -1. 调用 `Authority.AuthorizeOpen`。Authority 检查 grant、Runtime、当前 assignment 及其 epoch、资源 generation、允许的服务和访问权限,以及资源的 serve 授权是否当前有效。它返回 binding 身份、服务、lease、`ServiceFile` 的 export 和 `ServiceNetwork` 的 egress 规则。如果 Authority 裁决期间发生撤销,relay 会再次询问。 -2. 依次检查:每条 link 的 stream 上限(`LimitExceeded`);relay 见过的该资源最新 generation 不比 Open 中的更新(`StaleGeneration`);Open 所指 generation 的 serve peer 已连接并提供该服务(`ServiceUnavailable`),且支持 Open 的版本(`VersionMismatch`);非零的 `ExpectedServerInstanceID` 等于该 serve peer 的值(`InstanceChanged`);lease 尚未到期(`LeaseExpired`);relay 已以该 `AttachmentID` 持有的 attachment 具有完全相同的身份和 Runtime(`AttachmentConflict`)。 -3. 向 serve peer 打开 stream 并发送 `Bind`。serve peer 回复 `Bound`,或回复失败:对其不提供的服务回复 `ServiceUnavailable`,或回复 `VersionMismatch`,`ExpectedServerInstanceID` 不是自身值时回复 `InstanceChanged`,对最近关闭的 attachment 回复 `LeaseExpired`,或回复 `ProtocolViolation`。relay 将失败转交给 attach peer。`Bind` 已开始发送但未收到回复时,relay 回复带 `EffectPossible` 的 `ServiceUnavailable`。 -4. 回复 `Opened`,并将两个 stream 拼接起来。 +1. 占用该 link 的一个 stream [名额](#run-a-relay),并在 relay 未以 Open 的 `AttachmentID` 持有 attachment 时占用一个 attachment 名额;某个名额没有空闲时回复 `LimitExceeded`。 +2. 调用 `Authority.AuthorizeOpen`。Authority 检查 grant、Runtime、当前 assignment 及其 epoch、资源 generation、允许的服务和访问权限,以及资源的 serve 授权是否当前有效。它返回 binding 身份、服务、lease、`ServiceFile` 的 export 和 `ServiceNetwork` 的 egress 规则。如果 Authority 裁决期间发生撤销,relay 会再次询问。 +3. 依次检查:serve link 的 stream 上限(`LimitExceeded`);relay 为该资源持有的最新 generation 不比 Open 中的更新(`StaleGeneration`);Open 所指 generation 的 serve peer 已连接并提供该服务(`ServiceUnavailable`),且支持 Open 的版本(`VersionMismatch`);非零的 `ExpectedServerInstanceID` 等于该 serve peer 的值(`InstanceChanged`);lease 尚未到期,且 Open 到达时 relay 持有的 attachment 此后未被关闭(`LeaseExpired`);relay 已以该 `AttachmentID` 持有的 attachment 具有完全相同的身份和 Runtime(`AttachmentConflict`)。 +4. 向 serve peer 打开 stream 并发送 `Bind`。serve peer 回复 `Bound`,或回复失败:对其不提供的服务回复 `ServiceUnavailable`,或回复 `VersionMismatch`,`ExpectedServerInstanceID` 不是自身值时回复 `InstanceChanged`,对最近关闭的 attachment 回复 `LeaseExpired`,或回复 `ProtocolViolation`。relay 将失败转交给 attach peer。`Bind` 已开始发送但未收到回复时,relay 回复带 `EffectPossible` 的 `ServiceUnavailable`。 +5. 回复 `Opened`,并将两个 stream 拼接起来。 attachment 的 binding 身份由其 `AttachmentID`、`Resource`、`SessionID`、`AssignmentID` 和 `AssignmentEpoch` 组成。无论在同一 link 还是之后的 link 上重新打开 attachment,都要求同一 Runtime 提供完全相同的身份,并具有当前有效的授权。 @@ -233,7 +242,7 @@ relay 在内存中保存 link、attachment 和 lease。Authority 始终是持久 方法以 `*sandboxlink.Error` 表示类型化拒绝;其他任何错误都回复为 `ServiceUnavailable`。凭据 revision 和允许的访问来自 Authority,绝不来自 peer 的声明。 -重新创建的资源具有更高的 generation。相同或更高 generation 的 serve peer 替换资源当前的 serve peer;更高 generation 还会以 `CloseStaleGeneration` 关闭旧 generation 的所有 attachment。generation 比 relay 见过的最新 generation 更旧的 serve peer 或 Open 会以 `StaleGeneration` 被拒绝。 +重新创建的资源具有更高的 generation。相同或更高 generation 的 serve peer 替换资源当前的 serve peer;更高 generation 还会以 `CloseStaleGeneration` 关闭旧 generation 的所有 attachment。generation 比 relay 为该资源持有的最新 generation 更旧的 serve peer 或 Open 会以 `StaleGeneration` 被拒绝。relay 仅在[持有该资源](#run-a-relay)期间保留其 generation;此后由 Authority 单独拒绝更旧的 generation。 ## Lease、关闭与撤销 {#leases-closing-and-revocation} @@ -242,7 +251,7 @@ relay 在内存中保存 link、attachment 和 lease。Authority 始终是持久 - `CloseAttachment` 以 `CloseRequested` 关闭调用方的 attachment。关闭未知 attachment 会成功,关闭其他 Runtime 的 attachment 返回 `PermissionDenied`。 - `Relay.RevokeAttachment` 关闭一个 attachment。`Relay.RevokeResource` 关闭某个资源 generation 及更旧 generation 的所有 attachment,向 serve peer 写入相应的 `AttachmentClosed` 事件,然后断开它。两者都以 `CloseRevoked` 关闭。 -关闭 attachment 会重置其全部 stream,并向 serve peer 发送 `AttachmentClosed`;除 `CloseRequested` 外,也向其 attach peer 发送。relay 将每个 serve peer 未写出的事件保存在没有大小上限的集合中,仅在事件写出后才将其移除,因此已断开的 serve peer,或在事件写出前 link 断开的 serve peer,会在以相同 generation 重连时收到该事件。被关闭打断的 Open 根据原因以 `AttachmentConflict`、`LeaseExpired`、`PermissionDenied` 或 `StaleGeneration` 失败;其 `Bind` 可能已到达 serve peer 时带 `EffectPossible`。 +关闭 attachment 会重置其全部 stream,并向 serve peer 发送 `AttachmentClosed`;除 `CloseRequested` 外,也向其 attach peer 发送。relay 将每个 peer 未写出的事件保存在集合中,在事件写出后将其移除,因此已断开的 serve peer,或在事件写出前 link 断开的 serve peer,会在以相同 generation 重连时收到该事件。不再有 peer 能读取的事件会被丢弃:attach peer 的事件在其 link 结束时丢弃,serve peer 的事件在更新的 generation 连接或 `RevokeResource` 撤销其 generation 时丢弃。被关闭打断的 Open 根据原因以 `AttachmentConflict`、`LeaseExpired`、`PermissionDenied` 或 `StaleGeneration` 失败;其 `Bind` 可能已到达 serve peer 时带 `EffectPossible`。 丢失 link 会重置其承载的 stream,并保留其 attachment,直到 lease 到期或 attachment 被关闭。 @@ -271,7 +280,7 @@ relay 通过 32 KiB 缓冲区复制每个方向的数据,每个 stream 最多 | 8 | `InstanceChanged` | serve peer 的 `ServerInstanceID` 不是 `ExpectedServerInstanceID` | | 9 | `LeaseExpired` | lease 已过期,或 relay 不再持有该 attachment | | 10 | `AttachmentConflict` | 该 `AttachmentID` 以其他身份或 Runtime 被持有,或在 Open 期间被关闭 | -| 11 | `LimitExceeded` | 达到 link 的 stream 上限,或达到裁决中续期的数量上限 | +| 11 | `LimitExceeded` | Hello 或 Open 所需的[名额](#run-a-relay)没有空闲,serve link 已达 stream 上限,或 attach link 已达裁决中续期的数量上限 | | 12 | `ProtocolViolation` | 消息格式错误、出现在不允许的位置,或携带未递增的请求 ID | `ServiceUnavailable` 和 `LimitExceeded` 是临时失败:相同请求稍后可能成功,`Code.Retryable` 将它们报告为可重试。其他 code 都是最终失败:以相同凭据、attachment 和 generation 重复请求会再次失败。 @@ -280,4 +289,4 @@ relay 通过 32 KiB 缓冲区复制每个方向的数据,每个 stream 最多 ## 验证 {#verification} -`go test ./internal/sandboxlink/...` 覆盖 golden 帧、解码拒绝,以及 relay 的授权、generation、lease、撤销、续期上限和重连行为,包括有序结束与中止的传播。`go test -run '^$' -fuzz FuzzDecode ./internal/sandboxlink` 对解码器进行 fuzz 测试。 +`go test ./internal/sandboxlink/...` 覆盖 golden 帧、解码拒绝,以及 relay 的授权、generation、lease、撤销、容量与续期上限以及重连行为,包括有序结束与中止的传播。`go test -run '^$' -fuzz FuzzDecode ./internal/sandboxlink` 对解码器进行 fuzz 测试。 diff --git a/internal/sandboxlink/relay/export_test.go b/internal/sandboxlink/relay/export_test.go new file mode 100644 index 000000000..6b07e83ec --- /dev/null +++ b/internal/sandboxlink/relay/export_test.go @@ -0,0 +1,38 @@ +package relay + +import ( + "time" + + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" +) + +// The per-link and per-resource bounds the tests fill. +const ( + MaxStreams = maxStreams + MaxServeHellos = maxServeHellos +) + +// SetLimits lowers rl's capacity for resources, attachments and attach links. +func SetLimits(rl *Relay, resources, attachments, attachLinks int) { + rl.mu.Lock() + defer rl.mu.Unlock() + rl.maxResources, rl.maxAttachments, rl.maxAttachLinks = resources, attachments, attachLinks +} + +// Held reports the resources, attachment slots and attach links rl holds. +func Held(rl *Relay) (resources, attachments, attachLinks int) { + rl.mu.Lock() + defer rl.mu.Unlock() + return len(rl.resources), rl.held, rl.attachLinks +} + +// Lease reports the lease of the attachment rl holds under id, or the zero +// time when it holds none. +func Lease(rl *Relay, id sandboxwire.ID) time.Time { + rl.mu.Lock() + defer rl.mu.Unlock() + if a := rl.attachments[id]; a != nil { + return a.lease + } + return time.Time{} +} diff --git a/internal/sandboxlink/relay/relay.go b/internal/sandboxlink/relay/relay.go index 51ce1fb6c..ec671954d 100644 --- a/internal/sandboxlink/relay/relay.go +++ b/internal/sandboxlink/relay/relay.go @@ -20,7 +20,8 @@ import ( ) const ( - // maxStreams bounds the concurrent service streams of each link. + // maxStreams bounds the concurrent service streams of each link. An attach + // link's stream counts from its Open, through the Authority's decision. maxStreams = 256 // spliceBuffer bounds the bytes a splice holds per direction, on top of // one yamux window per stream. @@ -32,6 +33,26 @@ const ( answerQueue = 64 ) +// The relay's capacity, for one Core serving one deployment, which runs at +// most 1024 executions at once (the maximum OAC_EXECUTION_CONCURRENCY). +// Measured idle, the relay spends about 64 KiB on a serve link, 48 KiB on an +// attach link and 1 KiB on an attachment, so at capacity they hold under +// 500 MiB. +const ( + // maxResources bounds the resources held, and with them serve links, + // one per resource. It leaves room for idle sandboxes beside busy ones. + maxResources = 4096 + // maxServeHellos bounds the serve Hellos being decided for one resource, + // so a redial or a newer generation's peer can arrive while one is. + maxServeHellos = 2 + // maxAttachments bounds attachments, open or with AttachmentClosed events + // to write: four per resource. + maxAttachments = 16384 + // maxAttachLinks bounds attach links. An agent host dials one per + // attachment it uses, a few per execution. + maxAttachLinks = 4096 +) + // Relay accepts Link peers on ServeHTTP. It keeps the current serve peer of // each resource, the attachments it has opened and their leases in memory; the // Authority stays the durable judge of every grant. @@ -42,25 +63,41 @@ type Relay struct { mu sync.Mutex epoch uint64 // counts revocations, so a racing Open re-authorizes - generations map[resourceKey]uint64 - serves map[resourceKey]*serveLink - // closures holds the AttachmentClosed events not yet written to the serve - // peer of each resource's newest generation, connected or not. - closures map[resourceKey]closures + resources map[resourceKey]*resource attachments map[sandboxwire.ID]*attachment + // held counts attachment slots: open attachments, and closed ones whose + // AttachmentClosed events are not yet written or discarded. + held int + attachLinks int + // The capacity of each bounded kind. New sets the constants above; tests + // lower them. + maxResources, maxAttachments, maxAttachLinks int +} + +// resource is what the relay holds for one resource. It takes one of +// maxResources from the serve Hello that creates it until it holds nothing: +// no serve peer, Hello being decided, attachment or event. +type resource struct { + generation uint64 // the newest generation seen + serve *serveLink // nil while the serve peer is away + // closures holds the AttachmentClosed events not yet written to the serve + // peer of generation, connected or not. + closures closures + hellos int // serve Hellos being decided, at most maxServeHellos + attachments int // open attachments } -// closures is a set of AttachmentClosed events to write, by attachment. -type closures map[sandboxwire.ID]sandboxlink.CloseReason +// closures is a set of AttachmentClosed events to write, by attachment ID. +// Each event keeps its attachment's slot until it is written or discarded. +type closures map[sandboxwire.ID]*attachment // New returns a relay that asks auth to authenticate peers and authorize their // requests. func New(auth sandboxlink.Authority) *Relay { ctx, cancel := context.WithCancel(context.Background()) return &Relay{auth: auth, ctx: ctx, cancel: cancel, - generations: map[resourceKey]uint64{}, - serves: map[resourceKey]*serveLink{}, - closures: map[resourceKey]closures{}, + maxResources: maxResources, maxAttachments: maxAttachments, maxAttachLinks: maxAttachLinks, + resources: map[resourceKey]*resource{}, attachments: map[sandboxwire.ID]*attachment{}, } } @@ -89,8 +126,8 @@ func (rl *Relay) RevokeAttachment(id sandboxwire.ID) { // RevokeResource closes every attachment of ref's generation and older and // disconnects that serve peer after writing it their AttachmentClosed events. -// Events it could not write stay for a reconnect of the same generation. The -// caller withdraws the authority first. +// The caller withdraws the authority first, so the revoked generation cannot +// reconnect, and the events it could not read are discarded. func (rl *Relay) RevokeResource(ref sandboxlink.ResourceRef) { rl.mu.Lock() defer rl.mu.Unlock() @@ -101,10 +138,19 @@ func (rl *Relay) RevokeResource(ref sandboxlink.ResourceRef) { } } key := keyOf(ref) - if sl := rl.serves[key]; sl != nil && sl.hello.Resource.Generation <= ref.Generation { - delete(rl.serves, key) + r := rl.resources[key] + if r == nil || r.generation > ref.Generation { + return + } + if sl := r.serve; sl != nil { + // The link keeps the set to write before it ends. + r.serve = nil sl.end() + } else { + rl.dropLocked(r.closures) } + r.closures = closures{} + rl.releaseLocked(key, r) } type resourceKey struct { @@ -117,16 +163,17 @@ func keyOf(r sandboxlink.ResourceRef) resourceKey { } // link is one authenticated connection. Its writer sends answers from a -// bounded queue and AttachmentClosed events from an unbounded set, so the -// relay never blocks on a peer while holding its lock and never drops an -// event: an event leaves the set only once it is written. +// bounded queue and AttachmentClosed events from a set, so the relay never +// blocks on a peer while holding its lock: an event leaves the set once it is +// written, or with the set when no link can write it any more. type link struct { sess *yamux.Session ctl *yamux.Stream out chan outgoing wake chan struct{} // the set has events to write - // Under Relay.mu: the events to write and the open service streams. A - // serve link shares its resource's set, which outlives the link. + // Under Relay.mu: the events to write, nil once the link no longer + // writes them, and the open service streams. A serve link shares its + // resource's set, which outlives the link. closures closures streams uint32 } @@ -195,32 +242,66 @@ func (rl *Relay) write(l *link) { func (rl *Relay) writeClosures(l *link) error { for { rl.mu.Lock() - var c sandboxlink.AttachmentClosed - for id, reason := range l.closures { - c = sandboxlink.AttachmentClosed{AttachmentID: id, Reason: reason} + var a *attachment + for _, a = range l.closures { break } rl.mu.Unlock() - if c.Reason == 0 { + if a == nil { return nil } - if err := sandboxlink.WriteMessage(l.ctl, 0, c); err != nil { + id := a.identity.AttachmentID + if err := sandboxlink.WriteMessage(l.ctl, 0, sandboxlink.AttachmentClosed{AttachmentID: id, Reason: a.reason}); err != nil { return err } rl.mu.Lock() - if l.closures[c.AttachmentID] == c.Reason { - delete(l.closures, c.AttachmentID) + if l.closures[id] == a { + delete(l.closures, id) + rl.unrefLocked(a) } rl.mu.Unlock() } } -// closedLocked adds an AttachmentClosed event to l's set and wakes its writer. -func (l *link) closedLocked(id sandboxwire.ID, reason sandboxlink.CloseReason) { - if l.closures == nil { - l.closures = closures{} +// queueLocked adds a's AttachmentClosed event to set, unless set is nil. An +// unwritten event of an earlier attachment with the same ID gives way to it. +func (rl *Relay) queueLocked(set closures, a *attachment) { + if set == nil { + return + } + id := a.identity.AttachmentID + if old := set[id]; old != nil { + rl.unrefLocked(old) + } + set[id] = a + a.refs++ +} + +// dropLocked discards the events of a set that no link will write. +func (rl *Relay) dropLocked(set closures) { + for _, a := range set { + rl.unrefLocked(a) + } +} + +// unrefLocked drops one of a's references, the open attachment or one of its +// events. The last frees its slot. +func (rl *Relay) unrefLocked(a *attachment) { + a.refs-- + if a.refs == 0 { + rl.held-- } - l.closures[id] = reason +} + +// releaseLocked forgets r, freeing its slot, once it holds nothing. +func (rl *Relay) releaseLocked(key resourceKey, r *resource) { + if r.serve == nil && r.hellos == 0 && r.attachments == 0 && len(r.closures) == 0 { + delete(rl.resources, key) + } +} + +// wakeWriter tells l's writer that its set has events. +func (l *link) wakeWriter() { select { case l.wake <- struct{}{}: default: @@ -262,7 +343,9 @@ func linger(sess *yamux.Session) { } // attachment is an attachment the relay has opened. It outlives the links -// that carry its streams until its lease expires or it is closed. +// that carry its streams until its lease expires or it is closed, and keeps +// one of maxAttachments until its AttachmentClosed events are written or +// discarded. type attachment struct { identity sandboxlink.Identity runtime sandboxwire.ID @@ -271,6 +354,8 @@ type attachment struct { owner *attachLink // the link that last opened a stream on it bound bool // a Bind may have reached the serve peer splices map[*splice]struct{} + reason sandboxlink.CloseReason // why it closed; zero while open + refs int // its slot's holders: the open attachment and its unwritten events } // splice is one service stream from its Open to its end. @@ -353,17 +438,43 @@ func (rl *Relay) serve(l *link, id uint64, hello sandboxlink.ServeHello) { }() <-l.sess.CloseChan() rl.mu.Lock() - if rl.serves[key] == sl { - delete(rl.serves, key) + defer rl.mu.Unlock() + if r := rl.resources[key]; r != nil && r.serve == sl { + // The resource keeps the set for a reconnect. + r.serve = nil + rl.releaseLocked(key, r) + } else { + // A replaced link has no set; a revoked one takes it along. + rl.dropLocked(sl.closures) } sl.closures = nil - rl.mu.Unlock() } -// admitServe authenticates a serve Hello and installs the link. A revocation -// that lands while the Authority decides forces a fresh decision, so a -// withdrawn credential never installs a peer. +// admitServe authenticates a serve Hello and installs the link. Until it is +// decided, the Hello takes one of its resource's maxServeHellos, and a +// resource slot when the relay does not hold the resource. A revocation that +// lands while the Authority decides forces a fresh decision, so a withdrawn +// credential never installs a peer. func (rl *Relay) admitServe(l *link, id uint64, hello sandboxlink.ServeHello) (*serveLink, error) { + key := keyOf(hello.Resource) + rl.mu.Lock() + r := rl.resources[key] + if r == nil && len(rl.resources) >= rl.maxResources || r != nil && r.hellos >= maxServeHellos { + rl.mu.Unlock() + return nil, sandboxlink.Fail(sandboxlink.LimitExceeded) + } + if r == nil { + r = &resource{closures: closures{}} + rl.resources[key] = r + } + r.hellos++ + rl.mu.Unlock() + defer func() { + rl.mu.Lock() + r.hellos-- + rl.releaseLocked(key, r) + rl.mu.Unlock() + }() for range authorizeAttempts { rl.mu.Lock() epoch := rl.epoch @@ -382,7 +493,7 @@ func (rl *Relay) admitServe(l *link, id uint64, hello sandboxlink.ServeHello) (* rl.mu.Unlock() continue } - sl, old, err := rl.installServeLocked(l, id, hello) + sl, old, err := rl.installServeLocked(r, l, id, hello) rl.mu.Unlock() if old != nil { old.sess.Close() @@ -392,41 +503,62 @@ func (rl *Relay) admitServe(l *link, id uint64, hello sandboxlink.ServeHello) (* return nil, sandboxlink.Fail(sandboxlink.ServiceUnavailable) } -// installServeLocked makes l the resource's serve peer and queues its -// HelloAccepted while the decision is still current; its writer then writes -// the resource's pending AttachmentClosed events. It returns the replaced -// link for the caller to close. -func (rl *Relay) installServeLocked(l *link, id uint64, hello sandboxlink.ServeHello) (sl, old *serveLink, err error) { - key, generation := keyOf(hello.Resource), hello.Resource.Generation - if rl.generations[key] > generation { +// installServeLocked makes l r's serve peer and queues its HelloAccepted +// while the decision is still current; its writer then writes r's pending +// AttachmentClosed events. It returns the replaced link for the caller to +// close. +func (rl *Relay) installServeLocked(r *resource, l *link, id uint64, hello sandboxlink.ServeHello) (sl, old *serveLink, err error) { + generation := hello.Resource.Generation + if r.generation > generation { return nil, nil, sandboxlink.Fail(sandboxlink.StaleGeneration) } - if rl.generations[key] < generation { - rl.generations[key] = generation - delete(rl.closures, key) + if r.generation < generation { + r.generation = generation + rl.dropLocked(r.closures) + r.closures = closures{} for _, a := range rl.attachments { if a.identity.Resource.SameResource(hello.Resource) { rl.closeLocked(a, sandboxlink.CloseStaleGeneration) } } } - if rl.closures[key] == nil { - rl.closures[key] = closures{} - } sl = &serveLink{link: l, hello: hello} - sl.closures = rl.closures[key] - old = rl.serves[key] - if old != nil { + sl.closures = r.closures + if old = r.serve; old != nil { old.closures = nil } - rl.serves[key] = sl + r.serve = sl sl.send(id, sandboxlink.HelloAccepted{}) return sl, old, nil } // attach serves an attach peer's control requests and service streams until -// the link ends. Its attachments stay open until their leases expire. +// the link ends. The link takes an attach link slot first and keeps it until +// every stream and renewal it carried has finished, so the decisions of a +// dropped link stay bounded. Its attachments stay open until their leases +// expire. func (rl *Relay) attach(l *link, id uint64, hello sandboxlink.AttachHello) { + rl.mu.Lock() + full := rl.attachLinks >= rl.maxAttachLinks + if !full { + rl.attachLinks++ + l.closures = closures{} + } + rl.mu.Unlock() + if full { + l.fail(id, sandboxlink.OpHello, sandboxlink.LimitExceeded) + return + } + var serving sync.WaitGroup + defer func() { + <-l.sess.CloseChan() + serving.Wait() + rl.mu.Lock() + rl.attachLinks-- + rl.dropLocked(l.closures) + l.closures = nil + rl.mu.Unlock() + }() ctx, cancel := rl.authorityContext() peer, err := rl.auth.AuthenticateAttach(ctx, hello) cancel() @@ -441,15 +573,15 @@ func (rl *Relay) attach(l *link, id uint64, hello sandboxlink.AttachHello) { var seq sandboxwire.RequestSequence seq.Admit(id) // the Hello takes the first ID al.send(id, sandboxlink.HelloAccepted{}) - go func() { + serving.Go(func() { for { st, err := l.sess.AcceptStream() if err != nil { return } - go rl.open(al, st) + serving.Go(func() { rl.open(al, st) }) } - }() + }) for { id, m, err := sandboxlink.ReadMessage(l.ctl) if err != nil { @@ -468,10 +600,10 @@ func (rl *Relay) attach(l *link, id uint64, hello sandboxlink.AttachHello) { al.send(id, sandboxlink.FailureFor(sandboxlink.OpRenewAttachment, sandboxlink.Fail(sandboxlink.LimitExceeded))) continue } - go func() { + serving.Go(func() { defer al.inflight.Add(-1) rl.renew(al, id, r) - }() + }) case sandboxlink.CloseAttachment: if !admitted { l.fail(id, sandboxlink.OpCloseAttachment, sandboxlink.ProtocolViolation) @@ -555,27 +687,27 @@ func (rl *Relay) closeRequested(al *attachLink, id uint64, c sandboxlink.CloseAt // peer and the serve peer why. A serve peer that is away when its attachment // closes hears of it when the same generation reconnects. func (rl *Relay) closeLocked(a *attachment, reason sandboxlink.CloseReason) { - id := a.identity.AttachmentID - delete(rl.attachments, id) + delete(rl.attachments, a.identity.AttachmentID) a.timer.Stop() for sp := range a.splices { sp.abortLocked(abortCodes[reason]) } + a.reason = reason if reason != sandboxlink.CloseRequested { - a.owner.closedLocked(id, reason) - } - key, generation := keyOf(a.identity.Resource), a.identity.Resource.Generation - if !a.bound || rl.generations[key] != generation { - return - } - if sl := rl.serves[key]; sl != nil { - sl.closedLocked(id, reason) - return - } - if rl.closures[key] == nil { - rl.closures[key] = closures{} + rl.queueLocked(a.owner.closures, a) + a.owner.wakeWriter() + } + key := keyOf(a.identity.Resource) + r := rl.resources[key] // the attachment holds it + if a.bound && r.generation == a.identity.Resource.Generation { + rl.queueLocked(r.closures, a) + if r.serve != nil { + r.serve.wakeWriter() + } } - rl.closures[key][id] = reason + r.attachments-- + rl.releaseLocked(key, r) + rl.unrefLocked(a) } // abortCodes answers an Open that an attachment's close interrupts. @@ -661,15 +793,41 @@ func (rl *Relay) open(al *attachLink, st *yamux.Stream) { rl.splice(sp) } -// admit authorizes o and registers its splice. A revocation that lands while -// the Authority decides forces a fresh decision. -func (rl *Relay) admit(al *attachLink, st *yamux.Stream, o sandboxlink.Open) (*splice, sandboxlink.Authorization, error) { +// admit authorizes o and registers its splice. Before the Authority decides, +// the Open takes a stream slot of its link, which an admitted stream keeps +// until finish and a refused one gives back once it is decided, whether or +// not the peer reset it. An Open of an attachment the relay does not hold +// also takes an attachment slot. A revocation that lands while the Authority +// decides forces a fresh decision. +func (rl *Relay) admit(al *attachLink, st *yamux.Stream, o sandboxlink.Open) (sp *splice, auth sandboxlink.Authorization, err error) { + rl.mu.Lock() + prior := rl.attachments[o.AttachmentID] + reserved := prior == nil + if al.streams >= maxStreams || reserved && rl.held >= rl.maxAttachments { + rl.mu.Unlock() + return nil, auth, sandboxlink.Fail(sandboxlink.LimitExceeded) + } + al.streams++ + if reserved { + rl.held++ + } + rl.mu.Unlock() + defer func() { + rl.mu.Lock() + if sp == nil { + al.streams-- + } + if reserved { + rl.held-- + } + rl.mu.Unlock() + }() for range authorizeAttempts { rl.mu.Lock() epoch := rl.epoch rl.mu.Unlock() ctx, cancel := rl.authorityContext() - auth, err := rl.auth.AuthorizeOpen(ctx, al.peer, o) + auth, err = rl.auth.AuthorizeOpen(ctx, al.peer, o) cancel() if err != nil { return nil, auth, err @@ -682,16 +840,25 @@ func (rl *Relay) admit(al *attachLink, st *yamux.Stream, o sandboxlink.Open) (*s rl.mu.Unlock() continue } - sp, err := rl.admitLocked(al, st, o, auth) + sp, err = rl.admitLocked(al, st, o, auth, prior, &reserved) rl.mu.Unlock() return sp, auth, err } return nil, sandboxlink.Authorization{}, sandboxlink.Fail(sandboxlink.ServiceUnavailable) } -func (rl *Relay) admitLocked(al *attachLink, st *yamux.Stream, o sandboxlink.Open, auth sandboxlink.Authorization) (*splice, error) { +// admitLocked checks o against what the relay holds and registers its +// splice. A new attachment takes the slot reserved for it. The attachment +// prior, which the relay held when the Open arrived, must still be the one +// under its ID: once it closed, the Open fails even if another attachment +// took the ID since. +func (rl *Relay) admitLocked(al *attachLink, st *yamux.Stream, o sandboxlink.Open, auth sandboxlink.Authorization, prior *attachment, reserved *bool) (*splice, error) { key, generation := keyOf(o.Resource), o.Resource.Generation - sl := rl.serves[key] + r := rl.resources[key] + var sl *serveLink + if r != nil { + sl = r.serve + } var offered *sandboxlink.ServiceVersion if sl != nil { for i, s := range sl.hello.Services { @@ -703,9 +870,9 @@ func (rl *Relay) admitLocked(al *attachLink, st *yamux.Stream, o sandboxlink.Ope a := rl.attachments[o.AttachmentID] var code sandboxlink.Code switch { - case al.streams >= maxStreams || (sl != nil && sl.streams >= maxStreams): + case sl != nil && sl.streams >= maxStreams: code = sandboxlink.LimitExceeded - case rl.generations[key] > generation: + case r != nil && r.generation > generation: code = sandboxlink.StaleGeneration case sl == nil || sl.hello.Resource.Generation != generation || offered == nil: code = sandboxlink.ServiceUnavailable @@ -713,7 +880,7 @@ func (rl *Relay) admitLocked(al *attachLink, st *yamux.Stream, o sandboxlink.Ope code = sandboxlink.VersionMismatch case !o.ExpectedServerInstanceID.IsZero() && o.ExpectedServerInstanceID != sl.hello.ServerInstanceID: code = sandboxlink.InstanceChanged - case !auth.LeaseExpiresAt.After(time.Now()): + case !auth.LeaseExpiresAt.After(time.Now()) || prior != nil && a != prior: code = sandboxlink.LeaseExpired case a != nil && (a.identity != o.Identity() || a.runtime != al.peer.RuntimeID): code = sandboxlink.AttachmentConflict @@ -722,15 +889,16 @@ func (rl *Relay) admitLocked(al *attachLink, st *yamux.Stream, o sandboxlink.Ope return nil, sandboxlink.Fail(code) } if a == nil { - a = &attachment{identity: o.Identity(), runtime: al.peer.RuntimeID, splices: map[*splice]struct{}{}} + a = &attachment{identity: o.Identity(), runtime: al.peer.RuntimeID, splices: map[*splice]struct{}{}, refs: 1} rl.attachments[o.AttachmentID] = a + r.attachments++ + *reserved = false } a.owner = al a.bound = true rl.setLeaseLocked(a, auth.LeaseExpiresAt) sp := &splice{att: a, serve: sl, al: al, attach: st} a.splices[sp] = struct{}{} - al.streams++ sl.streams++ return sp, nil } diff --git a/internal/sandboxlink/relay/relay_test.go b/internal/sandboxlink/relay/relay_test.go index f7bff701b..9db74bad5 100644 --- a/internal/sandboxlink/relay/relay_test.go +++ b/internal/sandboxlink/relay/relay_test.go @@ -14,6 +14,7 @@ import ( "time" "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/relay" "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/sandboxlinktest" "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" ) @@ -50,12 +51,12 @@ func put[T any](ch chan T, v T) { } } -// hooked runs a test's hook after the Authority decides a serve Hello or a -// renewal, so a test can hold the decision. +// hooked runs a test's hook after the Authority decides a serve Hello, an +// Open or a renewal, so a test can hold the decision. type hooked struct { *sandboxlinktest.Authority - mu sync.Mutex - onServe, onRenew func() + mu sync.Mutex + onServe, onOpen, onRenew func() } func (h *hooked) hook(f *func()) { @@ -79,6 +80,12 @@ func (h *hooked) AuthenticateServe(ctx context.Context, hello sandboxlink.ServeH return peer, err } +func (h *hooked) AuthorizeOpen(ctx context.Context, peer sandboxlink.AttachPeer, o sandboxlink.Open) (sandboxlink.Authorization, error) { + auth, err := h.Authority.AuthorizeOpen(ctx, peer, o) + h.hook(&h.onOpen) + return auth, err +} + func (h *hooked) Renew(ctx context.Context, peer sandboxlink.AttachPeer, r sandboxlink.RenewAttachment) (sandboxlink.Authorization, error) { auth, err := h.Authority.Renew(ctx, peer, r) h.hook(&h.onRenew) @@ -193,6 +200,29 @@ func (f *fixture) startServe(generation uint64) *servePeer { return p } +// serveHello sends a serve Hello for ref with an unknown credential over a new +// link and returns the relay's answer. +func (f *fixture) serveHello(ref sandboxlink.ResourceRef) error { + ctx, cancel := context.WithTimeout(context.Background(), wait) + defer cancel() + conn, err := sandboxlink.DialWebSocket(ctx, f.srv.URL, f.srv.TLS) + if err != nil { + return err + } + sess, _ := sandboxlink.ClientSession(conn) + defer sess.Close() + ctl, err := sess.OpenStream(ctx) + if err == nil { + ctl.SetDeadline(time.Now().Add(wait)) + err = sandboxlink.WriteMessage(ctl, 1, sandboxlink.ServeHello{Version: sandboxlink.Version, Credential: []byte("unknown"), Resource: ref, + ServerInstanceID: sandboxwire.NewID(), Services: []sandboxlink.ServiceVersion{{Service: sandboxlink.ServiceFile, Version: 1}}}) + } + if err == nil { + _, err = sandboxlink.ReadReply(ctl, sandboxlink.OpHello, 1) + } + return err +} + // grant authorizes the fixture's Runtime for the resource's generation. func (f *fixture) grant(generation uint64) []byte { grant := []byte(fmt.Sprintf("grant %d", generation)) @@ -438,6 +468,238 @@ func TestClosuresReplayOnReconnect(t *testing.T) { } } +// A serve Hello for a resource the relay does not hold takes a slot before the +// Authority decides: beyond capacity even an unknown credential gets +// LimitExceeded. A refused Hello leaves no slot taken, and an admitted +// resource still reconnects at capacity. +func TestServesAreBounded(t *testing.T) { + f := newFixture(t) + relay.SetLimits(f.srv.Relay, 1, 16, 16) + other := resource(1) + other.ID = sandboxwire.NewID() + if err := f.serveHello(other); !errors.Is(err, sandboxlink.AuthenticationFailed) { + t.Fatalf("serve Hello with an unknown credential: %v, want AuthenticationFailed", err) + } + p := f.serve(1) + if err := f.serveHello(other); !errors.Is(err, sandboxlink.LimitExceeded) { + t.Fatalf("serve Hello beyond capacity: %v, want LimitExceeded", err) + } + if resources, _, _ := relay.Held(f.srv.Relay); resources != 1 { + t.Fatalf("relay holds %d resources, want 1", resources) + } + recv(t, p.conns).Close() + recv(t, p.connected) +} + +// Each Hello for a resource takes one of its MaxServeHellos while it is +// decided: past them, a Hello for a held resource gets LimitExceeded without +// reaching the Authority. +func TestServeHellosAreBounded(t *testing.T) { + f := newFixture(t) + f.serve(1) + entered, release := make(chan struct{}, relay.MaxServeHellos+1), make(chan struct{}) + f.auth.set(&f.auth.onServe, func() { + entered <- struct{}{} + <-release + }) + errs := make(chan error, relay.MaxServeHellos) + for range relay.MaxServeHellos { + go func() { errs <- f.serveHello(resource(1)) }() + recv(t, entered) + } + if err := f.serveHello(resource(1)); !errors.Is(err, sandboxlink.LimitExceeded) { + t.Fatalf("serve Hello beyond the bound: %v, want LimitExceeded", err) + } + close(release) + for range relay.MaxServeHellos { + if err := recv(t, errs); !errors.Is(err, sandboxlink.AuthenticationFailed) { + t.Fatalf("serve Hello within the bound: %v, want AuthenticationFailed", err) + } + } + if len(entered) != 0 { + t.Fatal("a serve Hello beyond the bound reached the Authority") + } +} + +// An Open takes one of its link's MaxStreams before the Authority decides and +// keeps it until the decision ends, even once the peer has reset the stream. +func TestOpenDecisionsAreBounded(t *testing.T) { + f := newFixture(t) + f.serve(1) + attachment := sandboxwire.NewID() + f.mustOpen(sandboxlink.ServiceFile, attachment, 1) // its stream takes one + grant := f.grant(1) + pending := relay.MaxStreams - 1 + entered, release := make(chan struct{}, pending+1), make(chan struct{}) + f.auth.set(&f.auth.onOpen, func() { + entered <- struct{}{} + <-release + }) + reset := make(chan error) + for range pending { + ctx, cancel := context.WithCancel(context.Background()) + go func() { + _, _, err := f.link.OpenService(ctx, sandboxlink.Open{Service: sandboxlink.ServiceFile, Version: 1, Resource: resource(1), + AttachmentID: attachment, SessionID: sessionID, AssignmentID: assignmentID, AssignmentEpoch: 1, AttachGrant: grant}) + reset <- err + }() + recv(t, entered) + cancel() // OpenService resets the stream + recv(t, reset) + } + if _, _, err := f.open(sandboxlink.ServiceFile, attachment, 1, grant); !errors.Is(err, sandboxlink.LimitExceeded) { + t.Fatalf("open beyond the stream bound: %v, want LimitExceeded", err) + } + if len(entered) != 0 { + t.Fatal("an Open beyond the bound reached the Authority") + } + f.auth.set(&f.auth.onOpen, nil) + close(release) + // The slots return as the decisions end; until then an Open is refused + // with LimitExceeded, which is retryable. + for { + _, _, err := f.open(sandboxlink.ServiceFile, attachment, 1, grant) + if err == nil { + break + } + if !errors.Is(err, sandboxlink.LimitExceeded) { + t.Fatalf("open after the decisions ended: %v", err) + } + } +} + +// An Open that finds an attachment fails with LeaseExpired when that +// attachment closes while the Open is decided, even if another Open has +// created a new attachment under the same ID since, and leaves the new +// attachment's lease alone. +func TestOpenOfAReplacedAttachment(t *testing.T) { + f := newFixture(t) + p := f.serve(1) + calls := make(chan chan struct{}) + f.auth.set(&f.auth.onOpen, func() { + release := make(chan struct{}) + calls <- release + <-release + }) + attachment := sandboxwire.NewID() + open := func(grant []byte) <-chan error { + errs := make(chan error, 1) + go func() { + _, _, err := f.open(sandboxlink.ServiceFile, attachment, 1, grant) + errs <- err + }() + return errs + } + grant := f.grant(1) + longer := []byte("longer grant") + f.auth.AddGrant(longer, sandboxlinktest.Grant{RuntimeID: f.runtime, Resource: resource(1), SessionID: sessionID, + AssignmentID: assignmentID, AssignmentEpoch: 1, Services: []sandboxlink.Service{sandboxlink.ServiceFile}, Lease: 2 * f.lease}) + + // Two Opens arrive while the relay holds no attachment; the first creates A. + first := open(grant) + releaseFirst := recv(t, calls) + second := open(grant) + releaseSecond := recv(t, calls) + close(releaseFirst) + if err := recv(t, first); err != nil { + t.Fatal(err) + } + // A third Open finds A, which closes while it is decided. + third := open(longer) + releaseThird := recv(t, calls) + if err := f.link.CloseAttachment(context.Background(), attachment); err != nil { + t.Fatal(err) + } + recv(t, p.closed) + // The second Open creates B under the same ID. The serve peer, which saw + // A close, refuses to bind it, but the relay holds B. + close(releaseSecond) + recv(t, second) + lease := relay.Lease(f.srv.Relay, attachment) + if lease.IsZero() { + t.Fatal("the second Open created no attachment") + } + close(releaseThird) + if err := recv(t, third); !errors.Is(err, sandboxlink.LeaseExpired) { + t.Fatalf("open of a replaced attachment: %v, want LeaseExpired", err) + } + if got := relay.Lease(f.srv.Relay, attachment); !got.Equal(lease) { + t.Fatalf("the new attachment's lease moved from %s to %s", lease, got) + } +} + +// An Open of an attachment the relay does not hold takes a slot before the +// Authority decides: beyond capacity even a forged grant gets LimitExceeded. +// A refused Open leaves no slot taken, and an open attachment still opens +// streams at capacity. +func TestAttachmentsAreBounded(t *testing.T) { + f := newFixture(t) + relay.SetLimits(f.srv.Relay, 16, 1, 16) + f.serve(1) + if _, _, err := f.open(sandboxlink.ServiceFile, sandboxwire.NewID(), 1, []byte("forged grant")); !errors.Is(err, sandboxlink.PermissionDenied) { + t.Fatalf("open with a forged grant: %v, want PermissionDenied", err) + } + attachment := sandboxwire.NewID() + f.mustOpen(sandboxlink.ServiceFile, attachment, 1) + if _, _, err := f.open(sandboxlink.ServiceFile, sandboxwire.NewID(), 1, []byte("forged grant")); !errors.Is(err, sandboxlink.LimitExceeded) { + t.Fatalf("open of a new attachment beyond capacity: %v, want LimitExceeded", err) + } + f.mustOpen(sandboxlink.ServiceFile, attachment, 1) + if _, attachments, _ := relay.Held(f.srv.Relay); attachments != 1 { + t.Fatalf("relay holds %d attachment slots, want 1", attachments) + } +} + +// A closed attachment keeps its slot until its AttachmentClosed event is +// written to the serve peer, which is away and then reconnects. +func TestClosuresKeepTheirSlots(t *testing.T) { + f := newFixture(t) + relay.SetLimits(f.srv.Relay, 16, 2, 16) + p := f.serve(1) + ids := []sandboxwire.ID{sandboxwire.NewID(), sandboxwire.NewID()} + for _, id := range ids { + f.mustOpen(sandboxlink.ServiceFile, id, 1) + } + held, release := make(chan struct{}), make(chan struct{}) + var once sync.Once + f.auth.set(&f.auth.onServe, func() { + once.Do(func() { close(held) }) + <-release + }) + recv(t, p.conns).Close() + recv(t, held) + for _, id := range ids { + if err := f.link.CloseAttachment(context.Background(), id); err != nil { + t.Fatal(err) + } + } + if _, _, err := f.open(sandboxlink.ServiceFile, sandboxwire.NewID(), 1, f.grant(1)); !errors.Is(err, sandboxlink.LimitExceeded) { + t.Fatalf("open while both closures are unwritten: %v, want LimitExceeded", err) + } + close(release) + for range ids { + if r := recv(t, p.closed); r != sandboxlink.CloseRequested { + t.Fatalf("serve peer saw close reason %d, want a requested close", r) + } + } + // The relay frees an event's slot before it writes the next, so one slot + // is free once both events have arrived. + f.mustOpen(sandboxlink.ServiceFile, sandboxwire.NewID(), 1) +} + +// An attach Hello takes a link slot before the Authority decides: beyond +// capacity even an unknown credential gets LimitExceeded. +func TestAttachLinksAreBounded(t *testing.T) { + f := newFixture(t) + relay.SetLimits(f.srv.Relay, 16, 16, 1) + ctx, cancel := context.WithTimeout(context.Background(), wait) + defer cancel() + _, err := sandboxlink.DialAttach(ctx, sandboxlink.AttachConfig{URL: f.srv.URL, TLS: f.srv.TLS, RuntimeID: sandboxwire.NewID(), Credential: []byte("unknown")}) + if !errors.Is(err, sandboxlink.LimitExceeded) { + t.Fatalf("attach Hello beyond capacity: %v, want LimitExceeded", err) + } +} + // A control call whose context ended before it was sent fails with no effect // and leaves the link up. func TestCancelledCallKeepsLink(t *testing.T) {