11 Commits
Author SHA1 Message Date
glm-5.3-flash d22cecd317 docs: ADR-049/050 in indexes; changelog link refs; ported-ADR range 001..048
- architecture README: index rows for ADR-049 (establishment phase) and
  ADR-050 (pump_bidi); stale ADR-001..045 range labels corrected
- channel-operations.md: design-decision table + references list the
  two new ADRs
- CHANGELOG: missing link references for 0.4.0/0.4.1/0.5.0/0.6.0
- AGENTS.md: ADR range 001..050
2026-09-07 09:03:10 +00:00
glm-5.3-flash d25deb5eaf docs(review 007): mark R-01/R-02/R-03 resolved; record the two sketch deviations
Status → Resolved with the three unit commits; both deviations from
the review's sketches (typed-opaque ChannelPlan over Option<Value>;
(u64, u64) over io::Result) recorded up top with their rationale, as
the ADRs carry them.
2026-09-07 08:48:06 +00:00
glm-5.3-flash 9c6fec17ca feat(review 007 Unit 3): pump_bidi two-pump helper (R-03, ADR-050)
Extracts the two-pump data-plane helper alknet ADR-078 deferred until
the shapes converged (they have: alktunnels POC pump_halves + alktty's
channels session). Purely additive.

- channels::pump::pump_bidi(channel, peer_read, peer_write) -> (u64, u64):
  two joined pumps, shutdown-on-completion per direction; copy counts
  for observability. The channel side is a single AsyncRead +
  AsyncWrite value (the accept_bi BiStream); the peer side takes split
  halves — the establisher's natural dial result (into_split).
- Return (u64, u64), not the review sketch's io::Result<(u64, u64)>:
  both pumps swallow copy errors by contract (mid-stream error =
  abrupt close, no error channel mid-stream per ADR-049 §6), so an
  Err state would be dead code. Deviation recorded in ADR-050.
- alktty's three-pump session does not fit (exit future as a third
  signal) and stays as-is, per the review's scope.
- Tests reproduce the POC's two-pump semantics through the helper:
  bidirectional flow with exact counts, EOF-from-one-side completes
  the other's shutdown (clean EOF at the far end), dead-source =
  EOF-shaped teardown.
- ADR-050 records the decision, deviations, and two-way door type.

Verification: cargo test (625 passed, +2), clippy -D warnings, fmt
--check, doc clean, test --all-features clean, wasm32-unknown-unknown
check clean.
2026-09-07 08:47:46 +00:00
glm-5.3-flash 8f122b0a38 feat(review 007 Unit 2): birth-teardown telemetry — unaccepted-stream debug hint (R-02)
The optional hardening half of R-02 (the doc notes landed with Unit 1,
dc4ad2b). The wrapper's teardown task now observes whether the pump
handler ever accepted the channel's stream:

- ChannelBidiStreamSource gains a shared acceptance flag
  (with_accepted_flag / accepted()); channel_source_with_accepted_flag
  threads it from run_open_wrapper.
- On handler exit without accept, the teardown task logs a debug!
  naming the contract: "the returned JoinHandle must track the
  data-plane lifetime ... await pumps inline" — the telemetry hint for
  the teardown-at-birth shape the alktunnels POC hit empirically.
- Telemetry only: no behavior change; the benign no-data close still
  tears down identically.

Verification: cargo test (623 passed, +2: accepted-flag flip on the
yield-once accept; accepting vs non-accepting handler teardown
behavior), clippy -D warnings, fmt.
2026-09-07 08:45:09 +00:00
glm-5.3-flash dc4ad2bc6d feat(review 007 Unit 1): Establishment carries the channel plan (R-01) + lifetime doc (R-02)
Implements ADR-049 amendment 2 — the reserved Establishment payload is
filled, and the OpenHandler lifetime contract is documented.

- Establishment { plan: Option<ChannelPlan> } with ChannelPlan =
  Arc<dyn Any + Send + Sync>: typed-opaque, because the payload an
  establisher hands the pump handler is a live handle (dialed socket,
  TTY handle), not JSON — the review's Option<Value> sketch could not
  satisfy its own verification gate. #[non_exhaustive] keeps a future
  carrier change from being another break. Construction:
  Establishment::new(plan) / Establishment::default().
- OpenHandler gains the plan parameter:
  Fn(Value, Option<ChannelPlan>, Connection, AuthContext) ->
  JoinHandle<()>. Separate parameter (not merged into input) — a
  typed payload cannot ride the JSON input; no schema collision.
  Wire surface unchanged: the plan is process-local (establisher ->
  wrapper -> handler).
- run_open_wrapper threads establishment.plan to the handler; None
  when no establisher is registered. Kills the alktunnels-POC
  side-channel handoff (resource-keyed slot + poll loop) whose
  concurrent same-resource race is now unreachable — each open's
  establisher result flows to its own handler.
- Lifetime contract documented (R-02, doc-only half): the returned
  JoinHandle must track the data-plane lifetime — the wrapper awaits
  it and its completion triggers teardown; early return = teardown
  at birth. Noted on the OpenHandler type docs and both registration
  entry points.
- Breaking at 0.6.0 (the point of landing it before alktunnels
  Phase 1): Ok(Establishment {}) sites become
  Ok(Establishment::default()) mechanically.

Verification: cargo test (621 passed, +4: plan-flows-to-handler,
concurrent same-resource opens get distinct plans, no-establisher
None plan, Establishment construction), clippy -D warnings, fmt
--check, doc clean, test --all-features clean.
2026-09-07 08:43:17 +00:00
glm-5.3-flash 6590ab005f docs: review 007 — establishment follow-ups (from the alktunnels UDP POC)
Findings filed to prevent a second fix->publish->update-dependents
cycle; each was reached by building working code against 0.5.0:

- R-01 [major] — Establishment is payloadless but the channel plan
  is exactly what establishers need to hand to the pump handler
  (ADR-049's own 'reserved for a channel plan' note). Costs verified
  in two consumers: alktunnels POC side-channel handoff (same-resource
  opens race the slot), alktty forced to keep backend allocate
  post-open in-band (allocate_failed stays an in-band frame — the
  shape ADR-049 eliminates, alive one layer down). Ask: fill the
  reserved field (plan: Option<Value>, process-local, wire unchanged)
  in a 0.6.0 sweep; the break is mechanical (Establishment::default).
- R-02 [minor] — the OpenHandler JoinHandle lifetime contract is
  undocumented and load-bearing: the wrapper's await of the returned
  handle IS the teardown trigger; a handler that returns before its
  pumps finish tears the channel down at birth (the POC found this
  empirically — every tunnel EOF'd instantly). Doc note + ADR-049
  amendment; optional debug warning.
- R-03 [minor, optional] — the ADR-078 two-pump helper's convergence
  test is satisfied (POC gives both shapes); extracting pump_bidi
  now is additive (no break) and pins the contract upstream. Decide
  in the same 0.6 sweep.

Non-findings: reverse-flow (-R) needs no upstream mechanism
(from_connection_with_serving + serving-side allocation, verified by
trace); establisher receives registry-validated input; typed
establishment errors complete on the wire; early-arrival cap
unchanged; EstablishmentError reason set sufficient.

Goal stated in the review: 0.6 is the last breaking sweep forced by
known work. alkcall tests: 617 passed (docs-only change; baseline
check).
2026-09-07 08:18:11 +00:00
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
24 changed files with 3543 additions and 61 deletions

No files matched your search

+3 -2
View File
@@ -215,13 +215,14 @@ cargo semver-checks check-release # before a release
## Architecture Context
- `docs/architecture/` — the authoritative spec. Read it before
non-trivial changes. ADRs are numbered 001..048; OQs (open questions)
non-trivial changes. ADRs are numbered 001..050; OQs (open questions)
track resolved/deferred decisions.
- This crate unifies `alknet-call` and `alknet-channels` from the
alknet mono-repo (`/workspace/@alkdev/alknet`). The source
architecture docs were ported from
`/workspace/@alkdev/alknet/docs/architecture/` and renumbered as
alkcall ADRs (001..048; ADR-048 was authored in this crate). The ALPN
alkcall ADRs (001..048; ADR-049 and later were authored in this
crate). The ALPN
strings (`alk/call`,
`alk/channels`) are wire-stable going forward (renamed from
`alknet/` to `alk/` in v0.1.1, before the first published consumer).
+105
View File
@@ -4,6 +4,107 @@ 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.6.0] - 2026-09-07
The establishment follow-ups sweep (review 007): the `Establishment`
plan payload lands (R-01), the `OpenHandler` lifetime contract is
documented (R-02), and the two-pump helper is extracted (R-03). One
breaking change (below); the wire surface is unchanged.
### Added
- **`pump_bidi` (review 007 R-03 / ADR-050).** The two-pump helper —
`channels::pump_bidi(channel, peer_read, peer_write)` — pumps a
channel stream (`BiStream`-shaped) against a peer's split
read/write halves: one pump per direction, each shutting down the
opposite sink on completion (alknet ADR-078's shutdown-on-
completion), joined, returning `(u64, u64)` copy counts. Errors
are EOF-shaped by design (the POC + ADR-078 semantics — no `Err`
state; the review sketch's `io::Result` return was dead code).
The helper pins the contract in one place instead of three
(alktunnels POC `pump_halves`, alktty's channels session, the next
consumer). Purely additive.
### Changed
- **`Establishment` carries the channel plan (review 007 R-01 —
breaking at 0.6.0).** The reserved field is filled:
`Establishment { plan: Option<ChannelPlan> }` with
`ChannelPlan = Arc<dyn Any + Send + Sync>` — typed-opaque, because
the payload an establisher hands the pump handler is a live handle
(a dialed socket, a TTY handle), not JSON. The wrapper threads
`establishment.plan` to the `OpenHandler`'s new second parameter
(`Fn(Value, Option<ChannelPlan>, Connection, AuthContext) ->
JoinHandle<()>`); `None` when no establisher is registered or it
returned `Establishment::default()`. Process-local: establisher →
wrapper → handler; nothing new crosses the transport. This kills
the side-channel handoff the alktunnels POC shipped (resource-keyed
slot + poll loop) with its concurrent same-resource race — each
open's establisher result flows to its own handler. Migration:
`Ok(Establishment {})` → `Ok(Establishment::default())` (or
`Establishment::new(handle)` to deliver a handle); handler closures
gain a `_plan` (or `plan`) parameter.
- **`OpenHandler` lifetime contract documented (review 007 R-02).**
Doc-only semantics note on the `OpenHandler` type and the
registration entry points: the returned `JoinHandle` must track the
data-plane lifetime — the wrapper awaits it and its completion
triggers channel teardown (drop of the demux sender = EOF to the
handler's read half); a handler that returns before its pumps
finish tears the channel down at birth (await pumps inline, never
spawn-and-forget). Plus a `debug!` telemetry line in
`run_open_wrapper` when a handler exits without having accepted the
channel's `BiStream` (the birth-teardown hint).
## [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
@@ -330,6 +431,10 @@ Vendored core types (`Connection`, `ProtocolHandler`, `BiStream`,
(ADR-046), the channels protocol with openable-ALPNs-as-operations
(ADR-047), and the `ChannelClient` transport-agnostic client.
[0.6.0]: https://git.alk.dev/alkdev/alkcall/releases/tag/v0.6.0
[0.5.0]: https://git.alk.dev/alkdev/alkcall/releases/tag/v0.5.0
[0.4.1]: https://git.alk.dev/alkdev/alkcall/releases/tag/v0.4.1
[0.4.0]: https://git.alk.dev/alkdev/alkcall/releases/tag/v0.4.0
[0.3.1]: https://git.alk.dev/alkdev/alkcall/releases/tag/v0.3.1
[0.3.0]: https://git.alk.dev/alkdev/alkcall/releases/tag/v0.3.0
[0.2.0]: https://git.alk.dev/alkdev/alkcall/releases/tag/v0.2.0
Generated
+1 -1
View File
@@ -27,7 +27,7 @@ dependencies = [
[[package]]
name = "alkcall"
version = "0.4.1"
version = "0.6.0"
dependencies = [
"async-trait",
"bytes",
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "alkcall"
version = "0.4.1"
version = "0.6.0"
edition = "2021"
rust-version = "1.85"
license = "MIT OR Apache-2.0"
+5 -3
View File
@@ -13,7 +13,7 @@ This crate unifies `alknet-call` and `alknet-channels` from the alknet
mono-repo, plus the vendored core types formerly in `alknet-core`. The
source architecture docs were ported from
`/workspace/@alkdev/alknet/docs/architecture/` and renumbered as alkcall
ADRs (ADR-001..045). The ALPN strings (`alk/call`, `alk/channels`)
ADRs (ADR-001..050). The ALPN strings (`alk/call`, `alk/channels`)
are wire-stable and unchanged — see ADR-004.
## Documents
@@ -101,6 +101,8 @@ are wire-stable and unchanged — see ADR-004.
| [046](decisions/046-publish-operation-type-and-handler-kind-sink.md) | Publish Operation Type and HandlerKind::Sink | `OperationType::Pub` (producer→consumer streaming); `SinkHandler` + `HandlerKind::Sink`; `call.published` wire event; `invoke_sink()` dispatch; `Subscription` renamed to `Sub` |
| [047](decisions/047-openable-alpns-are-operations.md) | Openable ALPNs Are Operations | `channel/open` dissolves into per-ALPN ops `channels/<alpn>/sub`/`pub`; `channel_open` marker on `OperationSpec`; `ChannelCore` wrapper; extension-trait `ChannelOperationEnv`; connection-owner allocates `channel_id`; opener ledger (Gap 2 fix); ALPNs are call apps |
| [048](decisions/048-dispatch-spine-gateway-module.md) | Dispatch Spine (feature-gated `gateway` module) | `alkcall::gateway` behind the `gateway` feature; `GatewayDispatch` invoke spine (deadline knob, re-rooted context) + `schema_disclosure_denial` (FORBIDDEN for ACL deny, spec-404 for Internal); promoted from alkhttp for hub/spoke reuse |
| [049](decisions/049-channel-open-establishment-phase.md) | Channel-Open Establishment Phase | `OpenEstablisher` + `register_openable_with_establisher` (awaited, bounded); typed `channel:open_failed` with `details.reason`; `Establishment.plan` (`ChannelPlan`) threaded to the `OpenHandler` (amendment 2); the `JoinHandle` data-plane lifetime contract |
| [050](decisions/050-pump-bidi-two-pump-helper.md) | `pump_bidi` Two-Pump Helper | `channels::pump_bidi` — shutdown-on-completion two-pump data plane pinned in one place (alknet ADR-078) |
## Relevant Open Questions
@@ -272,7 +274,7 @@ feature flags) or in the downstream alknet crate.
- `@alkdev/alknet: docs/architecture/` — the source architecture docs
these were ported from (renumbered from alknet ADR-001..094 to alkcall
ADR-001..045)
ADR-001..048; ADR-049 and later were authored in this crate)
- `@alkdev/alktype` — the binary struct engine; compiles BAST documents
(e.g. `chunk-header.bast.json`) into readers/writers/validators
- `@alkdev/pubsub` — the TypeScript EventEnvelope prior art the call
@@ -288,5 +290,5 @@ feature flags) or in the downstream alknet crate.
> TTY's wire format, ADR-082 for alknet-tls, ADR-086 for endpoint types).
> These are ADRs for sibling crates that are not part of alkcall. They
> retain their alknet numbering (052, 082, 086, etc.) — any ADR number
> outside the alkcall range 001..045 is an alknet source ADR, found at
> outside the alkcall range 001..048 is an alknet source ADR, found at
> `/workspace/@alkdev/alknet/docs/architecture/decisions/`.
+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
+32
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
@@ -443,6 +470,8 @@ All design decisions are documented as ADRs in [decisions/](decisions/).
| [035](decisions/035-channels-pure-channel-multiplexing.md) | Pure Channel Multiplexing | No stream_types; handler owns sub-mux |
| [021](decisions/021-streaming-handler-for-subscriptions.md) | StreamingHandler | The machinery `channel/resources/subscribe` uses |
| [046](decisions/046-publish-operation-type-and-handler-kind-sink.md) | Pub Operation Type | The `Pub`/`Sub` primitives the per-ALPN open ops build on |
| [049](decisions/049-channel-open-establishment-phase.md) | Channel-Open Establishment Phase | The establisher hook; typed `channel:open_failed`; `Establishment.plan` (amendment 2); the `JoinHandle` lifetime contract |
| [050](decisions/050-pump-bidi-two-pump-helper.md) | `pump_bidi` Two-Pump Helper | The data-plane helper openable-ALPN handlers await inline |
| [026](decisions/026-forwarded-for-identity.md) | Forwarded-For Identity | The auth chain for hub-relayed opens (and why the cap is per direct-caller, not per `forwarded_for`) |
| [011](decisions/011-dynamic-resource-ownership-for-runtime-spawned-resources.md) | Dynamic Resource Ownership | The parallel — a channel slot is a resource, the cap is a quota check; `resource_id_path` works again under per-ALPN ops |
@@ -450,6 +479,9 @@ All design decisions are documented as ADRs in [decisions/](decisions/).
- ADR-047: openable ALPNs are operations (the unifying ADR — per-ALPN
open ops, `channel_open` marker, `ChannelCore` wrapper, opener ledger)
- ADR-049: channel-open establishment phase (the establisher hook,
`channel:open_failed`, the plan payload)
- ADR-050: `pump_bidi` (the two-pump helper)
- ADR-037: channel lifecycle operations (amended by ADR-047)
- ADR-041: per-identity channel cap (amended by ADR-047 §7 — opener
ledger, every teardown path)
@@ -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,408 @@
# 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).
## Amendment 2 (plan payload, 2026-09-07 — review 007 R-01/R-02)
Review 007 (from the alktunnels UDP POC) filed two follow-ups on the
establishment surface; both landed in alkcall 0.6.0.
**1. `Establishment` carries the channel plan (R-01).** §1 reserved
the payload ("today: nothing") and the wrapper consulted only
success/failure — so an establisher whose backend produces a handle
(a dialed socket, a TTY allocation) had to cross it to the pump
handler through a per-crate side channel. The alktunnels POC shipped
a resource-keyed slot + poll loop whose concurrent same-resource race
is unfixable within that shape; alktty documented the same wall
(backend `allocate` cannot cross, so failure classes stayed in-band —
the phantom-channel shape ADR-049 removed, alive one layer down).
The plan is now real: `Establishment { plan: Option<ChannelPlan> }`
with `ChannelPlan = Arc<dyn Any + Send + Sync>` — **typed-opaque, not
`serde_json::Value`**. The review's `Option<Value>` sketch could not
satisfy its own verification gate ("establisher dials, `plan` carries
the handle"): the payloads establishers actually hand off are live
handles with no JSON representation. The establisher and the
`OpenHandler` agree on the concrete type; alkcall never inspects it.
The wrapper threads `establishment.plan` to the handler's new second
parameter (`OpenHandler = Fn(Value, Option<ChannelPlan>, Connection,
AuthContext) -> JoinHandle<()>`); the separate-parameter shape wins
over merging into `input` because a typed payload cannot ride the
JSON input without a downcast-side registry and the reserved-key
collision the review already anticipated. The plan is process-local
(establisher → wrapper → handler on the producing side); the wire
surface is unchanged — nothing crosses the transport that isn't
already the open op's input. `#[non_exhaustive]` on `Establishment`
keeps a future carrier change from being another breaking release.
Construction is `Establishment::new(plan)` /
`Establishment::default()`; the 0.5.0 `Ok(Establishment {})` sites
break mechanically at 0.6.0, which is the point of landing this now
(before alktunnels Phase 1 ships the side-channel shape into a real
crate and the payload lands later anyway as a second break).
**2. The `OpenHandler` lifetime contract is documented (R-02).** The
wrapper awaits the returned `JoinHandle` and its completion triggers
teardown — so the handle must track the data-plane lifetime: a
handler that returns before its pumps finish tears the channel down
at birth (the POC's first pump implementation hit exactly this: every
tunnel connected then instantly EOF'd). The contract was implemented
but never documented; the type docs now state it ("await the pumps
inline, never spawn-and-forget and return early") on `OpenHandler`
and the registration entry points, plus a `debug!` telemetry line in
`run_open_wrapper` when a handler exits without having accepted the
channel's `BiStream` (the birth-teardown hint; the accept is
observable in-process via the yield-once source). §6's pinned
EOF-shaped panic semantics are unchanged.
@@ -0,0 +1,105 @@
# ADR-050: `pump_bidi` — the two-pump helper, extracted upstream
## Status
Accepted — implemented in alkcall 0.6.0 (review 007 Unit 3 / R-03).
Pins alknet ADR-078's two-pump contract in one place.
## Context
alknet ADR-078 defined the two-pump data-plane contract for forwarding
channels (one pump per direction; each pump shuts the opposite sink
down on completion — `try_join!` alone deadlocks) and deferred helper
extraction until a second two-pump consumer existed and the shapes
converged. Review 007 (R-03) supplied both halves of that test from
the alktunnels UDP POC:
- **Producer side:** the POC's `pump_halves` — two `tokio::io::copy`
pumps over split channel-vs-substrate halves, each shutting down
the opposite sink on completion, joined.
- **Consumer side:** `TunnelSession::take_halves` + the assembly
layer's copy — the same shape modulo channel side.
alktty's channels session already implements the shape (its
`drive_session_pre_negotiated` awaited inline); alkhttp's ferry does
not pump. The shapes converged; the review asked this sweep to either
extract the helper or record the decision not to — an un-extracted
helper would surface during alktunnels Phase 1 as the same
fix→publish→update treadmill for a purely additive change.
## Decision
Extract the helper as `channels::pump_bidi`:
```rust
pub async fn pump_bidi<C, R, W>(channel: C, peer_read: R, mut peer_write: W) -> (u64, u64)
where
C: AsyncRead + AsyncWrite + Unpin,
R: AsyncRead + Unpin,
W: AsyncWrite + Unpin,
```
Two pumps, joined: `channel → peer_write` and `peer_read → channel`.
On each pump's completion the opposite sink is shut down
(shutdown-on-completion — the half-close the contract specifies).
Returns the two copy counts for observability.
Three deliberate deviations from the review's sketch, all
semantics-first:
1. **Return `(u64, u64)`, not `io::Result<(u64, u64)>`.** There is no
meaningful `Err` state: both pumps swallow copy errors by contract
— a mid-stream error is indistinguishable from an abrupt close
(no error channel exists mid-stream, per ADR-049 §6's pinned
EOF-shaped semantics), and shutdown-of-the-opposite-sink runs
either way. A returned `Err` would be dead code; the counts
(with `unwrap_or(0)` on error) are the observability.
2. **Peer side takes split halves (`peer_read: R`, `peer_write: W`),
not one `AsyncRead + AsyncWrite` value.** The dominant producer
shape is the establisher's dial result: `TcpStream::into_split`
halves (tunnels), pty pair halves (TTY). Taking halves matches the
POC (`pump_halves`) and avoids forcing consumers to re-join halves
just to satisfy a bound. The channel side stays a single value —
that is what the pump handler's `accept_bi` yields.
3. **Not `async fn` over `AsyncWriteExt` bounds on both sides.** The
channel side needs `tokio::io::split` (one value, two pumps), so
its bound is `AsyncRead + AsyncWrite + Unpin`; the `AsyncWriteExt`
bound the sketch had on every parameter is the extension trait
only where shutdown is called.
alktty's three-pump session does not fit the helper (the exit future
is a third signal) and stays as-is — the helper serves the two-pump
shape, exactly as the review scoped.
## Compatibility
Purely additive: a new pub fn in `channels::pump`, re-exported in the
module docs; no existing surface changes. Consumers may adopt it at
their own pace (alktunnels Phase 1 gets it for free; alktty's
channels session can migrate later if it wants; the POC's shape is
now canonical).
## Door type
**Two-way.** A helper function is trivially removable or re-shapable;
no wire surface, no trait, no state.
## Verification gate
The POC's two-pump semantics reproduced through the helper, as unit
tests in `src/channels/pump.rs`: data flows both directions with
exact copy counts; EOF from one side completes the other's
shutdown-on-completion (the far end observes clean EOF, not an
error); a dead source is EOF-shaped teardown, not an `Err`.
## References
- Review 007 R-03 (the extraction ask + the convergence evidence)
- alknet ADR-078 (the two-pump contract; shutdown-on-completion)
- alktunnels `poc-summary.md` §Issues Surfaced #4 (the producer-side
shape; `pump_halves` in the POC's producer.rs)
- alktty `src/channels.rs` (the channels session shape; three-pump
carve-out)
- ADR-049 amendment 2 (R-02 — the lifetime contract the helper's
callers must respect: await `pump_bidi(..).await` inline inside the
handler task)
+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).
@@ -0,0 +1,343 @@
# Review 007 — Establishment Follow-Ups (from the alktunnels UDP POC)
## Status
Resolved — all three units landed in alkcall 0.6.0 (2026-09-07):
Unit 1 `dc4ad2b` (R-01 + R-02's doc notes), Unit 2 `8f122b0` (R-02's
telemetry), Unit 3 `9c6fec1` (R-03 + ADR-050). Two deviations from
the sketches below, both recorded in the ADRs: (1) R-01's plan is
**typed-opaque** (`ChannelPlan = Arc<dyn Any + Send + Sync>`), not
`Option<Value>` — the sketch could not satisfy this review's own
verification gate, because the payloads establishers hand off are
live handles (dialed sockets, TTY handles) with no JSON
representation; (2) R-03's helper returns `(u64, u64)`, not
`io::Result<(u64, u64)>` — both pumps swallow copy errors by the
ADR-078 contract (mid-stream error = abrupt close), so an `Err` state
would be dead code. Original findings below, retained as filed
(verified against tree `36e74cd`, 0.5.0).
Open — findings filed from the alktunnels UDP POC pass
(`alktunnels-udp-poc`, 2026-09-06; summary at
`/workspace/@alkdev/alktunnels/docs/research/poc-summary.md`). This
review exists to prevent a second fix→publish→update-dependents cycle:
every finding below is either (a) guaranteed to be needed by
alktunnels Phase 1 implementation, or (b) a documented contract gap
that will bite the next `OpenHandler` author the way it bit the POC.
Nothing here is speculative — each finding was reached by writing
working code against 0.5.0 and finding the shape insufficient.
Findings continue the review numbering with prefix `R` (001–006 used
P/C/R/A/B/C/D/F/G/E — each review numbers independently).
> **Post-remediation note (2026-09-07):** the "Open —" paragraph
> above is the original filing state, retained for the record. All
> units landed; see Status above and ADR-049 amendment 2 + ADR-050.
## Scope
The ADR-049 establishment surface (`OpenEstablisher`, `Establishment`,
`EstablishmentError`, `register_openable_with_establisher`) and the
`OpenHandler` contract, reviewed from the alktunnels UDP POC's
consumer/producer implementations. Everything verified against source
at tree `36e74cd` (0.5.0). Cross-references: alkcall review 006
(E-01..E-04), alkcall ADR-049, alktunnels `poc-summary.md`
(§Issues Surfaced), alktty's channels establisher
(`make_tty_establisher`).
```
Verified against: alkcall 36e74cd (0.5.0)
Reading list: src/channels/operations.rs (OpenEstablisher, Establishment,
run_open_wrapper, teardown_failed_channel), src/channels/client.rs
(ChannelOpenError::establishment_reason), src/channels/manager.rs
(teardown_channel, adopt_channel), docs/architecture/decisions/
047 + 049, and the consumers: alktty src/channels.rs (make_tty_
establisher / make_tty_open_handler), alkhttp src/websocket/
upgrade.rs (OpenableAlpn ferry), alktunnels-udp-poc src/producer.rs
```
## Severity legend
Same scale as reviews 004–006.
---
# Part A — The findings
## R-01 [major] — `Establishment` is payloadless, but the channel plan is exactly what establishers need to hand to the pump handler
**Verified:** YES — by building against the API. The ADR-049 §1 text
itself anticipates this: *"Reserved for a channel plan — today the
wrapper consults only success/failure, so `()` carries no data"*
(`src/channels/operations.rs:341-345`, the empty
`pub struct Establishment {}`).
### The mechanism
The establisher is pre-data-plane: it cannot see the channel's
yield-once `BiStream` (ADR-049 amendment), so anything it
establishes — a dialed socket, an allocated handle — must cross to
the pump handler through a side channel the ALPN crate invents. The
alktunnels POC's workaround is a resource-keyed
`Mutex<HashMap<String, SubstrateHandle>>` + a poll-loop `take`
(`producer.rs::HandleHandoff`) that works only because the wrapper
guarantees establisher-before-handler ordering. The costs:
1. **Concurrent same-resource opens race the slot.** Two opens of the
same resource: the second establisher's `deliver` overwrites the
first handler's not-yet-taken handle. The POC documents this as a
simplification; a real crate cannot.
2. **alktty hit the same wall and documented it as a limitation.**
`alktty/src/channels.rs:253-257`: the establisher cannot carry the
allocated `TtyHandle` across, so *backend allocation* stays
post-open in the pump handler (`allocate_failed` remains an
in-band negotiation error frame — ADR-010) — exactly the
phantom-channel-ish shape ADR-049 exists to eliminate, still
alive one layer down. TTY's establisher can only validate/lookup/
ownership-check; the one thing that can actually fail with a
runtime error (allocate) is unreachable from the typed error path.
3. **Every ALPN crate pays the side-channel tax again.** Handoffs,
slots, poll loops or oneshots — per crate, per resource key, all
with the same concurrency caveat.
### The ask
Give `Establishment` its payload — the channel plan ADR-047 §3
originally described ("return a 'channel plan'... the wrapper consults
its result"). Minimal shape:
```rust
pub struct Establishment {
/// ALPN-defined plan data — the dialed handle, an allocation
/// ticket, whatever the pump phase needs. Opaque to alkcall.
pub plan: Option<Value>,
}
```
or, keeping it typed at the boundary:
```rust
pub struct Establishment { pub plan: Option<Value> }
```
with the wrapper passing `establishment.plan` (or `Null`) to the
`OpenHandler`'s `input` — e.g. as a well-known key the pump handler
reads, or as a third callback parameter. The wire surface is
**unchanged** (the plan is process-local: establisher → wrapper →
handler in the same process on the producing side; nothing crosses
the transport that isn't already the open op's input). Consumers'
`Ok(Establishment {})` construction sites (alktty has two; alkhttp's
ferry has none) break mechanically at 0.6 — `Establishment::new(plan)`
/ `Establishment::default()` make the migration one-liners.
**The break is the point of doing this now:** 0.5.0 published
yesterday with `Establishment {}` documented as reserved. The next
consumer (alktunnels) needs the payload in Phase 1 — implementing
Phase 1 without it means shipping the POC's side-channel handoff into
the real crate, with its same-resource race, and then the payload
lands later anyway as *another* breaking release. Filling the
reserved field now is the cheap moment; the alternative is paying the
breaking change twice.
### Severity
[major] — the capability gap is structural for any establisher whose
backend produces a handle (tunnels: dial; TTY: allocate), not a
polish item. Not [critical] because workarounds exist (the POC proves
one), but every workaround carries the same-resource race or forces
failure classes back into per-crate in-band frames — the exact cost
ADR-049 was written to remove.
---
## R-02 [minor] — The `OpenHandler` `JoinHandle` lifetime contract is undocumented and load-bearing (found empirically by the POC)
The POC's first pump implementation returned a wrapper task that
spawned the pump as a nested fire-and-forget task. The result:
**every tunnel connected then instantly EOF'd with zero bytes** — the
wrapper awaited the (already-complete) handler task, tore the channel
down at birth, and both pumps saw immediate EOF. Everything upstream
(establisher, open reply, channel routing) looked healthy; the bug
was purely in the handler's return-value semantics.
The contract, as implemented: the wrapper awaits the returned
`JoinHandle` and *that completion is the teardown trigger*
(`run_open_wrapper`'s spawned task: `let _ = raw_task.await;`
→ `teardown_channel(id)` → drop the demux sender → EOF to the
handler's read half). Therefore:
- **The returned `JoinHandle` must track the data-plane lifetime.** A
handler that returns before its pumps finish tears the channel down
at birth. The pump must be awaited *inline* inside the handler task
(`let _ = pump_halves(...).await;`), not spawned-and-forgotten.
- Half-open semantics fall out of this correctly (one pump finishing
shuts down the opposite sink per ADR-078; the handler's task
completes when both pumps finish — which is when teardown *should*
happen). The contract is right; only its documentation is missing.
The current type docs (`operations.rs:325-339`) say the handle is
"recorded for teardown (abort on `channel/close` / connection drop)"
— abort semantics — but never say **early return = teardown-at-birth**.
alktty got this right by accident of shape (`drive_session_pre_negotiated`
awaited inline in its single spawned task — `channels.rs:355-395`);
the POC got it wrong the natural way. The next consumer will write
the wrong shape too, because the natural reading of "spawn your
protocol and return the `JoinHandle`" is a task-spawner, not a
lifecycle promise.
**The ask:** doc note on `OpenHandler` and
`ChannelCore::register_openable*` — "the returned `JoinHandle` must
track the data-plane lifetime: the wrapper awaits it and tears the
channel down on completion; return a task that runs the handler to
completion, never a spawner that exits early" — plus the ADR-049
amendment (§1 or §6) recording the semantics. Optional hardening (not
required): a `debug!` in `run_open_wrapper` when the awaited handler
exits without the channel's `BiStream` having been accepted (a
telemetry hint for the birth-teardown pattern; the accept is
observable in-process). Doc-only; no break.
---
## R-03 [minor, optional unit] — The two-pump helper's convergence test is satisfied; extracting it now removes the last reason for a later sweep
alknet ADR-078 deferred helper extraction until a second two-pump
consumer existed and the shapes converged. The POC supplies both
halves of that test:
- **Producer side:** the alktunnels POC's `pump_halves` — two
`tokio::io::copy` pumps over split channel-vs-substrate halves,
each pump shutting down the opposite sink on completion, joined.
- **Consumer side:** `TunnelSession::take_halves` + the assembly
layer's copy — the same shape modulo channel side.
The shapes converged. The helper (`pump_bidi` or similar) is a
candidate for **this same 0.6 sweep** — additive (a new pub fn in
`channels` or `core`), zero breaking change, and it pins the
ADR-078 contract in one place instead of three. The natural signature
follows the POC:
```rust
pub async fn pump_bidi<A, B>(a: A, b: B) -> io::Result<(u64, u64)>
where
A: AsyncRead + AsyncWriteExt + Send + Unpin,
B: AsyncRead + AsyncWriteExt + Send + Unpin,
```
(join'd two-pump with shutdown-on-completion; returns the copy
counts for observability). alktty's three-pump session does not fit
it (the exit future is a third signal) — that is fine; the helper
serves the two-pump shape, TTY stays as-is.
**The ask:** decide in this sweep — either extract in 0.6 (additive;
recommended, since alktunnels Phase 1 will implement the pattern
anyway and an upstream helper makes the third consumer free), or
record the explicit decision to keep it per-crate. Do not leave it
half-decided: an un-extracted helper is not itself a break, but
finding this out during Phase 1 would be the same
fix→publish→update treadmill for a purely additive change.
---
# Part B — Verified non-issues (bounded so the next review doesn't re-check)
- **The reverse-flow (`-R`) story needs no upstream change.** A hub
proxy opening a channel *toward* a worker is the worker serving its
own open op on the connect side — supported by
`ChannelClient::from_connection_with_serving` (ADR-022 §2 both-sides
semantics) + `ChannelOperations::register_on` on the serving
registry + the wrapper allocating on the serving side (ADR-047 §5
"the side that holds the `ChannelManager` allocates" — both sides
do, with odd/even split per ADR-047 §5). Verified by trace; no API
gap. alktunnels' reverse-flow POC (OQ-TN-10 #2) will exercise this
end-to-end, but no upstream mechanism is missing.
- **The establisher receives registry-validated input** —
`invoke_streaming`'s `input_schema` check runs before the wrapper
(`registration.rs:397`); the establisher's `input.clone()`
(`operations.rs:718`) is post-schema. Confirmed; no gap.
- **Typed establishment errors are complete on the wire** —
`establishment_error_to_call_error` carries both `reason` and
`message` in `details` (`operations.rs:646-653`); the client-side
`establishment_reason()` branches on it (`client.rs:75-83`); alktty's
open spec declares the matching `ErrorDefinition` (ADR-016). The
vocabulary needs no extension for tunnels (`dial_failed` /
`unknown_resource` / `resource_shortage` / `handler_error` /
`timeout` cover every tunnel establisher failure class — verified by
the POC's three typed-error tests).
- **The early-arrival park cap (64)** — reviewed in 006 E-04; the POC
treated it as a sizing constraint, not a defect. Unchanged.
- **`EstablishmentError`'s reason set is sufficient** — the POC's
establisher used three of four reasons; `handler_error` covers
params-parse failures (the TTY pattern). No new reason needed.
---
# Remediation plan
## Unit 1 — `Establishment` carries the channel plan (R-01)
1. ADR-049 amendment (or ADR-050-adjacent amendment): `Establishment`
gains `plan: Option<Value>` (ALPN-opaque); the wrapper threads
`plan` into the pump handler (third `OpenHandler` parameter, or
merged into `input` under a reserved key — decide in the ADR; the
separate-parameter shape avoids colliding with the input schema).
2. Bump to 0.6.0 (breaks `Ok(Establishment {})` construction sites —
alktty, two sites; mechanical `Establishment::default()` or
`Establishment { plan: None }`).
3. Migration: alktty can then move backend `allocate` into its
establisher (killing the in-band `allocate_failed` frame on the
channels path — ADR-010's reason for that frame disappears);
alktunnels Phase 1 uses `plan` for the dialed handle; alkhttp's
ferry passes `Option` through unchanged.
## Unit 2 — `OpenHandler` lifetime doc note (R-02)
Doc-only (type docs + ADR-049 amendment). Optional debug warning.
No break.
## Unit 3 (optional) — `pump_bidi` helper extraction (R-03)
Additive pub fn + tests. No break. Do it in the same 0.6 or record
the decision not to.
## What must NOT ride along
Nothing else. The remaining alktunnels Phase 1 needs (UDP codec ADR,
reverse-flow lifecycle, access-policy details) are alktunnels-local
spec work — the POC validated the mechanics and found no further
upstream gaps. The deliberate goal of this review: **0.6 is the last
breaking sweep forced by known work**; anything after it should be a
genuinely new discovery, not a known-dangling item.
## Verification gates
- Unit 1: end-to-end test — establisher dials, `plan` carries the
handle, pump handler receives it; concurrent same-resource opens
each get their own handle (the race the POC documented is
unreachable); alktty migration test (`allocate` in the establisher
→ `allocate_failed` is a call error, no in-band frame).
- Unit 2: doc gate only (`cargo doc` clean; the note present on both
the type and the registration entry points).
- Unit 3: helper test — the POC's two-pump semantics (EOF from one
direction completes the other; shutdown-on-completion) reproduced
through the helper.
## References
- alkcall review 006 (E-01 establishment gap; R-01 is the payload
half of E-01's own remediation — ADR-047 §3's "channel plan"
through the reserved `Establishment` field) and ADR-049 (the
establishment phase; §1's original sketch already described the
plan-carrying shape).
- alkcall ADR-047 §3 ("return a 'channel plan'" — the original shape
this review restores), ADR-022 §2 (both-sides serving — the
reverse-flow non-finding), ADR-047 §5 (odd/even allocation —
reverse-flow allocation on the serving side).
- alknet ADR-078 (the two-pump contract; R-03's helper pins it
upstream).
- alktunnels `docs/research/poc-summary.md` (§Issues Surfaced #1/#2 —
the empirical findings; §#4 — the convergence input for R-03) and
the POC crate (`/workspace/alktunnels-udp-poc`, 17 tests — the
empirical basis for every finding here).
- alktty `src/channels.rs:253-257` — the existing documentation of
the R-01 cost on the TTY side (`allocate` cannot cross; in-band
`allocate_failed` remains), the second consumer confirming the gap
is not tunnel-specific.
+326 -9
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))
}
@@ -834,7 +897,7 @@ mod tests {
Arc::clone(&policy) as Arc<dyn ChannelLifecyclePolicy>;
let policy_for_assert = Arc::clone(&policy);
let open_handler: OpenHandler = Arc::new(|_input, _channel_conn, _auth| {
let open_handler: OpenHandler = Arc::new(|_input, _plan, _channel_conn, _auth| {
tokio::spawn(async move {
// No-op: the channel's BiStream is available via
// `_channel_conn.accept_bi()` if the test wanted
@@ -1157,7 +1220,7 @@ mod tests {
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 open_handler: OpenHandler = Arc::new(move |_input, _plan, channel_conn, _auth| {
let data_tx = data_tx.clone();
tokio::spawn(async move {
let mut bidi = channel_conn.accept_bi().await.expect("accept_bi");
@@ -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, _plan, _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, _plan, 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`
@@ -2016,7 +2333,7 @@ mod tests {
let policy_for_hook: Arc<dyn ChannelLifecyclePolicy> =
Arc::clone(&policy) as Arc<dyn ChannelLifecyclePolicy>;
let open_handler: OpenHandler = Arc::new(|_input, _channel_conn, _auth| {
let open_handler: OpenHandler = Arc::new(|_input, _plan, _channel_conn, _auth| {
tokio::spawn(async move {
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
})
+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;
+3
View File
@@ -25,6 +25,8 @@
//! ALPN crates via `ChannelCore` (ADR-047 §3).
//! - [`policy`]: `ChannelLifecyclePolicy` + `PerIdentityChannelPolicy`
//! (ADR-041, amended by ADR-047 §7 — opener ledger).
//! - [`pump`]: `pump_bidi` — the two-pump data-plane helper (alknet
//! ADR-078, pinned upstream by ADR-050).
//! - [`client`]: `ChannelClient` — transport-agnostic
//! `from_connection` (ADR-043).
//! - [`self::env`]: `ChannelOperationEnv` extension trait (ADR-047 §4 —
@@ -39,6 +41,7 @@ pub mod manager;
pub mod mux;
pub mod operations;
pub mod policy;
pub mod pump;
pub mod reassembly;
pub mod source;
pub mod wire;
+1012 -35
View File
File diff suppressed because it is too large. Load diff
+204
View File
@@ -0,0 +1,204 @@
//! `pump_bidi` — the two-pump data-plane helper (alknet ADR-078, as
//! pinned upstream by ADR-050 — review 007 R-03).
//!
//! A forwarding channel's data plane is two unidirectional pumps: one
//! copies channel→peer, the other peer→channel. The contract is
//! **shutdown-on-completion**: when one pump's copy source EOFs, it
//! shuts down the opposite sink so the peer sees a clean half-close;
//! the helper completes when both pumps finish. Copy errors are
//! EOF-shaped by design (the POC and ADR-078 semantics — there is no
//! error channel mid-stream; a pump error is an abrupt-close signal,
//! and shutdown-of-the-opposite-sink is the same either way). The
//! helper returns the two copy counts for observability.
//!
//! The shapes converged across three consumers (the extraction test
//! ADR-078 set): the alktunnels POC's `pump_halves`, alktty's channels
//! session, and the assembly-layer copies. This helper pins the
//! two-pump shape in one place; alktty's three-pump session (an exit
//! future as a third signal) does not fit and stays as-is.
use tokio::io::{AsyncRead, AsyncWrite, AsyncWriteExt};
/// Pump a channel stream against a peer's split halves (ADR-078, as
/// pinned by ADR-050). Two pumps, joined:
///
/// - `channel → peer_write`: copy until the channel EOFs, then
/// `shutdown()` the peer write half.
/// - `peer_read → channel`: copy until the peer EOFs, then
/// `shutdown()` the channel.
///
/// Returns `(channel_to_peer, peer_to_channel)` copy counts.
/// Copy errors are treated as EOF (shutdown still runs) — the
/// EOF-shaped teardown the ADR-078 contract specifies; there is no
/// `Err` state because mid-stream errors are indistinguishable from
/// abrupt closes by design.
///
/// The generic bounds follow the POC: the channel side is a single
/// `AsyncRead + AsyncWrite` value (the `BiStream` the pump handler
/// received), the peer side arrives as split read/write halves (a
/// dialed socket, a TTY pty pair — `into_split` is the natural
/// establisher result).
///
/// Await the returned future inline inside the `OpenHandler`'s task —
/// the handler's returned `JoinHandle` must track the data-plane
/// lifetime (ADR-049 amendment 2, R-02).
pub async fn pump_bidi<C, R, W>(channel: C, peer_read: R, mut peer_write: W) -> (u64, u64)
where
C: AsyncRead + AsyncWrite + Unpin,
R: AsyncRead + Unpin,
W: AsyncWrite + Unpin,
{
let (mut c_read, mut c_write) = tokio::io::split(channel);
let mut peer_read = peer_read;
let c2p = async {
let n = tokio::io::copy(&mut c_read, &mut peer_write)
.await
.unwrap_or(0);
// Shutdown-on-completion: the channel side is done writing;
// the peer must see the half-close.
let _ = peer_write.shutdown().await;
n
};
let p2c = async {
let n = tokio::io::copy(&mut peer_read, &mut c_write)
.await
.unwrap_or(0);
let _ = c_write.shutdown().await;
n
};
tokio::join!(c2p, p2c)
}
#[cfg(test)]
mod tests {
use super::*;
use tokio::io::{AsyncReadExt, AsyncWriteExt, DuplexStream};
/// Model the topology precisely: the helper gets the channel
/// (`BiStream`-shaped) and the peer\'s split halves. The test
/// drives the *other* ends: `channel_peer` (what the channel\'s
/// remote writes/reads) and the substrate far ends. `duplex`
/// pairs: a write on one half arrives on its counterpart.
struct Topology {
/// The channel value handed to the helper.
channel: DuplexStream,
/// The channel\'s remote end (test-driven).
channel_peer: DuplexStream,
/// The peer read half handed to the helper (helper reads it).
peer_read: DuplexStream,
/// The peer write half handed to the helper (helper writes it).
peer_write: DuplexStream,
/// The peer write half\'s remote (test reads what the helper
/// copied).
peer_write_far: DuplexStream,
/// The peer read half\'s remote (test writes what the helper
/// will copy).
peer_read_far: DuplexStream,
}
fn topology() -> Topology {
// The channel\'s two ends.
let (channel, channel_peer) = tokio::io::duplex(64);
// The substrate: peer_read\'s far end, peer_write\'s far end.
let (peer_read, peer_read_far) = tokio::io::duplex(64);
let (peer_write_far, peer_write) = tokio::io::duplex(64);
Topology {
channel,
channel_peer,
peer_read,
peer_write,
peer_write_far,
peer_read_far,
}
}
/// The two-pump semantics test (review 007 Unit 3 gate): the
/// POC\'s `pump_halves` behavior through the helper — data flows
/// both directions, shutdown-on-completion on EOF, exact counts.
#[tokio::test]
async fn pump_bidi_moves_data_both_ways_and_shuts_down_on_completion() {
let mut topo = topology();
let pumped = tokio::spawn(pump_bidi(topo.channel, topo.peer_read, topo.peer_write));
// channel -> substrate: write at the channel\'s remote, read
// at the substrate\'s far end.
topo.channel_peer
.write_all(b"outbound")
.await
.expect("write");
let mut buf = [0u8; 8];
topo.peer_write_far
.read_exact(&mut buf)
.await
.expect("helper copied channel->peer");
assert_eq!(&buf, b"outbound");
// substrate -> channel: write at the substrate\'s far end,
// read at the channel\'s remote.
topo.peer_read_far
.write_all(b"inbound!")
.await
.expect("write");
let mut buf2 = [0u8; 8];
topo.channel_peer
.read_exact(&mut buf2)
.await
.expect("helper copied peer->channel");
assert_eq!(&buf2, b"inbound!");
// EOF from the substrate\'s far end: the helper shuts the
// channel down (shutdown-on-completion).
topo.peer_read_far.shutdown().await.expect("shutdown");
let mut eof_buf = Vec::new();
let n = topo
.channel_peer
.read_to_end(&mut eof_buf)
.await
.expect("clean EOF at the channel remote");
assert_eq!(n, 0, "the channel sees EOF after the substrate finished");
// Finish the other pump (channel remote drops) and collect
// the counts.
drop(topo.channel_peer);
let (c2p, p2c) = pumped
.await
.expect("helper completes when both pumps finish");
assert_eq!(c2p, 8, "channel->peer count");
assert_eq!(p2c, 8, "peer->channel count");
// The substrate\'s far write end observes the helper\'s
// shutdown: read returns 0 (EOF), not an error.
let mut eof_buf2 = Vec::new();
let n2 = topo
.peer_write_far
.read_to_end(&mut eof_buf2)
.await
.expect("eof read");
assert_eq!(n2, 0, "peer sees clean EOF after the channel side finished");
}
/// Copy errors are EOF-shaped: a pump whose source dies at birth
/// still lets the helper complete and shut the opposite sink down
/// (no Err state — ADR-078\'s abrupt-close semantics).
#[tokio::test]
async fn pump_bidi_treats_errors_as_eof_shaped_teardown() {
let topo = topology();
// Drop the substrate\'s far write end: the helper\'s peer_read
// source EOFs immediately; the channel side must still be shut
// down cleanly once the channel itself finishes.
drop(topo.peer_read_far);
let pumped = tokio::spawn(pump_bidi(topo.channel, topo.peer_read, topo.peer_write));
// The channel remote first receives what nothing sends — EOF
// after the p2c pump finishes; drop the remote to finish the
// c2p pump too.
drop(topo.channel_peer);
let _counts = pumped
.await
.expect("helper completes despite the dead source");
}
}
+61 -1
View File
@@ -10,6 +10,8 @@
//! to the mux write half via `BiStream::from_joined`.
use std::net::SocketAddr;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use async_trait::async_trait;
use parking_lot::Mutex;
@@ -28,9 +30,16 @@ use crate::core::types::{BiStream, BidiStreamSource, StreamError};
/// yield-once contract (ADR-007) and the "split never crosses a crate
/// boundary as part of a constructor" rule (ADR-009) — the join happens
/// here, in the `BidiStreamSource` impl, not per-handler.
///
/// The `accepted` flag (R-02 telemetry) records whether `accept_bi`
/// ever yielded: the open wrapper's teardown task checks it when the
/// pump handler exits and logs a `debug!` hint when the channel's
/// stream was never accepted — the "teardown at birth" shape a
/// fire-and-forget handler produces (review 007 R-02).
pub struct ChannelBidiStreamSource {
stream: Mutex<Option<BiStream>>,
remote_addr: Option<SocketAddr>,
accepted: Arc<AtomicBool>,
}
impl ChannelBidiStreamSource {
@@ -38,11 +47,29 @@ impl ChannelBidiStreamSource {
/// half joined to the mux write half). The handler will call
/// `accept_bi()` once and receive this stream.
pub fn new(stream: BiStream, remote_addr: Option<SocketAddr>) -> Self {
Self::with_accepted_flag(stream, remote_addr, Arc::new(AtomicBool::new(false)))
}
/// Construct with a shared acceptance flag — the open wrapper
/// passes its own so the teardown task can observe whether the
/// handler ever accepted the channel's stream (R-02 telemetry).
pub fn with_accepted_flag(
stream: BiStream,
remote_addr: Option<SocketAddr>,
accepted: Arc<AtomicBool>,
) -> Self {
Self {
stream: Mutex::new(Some(stream)),
remote_addr,
accepted,
}
}
/// Whether `accept_bi` ever yielded the stream. `false` after
/// construction; flips once on the first successful accept.
pub fn accepted(&self) -> bool {
self.accepted.load(Ordering::SeqCst)
}
}
#[async_trait]
@@ -50,7 +77,10 @@ impl BidiStreamSource for ChannelBidiStreamSource {
async fn accept_bi(&self) -> Result<BiStream, StreamError> {
let mut guard = self.stream.lock();
match guard.take() {
Some(stream) => Ok(stream),
Some(stream) => {
self.accepted.store(true, Ordering::SeqCst);
Ok(stream)
}
None => Err(StreamError::ConnectionClosed),
}
}
@@ -89,11 +119,26 @@ pub fn channel_source(
ChannelBidiStreamSource::new(stream, remote_addr)
}
/// [`channel_source`] with a shared acceptance flag (R-02 telemetry):
/// the open wrapper passes the flag, hands the source to the
/// handler's `Connection`, and the teardown task reads it to log the
/// "handler exited without accepting the channel's stream" hint.
pub fn channel_source_with_accepted_flag(
recv: super::reassembly::MpscRecvStream,
send: super::reassembly::MpscSendStream,
remote_addr: Option<SocketAddr>,
accepted: Arc<AtomicBool>,
) -> ChannelBidiStreamSource {
let stream = BiStream::from_joined(recv, send);
ChannelBidiStreamSource::with_accepted_flag(stream, remote_addr, accepted)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::channels::reassembly::{MpscRecvStream, MpscSendStream};
use bytes::Bytes;
use std::sync::Arc;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
fn make_pair() -> (
@@ -166,4 +211,19 @@ mod tests {
bidi.read_exact(&mut buf).await.expect("read");
assert_eq!(&buf, b"inbound");
}
#[tokio::test]
async fn accepted_flag_tracks_the_yield_once_accept() {
let (send, _mux_recv, _demux_send, recv) = make_pair();
let accepted = Arc::new(AtomicBool::new(false));
let source = channel_source_with_accepted_flag(recv, send, None, Arc::clone(&accepted));
assert!(!source.accepted(), "false before the handler accepts");
let _bidi = source.accept_bi().await.expect("first accept");
assert!(source.accepted(), "flips on the accept");
assert!(matches!(
source.accept_bi().await,
Err(StreamError::ConnectionClosed)
));
assert!(source.accepted(), "stays true after exhaustion");
}
}
+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!(
+2 -1
View File
@@ -28,7 +28,8 @@
//! - **Channels** ([`channels`]): the channels protocol — 8-byte chunk
//! wire format, demux/mux, `ChannelManager`, `ChannelsAdapter`,
//! `ChannelBidiStreamSource`, `ChannelOperations`,
//! `ChannelLifecyclePolicy`, `ChannelClient`. Channel 0 is
//! `ChannelLifecyclePolicy`, `ChannelClient`, `pump_bidi` (the
//! two-pump helper, ADR-050). Channel 0 is
//! pre-negotiated as `alk/call`; channels 1..N are opened via
//! per-ALPN open ops (`channels/<alpn>/sub`, `channels/<alpn>/pub`)
//! on channel 0 (ADR-047).
+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(