Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
36e74cda11 | ||
|
|
f8dad9dbc8 | ||
|
|
2586c3b217 | ||
|
|
48ceeba55c | ||
|
|
88e3f5e9c3 |
No files matched your search
@@ -4,6 +4,55 @@ All notable changes to this crate are documented here. The format is
|
||||
based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and
|
||||
this crate adheres to [Semantic Versioning](https://semver.org/).
|
||||
|
||||
## [0.5.0] - 2026-09-06
|
||||
|
||||
The channel-open establishment phase (ADR-049 — review 006 E-01 +
|
||||
N-1). Additive establishment machinery; one breaking change (the
|
||||
`ChannelClient::open_channel` error type, below).
|
||||
|
||||
### Added
|
||||
|
||||
- **`OperationSpec.description` (review 006 E-02).** Additive
|
||||
`Option<String>` op description (builder: `with_description`),
|
||||
disclosed by `services/list` and `services/list-peers` (local
|
||||
listings) when set, carried in the `services/schema` wire shape
|
||||
(`spec_to_json_pub` emit / `rebuild_spec_for` parse — the field
|
||||
survives `from_call` discovery and `op/register` announcement).
|
||||
Describes the op, not the produced resource set (ADR-047 §6
|
||||
amendment: the live resource-enumeration half stays deferred,
|
||||
OQ-40). Absent on the wire when unset — additive for all consumers.
|
||||
|
||||
- **Channel-open establishment phase (ADR-049 — review 006 E-01).**
|
||||
`ChannelCore::register_openable_with_establisher` registers a
|
||||
per-ALPN open op with an establisher hook (`OpenEstablisher`): an
|
||||
awaited establishment phase — validate params semantically, dial the
|
||||
backend — bounded by the dispatch deadline or the registration's
|
||||
timeout override, else `ESTABLISHMENT_TIMEOUT` (10s; the earlier of
|
||||
the two). On establisher failure (error or deadline) the wrapper
|
||||
tears down the just-allocated channel (demux sender, opener-ledger
|
||||
take, `policy.on_close` un-increment — the allocation and teardown
|
||||
balance) and replies `channel:open_failed` with
|
||||
`details: { reason, message }`, reason ∈ `dial_failed` /
|
||||
`unknown_resource` / `resource_shortage` / `handler_error` /
|
||||
`timeout`. The SSH contract holds consumer-visibly: a failed open
|
||||
never returns a `channel_id`. `register_openable` is unchanged
|
||||
(no establisher = always-OK, existing registrations compile and
|
||||
behave identically); the establisher takes `(input, auth)` — the
|
||||
channel's yield-once `BiStream` belongs exclusively to the pump
|
||||
handler. Additive wire surface (new error-code string + `details`
|
||||
shape).
|
||||
|
||||
### Changed
|
||||
|
||||
- **`ChannelClient::open_channel` returns a typed error (ADR-049 §4 —
|
||||
review 006 N-1, breaking at 0.5.0).** The open op's `CallError` is
|
||||
carried verbatim in `ChannelOpenError::CallFailed` instead of being
|
||||
flattened into a debug-formatted string, so consumers branch on
|
||||
`channel:open_failed`'s typed reason (`establishment_reason()`).
|
||||
Other variants: `MissingChannelId` (malformed success reply),
|
||||
`AdoptFailed` (local adoption failure). Mechanical for consumers —
|
||||
the `String` was a debug-formatting wrapper.
|
||||
|
||||
## [0.4.1] - 2026-09-05
|
||||
|
||||
Bug-fix release: chunks arriving for a not-yet-adopted channel are
|
||||
|
||||
Generated
+1
-1
@@ -27,7 +27,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "alkcall"
|
||||
version = "0.4.1"
|
||||
version = "0.5.0"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"bytes",
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "alkcall"
|
||||
version = "0.4.1"
|
||||
version = "0.5.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.85"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -68,12 +68,19 @@ impl ChannelClient {
|
||||
/// the `MpscRecvStream` (read half). The caller can build a
|
||||
/// `Connection` from these via `channel_source` and
|
||||
/// `Connection::from_source`.
|
||||
///
|
||||
/// The error is typed (ADR-049 §4 — review 006 N-1):
|
||||
/// `ChannelOpenError::CallFailed` carries the wire `CallError`
|
||||
/// verbatim (branch on `establishment_reason()` for
|
||||
/// `channel:open_failed`'s `details.reason`); `MissingChannelId`
|
||||
/// covers a malformed success reply; `AdoptFailed` covers local
|
||||
/// adoption failure.
|
||||
pub async fn open_channel(
|
||||
&self,
|
||||
operation_id: &str,
|
||||
input: Value,
|
||||
alpn: &str,
|
||||
) -> Result<(u32, MpscSendStream, MpscRecvStream), String>;
|
||||
) -> Result<(u32, MpscSendStream, MpscRecvStream), ChannelOpenError>;
|
||||
|
||||
/// Take the `CallConnection` — used by the consumer to register
|
||||
/// imported ops (`from_call`) on the connection's overlay. After
|
||||
@@ -88,7 +95,11 @@ impl ChannelClient {
|
||||
ADR-047 dissolved the generic `channel/open` operation into per-ALPN
|
||||
open ops (`channels/<alpn>/sub`, `channels/<alpn>/pub`). Each ALPN
|
||||
crate registers its own open op via `ChannelCore::register_openable`
|
||||
(ADR-047 §3, as amended 2026-08-13 — per-connection registration).
|
||||
(ADR-047 §3, as amended 2026-08-13 — per-connection registration),
|
||||
optionally with an establisher via
|
||||
`register_openable_with_establisher` (ADR-049 — the awaited,
|
||||
bounded establishment phase; establishment failure resolves
|
||||
`Err` on the client with `channel:open_failed` + `details.reason`).
|
||||
The `ChannelClient` calls these ops by name on channel 0 via
|
||||
`call_open_op`; `open_channel` wraps `call_open_op` + `adopt_channel`.
|
||||
|
||||
@@ -97,6 +108,35 @@ and `subscribe_resources` method are deferred (OQ-40, OQ-41). The
|
||||
`ResourceEntry.access` preview was dropped by ADR-047 §6 — it is
|
||||
available via `services/schema` on the op spec.
|
||||
|
||||
## The early-arrival park bound (push-first producers)
|
||||
|
||||
The open-op response carrying `channel_id` races the producer's first
|
||||
data-plane writes: the producer's handler can start pumping before the
|
||||
consumer's `adopt_channel` runs. The demux parks those first chunks
|
||||
per-channel (FIFO) instead of dropping them, and `adopt_channel` drains
|
||||
the parked chunks into the new receiver — a push-first producer (a TTY
|
||||
backend's banner, a sub protocol's greeting) must not lose its first
|
||||
chunks.
|
||||
|
||||
The park is bounded: **up to 64 chunks per channel** are parked
|
||||
(`EARLY_ARRIVAL_CAP`, ADR-040's memory bounds); chunks arriving past
|
||||
the cap are dropped silently — at the consumer this presents as a
|
||||
truncated stream with clean framing everywhere else, not as an error.
|
||||
Consumer-facing consequences (review 006 E-04, noted for tunnel-style
|
||||
consumers):
|
||||
|
||||
- Size any pre-adopt buffering math (e.g. a UDP POC's MTU-vs-buffer
|
||||
sizing) against **64 parked chunks as the observable bound**, and
|
||||
adopt promptly — the open reply resolves *before* the first data
|
||||
arrives by design, so the adopt is the consumer's next step, not a
|
||||
slow path.
|
||||
- The producer side observes both halves of the race via
|
||||
`ChannelManager::early_arrival_count()` (chunks parked) and
|
||||
`ChannelManager::dropped_unknown_chunks()` (chunks lost to the cap) —
|
||||
a non-zero dropped counter under a push-first producer means the
|
||||
adopter was too slow for the producer's burst, and the affected
|
||||
streams were truncated at the park boundary.
|
||||
|
||||
## Transport-agnostic by construction
|
||||
|
||||
`ChannelClient` is the client side of the channels protocol. The channels
|
||||
|
||||
@@ -77,12 +77,39 @@ round-trip, so the open round-trip is not additive latency.
|
||||
| `channel:forbidden` | `AccessControl::check` denied the open | false |
|
||||
| `channel:allocation_failed` | Handler allocate failed (e.g., backend couldn't start) | true (often transient) |
|
||||
| `channel:too_many_channels` | Per-connection or per-identity channel limit hit (ADR-040, ADR-041) | false |
|
||||
| `channel:open_failed` | The establishment phase failed (ADR-049) — `details: { reason, message }`, reason ∈ `dial_failed` / `unknown_resource` / `resource_shortage` / `handler_error` / `timeout` | false |
|
||||
| `channel:no_channels_session` | The op was invoked outside a channels session (no `ChannelManager` — ADR-047 §2) | false |
|
||||
|
||||
`channel:unknown_alpn` and `channel:invalid_params` (ADR-037) are gone —
|
||||
an unregistered ALPN is an ordinary `NOT_FOUND` (op not registered); a
|
||||
bad `input` is ordinary schema rejection.
|
||||
|
||||
### The establishment phase (ADR-049)
|
||||
|
||||
An openable ALPN may register an **establisher** alongside its open
|
||||
handler (`ChannelCore::register_openable_with_establisher`). The
|
||||
establisher is the awaited preparation step ADR-047 §3 described
|
||||
("validate params, consult ownership, prepare the backend, return a
|
||||
channel plan"): the wrapper awaits it **bounded** (the dispatch
|
||||
deadline when the op carries one, else the registration's override or
|
||||
the 10s `ESTABLISHMENT_TIMEOUT` — the earlier of the two), **before
|
||||
replying and before spawning the pump handler**. On establisher
|
||||
failure (error or deadline) the wrapper tears down the just-allocated
|
||||
channel (demux sender, opener-ledger entry, `policy.on_close`
|
||||
un-increment — the same atomic ledger-take gate every teardown path
|
||||
runs) and replies `channel:open_failed` with
|
||||
`details: { reason, message }` (the table row above). The SSH
|
||||
contract holds consumer-visibly: a failed open never returns a
|
||||
`channel_id`.
|
||||
|
||||
Registrations without an establisher behave exactly as before (the
|
||||
open op cannot fail post-allocation; establishment work inside the
|
||||
handler is invisible to the open reply). The bound applies only to
|
||||
the establisher — the spawned pump handler's lifetime is governed by
|
||||
the existing teardown machinery, unchanged. Head-of-line safety: the
|
||||
serving loop spawns Once invocations as independent tasks, so a slow
|
||||
establisher on one open op does not block other calls on channel 0.
|
||||
|
||||
### `channels/<alpn>/pub` — publish a binary stream
|
||||
|
||||
`OperationType::Pub` (ADR-046). The initiator publishes a stream of
|
||||
|
||||
@@ -8,7 +8,40 @@ Accepted (amends ADR-037; refines ADR-044, ADR-046; §4 amended
|
||||
registration, 2026-08-13)"; amendment #2 (2026-09-03) — the
|
||||
per-connection registration mechanism is the **per-session fork of the
|
||||
base registry installed as the session's dispatch registry**, not the
|
||||
connection overlay — see "Amendment (§4 mechanism, 2026-09-03)" below)
|
||||
connection overlay — see "Amendment (§4 mechanism, 2026-09-03)" below;
|
||||
§6 amended 2026-09-06 — the static half of discovery gains an additive
|
||||
per-op `description` on the listing, the dynamic half stays deferred
|
||||
— see "Amendment (§6 listing enrichment, 2026-09-06)" below)
|
||||
|
||||
## Amendment (§6 listing enrichment, 2026-09-06)
|
||||
|
||||
Review 006 E-02 (from the alktunnels Phase 0 sweep) made §6's dynamic
|
||||
half load-bearing for the first time and asked for a decision on the
|
||||
static half. Decision: **both halves of the "static per-op" split gain
|
||||
what is cheap today; the dynamic half stays deferred** (OQ-40, now
|
||||
with a "load-bearing for alktunnels discovery UI" note in
|
||||
`open-questions.md`; alktunnels v1 uses config-known op names):
|
||||
|
||||
- `OperationSpec` gains `description: Option<String>` — a human-
|
||||
readable op description, set via `with_description` at registration.
|
||||
Additive (defaults `None`; no struct-literal construction sites
|
||||
exist — all sites use `OperationSpec::new`).
|
||||
- `services/list` (and the local-ops half of `services/list-peers`)
|
||||
emit `description` when set — one round-trip answers "which ops
|
||||
exist and what are they for" without the N+1 `services/schema`
|
||||
sweep. The output-schema docs on both listing specs and on
|
||||
`services/schema`'s `operation_spec_schema` advertise the field.
|
||||
- `spec_to_json_pub` emits it when set; `rebuild_spec_for` parses it
|
||||
back — the description survives discovery and peer announcement
|
||||
(`from_call`, `op/register`) like every other additive spec field
|
||||
(`resource_id_path`, `publish_schema` round-trip the same way).
|
||||
|
||||
Scope note from E-02 stands: the listing field describes **the op**,
|
||||
not **the produced resource set** (a tunnel producer registers one op
|
||||
and N resources). "Which tunnel resources may I open, live" remains
|
||||
`channel/resources/subscribe`'s job (§6's dynamic half, ADR-037 §
|
||||
`channel/resources/subscribe`) — deferred until a consumer needs live
|
||||
resource discovery.
|
||||
|
||||
## Amendment (§4 mechanism, 2026-09-03)
|
||||
|
||||
|
||||
@@ -0,0 +1,356 @@
|
||||
# ADR-049: Channel-Open Establishment Phase (`OpenEstablisher`)
|
||||
|
||||
## Status
|
||||
|
||||
Accepted — implemented in alkcall 0.5.0 (Unit 1: E-01 + N-1; see the
|
||||
"Amendment (Unit 1 implementation, 2026-09-06)" at the bottom).
|
||||
Amends ADR-047 §3 — the open-op wrapper gains an awaited
|
||||
establishment phase ahead of the spawned pump handler; resolves review
|
||||
006 E-01 and N-1
|
||||
|
||||
## Context
|
||||
|
||||
ADR-047 §3 made openable ALPNs operations: the ALPN crate supplies an
|
||||
open handler, and the channels wrapper does the channel machinery
|
||||
(allocation, ledger, policy, spawn). The §3 decision text describes the
|
||||
handler's job as "validate params, consult ownership, prepare the
|
||||
backend, return a 'channel plan'" — an awaited preparation step the
|
||||
wrapper consults **before** replying. The implemented `OpenHandler`
|
||||
type (`Arc<dyn Fn(Value, Connection, AuthContext) -> JoinHandle<()>`)
|
||||
collapsed that preparation into a fire-and-forget spawn: the wrapper
|
||||
collects the `JoinHandle`, records it for teardown, and writes
|
||||
`{ "channel_id": <id> }` to the wire the moment the handler task is
|
||||
*spawned* (`run_open_wrapper`, `src/channels/operations.rs`). The
|
||||
establishment phase the ADR described never became a thing the wrapper
|
||||
could consult.
|
||||
|
||||
The consequence (review 006 E-01, verified at tree `88e3f5e`): the open
|
||||
op **cannot fail after allocation**. Any establishment failure inside
|
||||
the handler — params valid at the schema level but semantically
|
||||
rejected, backend lookup failure, a target dial refused for a
|
||||
`direct-tcpip`-shaped tunnel, a resource no longer available — is
|
||||
invisible to the open reply. The consumer observes: the call op
|
||||
succeeds with `{channel_id}`, the channel is adopted, and then the
|
||||
channel EOFs (the handler exits without writing; the mux pump writes
|
||||
the implicit-EOF chunk; `MpscRecvStream::poll_read` returns clean EOF
|
||||
for both the sentinel and sender-drop arms). A dial failure is
|
||||
byte-for-byte indistinguishable from a target that closed immediately
|
||||
after connecting — the two most different failure/success stories map
|
||||
to the same consumer-visible event.
|
||||
|
||||
Every established tunnel/forwarding protocol puts establishment failure
|
||||
in the open reply, not in the data stream:
|
||||
|
||||
- **SSH** (RFC 4254 §5.1): `SSH_MSG_CHANNEL_OPEN_FAILURE` is a
|
||||
first-class reply carrying a reason code
|
||||
(`ADMINISTRATIVELY_PROHIBITED` / `CONNECT_FAILED` /
|
||||
`UNKNOWN_CHANNEL_TYPE` / `RESOURCE_SHORTAGE`) plus a description
|
||||
string; the channel never exists on the opener's side afterward.
|
||||
- **SOCKS5** (RFC 1928 §6): the reply carries REP codes 0x01–0x08;
|
||||
error-then-close, never "success then in-stream error."
|
||||
- **udpgw** (tun2proxy) is the counterexample: an opaque ERR bit with
|
||||
zero reason information — the vocabulary to avoid.
|
||||
|
||||
The per-crate workaround proves the gap is load-bearing: alktty's
|
||||
channels path answers establishment failures with a length-prefixed
|
||||
JSON error frame **on the channel stream**
|
||||
(`send_negotiation_error`, `alktty/src/adapter.rs`), disambiguated from
|
||||
data by a `0x00` first-byte peek. That is a per-crate reinvention of a
|
||||
protocol-level capability every ALPN crate will need: a structured,
|
||||
typed, **establishment-failure reply to the open op**. It also forces
|
||||
the phantom-opened channel to exist in the manager — allocation, ledger,
|
||||
and policy all fire for a channel that never carries data.
|
||||
|
||||
alktunnels (the next consumer, arbitrary TCP/UDP tunnels over channels)
|
||||
hits this on day one: its producer dials the tunnel target inside the
|
||||
open op, and dial failure is the *common* case, not the edge case
|
||||
(review 006 E-01; alktunnels OQ-TN-09).
|
||||
|
||||
This is the cheapest moment to fix upstream: three downstream crates
|
||||
(alktty, alkhttp, alktunnels-in-progress), alkcall 0.4.x, no published
|
||||
consumer depends on the phantom-open shape.
|
||||
|
||||
## Decision
|
||||
|
||||
### 1. The open-op wrapper gains an awaited establishment phase
|
||||
|
||||
`ChannelCore::register_openable` accepts an optional
|
||||
**establisher** alongside the existing `OpenHandler`:
|
||||
|
||||
```rust
|
||||
/// The establishment result. `Establishment` carries what the
|
||||
/// pump phase needs (today: nothing — reserved for a channel plan).
|
||||
/// `EstablishmentError` carries the reason code + message the
|
||||
/// wrapper puts in the open reply's `details`.
|
||||
pub type OpenEstablisher = Arc<
|
||||
dyn Fn(Value, Connection, AuthContext)
|
||||
-> BoxFuture<'static, Result<Establishment, EstablishmentError>>
|
||||
+ Send
|
||||
+ Sync,
|
||||
>;
|
||||
```
|
||||
|
||||
The establisher is **awaited by the wrapper, bounded** — before the
|
||||
reply is written, before the pump handler is spawned. The pump handler
|
||||
(the existing `OpenHandler`, unchanged) is spawned only on
|
||||
establishment success. The dial is the natural establisher step for
|
||||
tunnels; the pumps remain the spawned handler. This restores ADR-047
|
||||
§3's original "channel plan" shape: the establisher is the awaited
|
||||
preparation, the wrapper consults its result, the pumps are the
|
||||
spawned protocol.
|
||||
|
||||
**Why split-hook, not await-and-inspect** (the two candidates review
|
||||
006 proposed):
|
||||
|
||||
- The split preserves the `OpenHandler` type exactly. Await-and-inspect
|
||||
changes `OpenHandler`'s return type (`JoinHandle<()>` →
|
||||
`JoinHandle<OpenResult>`), breaking all three consumers' handlers and
|
||||
the alkhttp `OpenableAlpn` ferry for no compensating gain.
|
||||
- Establishment and pump are genuinely different lifecycles. The dial
|
||||
is synchronous with the open reply (SSH's semantics: the failure is
|
||||
the open's reply); the pumps outlive the reply. Coupling the
|
||||
establishment signal to the pump task's lifecycle (watching a
|
||||
`JoinHandle` for a first resolution) conflates them and makes the
|
||||
"established, continue" signal an out-of-band convention (a sentinel
|
||||
`Result` value, a oneshot the handler must remember to signal) —
|
||||
more protocol per crate, the thing this fix exists to remove.
|
||||
- The split is backward compatible by construction: no establisher
|
||||
registered = an always-OK establisher. Existing registrations compile
|
||||
and behave unchanged.
|
||||
|
||||
### 2. The deadline bounds the establisher, not the pumps
|
||||
|
||||
The establisher await is bounded by the dispatch deadline when the
|
||||
`OperationContext` carries one (`context.deadline`), else by a crate
|
||||
constant (`ESTABLISHMENT_TIMEOUT`, 10s default; overridable per
|
||||
registration via a `Duration` argument on the establisher-taking
|
||||
`register_openable` variant). On deadline expiry the wrapper treats it
|
||||
as establishment failure with reason `timeout`.
|
||||
|
||||
The bound applies **only** to the establisher. The spawned pump
|
||||
handler's lifetime is governed by the existing teardown machinery
|
||||
(`channel/close`, connection drop, handler exit) — unchanged.
|
||||
|
||||
Head-of-line safety is already proven: the serving loop spawns Once
|
||||
invocations as independent tasks
|
||||
(`Dispatcher::spawn_once_dispatch`), so a slow establisher on one open
|
||||
op does not block other calls on channel 0.
|
||||
|
||||
### 3. Establishment failure: teardown + typed `channel:open_failed`
|
||||
|
||||
On establishment failure (error or deadline), the wrapper:
|
||||
|
||||
1. Tears down the just-allocated channel
|
||||
(`teardown_channel` — drops the demux sender, returns the not-yet-
|
||||
installed handler task handle if any),
|
||||
2. Takes the opener-ledger entry and calls `policy.on_close(opener)`
|
||||
(the same un-increment path the allocation-failure arms already
|
||||
run — the ledger `take` is the atomic gate, ADR-047 §7),
|
||||
3. Replies with a new typed error:
|
||||
|
||||
```
|
||||
code: "channel:open_failed"
|
||||
message: human-readable establishment failure description
|
||||
retryable: false
|
||||
details: { "reason": <reason-code>, "message": <detail string> }
|
||||
```
|
||||
|
||||
The reason-code vocabulary maps 1:1 onto what an establisher can
|
||||
actually produce (per the SSH four; the survey's finding):
|
||||
|
||||
| reason | meaning |
|
||||
|---|---|
|
||||
| `dial_failed` | the backend/target could not be reached or refused |
|
||||
| `unknown_resource` | the requested resource does not exist |
|
||||
| `resource_shortage` | the backend is out of capacity (ports, fds, slots) |
|
||||
| `handler_error` | establisher-internal failure not covered above |
|
||||
| `timeout` | establishment exceeded the deadline |
|
||||
|
||||
Policy denial stays `channel:too_many_channels` (pre-allocation,
|
||||
unchanged); ACL denial stays `FORBIDDEN` (registry gate, unchanged).
|
||||
The new code is an additive wire addition (new error-code string +
|
||||
optional `details` shape); no existing consumer breaks. ALPN crates'
|
||||
open-op specs gain matching `ErrorDefinition` entries per ADR-016 so
|
||||
`services/schema` discloses the failure contract.
|
||||
|
||||
The SSH "channel never exists opener-side" property is the contract:
|
||||
the consumer's open resolves `Err` and no `channel_id` was ever
|
||||
returned. (The allocation still happened accept-side momentarily —
|
||||
that is invisible to the consumer and is what the teardown in step 1
|
||||
cleans up.)
|
||||
|
||||
### 4. `ChannelClient::open_channel` stops erasing the error (review 006 N-1)
|
||||
|
||||
`ChannelClient::open_channel` currently flattens the `CallError` into a
|
||||
`String` (`format!("open op failed: {e:?}")`), which would make the
|
||||
typed reason invisible to consumers — the E-01 fix would be unreachable
|
||||
end-to-end through the primary client path. It changes to return a
|
||||
typed error carrying the `CallError` (a new
|
||||
`ChannelOpenError { error: CallError }` or equivalent), so the
|
||||
consumer branches on `channel:open_failed` + `details.reason`.
|
||||
|
||||
This is a breaking change to a method signature introduced in this
|
||||
crate's 0.4.x — acceptable at 0.5.0 (see Consequences), and it is the
|
||||
point of the change: the reason must be consumer-usable.
|
||||
|
||||
### 5. Compatibility and migration
|
||||
|
||||
- `OpenHandler`'s type is unchanged. Existing registrations compile
|
||||
unchanged.
|
||||
- `ChannelCore::register_openable` keeps its current signature
|
||||
(no establisher = always-OK); a new
|
||||
`register_openable_with_establisher(spec, establisher, open_handler,
|
||||
registry, auth)` variant adds the hook. alkhttp's `OpenableAlpn`
|
||||
gains an optional `establisher` field (default `None`) — the ferry
|
||||
passes it through mechanically.
|
||||
- alktty migrates its **channels path** semantic failures (unknown
|
||||
backend, `carriage != "raw"`, `allocate_failed`, ownership denial —
|
||||
currently post-open error frames) into the establisher, resolving
|
||||
them as `channel:open_failed`. Its direct-ALPN path **keeps** the
|
||||
in-band error frame (two transports, two contracts; the direct path
|
||||
has no open op to fail). The `0x00`-peek disambiguation stays for
|
||||
the direct path only.
|
||||
- alkhttp is unaffected (no openable ops in the default surface; the
|
||||
`OpenableAlpn` change is additive).
|
||||
|
||||
### 6. Panicked pump handlers stay EOF-shaped (pinned as designed)
|
||||
|
||||
The wrapper's teardown task swallows the pump handler's `JoinError`
|
||||
(`let _ = raw_task.await`). A panicked pump = instant EOF, which is
|
||||
the correct consumer-visible outcome for a mid-stream handler crash
|
||||
(indistinguishable from an abrupt close — there is no error channel
|
||||
mid-stream by design; establishment errors are the only kind that
|
||||
belong in the open reply). This ADR pins that as intended; no change.
|
||||
The establisher, by contrast, runs pre-reply — its panic (a future
|
||||
that panics when polled) surfaces as the spawned Once task's panic,
|
||||
which the serving loop already tolerates (the call never resolves;
|
||||
the deadline / client timeout is the bound). Establisher
|
||||
implementations return `EstablishmentError` instead of panicking, per
|
||||
this crate's no-panic convention.
|
||||
|
||||
## Consequences
|
||||
|
||||
**Positive:**
|
||||
|
||||
- Establishment failure reaches the consumer as a typed, branchable
|
||||
call error — retry policy, client UX, and error reporting become
|
||||
possible for dial-refused, unknown-resource, and shortage cases
|
||||
(previously: instant-EOF ambiguity).
|
||||
- The SSH contract ("the channel never exists opener-side") holds
|
||||
consumer-visibly: a failed open never returns a `channel_id`.
|
||||
- No phantom channels: the ledger, policy count, and manager state are
|
||||
restored atomically on failure — allocation and teardown balance.
|
||||
- alktty's per-crate in-band error vocabulary is retired on the
|
||||
channels path; every future ALPN crate (alktunnels first) gets the
|
||||
establishment reply for free.
|
||||
- ADR-047 §3's "channel plan" shape is realized: awaited preparation
|
||||
before reply, spawned pumps after.
|
||||
|
||||
**Negative:**
|
||||
|
||||
- `channel:open_failed` + the reason vocabulary is a new wire-visible
|
||||
error surface — additive, but it joins the stable error set
|
||||
consumers may branch on (per ADR-016, `details` shapes are
|
||||
discoverable via `services/schema`).
|
||||
- `ChannelClient::open_channel`'s error type changes (breaking at
|
||||
0.5.0; mechanical for consumers — the `String` was a
|
||||
debug-formatting wrapper anyway).
|
||||
- `OpenableAlpn` (alkhttp) gains a field; its two construction sites
|
||||
add `None` (mechanical).
|
||||
- The establisher await adds a bounded latency to open-op replies
|
||||
where handlers previously replied instantly (the spawn). The 10s
|
||||
default is the worst case for a hung establisher; real establishers
|
||||
(dial, lookup) complete in dial-time. Consumers already tolerate
|
||||
call-op latency; the deadline is the bound.
|
||||
|
||||
## Door type
|
||||
|
||||
**One-way (wire-visible error surface).** `channel:open_failed` and its
|
||||
`details.reason` vocabulary join the stable error set: once consumers
|
||||
branch on reason codes, changing the vocabulary requires a migration
|
||||
(the same one-way-ness ADR-016 gives typed error details). The
|
||||
establisher hook shape itself — `OpenEstablisher`, the
|
||||
`register_openable_with_establisher` variant, the
|
||||
`Establishment`/`EstablishmentError` types — is a **two-way-door
|
||||
implementation detail** within the one-way decision (the wrapper shape,
|
||||
per ADR-047 §3's own door-type note). The `OpenHandler` type is
|
||||
untouched, which is what keeps the split cheap to revise.
|
||||
|
||||
## Implementation units
|
||||
|
||||
1. **alkcall 0.5.0** — `OpenEstablisher` +
|
||||
`register_openable_with_establisher`; wrapper flow (await bounded →
|
||||
teardown-on-failure → `channel:open_failed` with details);
|
||||
`ChannelClient::open_channel` typed error (N-1); tests:
|
||||
- establisher fails after allocation → consumer's `open_channel`
|
||||
resolves `Err(channel:open_failed)` + reason details; channel
|
||||
absent from `channel_ids()` afterward;
|
||||
- establisher never completes → `timeout`-reason failure within the
|
||||
deadline, channel torn down, ledger decremented;
|
||||
- no-establisher registration behaves exactly as today (compat
|
||||
gate);
|
||||
- establisher success spawns pumps and replies `{channel_id}`
|
||||
unchanged.
|
||||
2. **alktty migration** — channels-path semantic failures move into an
|
||||
establisher; `open_via_channels_surfaces_negotiation_rejected`
|
||||
resolves via call error; the direct-ALPN error-frame path is
|
||||
retained.
|
||||
3. **alkhttp pass** — `OpenableAlpn.establisher: Option<...>` (default
|
||||
`None`), threaded through the session fork (mechanical).
|
||||
|
||||
## References
|
||||
|
||||
- Review 006 E-01 (the establishment gap — findings and prior-art
|
||||
survey), N-1 (the client error-type gap this ADR also resolves),
|
||||
E-03/E-04 (adjacent teardown/early-arrival notes, filed separately
|
||||
from this ADR's scope)
|
||||
- ADR-047 §3 (openable ALPNs are operations — the "channel plan"
|
||||
wrapper shape this ADR restores; §7 opener ledger — the teardown
|
||||
un-increment path)
|
||||
- ADR-016 (typed error schemas — the `details` vehicle)
|
||||
- ADR-040/041 (backpressure/caps — untouched; the teardown path keeps
|
||||
the ledger `take` as the atomic gate)
|
||||
- alktty ADR-009 (the open op's input is the negotiation) + review 001
|
||||
L1/L3 — the in-band mechanism retired on the channels path
|
||||
- alktunnels OQ-TN-09 (dial-failure reporting — the first consumer of
|
||||
the new error) and `docs/research/ssh-socks5-survey.md`
|
||||
§"Open-failure path" (the reason-code prior art)
|
||||
- RFC 4254 §5.1, RFC 1928 §6 — SSH/SOCKS5 open-failure semantics
|
||||
|
||||
## Amendment (Unit 1 implementation, 2026-09-06)
|
||||
|
||||
Unit 1 landed in alkcall 0.5.0. Two implementation-shape notes, both
|
||||
within this ADR's two-way door (the hook shape is the revisable
|
||||
implementation detail; the wire surface is unchanged from §1/§3):
|
||||
|
||||
1. **The establisher does not receive the channel `Connection`.** §1's
|
||||
signature sketch passed `Connection` to the establisher, but the
|
||||
channel's `BiStream` is yield-once (`ChannelBidiStreamSource`) — it
|
||||
cannot be handed to both the establisher and the pump handler, and
|
||||
the establisher is pre-data-plane by design (its dial targets the
|
||||
backend, not the channel). The implemented signature is
|
||||
`Fn(Value, AuthContext) -> BoxFuture<'static,
|
||||
Result<Establishment, EstablishmentError>>`; the `Connection`
|
||||
belongs exclusively to the `OpenHandler` (unchanged).
|
||||
2. **The bound is the earlier of the dispatch deadline and the
|
||||
per-registration timeout.** §2 names the dispatch deadline "when
|
||||
the `OperationContext` carries one, else the crate constant";
|
||||
implemented as `min(deadline_remaining, timeout_override |
|
||||
ESTABLISHMENT_TIMEOUT)` — the registration override stays
|
||||
meaningful for `Query`/`Mutation`-typed open ops (whose dispatch
|
||||
carries a 30s deadline; `Sub` clears it), and a deadline already in
|
||||
the past yields a zero bound (immediate `timeout` reason).
|
||||
|
||||
Implemented surface: `OpenEstablisher`, `Establishment`,
|
||||
`EstablishmentError` (reasons `dial_failed`/`unknown_resource`/
|
||||
`resource_shortage`/`handler_error` + the wrapper's `timeout`),
|
||||
`ESTABLISHMENT_TIMEOUT` (10s), `CHANNEL_OPEN_FAILED`
|
||||
(`channel:open_failed`), `ChannelCore::register_openable_with_establisher`
|
||||
(`register_openable` delegates with `establisher: None`),
|
||||
`ChannelOpenError` (client-side typed error: `CallFailed { error:
|
||||
CallError }` / `MissingChannelId` / `AdoptFailed`, with
|
||||
`call_error()` + `establishment_reason()` accessors). All four
|
||||
verification gates from the review landed as tests (establisher
|
||||
failure e2e through a real channels connection with ledger
|
||||
un-increment + no-channel assertions, bounded timeout, no-establisher
|
||||
compat, establisher-success pump round-trip).
|
||||
@@ -83,7 +83,7 @@ is the load-bearing piece the broker composes on.
|
||||
| OQ-37 | `from_call` relay wrapper for marked ops | open | medium | ADR-047 §1 names it as a consumer (hub) concern; alkcall's `from_call` reconstructs the marker (Gap F resolved) so the consumer can branch on it |
|
||||
| OQ-38 | ALPN→path-segment mapping | resolved | low | ADR-047 §"Negative" — strip the `alk/` prefix; ALPNs without that prefix use the full ALPN string (rare, two-way-door) |
|
||||
| OQ-39 | `channel/control` control-handle surface | open | medium | The `channel/control` handler currently returns `channel:control_not_implemented`. The control-handle surface (per-channel control callbacks registered by ALPN crates, routing `message` to the handler's control handle for `channel_id`) is real design work — each ALPN crate needs a way to register a control callback, and the channels layer needs a control-handle registry keyed by `channel_id`. Deferred until an ALPN crate (TTY, tunnel) needs out-of-band control. |
|
||||
| OQ-40 | `channel/resources/subscribe` live subscription | open | medium | The `channel/resources/subscribe` handler currently returns `channel:resources_not_implemented`. The live subscription aggregated from ALPN-crate resource enumerators (ADR-047 §6) requires each ALPN crate to provide a resource enumerator, and the channels layer to aggregate them into a live `Stream` that emits on any change. The current stub is a one-shot error; the real implementation is deferred until a consumer (hub, dashboard) needs live resource discovery. |
|
||||
| OQ-40 | `channel/resources/subscribe` live subscription | open | medium | The `channel/resources/subscribe` handler currently returns `channel:resources_not_implemented`. The live subscription aggregated from ALPN-crate resource enumerators (ADR-047 §6) requires each ALPN crate to provide a resource enumerator, and the channels layer to aggregate them into a live `Stream` that emits on any change. The current stub is a one-shot error; the real implementation is deferred until a consumer (hub, dashboard) needs live resource discovery. **Load-bearing for alktunnels' discovery UI** (review 006 E-02): a consumer cannot distinguish produced tunnel resources by op name alone. The static half is covered (0.5.0: `OperationSpec.description` on `services/list` — describes the op, not the live resource set); the dynamic half stays deferred — alktunnels v1 uses config-known op names. |
|
||||
| OQ-41 | QUIC-native multi-stream substrate | open | medium | Only the in-line substrate mode is implemented (single bidi stream, header-demuxed N channels). The QUIC-native multi-stream substrate (accept remaining bidi streams, read headers off each — ADR-034 §substrate modes) is deferred to the downstream alknet crate. The wire format and demux loop are correct for both substrates; only the outer `accept_bi()` loop is missing. The alknet crate owns the QUIC dial/accept loop and is the natural place for the multi-stream accept loop. This crate stays transport-agnostic (no QUIC dependency, WASM-compatible). |
|
||||
|
||||
## Core Types
|
||||
|
||||
@@ -39,6 +39,9 @@ pub struct OperationSpec {
|
||||
pub output_schema: Value, // JSON Schema for output
|
||||
pub error_schemas: Vec<ErrorDefinition>, // Declared domain errors (ADR-016)
|
||||
pub access_control: AccessControl,
|
||||
/// Human-readable op description (review 006 E-02). Disclosed by
|
||||
/// `services/list` when set; `None` when the op declares none.
|
||||
pub description: Option<String>,
|
||||
/// JSON pointer into the input for the resource ID, when
|
||||
/// `access_control.resource_type` is set and the operation targets a
|
||||
/// specific runtime-spawned resource (ADR-011). e.g., `"$.containerId"`
|
||||
@@ -829,6 +832,12 @@ These are read-only — no admin operations are exposed through the call protoco
|
||||
}
|
||||
```
|
||||
|
||||
Each listing entry also carries `description` (review 006 E-02) when
|
||||
the op's spec declares one (`OperationSpec.description`, set via
|
||||
`with_description`) — the field is additive and absent otherwise. It
|
||||
describes the op, not the produced resource set: live resource
|
||||
discovery stays with `channel/resources/subscribe` (ADR-047 §6, OQ-40).
|
||||
|
||||
`services/schema` accepts `{ "name": "fs/readFile" }` (no leading slash —
|
||||
registry form, same as `OperationSpec.name`) and returns the full
|
||||
`OperationSpec` including input/output JSON Schemas and declared
|
||||
|
||||
@@ -0,0 +1,561 @@
|
||||
# Review 006 — Channel-Open Establishment Gap (from alktunnels Phase 0)
|
||||
|
||||
## Status
|
||||
|
||||
**Resolved-by-ADR (E-01, N-1) / Implemented (E-02, E-03, E-04, N-2).**
|
||||
Findings filed from the alktunnels Phase 0
|
||||
research pass (2026-09-06). This is a design review, not a code-defect
|
||||
review: the establishment gap (E-01) is real, POC-observable, and
|
||||
load-bearing for the next downstream crate; the remaining findings are
|
||||
smaller mechanism/coverage gaps noticed in the same sweep.
|
||||
|
||||
**2026-09-06 verification + remediation pass.** All four findings were
|
||||
independently re-verified against source at tree `88e3f5e` (0.4.1 +
|
||||
this review's own commit) — verdicts CONFIRMED for E-01..E-04, with
|
||||
one correction to E-02's cost estimate (see the verification appendix
|
||||
at the bottom of this file). Three additional findings filed from the
|
||||
same sweep: **N-1** (client error-type gap — blocking for E-01's
|
||||
consumer visibility), **N-2** (dead counter), **N-3** (panic posture,
|
||||
pinned). **E-01 + N-1 are resolved by ADR-049**
|
||||
(`docs/architecture/decisions/049-channel-open-establishment-phase.md`
|
||||
— split-hook `OpenEstablisher`, bounded await, `channel:open_failed`
|
||||
typed error). **Unit 1 (E-01 + N-1) is implemented** in alkcall 0.5.0
|
||||
(`OpenEstablisher` + `register_openable_with_establisher`, wrapper
|
||||
establishment phase with bounded await, `channel:open_failed` typed
|
||||
error, `ChannelOpenError` typed client error — all four verification
|
||||
gates landed as tests; two implementation-shape notes recorded in
|
||||
ADR-049's amendment: the establisher does not receive the channel
|
||||
`Connection` (yield-once BiStream belongs to the pump handler), and
|
||||
the bound is the earlier of dispatch deadline and per-registration
|
||||
timeout). **Unit 2 (E-02) is implemented** in alkcall 0.5.0
|
||||
(`OperationSpec.description: Option<String>` +
|
||||
`with_description`, `spec_to_json_pub` emits it when set,
|
||||
`rebuild_spec_for` parses it (shared by `from_call` and `op/register`),
|
||||
`services/list` and `services/list-peers` local listings emit it when
|
||||
set; the ADR-047 §6 amendment below records the discovery decision).
|
||||
**Unit 3 (E-03, E-04, N-2) is implemented** in alkcall 0.5.0 (E-03: the
|
||||
wrapper's handler-exit teardown discards `UnknownChannel` through a
|
||||
debug log + a benign-race pinning comment; E-04: the 64-parked-chunks
|
||||
observable bound is documented consumer-side in `channel-client.md` §
|
||||
"The early-arrival park bound (push-first producers)" and on the
|
||||
`EARLY_ARRIVAL_CAP` const; N-2: `early_arrival_count` is exposed via
|
||||
`ChannelManager::early_arrival_count()` — the observability choice,
|
||||
paired with `dropped_unknown_chunks` — with a monotonicity test).
|
||||
The original remediation sketch below is superseded by the
|
||||
"Remediation plan (post-verification)" section; the original sketch
|
||||
is retained for the record.
|
||||
|
||||
Findings continue the review numbering with prefix `E` (001–005 used
|
||||
P/C/R/A/B/C/D/F/G — each review numbers independently).
|
||||
|
||||
## Scope
|
||||
|
||||
The channels open-op path (`run_open_wrapper` and the `OpenHandler`
|
||||
contract) reviewed from the perspective of the next consumer crate
|
||||
(alktunnels — arbitrary TCP/UDP tunnels over channels), cross-checked
|
||||
against the two existing consumers (alktty, alkhttp) and the SSH
|
||||
forwarding prior art recorded in
|
||||
`/workspace/@alkdev/alktunnels/docs/research/ssh-socks5-survey.md`.
|
||||
|
||||
Everything below was verified directly in source at tree `a22b2b8`
|
||||
(0.4.1 + the early-arrival park fix). No code changes were made in
|
||||
this repo by this review. The 2026-09-06 re-verification pass
|
||||
(verification appendix) re-traced all findings at tree `88e3f5e` —
|
||||
no source changes touched the reviewed paths between the two trees
|
||||
(the only delta is this review's own commit).
|
||||
|
||||
```
|
||||
Verified against: alkcall a22b2b8 (0.4.1); re-verified at 88e3f5e
|
||||
Reading list: src/channels/operations.rs (run_open_wrapper, OpenHandler,
|
||||
make_open_handler_once/stream/sink), src/channels/manager.rs
|
||||
(open_channel, teardown_channel, route_payload, early-arrival park),
|
||||
src/channels/reassembly.rs (MpscSendStream::poll_shutdown, Drop,
|
||||
MpscRecvStream::poll_read EOF arms), src/channels/mux.rs (implicit-EOF
|
||||
pump path), src/protocol/wire.rs (CallError), src/registry/
|
||||
registration.rs (invoke/invoke_streaming gates), src/registry/
|
||||
discovery.rs (services/list, spec_to_json_pub), src/client/from_call.rs
|
||||
(rebuild_spec_for), docs/architecture/decisions/047,
|
||||
docs/architecture/channel-operations.md
|
||||
```
|
||||
|
||||
## Severity legend
|
||||
|
||||
Same scale as reviews 004/005:
|
||||
|
||||
- **[critical]** — a decided spec invariant is violated in a way that
|
||||
makes a promised capability unreachable end-to-end; or corrupts data.
|
||||
- **[major]** — a core protocol path cannot serve a decided behavior;
|
||||
works only via shapes the spec does not describe.
|
||||
- **[minor]** — drift, doc/spec inconsistency, or a missing
|
||||
convenience with no correctness impact.
|
||||
|
||||
---
|
||||
|
||||
# Part A — The open-op establishment gap
|
||||
|
||||
## E-01 [major] — The open op cannot fail after allocation: the consumer receives a live channel for a tunnel whose establishment failed, with no error channel
|
||||
|
||||
**Verified:** YES, by code trace and mechanism analysis.
|
||||
|
||||
### The mechanics
|
||||
|
||||
`OpenHandler` (ADR-047 §3) is the ALPN crate's hook for "validate
|
||||
params, consult ownership, prepare the backend" — i.e., the
|
||||
establishment phase of a channel. But the wrapper does not await it:
|
||||
|
||||
1. `run_open_wrapper` (`src/channels/operations.rs:483-528`):
|
||||
`policy.check_open` → `manager.open_channel(alpn, opener_id, None)`
|
||||
→ build the channel `Connection` → `open_handler(input, channel_conn,
|
||||
auth)` **spawns** the handler and collects its `JoinHandle` →
|
||||
`ResponseEnvelope::ok(request_id, json!({ "channel_id": channel_id }))`
|
||||
(`operations.rs:503-527, 537`). The reply is written to the wire the
|
||||
moment the handler task is *spawned*, not when the handler has done
|
||||
its establishment work.
|
||||
2. The `OpenHandler` type (`operations.rs:333-334`) returns
|
||||
`tokio::task::JoinHandle<()>` — there is no result, no error variant,
|
||||
no establishment phase the wrapper can consult. Any failure inside
|
||||
the handler (params valid at the schema level but semantically
|
||||
rejected, backend lookup failure, target dial failure for a
|
||||
`direct-tcpip`-shaped tunnel, resource no longer available) is
|
||||
invisible to the open op's reply.
|
||||
3. What the consumer observes on handler-side failure: the call op
|
||||
**succeeds** with `{channel_id}`, the channel is adopted
|
||||
(`ChannelClient::open_channel`, `src/channels/client.rs:227-250`),
|
||||
and then the channel EOFs — the handler wrote nothing before
|
||||
exiting, the mux pump writes the implicit-EOF chunk on receiver end
|
||||
(`src/channels/mux.rs:78-81`, REQ-CH-01 implicit-EOF path), and
|
||||
`MpscRecvStream::poll_read` returns clean EOF for both the sentinel
|
||||
and sender-drop arms (`src/channels/reassembly.rs:111-135`). A
|
||||
dial-failure EOF is byte-for-byte indistinguishable from a target
|
||||
that closed immediately after connecting — the two most different
|
||||
failure/success stories map to the same consumer-visible event.
|
||||
|
||||
### Why this is the wrong shape (SSH's semantics, prior-art checked)
|
||||
|
||||
Every established tunnel/forwarding protocol puts establishment
|
||||
failure in the open reply, not in the data stream:
|
||||
|
||||
- **SSH** (RFC 4254 §5.1): `SSH_MSG_CHANNEL_OPEN_FAILURE` is a
|
||||
first-class reply to the open, carrying a reason code
|
||||
(`ADMINISTRATIVELY_PROHIBITED` / `CONNECT_FAILED` /
|
||||
`UNKNOWN_CHANNEL_TYPE` / `RESOURCE_SHORTAGE`) plus a description
|
||||
string; the channel never exists on the opener's side afterward
|
||||
(verified end-to-end in russh:
|
||||
`src/server/encrypted.rs:1244-1278`, `src/client/encrypted.rs:414-434`
|
||||
— see `/workspace/@alkdev/alktunnels/docs/research/ssh-socks5-survey.md`
|
||||
§"Open-failure path").
|
||||
- **SOCKS5** (RFC 1928 §6): the reply carries REP codes 0x01–0x08 and
|
||||
the connection closes within 10s on failure — error-then-close, never
|
||||
"success then in-stream error."
|
||||
- **udpgw** (`tun2proxy/src/udpgw.rs:21-26`) is the counterexample: an
|
||||
opaque ERR bit with zero reason information — the survey flags it as
|
||||
the vocabulary to avoid.
|
||||
- **alktty was forced to reinvent the missing mechanism in-band.** The
|
||||
direct-ALPN path answers the negotiation with a length-prefixed JSON
|
||||
error frame *on the stream* (`send_negotiation_error`,
|
||||
`alktty/src/adapter.rs:177-187`; the `0x00`-prefix first-byte
|
||||
disambiguation trick, `alktty/docs/architecture/tty-adapter.md:186-197`),
|
||||
and the channels path retains the same error-frame shape even though
|
||||
the open op's `input` is the negotiation (alktty ADR-009) — because
|
||||
post-allocation failures have nowhere better to go. That is
|
||||
per-crate reinvention of a protocol-level capability every ALPN
|
||||
crate will need: a structured, typed, **establishment-failure reply
|
||||
to the open op**.
|
||||
|
||||
### The cost today, concretely
|
||||
|
||||
- A consumer cannot distinguish "ACL denied" (call error,
|
||||
`channel:forbidden` — never allocated) from "dial refused" (open
|
||||
succeeded, instant EOF) from "target accepted then instantly closed"
|
||||
(open succeeded, instant EOF). Retry policy, client UX, and error
|
||||
reporting are impossible on the second and third.
|
||||
- ADR-016's typed error details (`CallError.details`) — which the
|
||||
wrapper already uses for `channel:too_many_channels` with
|
||||
`{count, max}` details (`operations.rs:530-545`) — cannot carry
|
||||
dial-failure reasons.
|
||||
- The alktty in-band error-frame path exists *only* because the
|
||||
wrapper can't fail the open post-allocation; every future ALPN crate
|
||||
faces the same fork: reinvent an in-band error vocabulary or
|
||||
silently-EOF.
|
||||
|
||||
### The ask (proposed shape, for the ADR — not a prescriptive API)
|
||||
|
||||
Give the open-op wrapper an **establishment phase it awaits before
|
||||
replying**. Minimal, backward-compatible shape:
|
||||
|
||||
1. `OpenHandler` gains an establishment result. Two candidate shapes:
|
||||
- **Split the hook**: `OpenEstablisher` (async, awaited by the
|
||||
wrapper — validates params semantically, prepares/dials the
|
||||
backend, returns `Result<Establishment, EstablishmentError>`)
|
||||
followed by `OpenHandler` (spawned on the established channel, as
|
||||
today). The dial is the natural establisher step for tunnels; the
|
||||
pumps remain the spawned handler.
|
||||
- **Await-and-inspect**: keep the single `OpenHandler` returning
|
||||
`JoinHandle<OpenResult>`; the wrapper awaits a bounded
|
||||
establishment phase (a `JoinHandle::timeout` equivalent — select
|
||||
on the handle vs an establishment deadline) before replying. The
|
||||
handler signals "established, continue" via an agreed value
|
||||
(e.g. the handler resolves a first `Result<(), HandlerError>`
|
||||
promptly, or the wrapper watches a oneshot the handler signals).
|
||||
2. On establishment failure: the wrapper tears down the just-allocated
|
||||
channel (`teardown_channel` — the ledger/policy paths already
|
||||
handle this atomically) and replies with a **new typed CallError**,
|
||||
e.g. `channel:open_failed`, with `details` carrying a reason code +
|
||||
message. Reason-code vocabulary per the SSH four (the survey's
|
||||
finding): policy-denied (already distinct — `channel:forbidden`),
|
||||
dial-failed, unknown-resource-or-substrate, resource-shortage —
|
||||
mapping 1:1 onto what an open handler can actually produce. Wire
|
||||
addition is additive (new error code string + optional details
|
||||
shape), no existing consumer breaks.
|
||||
3. Backward compatibility: the existing `OpenHandler` signature is
|
||||
preserved if the split shape is chosen (old handlers still compile —
|
||||
the establisher is a new, separately-registered hook, defaulting to
|
||||
an always-OK establisher for the no-establishment-work case).
|
||||
|
||||
### Alternative considered and rejected
|
||||
|
||||
An in-band establishment/error frame (alktty-style, on the channel
|
||||
stream, alktunnels' original OQ-TN-09 direction) works without an
|
||||
upstream change, but: (a) it forces every ALPN crate to define a frame
|
||||
vocabulary and a disambiguation scheme (alktty's `0x00` peek is
|
||||
exactly this cost, paid once per crate); (b) it cannot carry typed
|
||||
`CallError.details` or participate in ADR-016 error schemas; (c) it
|
||||
leaves the phantom-opened channel in the manager (allocation/ledger/
|
||||
policy all fire for a channel that never carried data); (d) SSH's
|
||||
semantics — the failure is the open's reply, the channel never exists
|
||||
opener-side — are the cleaner contract, and channels is at exactly the
|
||||
maturity point (three downstream dependents, alkcall 0.4.x) to fix it
|
||||
upstream cheaply.
|
||||
|
||||
### Severity
|
||||
|
||||
[major] — not [critical]: no data corruption, and the capability is
|
||||
reachable via per-crate workarounds (alktty proves it). But it is the
|
||||
shape every future protocol crate will fight, the fix is cheaper now
|
||||
than after the next consumer, and the workaround path (in-band frames)
|
||||
bakes in a wire format that would then need to stay stable per crate.
|
||||
|
||||
---
|
||||
|
||||
# Part B — Smaller findings from the same sweep
|
||||
|
||||
## E-02 [minor] — `services/list` discloses no per-op metadata; openable resources are indistinguishable by name alone
|
||||
|
||||
`services_list_handler` (`src/registry/discovery.rs:247-269`) maps
|
||||
`list_operations()` to `{name, namespace, op_type}` only. For a tunnel
|
||||
registry, the consumer's discovery question is "which tunnel resources
|
||||
may I open" — and the answer arrives as a bare list of op names
|
||||
(`channels/tunnel/sub`...). Per-resource identity (which target, which
|
||||
substrate, human description) has no field to live in:
|
||||
|
||||
- `services/schema` (`spec_to_json_pub`, `discovery.rs:211-236`) can
|
||||
disclose it *per op* (input schema, access_control), so the data
|
||||
path exists — but it is N+1 round-trips, and the input schema
|
||||
describes the *open params contract*, not the *set of produced
|
||||
resources* (a tunnel producer registers one op and N resources).
|
||||
- ADR-047 §6 already anticipated the dynamic half:
|
||||
`channel/resources/subscribe` aggregates per-ALPN resource
|
||||
enumerators — but the handler is a not-implemented stub
|
||||
(`channel:resources_not_implemented`, OQ-40,
|
||||
`operations.rs:274-299`).
|
||||
|
||||
This is not an alkcall defect — OQ-40 is decided-deferred and the stub
|
||||
fails loudly by design. Filing it because alktunnels' discovery need
|
||||
(OQ-TN-08 resolution: "the ops listing IS tunnel-resource discovery")
|
||||
makes OQ-40 load-bearing for the first time: a consumer UI cannot
|
||||
distinguish produced tunnel resources without either the enumerator
|
||||
aggregation or a listing enrichment. Recommend deciding (small ADR or
|
||||
OQ-40 update) whether:
|
||||
|
||||
1. `channel/resources/subscribe` lands (per-ALPN enumerators — the
|
||||
decided shape), or
|
||||
2. `services/list` gains an additive per-op `description`/`metadata`
|
||||
field (static, registry-side — cheap, but describes the op, not the
|
||||
resource set), or
|
||||
3. Both: subscribe for live resource sets, list enriched for op
|
||||
descriptions.
|
||||
|
||||
alktunnels Phase 1 can proceed with (3)'s shape assumed; the minimal
|
||||
v1 consumer can also just know the op name out-of-band (config), which
|
||||
is why this is [minor] today.
|
||||
|
||||
## E-03 [minor] — `OpenHandler`-exit teardown races `channel/close`'s awaited teardown, but the ledger `take` makes it benign (verify + document)
|
||||
|
||||
`run_open_wrapper`'s spawned teardown task (`operations.rs:505-518`)
|
||||
and the `channel/close` handler's teardown path
|
||||
(`make_close_handler`, `operations.rs:170-238` — 5s await on the
|
||||
handler task, then ledger `take` + `on_close`) both walk
|
||||
`opener_ledger().take(id)` → `policy.on_close(opener)`. The `take` is
|
||||
atomic (first caller removes the entry), so no double-decrement — the
|
||||
cap cannot drift upward. But the loser of the race:
|
||||
- The wrapper task's `teardown_channel(id)` returns
|
||||
`Err(UnknownChannel)` (logged? no — `let _ =` discard,
|
||||
`operations.rs:509`) and skips the ledger/policy half silently.
|
||||
- The close handler's `Ok(task)` arm then awaits a task that is
|
||||
already exiting.
|
||||
|
||||
Verified benign for the cap invariant (`take` is the gate —
|
||||
`operations.rs:510`, `operations.rs:210`), but the `let _ =` discard
|
||||
at `operations.rs:509` means a real `UnknownChannel` after a
|
||||
`set_handler_task` failure (`operations.rs:520-527` — the "channel
|
||||
vanished between open and handler-task install" warn) is
|
||||
indistinguishable from a benign race. Recommend either a debug log on
|
||||
the discard or a comment pinning the race as designed. No correctness
|
||||
impact; filing for the record since it was checked during the E-01
|
||||
trace.
|
||||
|
||||
## E-04 [minor] — Early-arrival park cap (64) interacts with push-first producers under slow adopters — cap-drop is silent at the consumer
|
||||
|
||||
The early-arrival buffer (`EARLY_ARRIVAL_CAP = 64`,
|
||||
`manager.rs:111`, `park_early_arrival` `manager.rs:440-462`) is
|
||||
per-channel and drops past the cap with a debug log + counter —
|
||||
correct per the adopt-race design. But for a *tunnel* producer whose
|
||||
handler dials then immediately pumps (the common case), 64 chunks can
|
||||
arrive before the consumer's `adopt_channel` runs if the open-op
|
||||
response is delayed (relay hops, scheduling). The drops are counted
|
||||
(`dropped_unknown_chunks`) but not visible to the channel's
|
||||
consumer — data loss presents as a truncated stream with clean framing
|
||||
elsewhere. Not a correctness bug (the cap is the documented behavior;
|
||||
the fix shipped in `a22b2b8` is the right shape), but worth a note in
|
||||
the tunnel-crate's POC checklist: the UDP POC's MTU-vs-buffer sizing
|
||||
should treat 64 parked chunks as the observable bound. No alkcall
|
||||
change requested; filed so the constraint is visible to the next
|
||||
consumer.
|
||||
|
||||
## Non-findings (verified correct, recorded to bound the re-review)
|
||||
|
||||
- **The open-op ACL path is complete**: `invoke_streaming` runs the
|
||||
same visibility + `AccessControl::check` + `input_schema` gates as
|
||||
`invoke` (`registration.rs:380-404` vs `:347+`); the Sub-typed open
|
||||
op's ACL failure is a `call.error` on the open, no channel
|
||||
allocated. Verified against the `invoke_streaming_acl_denied_yields_
|
||||
forbidden` test (`registration.rs:1628`).
|
||||
- **`input_schema` enforcement covers open params** (0.4.0, the
|
||||
alktty-review-L1 fix): `check_input_schema` runs before the handler
|
||||
in all three dispatch entry points; schema-invalid open params are
|
||||
`INVALID_INPUT` call errors — never a phantom channel. The
|
||||
semantic-beyond-schema gap is E-01's subject and is properly
|
||||
post-schema.
|
||||
- **EOF arms are unified and clean**: both the sentinel (`Bytes::new`)
|
||||
and sender-drop arms of `MpscRecvStream::poll_read` yield clean EOF
|
||||
(`reassembly.rs:111-135`); the mux pump writes the implicit-EOF
|
||||
chunk for handler-drop-without-shutdown (`mux.rs:78-81` +
|
||||
`mux_pump_writes_eof_on_implicit_close` test). The two-pump
|
||||
shutdown contract (alknet ADR-078) has its upstream half in place.
|
||||
- **The `channel_open` marker survives the wire round-trip**
|
||||
(`spec_to_json_pub` emits the boolean; `rebuild_spec_for` re-derives
|
||||
the ALPN from the op name — `from_call.rs:265-280`, `:298-314`),
|
||||
including the `custom/proto` multi-segment case. The G-04
|
||||
`resource_id_path` fix is present (`from_call.rs:260`,
|
||||
round-trip test `from_call.rs:610-617`).
|
||||
- **The opener-ledger cap decrement is atomic with removal**
|
||||
(`opener_ledger().take` gates every `on_close` path); the
|
||||
double-decrement concern from ADR-047 §7 is closed at every call
|
||||
site traced.
|
||||
|
||||
---
|
||||
|
||||
# Remediation sketch (original, superseded — see "Remediation plan (post-verification)")
|
||||
|
||||
**Unit 1 — E-01 (the establishment phase).** Decided-shape ADR first
|
||||
(this is ADR-047 §3 contract territory — the `OpenHandler` type shape
|
||||
is a cross-crate API surface, and both existing consumers' handlers
|
||||
must be considered; alktty's `TtyOpenHandler` is the migration
|
||||
prototype). Implementation sketch: `ChannelCore::register_openable`
|
||||
gains an optional establisher hook; `run_open_wrapper` awaits it
|
||||
bounded (establishment deadline — a constant, e.g. 10s, or per-spec),
|
||||
tears down on failure, replies `channel:open_failed` with
|
||||
`{reason, message}` details on failure and `{channel_id}` on success.
|
||||
alktty migrates its channels-path error-frame to the call-error path
|
||||
where applicable (its direct-ALPN path keeps the in-band frame — two
|
||||
transports, two contracts). alkhttp unaffected (no openable ops).
|
||||
|
||||
**Unit 2 — E-02 (discovery enrichment).** Decide via OQ-40 update +
|
||||
small ADR amendment: recommend (3) — subscribe for live resource sets
|
||||
(the decided shape, now load-bearing), plus an additive
|
||||
`description` field on the listing (cheap, immediately useful).
|
||||
**Implemented 2026-09-06 (alkcall 0.5.0):** the listing half landed
|
||||
(`OperationSpec.description`, four touchpoints as the appendix
|
||||
corrected); the subscribe half stays deferred (OQ-40, now with the
|
||||
load-bearing note).
|
||||
|
||||
**Unit 3 — E-03/E-04.** E-03: log-or-comment; E-04: doc note. Trivial.
|
||||
|
||||
# Remediation plan (post-verification)
|
||||
|
||||
Firm ordering, decided 2026-09-06. E-01's ADR is **ADR-049**
|
||||
(`docs/architecture/decisions/049-channel-open-establishment-phase.md`)
|
||||
— decided shape: the **split hook** (`OpenEstablisher` awaited bounded
|
||||
by the wrapper, `OpenHandler` spawned unchanged after success), chosen
|
||||
over await-and-inspect because it preserves the `OpenHandler` type
|
||||
(backward compatible by construction), keeps establishment and pump
|
||||
lifecycles separate (SSH semantics: the dial is synchronous with the
|
||||
open reply), and restores ADR-047 §3's original "channel plan" wrapper
|
||||
shape. N-1 rides the same unit — without the client error-type fix,
|
||||
`channel:open_failed`'s typed reason is unreachable through the
|
||||
primary client path.
|
||||
|
||||
**Unit 1 — E-01 + N-1 (alkcall 0.5.0).** `OpenEstablisher` +
|
||||
`register_openable_with_establisher` (no-establisher = always-OK,
|
||||
existing registrations compile unchanged); wrapper flow: `check_open`
|
||||
→ `open_channel` → await establisher bounded (dispatch deadline when
|
||||
`Some`, else `ESTABLISHMENT_TIMEOUT` = 10s) → on failure:
|
||||
`teardown_channel` + ledger `take` + `policy.on_close` + reply
|
||||
`channel:open_failed` with `details: {reason, message}` (reason ∈
|
||||
`dial_failed` / `unknown_resource` / `resource_shortage` /
|
||||
`handler_error` / `timeout`); on success: spawn pumps,
|
||||
`set_handler_task`, reply `{channel_id}`. `ChannelClient::open_channel`
|
||||
returns a typed error carrying the `CallError`. Concurrency is safe:
|
||||
the dispatcher spawns Once invocations as independent tasks
|
||||
(`dispatch.rs` `spawn_once_dispatch`), so an awaited establisher does
|
||||
not head-of-line-block channel 0. Version 0.5.0 (the `open_channel`
|
||||
error-type change is semver-relevant). **Status: IMPLEMENTED
|
||||
(2026-09-06)** — with two implementation-shape notes recorded in
|
||||
ADR-049's amendment: the establisher takes `(input, auth)` only (the
|
||||
yield-once channel `BiStream` belongs exclusively to the pump
|
||||
handler), and the effective bound is `min(dispatch deadline,
|
||||
registration timeout | ESTABLISHMENT_TIMEOUT)`. All four verification
|
||||
gates below are tests in the crate.
|
||||
|
||||
**Unit 2 — E-02 (ride the 0.5.0 release).** Additive
|
||||
`description: Option<String>` on `OperationSpec` — note this is four
|
||||
touchpoints, not one (see the verification appendix's E-02
|
||||
correction): struct field + `spec_to_json_pub` emit +
|
||||
`rebuild_spec_for` parse + the `services/list` output-schema doc.
|
||||
`services/list` emits it when set. OQ-40 stays deferred but gains a
|
||||
"load-bearing for alktunnels discovery UI" note; the
|
||||
`channel/resources/subscribe` half stays deferred (alktunnels v1 uses
|
||||
config-known op names). **Status: IMPLEMENTED (2026-09-06)** — the
|
||||
discovery decision (3: enriched listing now, subscribe deferred) is
|
||||
recorded as an amendment to ADR-047 §6.
|
||||
|
||||
**Unit 3 — E-03/E-04/N-2 (trivial batch, same PR series as Unit 1).**
|
||||
E-03: debug log (or pinning comment) on the discarded `UnknownChannel`
|
||||
at the wrapper's teardown discard. E-04: doc note on `EARLY_ARRIVAL_CAP`
|
||||
(the 64-parked-chunks observable bound) for tunnel-crate-facing
|
||||
consumers. N-2: expose an accessor for `early_arrival_count` or remove
|
||||
the write-only counter. **Status: IMPLEMENTED (2026-09-06)** — E-03:
|
||||
debug log + benign-race pinning comment on the wrapper's handler-exit
|
||||
teardown (mirroring the sibling log the Unit 1 establisher-failure
|
||||
teardown already carries); E-04: consumer-side doc note in
|
||||
`channel-client.md` plus the const doc; N-2: accessor
|
||||
(`ChannelManager::early_arrival_count()`, documented as monotonic —
|
||||
adopt-drain does not decrement) + test.
|
||||
|
||||
**Sequencing:** ADR-049 (landed) → alkcall 0.5.0 (Units 1+2+3) →
|
||||
alktty migration (channels-path semantic failures move into an
|
||||
establisher; the direct-ALPN in-band error frame is retained — two
|
||||
transports, two contracts) → alkhttp mechanical pass
|
||||
(`OpenableAlpn.establisher: Option<...>`, default `None`). alktunnels
|
||||
Phase 1 blocks only on the ADR decision (landed), not the
|
||||
implementation: its OQ-TN-09 direction ("self-contained control frame")
|
||||
shrinks to *post-establishment* control only once `channel:open_failed`
|
||||
exists.
|
||||
|
||||
## Verification gates for the E-01 remediation
|
||||
|
||||
- A channels end-to-end test: producer handler whose establisher fails
|
||||
after channel allocation → consumer's `open_channel` resolves
|
||||
`Err` with `channel:open_failed` + reason details; no channel in the
|
||||
manager's `channel_ids()` afterward (the SSH "channel never exists
|
||||
opener-side" property, consumer-visible as "no `channel_id` was
|
||||
ever returned").
|
||||
- A bounded-establishment test: establisher that never completes →
|
||||
open op fails with a timeout-flavored error within the deadline,
|
||||
channel torn down, ledger decremented.
|
||||
- An alktty-migration test: the existing
|
||||
`open_via_channels_surfaces_negotiation_rejected` scenario resolves
|
||||
via call error (or retained error-frame — per the ADR's chosen
|
||||
compat shape) unchanged in behavior.
|
||||
- (Added post-verification) A compat gate: a no-establisher
|
||||
registration behaves exactly as today — `{channel_id}` reply,
|
||||
handler spawned, teardown on handler exit.
|
||||
|
||||
## Verification appendix (2026-09-06 re-verification pass)
|
||||
|
||||
All findings re-verified by independent source trace at tree `88e3f5e`.
|
||||
Verdicts, corrections, and the additional findings:
|
||||
|
||||
- **E-01 CONFIRMED — and strengthened.** The wrapper's
|
||||
spawn-then-reply flow is as described
|
||||
(`src/channels/operations.rs` `run_open_wrapper`: handler spawned,
|
||||
reply written before the handler performs any work). New supporting
|
||||
evidence the original pass missed: **ADR-047 §3's decision text**
|
||||
describes the ALPN open handler as "validate params, consult
|
||||
ownership, prepare the backend, return a 'channel plan'" — an
|
||||
awaited preparation the wrapper consults before replying. The
|
||||
implemented `OpenHandler` (`JoinHandle<()>`) collapsed that phase
|
||||
into a fire-and-forget spawn. E-01 is therefore **drift from
|
||||
ADR-047 §3's own wrapper shape**, not merely a missing convenience —
|
||||
the remediation restores the decided design, which lowers the ADR's
|
||||
contention cost. Also verified: the serving loop spawns Once
|
||||
invocations as independent tasks, so an awaited establishment phase
|
||||
does not head-of-line-block other calls on channel 0 (a fix-feasibility
|
||||
question the original pass did not address).
|
||||
- **E-02 CONFIRMED — one correction.** `services_list_handler` emits
|
||||
`{name, namespace, op_type}` only, as described. But the "cheap,
|
||||
additive `description` field" framing understates the work:
|
||||
`OperationSpec` has **no `description` field at all** (the
|
||||
`description` in `spec.rs` is on `ErrorDefinition`). The change is
|
||||
four touchpoints: struct field, `spec_to_json_pub` emit,
|
||||
`rebuild_spec_for` parse, and the `services/list` output-schema doc.
|
||||
Still small; still [minor].
|
||||
- **E-03 CONFIRMED.** Both teardown paths gate on the atomic
|
||||
`opener_ledger().take(id)`; the race loser gets `UnknownChannel` /
|
||||
`take → None` — no double-decrement. The `let _ =` discard on the
|
||||
wrapper's `teardown_channel` result is as described. Also traced:
|
||||
the `set_handler_task` failure path still runs the wrapper task to
|
||||
completion (self-teardown), so no leak in that arm either. The
|
||||
log-or-comment remedy stands.
|
||||
- **E-04 CONFIRMED.** `EARLY_ARRIVAL_CAP = 64`, per-channel, FIFO,
|
||||
drop-past-cap with debug log + `dropped_unknown_chunks` counter.
|
||||
Doc-note-only remedy stands.
|
||||
- **N-1 [minor today, blocking for E-01's consumer visibility] —
|
||||
`ChannelClient::open_channel` erases the typed error.**
|
||||
(`src/channels/client.rs`) It flattens the `CallError` into a
|
||||
`String` via `format!("open op failed: {e:?}")`. Even once the
|
||||
wrapper replies `channel:open_failed` with reason details, a
|
||||
consumer cannot branch on the reason — the error type destroys it.
|
||||
E-01's verification gate ("consumer's `open_channel` resolves `Err`
|
||||
with `channel:open_failed` + reason details") is unreachable without
|
||||
changing this error type, so the fix is in scope for the E-01 unit
|
||||
(ADR-049 §4), not deferred.
|
||||
- **N-2 [trivial] — `early_arrival_count` is write-only.** Incremented
|
||||
in `park_early_arrival`, never read anywhere (no accessor; only
|
||||
`dropped_unknown_chunks` is exposed). Either expose an accessor
|
||||
(observability for the E-04 bound) or remove the counter.
|
||||
- **N-3 [observation, pinned by ADR-049 §6] — panicked pump handlers
|
||||
are EOF-shaped by design.** The wrapper's teardown task swallows the
|
||||
pump handler's `JoinError` (`let _ = raw_task.await`). A panicked
|
||||
handler = phantom channel + instant EOF, indistinguishable from a
|
||||
clean short-lived channel. With the establisher split, the pump-phase
|
||||
panic stays in this category (correct — there is no mid-stream error
|
||||
channel by design; establishment errors are the only kind that
|
||||
belong in the open reply). ADR-049 §6 pins this posture; no change.
|
||||
|
||||
## References
|
||||
|
||||
- alknet ADR-078 / channels two-pump contract — the teardown half this
|
||||
review's EOF-arms non-finding confirms upstream.
|
||||
- alkcall ADR-047 (openable ALPNs are operations), ADR-016 (typed
|
||||
error schemas — the vehicle for `channel:open_failed` details),
|
||||
ADR-040/041 (backpressure/caps — untouched by E-01).
|
||||
- alktty ADR-009 (the open op's input is the negotiation) + review
|
||||
#001 L1/L3 resolution — the per-crate workaround E-01 obsoletes;
|
||||
`send_negotiation_error` (`alktty/src/adapter.rs:177-187`) and the
|
||||
`0x00`-peek (`alktty/docs/architecture/tty-adapter.md:186-197`) are
|
||||
the in-band mechanism the upstream establisher replaces for the
|
||||
channels path.
|
||||
- alktunnels Phase 0 (`docs/research/phase-0-findings.md` OQ-TN-09)
|
||||
and the SSH/SOCKS5 survey (`docs/research/ssh-socks5-survey.md`
|
||||
§"Open-failure path", §"Comparison") — the prior art motivating
|
||||
E-01's reason-code vocabulary.
|
||||
- The consumer-findings ledger convention (`docs/reviews/consumer-
|
||||
findings-ledger.md`) — findings here follow the same spirit
|
||||
(downstream-discovered, filed for upstream action); E-01..E-04 are
|
||||
alktunnels-discovered but numbered in alkcall's review series since
|
||||
they are alkcall findings.
|
||||
- **ADR-049** (`docs/architecture/decisions/049-channel-open-
|
||||
establishment-phase.md`) — the E-01/N-1 resolution (split-hook
|
||||
`OpenEstablisher`, bounded establishment, `channel:open_failed`
|
||||
typed error, client error-type fix).
|
||||
+323
-6
@@ -25,13 +25,69 @@ use tokio::sync::Mutex;
|
||||
|
||||
use crate::core::types::{Connection, StreamError};
|
||||
use crate::protocol::connection::CallConnection;
|
||||
use crate::protocol::wire::ResponseEnvelope;
|
||||
use crate::protocol::wire::{CallError, ResponseEnvelope};
|
||||
use crate::registry::registration::OperationRegistry;
|
||||
|
||||
use super::manager::{ChannelManager, ChannelSide};
|
||||
use super::mux::MuxRunner;
|
||||
use super::reassembly::{MpscRecvStream, MpscSendStream};
|
||||
|
||||
/// The typed error from [`ChannelClient::open_channel`] (ADR-049 §4 —
|
||||
/// review 006 N-1). The open op's `CallError` — including the
|
||||
/// `channel:open_failed` code with `details: { reason, message }`
|
||||
/// (ADR-049 §3) — is carried verbatim instead of flattened into a
|
||||
/// string, so the consumer can branch on the establishment-failure
|
||||
/// reason (`dial_failed`, `unknown_resource`, `resource_shortage`,
|
||||
/// `handler_error`, `timeout`).
|
||||
#[derive(Debug, Clone, PartialEq, thiserror::Error)]
|
||||
pub enum ChannelOpenError {
|
||||
/// The open op failed — the `CallError` carries the wire code and
|
||||
/// typed `details` (`channel:open_failed` + `details.reason` per
|
||||
/// ADR-049 §3, `channel:too_many_channels`, `FORBIDDEN`, …).
|
||||
#[error("open op failed: {error:?}")]
|
||||
CallFailed { error: CallError },
|
||||
/// The open op succeeded but the response carried no `channel_id`
|
||||
/// (malformed responder).
|
||||
#[error("open op response missing channel_id")]
|
||||
MissingChannelId,
|
||||
/// The open op succeeded but local channel adoption failed (the
|
||||
/// per-connection cap, or an ID collision).
|
||||
#[error("adopt_channel failed: {0}")]
|
||||
AdoptFailed(#[from] super::manager::ManagerError),
|
||||
}
|
||||
|
||||
impl ChannelOpenError {
|
||||
/// The wire `CallError` behind this failure, when the open op
|
||||
/// itself failed. `None` for the local-only variants
|
||||
/// (`MissingChannelId` comes from a success envelope;
|
||||
/// `AdoptFailed` is post-reply).
|
||||
pub fn call_error(&self) -> Option<&CallError> {
|
||||
match self {
|
||||
ChannelOpenError::CallFailed { error } => Some(error),
|
||||
ChannelOpenError::MissingChannelId | ChannelOpenError::AdoptFailed(_) => None,
|
||||
}
|
||||
}
|
||||
|
||||
/// The establishment-failure reason code — `details.reason` of a
|
||||
/// `channel:open_failed` error (`dial_failed`, `unknown_resource`,
|
||||
/// `resource_shortage`, `handler_error`, `timeout`). `None` for
|
||||
/// any other failure shape.
|
||||
pub fn establishment_reason(&self) -> Option<&str> {
|
||||
match self {
|
||||
ChannelOpenError::CallFailed { error }
|
||||
if error.code == super::operations::CHANNEL_OPEN_FAILED =>
|
||||
{
|
||||
error
|
||||
.details
|
||||
.as_ref()
|
||||
.and_then(|d| d.get("reason"))
|
||||
.and_then(|r| r.as_str())
|
||||
}
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// The client-side handle for a channels connection. Constructed via
|
||||
/// [`ChannelClient::from_connection`] from an established
|
||||
/// `Connection`. Holds the `ChannelManager` (for the demux/mux state)
|
||||
@@ -224,26 +280,33 @@ impl ChannelClient {
|
||||
///
|
||||
/// `alpn` is the data-plane ALPN (e.g. `alk/tty`), used for
|
||||
/// observability in the local manager.
|
||||
///
|
||||
/// The error is typed (ADR-049 §4 — review 006 N-1): a failed open
|
||||
/// resolves [`ChannelOpenError::CallFailed`] carrying the
|
||||
/// wire `CallError` verbatim — branch on
|
||||
/// [`ChannelOpenError::establishment_reason`] for
|
||||
/// `channel:open_failed`'s reason code (`dial_failed`,
|
||||
/// `unknown_resource`, `resource_shortage`, `handler_error`,
|
||||
/// `timeout`).
|
||||
pub async fn open_channel(
|
||||
&self,
|
||||
operation_id: &str,
|
||||
input: Value,
|
||||
alpn: &str,
|
||||
) -> Result<(u32, MpscSendStream, MpscRecvStream), String> {
|
||||
) -> Result<(u32, MpscSendStream, MpscRecvStream), ChannelOpenError> {
|
||||
let response = self.call_open_op(operation_id, input).await;
|
||||
let out = response
|
||||
.result
|
||||
.map_err(|e| format!("open op failed: {e:?}"))?;
|
||||
.map_err(|error| ChannelOpenError::CallFailed { error })?;
|
||||
let channel_id = out
|
||||
.get("channel_id")
|
||||
.and_then(|v| v.as_u64())
|
||||
.ok_or_else(|| "open op response missing channel_id".to_string())?
|
||||
as u32;
|
||||
.ok_or(ChannelOpenError::MissingChannelId)? as u32;
|
||||
|
||||
self.manager
|
||||
.adopt_channel(channel_id, alpn, None)
|
||||
.await
|
||||
.map_err(|e| format!("adopt_channel failed: {e}"))
|
||||
.map_err(ChannelOpenError::AdoptFailed)
|
||||
.map(|(send, recv)| (channel_id, send, recv))
|
||||
}
|
||||
|
||||
@@ -1256,6 +1319,260 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
/// ADR-049 Unit 1 acceptance gate (review 006 E-01/N-1) — the
|
||||
/// establisher fails after channel allocation: the consumer's
|
||||
/// `open_channel` resolves `Err` carrying the wire `CallError`
|
||||
/// (`channel:open_failed` + `details.reason`), no `channel_id` was
|
||||
/// ever returned (the SSH "channel never exists opener-side"
|
||||
/// property, consumer-visible), and the accept side's manager has
|
||||
/// no channel afterward (teardown + ledger un-increment).
|
||||
#[tokio::test]
|
||||
async fn establisher_failure_resolves_typed_open_failed_on_the_client() {
|
||||
use crate::channels::operations::{
|
||||
ChannelCore, EstablishmentError, OpenEstablisher, OpenHandler,
|
||||
};
|
||||
use crate::channels::policy::PerIdentityChannelPolicy;
|
||||
use crate::registry::spec::ChannelOpenSpec;
|
||||
|
||||
let policy: Arc<PerIdentityChannelPolicy> = Arc::new(PerIdentityChannelPolicy::new(256));
|
||||
let policy_for_hook: Arc<dyn ChannelLifecyclePolicy> =
|
||||
Arc::clone(&policy) as Arc<dyn ChannelLifecyclePolicy>;
|
||||
let policy_for_assert = Arc::clone(&policy);
|
||||
|
||||
let open_handler: OpenHandler =
|
||||
Arc::new(|_input, _channel_conn, _auth| tokio::spawn(async {}));
|
||||
let establisher: OpenEstablisher = Arc::new(|_input, _auth| {
|
||||
Box::pin(async {
|
||||
Err(EstablishmentError::DialFailed {
|
||||
message: "target refused the connection".to_string(),
|
||||
})
|
||||
})
|
||||
});
|
||||
|
||||
let install_hook: crate::channels::adapter::InstallChannelZero =
|
||||
Arc::new(move |manager, channel0_conn, auth| {
|
||||
let open_handler = Arc::clone(&open_handler);
|
||||
let establisher = Arc::clone(&establisher);
|
||||
let policy = Arc::clone(&policy_for_hook);
|
||||
tokio::spawn(async move {
|
||||
let channel0_bidi = match channel0_conn.accept_bi().await {
|
||||
Ok(s) => s,
|
||||
Err(_) => return,
|
||||
};
|
||||
let (writer, reader) = split_single_stream(channel0_bidi);
|
||||
let core = ChannelCore::new(manager, policy);
|
||||
let registry = crate::registry::registration::OperationRegistry::new();
|
||||
let spec = OperationSpec::new(
|
||||
"channels/tty/sub",
|
||||
OperationType::Sub,
|
||||
Visibility::External,
|
||||
serde_json::json!({}),
|
||||
serde_json::json!({
|
||||
"type": "object",
|
||||
"properties": { "channel_id": { "type": "integer" } }
|
||||
}),
|
||||
vec![],
|
||||
AccessControl::default(),
|
||||
None,
|
||||
)
|
||||
.with_channel_open(ChannelOpenSpec::new("alk/tty"));
|
||||
core.register_openable_with_establisher(
|
||||
spec,
|
||||
Some(establisher),
|
||||
open_handler,
|
||||
®istry,
|
||||
auth.clone(),
|
||||
None,
|
||||
)
|
||||
.expect("register_openable_with_establisher");
|
||||
let registry = Arc::new(registry);
|
||||
let provider: Arc<dyn IdentityProvider> = Arc::new(NoopIdProvider);
|
||||
let call_connection = Arc::new(CallConnection::new_single_stream(
|
||||
channel0_conn,
|
||||
Arc::clone(&writer),
|
||||
));
|
||||
let dp = Dispatcher::new(registry, provider);
|
||||
dp.run_loop_single_stream(call_connection, reader, writer)
|
||||
.await;
|
||||
})
|
||||
});
|
||||
|
||||
let (client_end, server_end) = tokio::io::duplex(64 * 1024);
|
||||
let client_conn =
|
||||
Connection::from_bidi(client_end, b"alk/channels".to_vec(), Some(TEST_ADDR));
|
||||
let server_conn =
|
||||
Connection::from_bidi(server_end, b"alk/channels".to_vec(), Some(TEST_ADDR));
|
||||
|
||||
let adapter = ChannelsAdapter::new(install_hook, Arc::new(NoCap));
|
||||
let auth = AuthContext::anonymous(b"alk/channels");
|
||||
let _server_handle = tokio::spawn(async move {
|
||||
let _ = crate::core::types::ProtocolHandler::handle(&adapter, server_conn, &auth).await;
|
||||
});
|
||||
|
||||
let client = ChannelClient::from_connection(client_conn)
|
||||
.await
|
||||
.expect("channel client init");
|
||||
|
||||
let err = tokio::time::timeout(
|
||||
std::time::Duration::from_secs(10),
|
||||
client.open_channel(
|
||||
"channels/tty/sub",
|
||||
serde_json::json!({ "container": "abc" }),
|
||||
"alk/tty",
|
||||
),
|
||||
)
|
||||
.await
|
||||
.expect("open_channel timed out");
|
||||
let err = match err {
|
||||
Err(e) => e,
|
||||
Ok(_) => panic!("establisher failure must resolve Err through the client"),
|
||||
};
|
||||
|
||||
// N-1 gate: the typed error carries the wire CallError verbatim.
|
||||
let call_error = err.call_error().expect("CallFailed carries the CallError");
|
||||
assert_eq!(call_error.code, "channel:open_failed");
|
||||
assert_eq!(
|
||||
err.establishment_reason(),
|
||||
Some("dial_failed"),
|
||||
"consumer branches on the typed reason code"
|
||||
);
|
||||
|
||||
// The SSH property, consumer-visible: no channel_id was ever
|
||||
// returned, and the accept side holds no channel for the failed
|
||||
// open (channel 0 is the pre-negotiated call channel — always
|
||||
// present, never the failed open's).
|
||||
assert!(
|
||||
client.manager().channel_ids().into_iter().all(|id| id == 0),
|
||||
"no data channel exists for the failed open"
|
||||
);
|
||||
let anonymous = crate::core::auth::Identity {
|
||||
id: "anonymous".to_string(),
|
||||
scopes: vec![],
|
||||
resources: Default::default(),
|
||||
};
|
||||
assert_eq!(
|
||||
policy_for_assert.count_for(&anonymous),
|
||||
0,
|
||||
"ledger un-incremented on the accept side after establishment failure"
|
||||
);
|
||||
}
|
||||
|
||||
/// ADR-049 Unit 1 acceptance gate (companion to the E-01 gate) —
|
||||
/// establisher success through a real channels connection: the
|
||||
/// client's `open_channel` adopts and the pumps flow (the happy
|
||||
/// path is unchanged by the establishment phase).
|
||||
#[tokio::test]
|
||||
async fn establisher_success_opens_and_pumps_data_end_to_end() {
|
||||
use crate::channels::operations::{
|
||||
ChannelCore, Establishment, OpenEstablisher, OpenHandler,
|
||||
};
|
||||
use crate::channels::policy::NoCap;
|
||||
use crate::registry::spec::ChannelOpenSpec;
|
||||
use tokio::io::AsyncReadExt;
|
||||
|
||||
let (data_tx, mut data_rx) = tokio::sync::mpsc::channel::<Vec<u8>>(1);
|
||||
|
||||
let open_handler: OpenHandler = Arc::new(move |_input, channel_conn, _auth| {
|
||||
let data_tx = data_tx.clone();
|
||||
tokio::spawn(async move {
|
||||
let mut bidi = channel_conn.accept_bi().await.expect("accept_bi");
|
||||
let mut buf = [0u8; 4];
|
||||
bidi.read_exact(&mut buf).await.expect("read");
|
||||
data_tx.send(buf.to_vec()).await.expect("send to channel");
|
||||
})
|
||||
});
|
||||
let establisher: OpenEstablisher =
|
||||
Arc::new(|_input, _auth| Box::pin(async { Ok(Establishment::default()) }));
|
||||
|
||||
let establisher_for_hook = Arc::clone(&establisher);
|
||||
let open_handler_for_hook = Arc::clone(&open_handler);
|
||||
let install_hook: crate::channels::adapter::InstallChannelZero =
|
||||
Arc::new(move |manager, channel0_conn, auth| {
|
||||
let establisher = Arc::clone(&establisher_for_hook);
|
||||
let open_handler = Arc::clone(&open_handler_for_hook);
|
||||
tokio::spawn(async move {
|
||||
let channel0_bidi = match channel0_conn.accept_bi().await {
|
||||
Ok(s) => s,
|
||||
Err(_) => return,
|
||||
};
|
||||
let (writer, reader) = split_single_stream(channel0_bidi);
|
||||
let core = ChannelCore::new(manager, Arc::new(NoCap));
|
||||
let registry = crate::registry::registration::OperationRegistry::new();
|
||||
let spec = OperationSpec::new(
|
||||
"channels/tty/sub",
|
||||
OperationType::Sub,
|
||||
Visibility::External,
|
||||
serde_json::json!({
|
||||
"type": "object",
|
||||
"properties": { "container": { "type": "string" } },
|
||||
"required": ["container"]
|
||||
}),
|
||||
serde_json::json!({
|
||||
"type": "object",
|
||||
"properties": { "channel_id": { "type": "integer" } }
|
||||
}),
|
||||
vec![],
|
||||
AccessControl::default(),
|
||||
None,
|
||||
)
|
||||
.with_channel_open(ChannelOpenSpec::new("alk/tty"));
|
||||
core.register_openable_with_establisher(
|
||||
spec,
|
||||
Some(establisher),
|
||||
open_handler,
|
||||
®istry,
|
||||
auth.clone(),
|
||||
None,
|
||||
)
|
||||
.expect("register_openable_with_establisher");
|
||||
let registry = Arc::new(registry);
|
||||
let provider: Arc<dyn IdentityProvider> = Arc::new(NoopIdProvider);
|
||||
let call_connection = Arc::new(CallConnection::new_single_stream(
|
||||
channel0_conn,
|
||||
Arc::clone(&writer),
|
||||
));
|
||||
let dp = Dispatcher::new(registry, provider);
|
||||
dp.run_loop_single_stream(call_connection, reader, writer)
|
||||
.await;
|
||||
})
|
||||
});
|
||||
|
||||
let (client_end, server_end) = tokio::io::duplex(64 * 1024);
|
||||
let client_conn =
|
||||
Connection::from_bidi(client_end, b"alk/channels".to_vec(), Some(TEST_ADDR));
|
||||
let server_conn =
|
||||
Connection::from_bidi(server_end, b"alk/channels".to_vec(), Some(TEST_ADDR));
|
||||
|
||||
let adapter = ChannelsAdapter::new(install_hook, Arc::new(NoCap));
|
||||
let auth = AuthContext::anonymous(b"alk/channels");
|
||||
let _server_handle = tokio::spawn(async move {
|
||||
let _ = crate::core::types::ProtocolHandler::handle(&adapter, server_conn, &auth).await;
|
||||
});
|
||||
|
||||
let client = ChannelClient::from_connection(client_conn)
|
||||
.await
|
||||
.expect("channel client init");
|
||||
|
||||
let (channel_id, mut send, _recv) = client
|
||||
.open_channel(
|
||||
"channels/tty/sub",
|
||||
serde_json::json!({ "container": "abc" }),
|
||||
"alk/tty",
|
||||
)
|
||||
.await
|
||||
.expect("open_channel with establisher success");
|
||||
|
||||
assert!(channel_id > 0, "channel_id should be non-zero");
|
||||
send.write_all(b"ping").await.expect("write ping");
|
||||
drop(send);
|
||||
|
||||
let data = tokio::time::timeout(std::time::Duration::from_secs(5), data_rx.recv())
|
||||
.await
|
||||
.expect("timed out waiting for handler data")
|
||||
.expect("handler should receive data");
|
||||
assert_eq!(&data, b"ping", "handler received ping post-establishment");
|
||||
}
|
||||
|
||||
// --- review 004 Unit 3 acceptance gates (F-04 serving half) -----------
|
||||
|
||||
/// F-04 gate 1: hub→consumer call over an existing `ChannelClient`
|
||||
|
||||
@@ -108,6 +108,16 @@ struct Inner {
|
||||
/// the open-op response/first-data race). A channel whose adopter never
|
||||
/// arrives leaks its parked chunks until `clear_all`; the cap bounds
|
||||
/// that leak per channel.
|
||||
///
|
||||
/// Consumer-facing bound (review 006 E-04): for a push-first producer
|
||||
/// (dial-then-pump is the common tunnel shape), up to `EARLY_ARRIVAL_CAP`
|
||||
/// chunks per channel can be parked before the consumer's
|
||||
/// `adopt_channel` runs — and chunks past the cap drop silently at the
|
||||
/// consumer (data loss presents as a truncated stream with clean
|
||||
/// framing elsewhere; `dropped_unknown_chunks` /
|
||||
/// `early_arrival_count` are the producer-side observability). A
|
||||
/// tunnel crate's sizing (e.g. a UDP POC's MTU-vs-buffer math) should
|
||||
/// treat 64 parked chunks as the observable bound and adopt promptly.
|
||||
const EARLY_ARRIVAL_CAP: usize = 64;
|
||||
|
||||
impl ChannelManager {
|
||||
@@ -192,6 +202,19 @@ impl ChannelManager {
|
||||
self.inner.dropped_unknown_chunks.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// The number of chunks parked in the early-arrival buffer since
|
||||
/// the connection started (monotonic — parked chunks handed to the
|
||||
/// adopter are not subtracted). Together with
|
||||
/// [`Self::dropped_unknown_chunks`] this bounds the two observable
|
||||
/// halves of the open-op-response/first-data race: how many of a
|
||||
/// push-first producer's first chunks were buffered
|
||||
/// (`early_arrival_count`) versus lost to the per-channel cap
|
||||
/// (`dropped_unknown_chunks` — see `EARLY_ARRIVAL_CAP`, review 006
|
||||
/// E-04). Observability only — neither counter drives control flow.
|
||||
pub fn early_arrival_count(&self) -> u64 {
|
||||
self.inner.early_arrival_count.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// The mux handle — for registering new channels' write halves.
|
||||
pub fn mux(&self) -> &MuxHandle {
|
||||
&self.inner.mux
|
||||
@@ -712,6 +735,42 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
/// `early_arrival_count` observes parked chunks (review 006 N-2 —
|
||||
/// the counter was write-only before the accessor). Monotonic by
|
||||
/// design: it counts chunks parked since connection start; draining
|
||||
/// on adopt does not subtract (paired with `dropped_unknown_chunks`
|
||||
/// it bounds the open-response/first-data race's two halves).
|
||||
#[tokio::test]
|
||||
async fn early_arrival_count_tracks_parked_chunks_monotonically() {
|
||||
let manager = make_manager_with_runner().await;
|
||||
assert_eq!(manager.early_arrival_count(), 0);
|
||||
for _ in 0..3 {
|
||||
manager
|
||||
.route_payload(7, Bytes::from_static(b"parked"))
|
||||
.await;
|
||||
}
|
||||
assert_eq!(
|
||||
manager.early_arrival_count(),
|
||||
3,
|
||||
"each parked chunk increments the counter"
|
||||
);
|
||||
// Adoption drains the buffer into the receiver but does not
|
||||
// decrement (the accessor is an observability counter, not a
|
||||
// live-depth gauge).
|
||||
let (_send, mut recv) = manager
|
||||
.adopt_channel(7, "alk/tty", None)
|
||||
.await
|
||||
.expect("adopt");
|
||||
assert_eq!(
|
||||
manager.early_arrival_count(),
|
||||
3,
|
||||
"adopt-drain does not decrement the monotonic counter"
|
||||
);
|
||||
use tokio::io::AsyncReadExt;
|
||||
let mut buf = [0u8; 6 * 3];
|
||||
recv.read_exact(&mut buf).await.expect("read parked chunks");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn route_payload_to_known_channel_does_not_increment_dropped_counter() {
|
||||
let manager = make_manager_with_runner().await;
|
||||
|
||||
+659
-16
@@ -14,7 +14,9 @@
|
||||
//! See `docs/architecture/channel-operations.md` for the spec.
|
||||
|
||||
use std::sync::Arc;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
use futures::future::BoxFuture;
|
||||
use serde_json::{json, Value};
|
||||
|
||||
use crate::core::auth::{AuthContext, Identity};
|
||||
@@ -295,7 +297,10 @@ fn make_resources_subscribe_handler(manager: ChannelManager) -> StreamingHandler
|
||||
}
|
||||
|
||||
/// `ChannelCore` — the channel machinery the ALPN crate's open-op
|
||||
/// wrapper uses (ADR-047 §3). Provides [`ChannelCore::register_openable`],
|
||||
/// wrapper uses (ADR-047 §3, as amended by ADR-049 — the wrapper
|
||||
/// gains an awaited, bounded establishment phase via
|
||||
/// [`ChannelCore::register_openable_with_establisher`]). Provides
|
||||
/// [`ChannelCore::register_openable`],
|
||||
/// which wraps the ALPN's open handler with channel-id allocation,
|
||||
/// `ChannelManager` integration, opener-ledger recording,
|
||||
/// `ChannelLifecyclePolicy` consultation, and teardown hooks.
|
||||
@@ -333,6 +338,85 @@ pub struct ChannelCore {
|
||||
pub type OpenHandler =
|
||||
Arc<dyn Fn(Value, Connection, AuthContext) -> tokio::task::JoinHandle<()> + Send + Sync>;
|
||||
|
||||
/// The establishment-phase result (ADR-049 §1). Reserved for a channel
|
||||
/// plan — today the wrapper consults only success/failure, so `()`
|
||||
/// carries no data.
|
||||
#[derive(Debug, Clone, Default)]
|
||||
pub struct Establishment {}
|
||||
|
||||
/// The establishment-failure reason code the wrapper puts in the open
|
||||
/// reply's `details.reason` (ADR-049 §3). The vocabulary maps 1:1 onto
|
||||
/// what an establisher can actually produce (per the SSH four;
|
||||
/// review 006's survey) plus the wrapper's own deadline.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
|
||||
pub enum EstablishmentError {
|
||||
/// The backend/target could not be reached or refused the
|
||||
/// connection.
|
||||
#[error("dial failed: {message}")]
|
||||
DialFailed { message: String },
|
||||
/// The requested resource does not exist.
|
||||
#[error("unknown resource: {message}")]
|
||||
UnknownResource { message: String },
|
||||
/// The backend is out of capacity (ports, fds, slots).
|
||||
#[error("resource shortage: {message}")]
|
||||
ResourceShortage { message: String },
|
||||
/// Establisher-internal failure not covered above.
|
||||
#[error("handler error: {message}")]
|
||||
HandlerError { message: String },
|
||||
}
|
||||
|
||||
impl EstablishmentError {
|
||||
/// The wire-visible reason code (`details.reason` in
|
||||
/// `channel:open_failed`).
|
||||
pub fn reason(&self) -> &'static str {
|
||||
match self {
|
||||
EstablishmentError::DialFailed { .. } => "dial_failed",
|
||||
EstablishmentError::UnknownResource { .. } => "unknown_resource",
|
||||
EstablishmentError::ResourceShortage { .. } => "resource_shortage",
|
||||
EstablishmentError::HandlerError { .. } => "handler_error",
|
||||
}
|
||||
}
|
||||
|
||||
/// The human-readable detail message.
|
||||
pub fn message(&self) -> &str {
|
||||
match self {
|
||||
EstablishmentError::DialFailed { message }
|
||||
| EstablishmentError::UnknownResource { message }
|
||||
| EstablishmentError::ResourceShortage { message }
|
||||
| EstablishmentError::HandlerError { message } => message,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// The establisher hook (ADR-049 §1): the awaited establishment phase
|
||||
/// of a channel open — validate params semantically, consult ownership,
|
||||
/// prepare/dial the backend. **Awaited by the open-op wrapper,
|
||||
/// bounded** (dispatch deadline or `ESTABLISHMENT_TIMEOUT`, whichever
|
||||
/// is earlier) — before the reply is written, before the pump handler
|
||||
/// is spawned. On failure the wrapper tears down the just-allocated
|
||||
/// channel and replies `channel:open_failed` with `details.reason`.
|
||||
///
|
||||
/// The establisher takes the open op's `input` (params) and the peer's
|
||||
/// `AuthContext`. It deliberately does **not** receive the channel's
|
||||
/// `Connection`: the channel `BiStream` is yield-once
|
||||
/// (`ChannelBidiStreamSource`) and belongs exclusively to the pump
|
||||
/// handler (the `OpenHandler`); the establisher is pre-data-plane (its
|
||||
/// dial targets the backend, not the channel). Amendment to the
|
||||
/// ADR-049 §1 signature sketch.
|
||||
pub type OpenEstablisher = Arc<
|
||||
dyn Fn(Value, AuthContext) -> BoxFuture<'static, Result<Establishment, EstablishmentError>>
|
||||
+ Send
|
||||
+ Sync,
|
||||
>;
|
||||
|
||||
/// The default bound on the establisher await (ADR-049 §2) when the
|
||||
/// `OperationContext` carries no deadline and the registration set no
|
||||
/// override.
|
||||
pub const ESTABLISHMENT_TIMEOUT: Duration = Duration::from_secs(10);
|
||||
|
||||
/// The wire-visible error code for establishment failure (ADR-049 §3).
|
||||
pub const CHANNEL_OPEN_FAILED: &str = "channel:open_failed";
|
||||
|
||||
/// The error code prefix for channel-open failures (ADR-047 §3). The
|
||||
/// wrapper maps `ChannelError` / `ManagerError` to `CallError` with
|
||||
/// codes prefixed by `channel:` so the initiator can branch on the
|
||||
@@ -374,10 +458,17 @@ impl ChannelCore {
|
||||
/// spawn the protocol on the channel's `BiStream`). This wrapper
|
||||
/// does the channel machinery: ACL check (run by the registry's
|
||||
/// `invoke`/`invoke_streaming`/`invoke_sink` before the wrapper) →
|
||||
/// `check_open(identity)` → `manager.open_channel(alpn, opener)`
|
||||
/// → spawn the ALPN handler on the channel's `Connection` →
|
||||
/// `check_open(identity)` → `manager.open_channel(alpn, opener)` →
|
||||
/// spawn the ALPN handler on the channel's `Connection` →
|
||||
/// respond with `{ "channel_id": <id> }`.
|
||||
///
|
||||
/// **No establisher**: the open op cannot fail after allocation —
|
||||
/// establishment work inside the handler is invisible to the open
|
||||
/// reply (review 006 E-01's subject). Use
|
||||
/// [`ChannelCore::register_openable_with_establisher`] when the
|
||||
/// establishment phase can fail and the failure belongs in the
|
||||
/// open reply (ADR-049).
|
||||
///
|
||||
/// The op is registered on the given `registry` (the connection
|
||||
/// overlay registry, Layer 2 per ADR-019 — this is the
|
||||
/// per-connection registration the §4 amendment blesses). The
|
||||
@@ -402,6 +493,38 @@ impl ChannelCore {
|
||||
open_handler: OpenHandler,
|
||||
registry: &OperationRegistry,
|
||||
auth: AuthContext,
|
||||
) -> Result<(), String> {
|
||||
self.register_openable_with_establisher(spec, None, open_handler, registry, auth, None)
|
||||
}
|
||||
|
||||
/// Register a per-ALPN open op with an establishment phase
|
||||
/// (ADR-049 §1/§5). The `establisher` is **awaited by the wrapper,
|
||||
/// bounded** — after allocation, before the reply and before the
|
||||
/// pump handler is spawned. On establisher failure (error or
|
||||
/// deadline) the wrapper tears down the just-allocated channel and
|
||||
/// replies `channel:open_failed` with
|
||||
/// `details: { reason, message }`; on success the pump handler is
|
||||
/// spawned exactly as [`ChannelCore::register_openable`] does.
|
||||
///
|
||||
/// `timeout` bounds the establisher when the dispatch carries no
|
||||
/// deadline (`None` = [`ESTABLISHMENT_TIMEOUT`], 10s). When the
|
||||
/// `OperationContext` carries a deadline, the effective bound is
|
||||
/// the earlier of the two (ADR-049 §2). The bound applies only to
|
||||
/// the establisher — the spawned pump's lifetime is governed by
|
||||
/// the existing teardown machinery, unchanged.
|
||||
///
|
||||
/// Backward compatible by construction: `establisher: None` behaves
|
||||
/// exactly like [`ChannelCore::register_openable`] (the
|
||||
/// no-establisher shape), and existing `OpenHandler`s compile
|
||||
/// unchanged.
|
||||
pub fn register_openable_with_establisher(
|
||||
&self,
|
||||
spec: OperationSpec,
|
||||
establisher: Option<OpenEstablisher>,
|
||||
open_handler: OpenHandler,
|
||||
registry: &OperationRegistry,
|
||||
auth: AuthContext,
|
||||
timeout: Option<Duration>,
|
||||
) -> Result<(), String> {
|
||||
let op_type = spec.op_type;
|
||||
let channel_open = spec.channel_open.clone();
|
||||
@@ -417,17 +540,22 @@ impl ChannelCore {
|
||||
.alpn
|
||||
.to_string();
|
||||
|
||||
let establisher = establisher.map(|establisher| EstablisherHook {
|
||||
establisher,
|
||||
timeout,
|
||||
});
|
||||
let manager = self.manager.clone();
|
||||
let policy = Arc::clone(&self.policy);
|
||||
let open_handler = Arc::clone(&open_handler);
|
||||
|
||||
let handler_kind = match op_type {
|
||||
OperationType::Query | OperationType::Mutation => HandlerKind::Once(
|
||||
make_open_handler_once(manager, policy, open_handler, auth, alpn),
|
||||
make_open_handler_once(manager, policy, establisher, open_handler, auth, alpn),
|
||||
),
|
||||
OperationType::Sub => HandlerKind::Stream(make_open_handler_stream(
|
||||
manager,
|
||||
policy,
|
||||
establisher,
|
||||
open_handler,
|
||||
auth,
|
||||
alpn,
|
||||
@@ -435,6 +563,7 @@ impl ChannelCore {
|
||||
OperationType::Pub => HandlerKind::Sink(make_open_handler_sink(
|
||||
manager,
|
||||
policy,
|
||||
establisher,
|
||||
open_handler,
|
||||
auth,
|
||||
alpn,
|
||||
@@ -467,22 +596,103 @@ fn opener_identity_from_context(ctx: &OperationContext) -> Identity {
|
||||
})
|
||||
}
|
||||
|
||||
/// The establisher hook + its per-registration bound override,
|
||||
/// threaded from `ChannelCore::register_openable_with_establisher`
|
||||
/// into the wrapper factories. `timeout: None` = `ESTABLISHMENT_TIMEOUT`.
|
||||
#[derive(Clone)]
|
||||
struct EstablisherHook {
|
||||
establisher: OpenEstablisher,
|
||||
timeout: Option<Duration>,
|
||||
}
|
||||
|
||||
/// The effective bound on the establisher await (ADR-049 §2): the
|
||||
/// earlier of the dispatch deadline (when the `OperationContext`
|
||||
/// carries one) and the establishment timeout (the per-registration
|
||||
/// override, else `ESTABLISHMENT_TIMEOUT`). A deadline already in the
|
||||
/// past yields a zero bound — the timeout fires immediately.
|
||||
fn establishment_bound(timeout_override: Option<Duration>, deadline: Option<Instant>) -> Duration {
|
||||
let timeout = timeout_override.unwrap_or(ESTABLISHMENT_TIMEOUT);
|
||||
deadline
|
||||
.map(|d| d.saturating_duration_since(Instant::now()).min(timeout))
|
||||
.unwrap_or(timeout)
|
||||
}
|
||||
|
||||
/// Tear down a just-allocated channel after establishment failure
|
||||
/// (ADR-049 §3): drop the demux sender (EOF to nothing — no handler
|
||||
/// was ever spawned), take the opener-ledger entry, and un-reserve the
|
||||
/// per-identity count (`policy.on_close` — the same un-increment path
|
||||
/// the allocation-failure arms run; the ledger `take` is the atomic
|
||||
/// gate, ADR-047 §7).
|
||||
fn teardown_failed_channel(
|
||||
manager: &ChannelManager,
|
||||
policy: &Arc<dyn ChannelLifecyclePolicy>,
|
||||
channel_id: u32,
|
||||
opener: &Identity,
|
||||
) {
|
||||
if let Err(e) = manager.teardown_channel(channel_id) {
|
||||
tracing::debug!(
|
||||
channel_id,
|
||||
error = %e,
|
||||
"open wrapper: establisher-failure teardown found no channel to remove"
|
||||
);
|
||||
}
|
||||
if manager.opener_ledger().take(channel_id).is_some() {
|
||||
policy.on_close(opener);
|
||||
}
|
||||
}
|
||||
|
||||
/// The `channel:open_failed` reply for an establisher error
|
||||
/// (ADR-049 §3): `details` carries the reason code + message.
|
||||
fn establishment_error_to_call_error(err: &EstablishmentError) -> CallError {
|
||||
CallError::new(
|
||||
CHANNEL_OPEN_FAILED,
|
||||
format!("channel establishment failed: {}", err.message()),
|
||||
false,
|
||||
)
|
||||
.with_details(json!({ "reason": err.reason(), "message": err.message() }))
|
||||
}
|
||||
|
||||
/// The `channel:open_failed` reply for a deadline expiry — same
|
||||
/// `details` shape, reason `timeout`.
|
||||
fn establishment_timeout_call_error() -> CallError {
|
||||
CallError::new(
|
||||
CHANNEL_OPEN_FAILED,
|
||||
"channel establishment exceeded the deadline",
|
||||
false,
|
||||
)
|
||||
.with_details(json!({
|
||||
"reason": "timeout",
|
||||
"message": "establishment exceeded the deadline",
|
||||
}))
|
||||
}
|
||||
|
||||
/// Run the open-op wrapper's shared post-ACL steps: `check_open` →
|
||||
/// `open_channel` → build the channel `Connection` → spawn the ALPN
|
||||
/// handler → return `{ channel_id }` on success, or a `CallError` on
|
||||
/// any failure. The ACL check has already run (by the registry's
|
||||
/// `invoke`/`invoke_streaming`/`invoke_sink` before the wrapper); this
|
||||
/// function is the `check_open` → `allocate` → `ledger` → `spawn` →
|
||||
/// `respond` tail of the ADR-047 §3 flow.
|
||||
/// `open_channel` → **await the establisher bounded** (ADR-049 §1/§2 —
|
||||
/// no-op when no establisher is registered) → build the channel
|
||||
/// `Connection` → spawn the ALPN handler → return `{ channel_id }` on
|
||||
/// success, or a `CallError` on any failure. The ACL check has already
|
||||
/// run (by the registry's `invoke`/`invoke_streaming`/`invoke_sink`
|
||||
/// before the wrapper); this function is the `check_open` →
|
||||
/// `allocate` → `establish` → `ledger` → `spawn` → `respond` tail of
|
||||
/// the ADR-047 §3 flow (as amended by ADR-049).
|
||||
///
|
||||
/// On establishment failure (error or deadline) the wrapper tears down
|
||||
/// the just-allocated channel (demux sender, opener-ledger entry,
|
||||
/// `policy.on_close` un-increment) and replies `channel:open_failed`
|
||||
/// with `details: { reason, message }` — the SSH "channel never exists
|
||||
/// opener-side" property: the consumer's open resolves `Err` and no
|
||||
/// `channel_id` was ever returned.
|
||||
///
|
||||
/// `input` is the open op's input (params), passed through to the
|
||||
/// ALPN's `OpenHandler` so it can validate params and prepare the
|
||||
/// backend. `opener_id` is the `PeerId` recorded in the opener ledger
|
||||
/// (the direct caller's identity id).
|
||||
/// establisher (when registered) and the ALPN's `OpenHandler`.
|
||||
/// `opener_id` is the `PeerId` recorded in the opener ledger (the
|
||||
/// direct caller's identity id). `opener_identity` is the identity the
|
||||
/// cap was reserved against (un-reserved on failure paths).
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
async fn run_open_wrapper(
|
||||
manager: &ChannelManager,
|
||||
policy: &Arc<dyn ChannelLifecyclePolicy>,
|
||||
establisher: Option<&EstablisherHook>,
|
||||
open_handler: &OpenHandler,
|
||||
auth: &AuthContext,
|
||||
alpn: &str,
|
||||
@@ -490,6 +700,7 @@ async fn run_open_wrapper(
|
||||
opener_id: String,
|
||||
opener_identity: Identity,
|
||||
request_id: String,
|
||||
deadline: Option<Instant>,
|
||||
) -> ResponseEnvelope {
|
||||
if let Err(channel_err) = policy.check_open(&opener_identity) {
|
||||
return ResponseEnvelope::error(request_id, map_channel_error_to_call_error(&channel_err));
|
||||
@@ -497,6 +708,39 @@ async fn run_open_wrapper(
|
||||
|
||||
let channel_id = match manager.open_channel(alpn, opener_id, None).await {
|
||||
Ok((id, send, recv)) => {
|
||||
// The establishment phase (ADR-049 §1): awaited bounded,
|
||||
// before the reply and before the pump handler is spawned.
|
||||
// No establisher registered = an always-OK establisher
|
||||
// (the pre-ADR-049 shape, unchanged).
|
||||
if let Some(hook) = establisher {
|
||||
let bound = establishment_bound(hook.timeout, deadline);
|
||||
let establishment =
|
||||
tokio::time::timeout(bound, (hook.establisher)(input.clone(), auth.clone()))
|
||||
.await;
|
||||
match establishment {
|
||||
Err(_elapsed) => {
|
||||
teardown_failed_channel(manager, policy, id, &opener_identity);
|
||||
return ResponseEnvelope::error(
|
||||
request_id,
|
||||
establishment_timeout_call_error(),
|
||||
);
|
||||
}
|
||||
Ok(Err(e)) => {
|
||||
tracing::debug!(
|
||||
channel_id = id,
|
||||
reason = e.reason(),
|
||||
"open wrapper: establisher failed; tearing down channel"
|
||||
);
|
||||
teardown_failed_channel(manager, policy, id, &opener_identity);
|
||||
return ResponseEnvelope::error(
|
||||
request_id,
|
||||
establishment_error_to_call_error(&e),
|
||||
);
|
||||
}
|
||||
Ok(Ok(_establishment)) => {}
|
||||
}
|
||||
}
|
||||
|
||||
let remote_addr = manager.remote_addr();
|
||||
let source = super::source::channel_source(recv, send, remote_addr);
|
||||
let channel_conn = Connection::from_source(source, alpn.as_bytes().to_vec());
|
||||
@@ -506,7 +750,23 @@ async fn run_open_wrapper(
|
||||
let teardown_policy = Arc::clone(policy);
|
||||
let task = tokio::spawn(async move {
|
||||
let _ = raw_task.await;
|
||||
let _ = teardown_manager.teardown_channel(id);
|
||||
// Benign-race pin (review 006 E-03): `channel/close`'s
|
||||
// awaited teardown and this handler-exit teardown both
|
||||
// gate on the atomic `opener_ledger().take(id)`, so the
|
||||
// loser of the race gets `UnknownChannel` here and skips
|
||||
// the ledger/policy half — no double-decrement. A
|
||||
// vanished-after-install channel (transport EOF
|
||||
// `clear_all`) is equally benign: `clear_all` drained
|
||||
// the ledger and decremented the policy itself. The
|
||||
// vanished-between-open-and-install case surfaces as the
|
||||
// `set_handler_task` warn below.
|
||||
if let Err(e) = teardown_manager.teardown_channel(id) {
|
||||
tracing::debug!(
|
||||
channel_id = id,
|
||||
error = %e,
|
||||
"open wrapper: handler-exit teardown found no channel to remove (benign when racing channel/close)"
|
||||
);
|
||||
}
|
||||
if let Some(opener_id) = teardown_manager.opener_ledger().take(id) {
|
||||
let opener = Identity {
|
||||
id: opener_id,
|
||||
@@ -584,6 +844,7 @@ fn map_channel_error_to_call_error(err: &ChannelError) -> CallError {
|
||||
fn make_open_handler_once(
|
||||
manager: ChannelManager,
|
||||
policy: Arc<dyn ChannelLifecyclePolicy>,
|
||||
establisher: Option<EstablisherHook>,
|
||||
open_handler: OpenHandler,
|
||||
auth: AuthContext,
|
||||
alpn: String,
|
||||
@@ -591,6 +852,7 @@ fn make_open_handler_once(
|
||||
Arc::new(move |input: Value, ctx: OperationContext| {
|
||||
let manager = manager.clone();
|
||||
let policy = Arc::clone(&policy);
|
||||
let establisher = establisher.clone();
|
||||
let open_handler = Arc::clone(&open_handler);
|
||||
let auth = auth.clone();
|
||||
let alpn = alpn.clone();
|
||||
@@ -601,6 +863,7 @@ fn make_open_handler_once(
|
||||
run_open_wrapper(
|
||||
&manager,
|
||||
&policy,
|
||||
establisher.as_ref(),
|
||||
&open_handler,
|
||||
&auth,
|
||||
&alpn,
|
||||
@@ -608,6 +871,7 @@ fn make_open_handler_once(
|
||||
opener_id,
|
||||
opener_identity,
|
||||
request_id,
|
||||
ctx.deadline,
|
||||
)
|
||||
.await
|
||||
})
|
||||
@@ -625,6 +889,7 @@ fn make_open_handler_once(
|
||||
fn make_open_handler_stream(
|
||||
manager: ChannelManager,
|
||||
policy: Arc<dyn ChannelLifecyclePolicy>,
|
||||
establisher: Option<EstablisherHook>,
|
||||
open_handler: OpenHandler,
|
||||
auth: AuthContext,
|
||||
alpn: String,
|
||||
@@ -632,6 +897,7 @@ fn make_open_handler_stream(
|
||||
Arc::new(move |input: Value, ctx: OperationContext| {
|
||||
let manager = manager.clone();
|
||||
let policy = Arc::clone(&policy);
|
||||
let establisher = establisher.clone();
|
||||
let open_handler = Arc::clone(&open_handler);
|
||||
let auth = auth.clone();
|
||||
let alpn = alpn.clone();
|
||||
@@ -642,6 +908,7 @@ fn make_open_handler_stream(
|
||||
run_open_wrapper(
|
||||
&manager,
|
||||
&policy,
|
||||
establisher.as_ref(),
|
||||
&open_handler,
|
||||
&auth,
|
||||
&alpn,
|
||||
@@ -649,6 +916,7 @@ fn make_open_handler_stream(
|
||||
opener_id,
|
||||
opener_identity,
|
||||
request_id,
|
||||
ctx.deadline,
|
||||
)
|
||||
.await
|
||||
}))
|
||||
@@ -675,6 +943,7 @@ fn make_open_handler_stream(
|
||||
fn make_open_handler_sink(
|
||||
_manager: ChannelManager,
|
||||
_policy: Arc<dyn ChannelLifecyclePolicy>,
|
||||
_establisher: Option<EstablisherHook>,
|
||||
_open_handler: OpenHandler,
|
||||
_auth: AuthContext,
|
||||
_alpn: String,
|
||||
@@ -1039,6 +1308,7 @@ mod tests {
|
||||
let handler = make_open_handler_once(
|
||||
manager.clone(),
|
||||
policy,
|
||||
None,
|
||||
open_handler,
|
||||
auth,
|
||||
"alk/tty".to_string(),
|
||||
@@ -1069,6 +1339,7 @@ mod tests {
|
||||
let handler = make_open_handler_stream(
|
||||
manager.clone(),
|
||||
policy,
|
||||
None,
|
||||
open_handler,
|
||||
auth,
|
||||
"alk/tty".to_string(),
|
||||
@@ -1096,8 +1367,14 @@ mod tests {
|
||||
let policy = super::super::policy::default_policy();
|
||||
let open_handler: OpenHandler = Arc::new(|_input, _conn, _auth| tokio::spawn(async {}));
|
||||
let auth = AuthContext::anonymous(b"alk/call");
|
||||
let handler =
|
||||
make_open_handler_sink(manager, policy, open_handler, auth, "alk/tty".to_string());
|
||||
let handler = make_open_handler_sink(
|
||||
manager,
|
||||
policy,
|
||||
None,
|
||||
open_handler,
|
||||
auth,
|
||||
"alk/tty".to_string(),
|
||||
);
|
||||
let publish_stream: crate::registry::registration::PublishStream =
|
||||
Box::pin(futures::stream::empty());
|
||||
let env = handler(json!({}), test_context("open-sink-1"), publish_stream).await;
|
||||
@@ -1265,6 +1542,7 @@ mod tests {
|
||||
let env = run_open_wrapper(
|
||||
&manager,
|
||||
&policy,
|
||||
None,
|
||||
&open_handler,
|
||||
&auth,
|
||||
"alk/tty",
|
||||
@@ -1272,6 +1550,7 @@ mod tests {
|
||||
opener_id.id.clone(),
|
||||
opener_id,
|
||||
"req-cap-deny".to_string(),
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
match env.result {
|
||||
@@ -1279,4 +1558,368 @@ mod tests {
|
||||
Ok(_) => panic!("cap 0 should deny open"),
|
||||
}
|
||||
}
|
||||
|
||||
// --- ADR-049 Unit 1: the establishment phase -------------------------
|
||||
|
||||
fn failing_establisher() -> OpenEstablisher {
|
||||
Arc::new(|_input, _auth| {
|
||||
Box::pin(async {
|
||||
Err(EstablishmentError::DialFailed {
|
||||
message: "connection refused".to_string(),
|
||||
})
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
fn ok_establisher() -> OpenEstablisher {
|
||||
Arc::new(|_input, _auth| Box::pin(async { Ok(Establishment::default()) }))
|
||||
}
|
||||
|
||||
fn hook(establisher: OpenEstablisher, timeout: Option<Duration>) -> Option<EstablisherHook> {
|
||||
Some(EstablisherHook {
|
||||
establisher,
|
||||
timeout,
|
||||
})
|
||||
}
|
||||
|
||||
fn identity(id: &str) -> Identity {
|
||||
Identity {
|
||||
id: id.to_string(),
|
||||
scopes: vec![],
|
||||
resources: HashMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn run_open_wrapper_establisher_failure_tears_down_and_replies_open_failed() {
|
||||
let manager = make_manager().await;
|
||||
let concrete_policy = Arc::new(super::super::policy::PerIdentityChannelPolicy::new(8));
|
||||
let policy: Arc<dyn ChannelLifecyclePolicy> =
|
||||
Arc::clone(&concrete_policy) as Arc<dyn ChannelLifecyclePolicy>;
|
||||
let spawned = Arc::new(AtomicBool::new(false));
|
||||
let spawned_clone = Arc::clone(&spawned);
|
||||
let open_handler: OpenHandler = Arc::new(move |_input, _conn, _auth| {
|
||||
spawned_clone.store(true, Ordering::SeqCst);
|
||||
tokio::spawn(async {})
|
||||
});
|
||||
let auth = AuthContext::anonymous(b"alk/call");
|
||||
let env = run_open_wrapper(
|
||||
&manager,
|
||||
&policy,
|
||||
hook(failing_establisher(), None).as_ref(),
|
||||
&open_handler,
|
||||
&auth,
|
||||
"alk/tty",
|
||||
json!({}),
|
||||
"alice".to_string(),
|
||||
identity("alice"),
|
||||
"req-est-fail".to_string(),
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
match env.result {
|
||||
Err(e) => {
|
||||
assert_eq!(e.code, "channel:open_failed");
|
||||
assert!(!e.retryable);
|
||||
let details = e.details.expect("details carry the reason");
|
||||
assert_eq!(details["reason"], "dial_failed");
|
||||
assert_eq!(details["message"], "connection refused");
|
||||
}
|
||||
Ok(_) => panic!("establisher failure must fail the open"),
|
||||
}
|
||||
// The SSH "channel never exists opener-side" property: the
|
||||
// channel was torn down and no `channel_id` was ever returned.
|
||||
assert!(
|
||||
manager.channel_ids().is_empty(),
|
||||
"no channel survives a failed establishment"
|
||||
);
|
||||
assert_eq!(manager.open_count(), 0);
|
||||
// The per-identity cap reservation was un-incremented.
|
||||
assert_eq!(
|
||||
concrete_policy.count_for(&identity("alice")),
|
||||
0,
|
||||
"ledger take + policy.on_close restored the cap count"
|
||||
);
|
||||
// The pump handler was never spawned.
|
||||
assert!(
|
||||
!spawned.load(Ordering::SeqCst),
|
||||
"pump handler must not spawn on establishment failure"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn run_open_wrapper_establisher_timeout_replies_timeout_reason() {
|
||||
let manager = make_manager().await;
|
||||
let concrete_policy = Arc::new(super::super::policy::PerIdentityChannelPolicy::new(8));
|
||||
let policy: Arc<dyn ChannelLifecyclePolicy> =
|
||||
Arc::clone(&concrete_policy) as Arc<dyn ChannelLifecyclePolicy>;
|
||||
let open_handler: OpenHandler = Arc::new(|_input, _conn, _auth| tokio::spawn(async {}));
|
||||
let auth = AuthContext::anonymous(b"alk/call");
|
||||
let hanging: OpenEstablisher = Arc::new(|_input, _auth| {
|
||||
Box::pin(async {
|
||||
tokio::time::sleep(Duration::from_secs(5)).await;
|
||||
Ok(Establishment::default())
|
||||
})
|
||||
});
|
||||
let env = run_open_wrapper(
|
||||
&manager,
|
||||
&policy,
|
||||
hook(hanging, Some(Duration::from_millis(50))).as_ref(),
|
||||
&open_handler,
|
||||
&auth,
|
||||
"alk/tty",
|
||||
json!({}),
|
||||
"alice".to_string(),
|
||||
identity("alice"),
|
||||
"req-est-timeout".to_string(),
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
match env.result {
|
||||
Err(e) => {
|
||||
assert_eq!(e.code, "channel:open_failed");
|
||||
let details = e.details.expect("details carry the reason");
|
||||
assert_eq!(details["reason"], "timeout");
|
||||
}
|
||||
Ok(_) => panic!("a hung establisher must fail the open within the bound"),
|
||||
}
|
||||
assert!(manager.channel_ids().is_empty(), "channel torn down");
|
||||
assert_eq!(
|
||||
concrete_policy.count_for(&identity("alice")),
|
||||
0,
|
||||
"ledger decremented on timeout"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn run_open_wrapper_dispatch_deadline_bounds_establisher() {
|
||||
let manager = make_manager().await;
|
||||
let policy = super::super::policy::default_policy();
|
||||
let open_handler: OpenHandler = Arc::new(|_input, _conn, _auth| tokio::spawn(async {}));
|
||||
let auth = AuthContext::anonymous(b"alk/call");
|
||||
let hanging: OpenEstablisher = Arc::new(|_input, _auth| {
|
||||
Box::pin(async {
|
||||
tokio::time::sleep(Duration::from_secs(5)).await;
|
||||
Ok(Establishment::default())
|
||||
})
|
||||
});
|
||||
// No per-registration override; the dispatch deadline (50ms) is
|
||||
// the bound — the earlier of the two per ADR-049 §2.
|
||||
let env = run_open_wrapper(
|
||||
&manager,
|
||||
&policy,
|
||||
hook(hanging, None).as_ref(),
|
||||
&open_handler,
|
||||
&auth,
|
||||
"alk/tty",
|
||||
json!({}),
|
||||
"alice".to_string(),
|
||||
identity("alice"),
|
||||
"req-est-deadline".to_string(),
|
||||
Some(Instant::now() + Duration::from_millis(50)),
|
||||
)
|
||||
.await;
|
||||
match env.result {
|
||||
Err(e) => {
|
||||
assert_eq!(e.code, "channel:open_failed");
|
||||
assert_eq!(e.details.as_ref().expect("details")["reason"], "timeout");
|
||||
}
|
||||
Ok(_) => panic!("deadline must bound the establisher"),
|
||||
}
|
||||
assert!(manager.channel_ids().is_empty());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn run_open_wrapper_establisher_success_spawns_pumps_and_replies_channel_id() {
|
||||
let manager = make_manager().await;
|
||||
let policy = super::super::policy::default_policy();
|
||||
let spawned = Arc::new(AtomicBool::new(false));
|
||||
let spawned_clone = Arc::clone(&spawned);
|
||||
let open_handler: OpenHandler = Arc::new(move |_input, _conn, _auth| {
|
||||
spawned_clone.store(true, Ordering::SeqCst);
|
||||
tokio::spawn(async {})
|
||||
});
|
||||
let auth = AuthContext::anonymous(b"alk/call");
|
||||
let env = run_open_wrapper(
|
||||
&manager,
|
||||
&policy,
|
||||
hook(ok_establisher(), None).as_ref(),
|
||||
&open_handler,
|
||||
&auth,
|
||||
"alk/tty",
|
||||
json!({}),
|
||||
"alice".to_string(),
|
||||
identity("alice"),
|
||||
"req-est-ok".to_string(),
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
match env.result {
|
||||
Ok(v) => {
|
||||
let channel_id = v["channel_id"].as_u64().expect("channel_id");
|
||||
assert!(channel_id > 0);
|
||||
assert!(manager.has_channel(channel_id as u32));
|
||||
}
|
||||
Err(e) => panic!("successful establishment should open, got: {e:?}"),
|
||||
}
|
||||
assert!(spawned.load(Ordering::SeqCst));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn establishment_error_reason_mapping_covers_the_vocabulary() {
|
||||
let cases: Vec<(EstablishmentError, &str)> = vec![
|
||||
(
|
||||
EstablishmentError::DialFailed {
|
||||
message: "refused".into(),
|
||||
},
|
||||
"dial_failed",
|
||||
),
|
||||
(
|
||||
EstablishmentError::UnknownResource {
|
||||
message: "gone".into(),
|
||||
},
|
||||
"unknown_resource",
|
||||
),
|
||||
(
|
||||
EstablishmentError::ResourceShortage {
|
||||
message: "no fds".into(),
|
||||
},
|
||||
"resource_shortage",
|
||||
),
|
||||
(
|
||||
EstablishmentError::HandlerError {
|
||||
message: "boom".into(),
|
||||
},
|
||||
"handler_error",
|
||||
),
|
||||
];
|
||||
for (err, reason) in cases {
|
||||
let call_err = establishment_error_to_call_error(&err);
|
||||
assert_eq!(call_err.code, "channel:open_failed");
|
||||
let details = call_err.details.expect("details");
|
||||
assert_eq!(details["reason"], reason);
|
||||
assert_eq!(details["message"], err.message());
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn establishment_timeout_error_carries_timeout_reason() {
|
||||
let err = establishment_timeout_call_error();
|
||||
assert_eq!(err.code, "channel:open_failed");
|
||||
assert_eq!(err.details.as_ref().expect("details")["reason"], "timeout");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn establishment_bound_is_the_earlier_of_deadline_and_timeout() {
|
||||
// No deadline → the override (or the 10s default).
|
||||
assert_eq!(establishment_bound(None, None), ESTABLISHMENT_TIMEOUT);
|
||||
assert_eq!(
|
||||
establishment_bound(Some(Duration::from_secs(3)), None),
|
||||
Duration::from_secs(3)
|
||||
);
|
||||
// Deadline earlier than the timeout → the deadline.
|
||||
let soon = Instant::now() + Duration::from_millis(100);
|
||||
let bound = establishment_bound(None, Some(soon));
|
||||
assert!(bound <= Duration::from_millis(100));
|
||||
// Deadline later than the timeout → the timeout wins.
|
||||
let late = Instant::now() + Duration::from_secs(60);
|
||||
assert_eq!(
|
||||
establishment_bound(Some(Duration::from_secs(3)), Some(late)),
|
||||
Duration::from_secs(3)
|
||||
);
|
||||
// Deadline already past → zero bound (timeout fires immediately).
|
||||
let past = Instant::now() - Duration::from_secs(1);
|
||||
assert_eq!(establishment_bound(None, Some(past)), Duration::ZERO);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn register_openable_with_establisher_fails_via_invoke_streaming() {
|
||||
let manager = make_manager().await;
|
||||
let concrete_policy = Arc::new(super::super::policy::PerIdentityChannelPolicy::new(8));
|
||||
let policy: Arc<dyn ChannelLifecyclePolicy> =
|
||||
Arc::clone(&concrete_policy) as Arc<dyn ChannelLifecyclePolicy>;
|
||||
let core = ChannelCore::new(manager.clone(), policy);
|
||||
let spec = OperationSpec::new(
|
||||
"channels/tty/sub",
|
||||
OperationType::Sub,
|
||||
Visibility::External,
|
||||
json!({}),
|
||||
json!({ "type": "object", "properties": { "channel_id": { "type": "integer" } } }),
|
||||
vec![],
|
||||
AccessControl::default(),
|
||||
None,
|
||||
)
|
||||
.with_channel_open(ChannelOpenSpec::new("alk/tty"));
|
||||
let open_handler: OpenHandler = Arc::new(|_input, _conn, _auth| tokio::spawn(async {}));
|
||||
let registry = OperationRegistry::new();
|
||||
core.register_openable_with_establisher(
|
||||
spec,
|
||||
Some(failing_establisher()),
|
||||
open_handler,
|
||||
®istry,
|
||||
AuthContext::anonymous(b"alk/call"),
|
||||
None,
|
||||
)
|
||||
.expect("register");
|
||||
let mut stream =
|
||||
registry.invoke_streaming("channels/tty/sub", json!({}), test_context("reg-est-fail"));
|
||||
let env = stream.next().await.expect("one envelope");
|
||||
match env.result {
|
||||
Err(e) => {
|
||||
assert_eq!(e.code, "channel:open_failed");
|
||||
assert_eq!(
|
||||
e.details.as_ref().expect("details")["reason"],
|
||||
"dial_failed"
|
||||
);
|
||||
}
|
||||
Ok(_) => panic!("establisher failure must surface as a call error"),
|
||||
}
|
||||
assert!(manager.channel_ids().is_empty());
|
||||
assert_eq!(
|
||||
concrete_policy.count_for(&identity("alice")),
|
||||
0,
|
||||
"no-establisher compat: registry-level failure un-reserves the ledger"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn register_openable_no_establisher_behaves_as_today() {
|
||||
let manager = make_manager().await;
|
||||
let policy = super::super::policy::default_policy();
|
||||
let core = ChannelCore::new(manager.clone(), policy);
|
||||
let spec = OperationSpec::new(
|
||||
"channels/tty/sub",
|
||||
OperationType::Query,
|
||||
Visibility::External,
|
||||
json!({}),
|
||||
json!({ "type": "object", "properties": { "channel_id": { "type": "integer" } } }),
|
||||
vec![],
|
||||
AccessControl::default(),
|
||||
None,
|
||||
)
|
||||
.with_channel_open(ChannelOpenSpec::new("alk/tty"));
|
||||
let spawned = Arc::new(AtomicBool::new(false));
|
||||
let spawned_clone = Arc::clone(&spawned);
|
||||
let open_handler: OpenHandler = Arc::new(move |_input, _conn, _auth| {
|
||||
spawned_clone.store(true, Ordering::SeqCst);
|
||||
tokio::spawn(async {})
|
||||
});
|
||||
let registry = OperationRegistry::new();
|
||||
core.register_openable(
|
||||
spec,
|
||||
open_handler,
|
||||
®istry,
|
||||
AuthContext::anonymous(b"alk/call"),
|
||||
)
|
||||
.expect("register");
|
||||
let env = registry
|
||||
.invoke("channels/tty/sub", json!({}), test_context("reg-compat"))
|
||||
.await;
|
||||
match env.result {
|
||||
Ok(v) => assert!(v["channel_id"].as_u64().is_some()),
|
||||
Err(e) => panic!("no-establisher open should succeed, got: {e:?}"),
|
||||
}
|
||||
assert!(spawned.load(Ordering::SeqCst));
|
||||
assert_eq!(manager.open_count(), 1, "channel stays open (compat)");
|
||||
}
|
||||
}
|
||||
@@ -285,6 +285,10 @@ pub(crate) fn rebuild_spec_for(
|
||||
}
|
||||
}
|
||||
|
||||
if let Some(description) = schema_json.get("description").and_then(|v| v.as_str()) {
|
||||
spec = spec.with_description(description);
|
||||
}
|
||||
|
||||
Ok(spec)
|
||||
}
|
||||
|
||||
@@ -768,6 +772,64 @@ mod tests {
|
||||
assert_eq!(rebuilt.publish_schema.as_ref(), Some(&publish_schema));
|
||||
}
|
||||
|
||||
/// E-02 gate (review 006): `description` survives the spec wire
|
||||
/// round-trip — `spec_to_json_pub` serializes it when set,
|
||||
/// `rebuild_spec_for` parses it back (the `op/register` announced-spec
|
||||
/// path and the `from_call` import path both parse through here).
|
||||
#[test]
|
||||
fn spec_round_trips_description() {
|
||||
use crate::registry::discovery::spec_to_json_pub;
|
||||
|
||||
let spec = OperationSpec::new(
|
||||
"channels/tty/sub",
|
||||
OperationType::Sub,
|
||||
Visibility::External,
|
||||
json!({}),
|
||||
json!({}),
|
||||
vec![],
|
||||
crate::registry::spec::AccessControl::default(),
|
||||
None,
|
||||
)
|
||||
.with_description("Interactive TTY sessions");
|
||||
let wire = spec_to_json_pub(&spec);
|
||||
assert_eq!(
|
||||
wire.get("description").and_then(|v| v.as_str()),
|
||||
Some("Interactive TTY sessions"),
|
||||
"description serialized"
|
||||
);
|
||||
|
||||
let rebuilt = rebuild_spec_for(&wire, "channels/tty/sub", &None).expect("rebuild");
|
||||
assert_eq!(
|
||||
rebuilt.description.as_deref(),
|
||||
Some("Interactive TTY sessions"),
|
||||
"description survives the round-trip"
|
||||
);
|
||||
}
|
||||
|
||||
/// E-02 companion: a spec without `description` serializes no
|
||||
/// `description` key and rebuilds with `None` (additive optional
|
||||
/// field — absent stays absent, old producers stay parseable).
|
||||
#[test]
|
||||
fn spec_without_description_stays_absent_through_round_trip() {
|
||||
use crate::registry::discovery::spec_to_json_pub;
|
||||
|
||||
let spec = OperationSpec::new(
|
||||
"fs/readFile",
|
||||
OperationType::Query,
|
||||
Visibility::External,
|
||||
json!({}),
|
||||
json!({}),
|
||||
vec![],
|
||||
crate::registry::spec::AccessControl::default(),
|
||||
None,
|
||||
);
|
||||
let wire = spec_to_json_pub(&spec);
|
||||
assert!(wire.get("description").is_none());
|
||||
|
||||
let rebuilt = rebuild_spec_for(&wire, "fs/readFile", &None).expect("rebuild");
|
||||
assert_eq!(rebuilt.description, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn derive_alpn_from_op_name_strips_channels_prefix() {
|
||||
assert_eq!(
|
||||
|
||||
+126
-4
@@ -33,6 +33,10 @@ pub fn services_list_spec() -> OperationSpec {
|
||||
"op_type": {
|
||||
"type": "string",
|
||||
"enum": ["query", "mutation", "sub", "pub"]
|
||||
},
|
||||
"description": {
|
||||
"type": "string",
|
||||
"description": "Human-readable op description (review 006 E-02). Absent when the producer declares none."
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -87,6 +91,10 @@ pub fn services_list_peers_spec() -> OperationSpec {
|
||||
"op_type": {
|
||||
"type": "string",
|
||||
"enum": ["query", "mutation", "sub", "pub"]
|
||||
},
|
||||
"description": {
|
||||
"type": "string",
|
||||
"description": "Human-readable op description. Absent when the op declares none."
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -116,6 +124,9 @@ fn operation_spec_schema() -> Value {
|
||||
"type": "string",
|
||||
"enum": ["external", "internal"]
|
||||
},
|
||||
"description": {
|
||||
"description": "Human-readable op description (review 006 E-02). Absent when the op declares none; disclosed verbatim in `services/list` when set."
|
||||
},
|
||||
"input_schema": {},
|
||||
"output_schema": {},
|
||||
"error_schemas": {
|
||||
@@ -227,6 +238,9 @@ pub fn spec_to_json_pub(spec: &OperationSpec) -> Value {
|
||||
if let Some(resource_id_path) = &spec.resource_id_path {
|
||||
json["resource_id_path"] = json!(resource_id_path);
|
||||
}
|
||||
if let Some(description) = &spec.description {
|
||||
json["description"] = json!(description);
|
||||
}
|
||||
if spec.channel_open.is_some() {
|
||||
json["channel_open"] = json!(true);
|
||||
}
|
||||
@@ -259,11 +273,15 @@ pub fn services_list_handler(registry: Arc<OperationRegistry>) -> Handler {
|
||||
.is_allowed()
|
||||
})
|
||||
.map(|s| {
|
||||
json!({
|
||||
let mut listing = json!({
|
||||
"name": s.name,
|
||||
"namespace": s.namespace,
|
||||
"op_type": op_type_str(s.op_type),
|
||||
})
|
||||
});
|
||||
if let Some(description) = &s.description {
|
||||
listing["description"] = json!(description);
|
||||
}
|
||||
listing
|
||||
})
|
||||
.collect();
|
||||
ResponseEnvelope::ok(ctx.request_id, json!({ "operations": ops }))
|
||||
@@ -333,11 +351,15 @@ pub fn services_list_peers_handler(registry: Arc<OperationRegistry>) -> Handler
|
||||
.is_allowed()
|
||||
})
|
||||
.map(|s| {
|
||||
json!({
|
||||
let mut listing = json!({
|
||||
"name": s.name,
|
||||
"namespace": s.namespace,
|
||||
"op_type": op_type_str(s.op_type),
|
||||
})
|
||||
});
|
||||
if let Some(description) = &s.description {
|
||||
listing["description"] = json!(description);
|
||||
}
|
||||
listing
|
||||
})
|
||||
.collect();
|
||||
let mut peers: Vec<Value> = Vec::new();
|
||||
@@ -1057,6 +1079,106 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
// --- review 006 E-02: `description` on OperationSpec -------------------
|
||||
|
||||
#[test]
|
||||
fn spec_to_json_emits_description_when_set() {
|
||||
let spec = external_spec("fs/readFile").with_description("Read a file");
|
||||
let json_val = spec_to_json(&spec);
|
||||
assert_eq!(json_val.get("description"), Some(&json!("Read a file")));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn spec_to_json_omits_description_when_absent() {
|
||||
let spec = external_spec("fs/readFile");
|
||||
let json_val = spec_to_json(&spec);
|
||||
assert!(
|
||||
json_val.get("description").is_none(),
|
||||
"description must be absent when not set (additive optional field)"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn services_list_emits_description_when_set() {
|
||||
let registry = Arc::new(OperationRegistry::new());
|
||||
registry
|
||||
.register(HandlerRegistration::new(
|
||||
external_spec("channels/tty/sub").with_description("Interactive TTY sessions"),
|
||||
HandlerKind::Once(echo_handler()),
|
||||
OperationProvenance::Local,
|
||||
None,
|
||||
None,
|
||||
Capabilities::new(),
|
||||
))
|
||||
.unwrap();
|
||||
registry
|
||||
.register(HandlerRegistration::new(
|
||||
external_spec("fs/readFile"),
|
||||
HandlerKind::Once(echo_handler()),
|
||||
OperationProvenance::Local,
|
||||
None,
|
||||
None,
|
||||
Capabilities::new(),
|
||||
))
|
||||
.unwrap();
|
||||
let handler = services_list_handler(Arc::clone(®istry));
|
||||
let response = handler(json!({}), root_context("req-e02-1")).await;
|
||||
let output = response.result.expect("ok response");
|
||||
let ops = output
|
||||
.get("operations")
|
||||
.and_then(|v| v.as_array())
|
||||
.expect("operations array")
|
||||
.iter()
|
||||
.map(|o| {
|
||||
(
|
||||
o.get("name").and_then(|n| n.as_str()).unwrap_or(""),
|
||||
o.get("description").and_then(|d| d.as_str()),
|
||||
)
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
assert!(
|
||||
ops.contains(&("channels/tty/sub", Some("Interactive TTY sessions"))),
|
||||
"described op carries its description: {ops:?}"
|
||||
);
|
||||
assert!(
|
||||
ops.contains(&("fs/readFile", None)),
|
||||
"undescribed op omits the description key: {ops:?}"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn services_schema_discloses_description() {
|
||||
let registry = registry_with_ops();
|
||||
let handler = services_schema_handler(Arc::clone(®istry));
|
||||
let described = external_spec("fs/readFile").with_description("Read a file");
|
||||
registry
|
||||
.register(HandlerRegistration::new(
|
||||
described,
|
||||
HandlerKind::Once(echo_handler()),
|
||||
OperationProvenance::Local,
|
||||
None,
|
||||
None,
|
||||
Capabilities::new(),
|
||||
))
|
||||
.unwrap();
|
||||
let response = handler(json!({ "name": "fs/readFile" }), root_context("req-e02-2")).await;
|
||||
let spec = response.result.expect("ok response");
|
||||
assert_eq!(spec.get("description"), Some(&json!("Read a file")));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn operation_spec_schema_documents_description_property() {
|
||||
let schema = operation_spec_schema();
|
||||
let props = schema
|
||||
.get("properties")
|
||||
.and_then(|v| v.as_object())
|
||||
.expect("properties object");
|
||||
assert!(
|
||||
props.contains_key("description"),
|
||||
"operation_spec_schema must advertise description"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn services_list_filters_by_access_control_authorized_peer() {
|
||||
let registry = registry_with_access_controlled_ops();
|
||||
|
||||
@@ -181,6 +181,13 @@ pub struct OperationSpec {
|
||||
pub output_schema: Value,
|
||||
pub error_schemas: Vec<ErrorDefinition>,
|
||||
pub access_control: AccessControl,
|
||||
/// Human-readable op description, disclosed by `services/list` when
|
||||
/// set and by `services/schema` (`spec_to_json_pub`) — review 006
|
||||
/// E-02. `None` (absent on the wire) for ops that declare none; the
|
||||
/// field is additive and registry-side (it describes the op, not the
|
||||
/// produced resource set — ADR-047 §6 keeps the live half of
|
||||
/// discovery in `channel/resources/subscribe`).
|
||||
pub description: Option<String>,
|
||||
/// JSON pointer into the input for the resource ID, when
|
||||
/// `access_control.resource_type` is set and the operation targets a
|
||||
/// specific runtime-spawned resource (ADR-011). e.g. `"$.containerId"`
|
||||
@@ -234,12 +241,22 @@ impl OperationSpec {
|
||||
output_schema,
|
||||
error_schemas,
|
||||
access_control,
|
||||
description: None,
|
||||
resource_id_path,
|
||||
publish_schema: None,
|
||||
channel_open: None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Set the op's `description` (review 006 E-02). Disclosed by
|
||||
/// `services/list` when set; carried in the `services/schema` wire
|
||||
/// shape (`spec_to_json_pub` / `rebuild_spec_for`). Builder-style;
|
||||
/// returns `self` for chaining at registration sites.
|
||||
pub fn with_description(mut self, description: impl Into<String>) -> Self {
|
||||
self.description = Some(description.into());
|
||||
self
|
||||
}
|
||||
|
||||
/// Set the `publish_schema` (Pub ops only, ADR-046). Validates each
|
||||
/// `call.published` chunk's `input`. Builder-style; returns `self`
|
||||
/// for chaining at registration sites.
|
||||
@@ -359,6 +376,27 @@ mod tests {
|
||||
assert_eq!(spec.channel_open, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn description_defaults_to_none_and_builder_sets_it() {
|
||||
let spec = OperationSpec::new(
|
||||
"channels/tty/sub",
|
||||
OperationType::Sub,
|
||||
Visibility::External,
|
||||
serde_json::json!({}),
|
||||
serde_json::json!({}),
|
||||
vec![],
|
||||
AccessControl::default(),
|
||||
None,
|
||||
);
|
||||
assert_eq!(spec.description, None);
|
||||
|
||||
let described = spec.with_description("Interactive TTY sessions");
|
||||
assert_eq!(
|
||||
described.description.as_deref(),
|
||||
Some("Interactive TTY sessions")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn with_channel_open_sets_marker() {
|
||||
let spec = OperationSpec::new(
|
||||
|
||||
Reference in new issue
Block a user