5 Commits
Author SHA1 Message Date
glm-5.3-flash 36e74cda11 feat(review 006 Unit 3): teardown-race log, early-arrival bound docs, count accessor (E-03, E-04, N-2)
- E-03: the open wrapper's handler-exit teardown no longer discards
  UnknownChannel silently — debug log + benign-race pinning comment
  (ledger take is the atomic gate; no double-decrement)
- E-04: consumer-facing doc note on the 64-parked-chunks observable
  bound (channel-client.md + EARLY_ARRIVAL_CAP const doc) for
  tunnel-style push-first producers
- N-2: ChannelManager::early_arrival_count() accessor (the
  observability choice over removing the write-only counter),
  documented monotonic, with a park/adopt-drain monotonicity test
- review 006: Unit 3 marked implemented in Status and remediation plan

Verification: 617 tests pass; clippy -D warnings clean (host +
wasm32 check); fmt clean; doc clean
2026-09-06 19:35:18 +00:00
glm-5.3-flash f8dad9dbc8 feat(review 006 Unit 2): additive OperationSpec.description disclosed via discovery (E-02)
- OperationSpec gains description: Option<String> (builder
  with_description, defaults None; no struct-literal construction
  sites exist, so additive by construction)
- spec_to_json_pub emits description when set; rebuild_spec_for
  parses it back — the field survives from_call discovery and
  op/register announcement (same round-trip pattern as
  resource_id_path / publish_schema)
- services/list and the local-ops half of services/list-peers emit
  description when set; output-schema docs on both listing specs and
  operation_spec_schema advertise the field
- Tests: builder/default, emit/omit, listing emission, schema
  disclosure, schema-doc presence, round-trip + absent-stays-absent
  (8 new; 616 total)
- Docs: review 006 Unit 2 marked IMPLEMENTED; ADR-047 §6 amendment
  records the E-02 discovery decision (listing enrichment lands, the
  channel/resources/subscribe half stays deferred); OQ-40 gains the
  load-bearing note; operation-registry.md struct + listing docs;
  CHANGELOG

Verification: cargo test (616 pass), clippy -D warnings (host +
wasm32), fmt --check, cargo doc --no-deps, wasm32 check — all clean
2026-09-06 19:26:51 +00:00
glm-5.3-flash 2586c3b217 feat(review 006 Unit 1): channel-open establishment phase + typed client error (E-01, N-1)
Implements ADR-049 Unit 1 — the open-op wrapper gains an awaited,
bounded establishment phase, and the client stops erasing the error.

- OpenEstablisher hook + Establishment/EstablishmentError types:
  register_openable_with_establisher awaits the establisher bounded
  (earlier of dispatch deadline and per-registration timeout, else
  ESTABLISHMENT_TIMEOUT = 10s) after allocation, before the reply and
  before the pump handler is spawned (ADR-049 §1/§2). Implementation
  note: the establisher takes (input, auth) only — the channel's
  yield-once BiStream belongs exclusively to the pump handler
  (amendment recorded in ADR-049).
- Establishment failure: teardown_channel + opener-ledger take +
  policy.on_close un-increment (allocation and teardown balance;
  the ledger take is the atomic gate, ADR-047 §7), reply
  channel:open_failed with details {reason, message} — reason ∈
  dial_failed / unknown_resource / resource_shortage / handler_error
  / timeout (ADR-049 §3). SSH contract consumer-visible: a failed
  open never returns a channel_id.
- register_openable unchanged (no establisher = always-OK; existing
  registrations compile and behave identically — compat gate test).
- ChannelClient::open_channel returns ChannelOpenError (breaking at
  0.5.0): CallFailed { error: CallError } carries the wire error
  verbatim (establishment_reason() branches on details.reason);
  MissingChannelId / AdoptFailed cover the local-only shapes
  (ADR-049 §4, review 006 N-1).
- Tests cover all four verification gates from the review: e2e
  establisher failure through a real channels connection (typed
  reason + no-channel + ledger un-increment), bounded timeout,
  no-establisher compat, establisher-success pump round-trip; plus
  reason-vocabulary mapping and bound arithmetic.
- Bump to 0.5.0 (open_channel error-type change is semver-relevant).

Verification: cargo test (608 passed), clippy --all-targets -D
warnings, fmt --check, doc --no-deps, wasm32 check — all clean.
2026-09-06 19:00:04 +00:00
glm-5.3-flash 48ceeba55c docs: ADR-049 channel-open establishment phase; verify review 006
Verify review 006's findings against source at 88e3f5e (E-01..E-04
all confirmed; E-02 cost corrected — OperationSpec has no description
field, four touchpoints) and file three additional findings from the
same sweep (N-1 client error-type gap, N-2 write-only early-arrival
counter, N-3 pump-panic posture).

ADR-049 resolves E-01 + N-1: split-hook OpenEstablisher awaited
bounded by the open-op wrapper (restoring ADR-047 §3's "channel
plan" shape), teardown + typed channel:open_failed reply on
establishment failure, ChannelClient::open_channel typed error.
Review 006 gains the post-verification remediation plan and verdict
appendix.

Verification: cargo test (597 passed), cargo doc --no-deps clean.
2026-09-06 10:57:20 +00:00
glm-5.3-flash 88e3f5e9c3 docs: review 006 — channel-open establishment gap (from alktunnels phase 0)
Design review from the alktunnels Phase 0 research pass, verified
against tree a22b2b8 (0.4.1). Findings numbered E-01..E-04:

- E-01 [major] — the open op cannot fail after allocation: the
  wrapper replies {channel_id} the moment the OpenHandler is spawned;
  establishment failures (params-valid-but-rejected, backend lookup
  failure, target dial failure) present to the consumer as a
  successful open followed by an instant, indistinguishable clean
  EOF (implicit-EOF mux path + unified poll_read EOF arms). SSH
  semantics (RFC 4254 §5.1 open-failure reply with reason codes;
  channel never exists opener-side), SOCKS5 reply codes, and
  udpgw's opaque ERR bit (counterexample) surveyed in
  alktunnels/docs/research/ssh-socks5-survey.md. alktty's in-band
  error-frame mechanism (send_negotiation_error, 0x00-peek) is the
  per-crate workaround this upstream establisher obsoletes for the
  channels path. Proposed shape: an awaited establishment hook
  (OpenEstablisher) or await-and-inspect OpenHandler, tearing down on
  failure and replying channel:open_failed with SSH-four reason codes
  in ADR-016 details. Remediation sketch + verification gates
  included.
- E-02 [minor] — services/list discloses no per-op metadata; OQ-40
  (channel/resources/subscribe) becomes load-bearing for the first
  time via the alktunnels discovery resolution (OQ-TN-08).
- E-03 [minor] — OpenHandler-exit vs channel/close teardown race is
  benign (ledger take is the gate) but the let _ = discard at
  operations.rs:509 is silent; recommend log-or-comment.
- E-04 [minor] — early-arrival park cap (64) is an observable bound
  for push-first producers under slow adopters; no change requested,
  filed so the constraint is visible to the next consumer.

Non-findings recorded: open-op ACL path complete across all three
dispatch entry points; input_schema enforcement covers open params;
EOF arms unified; channel_open marker + resource_id_path wire
round-trip intact; opener-ledger decrement atomic at every call site.

alkcall tests: 597 passed (docs-only change; baseline check).
2026-09-06 09:54:13 +00:00
16 changed files with 2348 additions and 32 deletions

No files matched your search

+49
View File
@@ -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
View File
@@ -27,7 +27,7 @@ dependencies = [
[[package]]
name = "alkcall"
version = "0.4.1"
version = "0.5.0"
dependencies = [
"async-trait",
"bytes",
+1 -1
View File
@@ -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"
+42 -2
View File
@@ -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
+27
View File
@@ -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).
+1 -1
View File
@@ -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
+9
View File
@@ -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
View File
@@ -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,
&registry,
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,
&registry,
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`
+59
View File
@@ -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
View File
@@ -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,
&registry,
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,
&registry,
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)");
}
}
+62
View File
@@ -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
View File
@@ -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(&registry));
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(&registry));
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();
+38
View File
@@ -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(