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.
This commit is contained in:
1 parent
48ceeba55c
commit
2586c3b217
9 files changed
+1105
-30
No files matched your search
@@ -4,6 +4,45 @@ All notable changes to this crate are documented here. The format is
|
||||
based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and
|
||||
this crate adheres to [Semantic Versioning](https://semver.org/).
|
||||
|
||||
## [0.5.0] - 2026-09-06
|
||||
|
||||
The channel-open establishment phase (ADR-049 — review 006 E-01 +
|
||||
N-1). Additive establishment machinery; one breaking change (the
|
||||
`ChannelClient::open_channel` error type, below).
|
||||
|
||||
### Added
|
||||
|
||||
- **Channel-open establishment phase (ADR-049 — review 006 E-01).**
|
||||
`ChannelCore::register_openable_with_establisher` registers a
|
||||
per-ALPN open op with an establisher hook (`OpenEstablisher`): an
|
||||
awaited establishment phase — validate params semantically, dial the
|
||||
backend — bounded by the dispatch deadline or the registration's
|
||||
timeout override, else `ESTABLISHMENT_TIMEOUT` (10s; the earlier of
|
||||
the two). On establisher failure (error or deadline) the wrapper
|
||||
tears down the just-allocated channel (demux sender, opener-ledger
|
||||
take, `policy.on_close` un-increment — the allocation and teardown
|
||||
balance) and replies `channel:open_failed` with
|
||||
`details: { reason, message }`, reason ∈ `dial_failed` /
|
||||
`unknown_resource` / `resource_shortage` / `handler_error` /
|
||||
`timeout`. The SSH contract holds consumer-visibly: a failed open
|
||||
never returns a `channel_id`. `register_openable` is unchanged
|
||||
(no establisher = always-OK, existing registrations compile and
|
||||
behave identically); the establisher takes `(input, auth)` — the
|
||||
channel's yield-once `BiStream` belongs exclusively to the pump
|
||||
handler. Additive wire surface (new error-code string + `details`
|
||||
shape).
|
||||
|
||||
### Changed
|
||||
|
||||
- **`ChannelClient::open_channel` returns a typed error (ADR-049 §4 —
|
||||
review 006 N-1, breaking at 0.5.0).** The open op's `CallError` is
|
||||
carried verbatim in `ChannelOpenError::CallFailed` instead of being
|
||||
flattened into a debug-formatted string, so consumers branch on
|
||||
`channel:open_failed`'s typed reason (`establishment_reason()`).
|
||||
Other variants: `MissingChannelId` (malformed success reply),
|
||||
`AdoptFailed` (local adoption failure). Mechanical for consumers —
|
||||
the `String` was a debug-formatting wrapper.
|
||||
|
||||
## [0.4.1] - 2026-09-05
|
||||
|
||||
Bug-fix release: chunks arriving for a not-yet-adopted channel are
|
||||
|
||||
Generated
+1
-1
@@ -27,7 +27,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "alkcall"
|
||||
version = "0.4.1"
|
||||
version = "0.5.0"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"bytes",
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "alkcall"
|
||||
version = "0.4.1"
|
||||
version = "0.5.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.85"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -68,12 +68,19 @@ impl ChannelClient {
|
||||
/// the `MpscRecvStream` (read half). The caller can build a
|
||||
/// `Connection` from these via `channel_source` and
|
||||
/// `Connection::from_source`.
|
||||
///
|
||||
/// The error is typed (ADR-049 §4 — review 006 N-1):
|
||||
/// `ChannelOpenError::CallFailed` carries the wire `CallError`
|
||||
/// verbatim (branch on `establishment_reason()` for
|
||||
/// `channel:open_failed`'s `details.reason`); `MissingChannelId`
|
||||
/// covers a malformed success reply; `AdoptFailed` covers local
|
||||
/// adoption failure.
|
||||
pub async fn open_channel(
|
||||
&self,
|
||||
operation_id: &str,
|
||||
input: Value,
|
||||
alpn: &str,
|
||||
) -> Result<(u32, MpscSendStream, MpscRecvStream), String>;
|
||||
) -> Result<(u32, MpscSendStream, MpscRecvStream), ChannelOpenError>;
|
||||
|
||||
/// Take the `CallConnection` — used by the consumer to register
|
||||
/// imported ops (`from_call`) on the connection's overlay. After
|
||||
@@ -88,7 +95,11 @@ impl ChannelClient {
|
||||
ADR-047 dissolved the generic `channel/open` operation into per-ALPN
|
||||
open ops (`channels/<alpn>/sub`, `channels/<alpn>/pub`). Each ALPN
|
||||
crate registers its own open op via `ChannelCore::register_openable`
|
||||
(ADR-047 §3, as amended 2026-08-13 — per-connection registration).
|
||||
(ADR-047 §3, as amended 2026-08-13 — per-connection registration),
|
||||
optionally with an establisher via
|
||||
`register_openable_with_establisher` (ADR-049 — the awaited,
|
||||
bounded establishment phase; establishment failure resolves
|
||||
`Err` on the client with `channel:open_failed` + `details.reason`).
|
||||
The `ChannelClient` calls these ops by name on channel 0 via
|
||||
`call_open_op`; `open_channel` wraps `call_open_op` + `adopt_channel`.
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -2,9 +2,11 @@
|
||||
|
||||
## Status
|
||||
|
||||
Accepted (amends ADR-047 §3 — the open-op wrapper gains an awaited
|
||||
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)
|
||||
006 E-01 and N-1
|
||||
|
||||
## Context
|
||||
|
||||
@@ -314,3 +316,41 @@ untouched, which is what keeps the split cheap to revise.
|
||||
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).
|
||||
@@ -19,8 +19,16 @@ 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). The original remediation sketch below is superseded by
|
||||
the "Remediation plan (post-verification)" section; the original sketch
|
||||
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). 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
|
||||
@@ -383,7 +391,13 @@ 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).
|
||||
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
|
||||
|
||||
+323
-6
@@ -25,13 +25,69 @@ use tokio::sync::Mutex;
|
||||
|
||||
use crate::core::types::{Connection, StreamError};
|
||||
use crate::protocol::connection::CallConnection;
|
||||
use crate::protocol::wire::ResponseEnvelope;
|
||||
use crate::protocol::wire::{CallError, ResponseEnvelope};
|
||||
use crate::registry::registration::OperationRegistry;
|
||||
|
||||
use super::manager::{ChannelManager, ChannelSide};
|
||||
use super::mux::MuxRunner;
|
||||
use super::reassembly::{MpscRecvStream, MpscSendStream};
|
||||
|
||||
/// The typed error from [`ChannelClient::open_channel`] (ADR-049 §4 —
|
||||
/// review 006 N-1). The open op's `CallError` — including the
|
||||
/// `channel:open_failed` code with `details: { reason, message }`
|
||||
/// (ADR-049 §3) — is carried verbatim instead of flattened into a
|
||||
/// string, so the consumer can branch on the establishment-failure
|
||||
/// reason (`dial_failed`, `unknown_resource`, `resource_shortage`,
|
||||
/// `handler_error`, `timeout`).
|
||||
#[derive(Debug, Clone, PartialEq, thiserror::Error)]
|
||||
pub enum ChannelOpenError {
|
||||
/// The open op failed — the `CallError` carries the wire code and
|
||||
/// typed `details` (`channel:open_failed` + `details.reason` per
|
||||
/// ADR-049 §3, `channel:too_many_channels`, `FORBIDDEN`, …).
|
||||
#[error("open op failed: {error:?}")]
|
||||
CallFailed { error: CallError },
|
||||
/// The open op succeeded but the response carried no `channel_id`
|
||||
/// (malformed responder).
|
||||
#[error("open op response missing channel_id")]
|
||||
MissingChannelId,
|
||||
/// The open op succeeded but local channel adoption failed (the
|
||||
/// per-connection cap, or an ID collision).
|
||||
#[error("adopt_channel failed: {0}")]
|
||||
AdoptFailed(#[from] super::manager::ManagerError),
|
||||
}
|
||||
|
||||
impl ChannelOpenError {
|
||||
/// The wire `CallError` behind this failure, when the open op
|
||||
/// itself failed. `None` for the local-only variants
|
||||
/// (`MissingChannelId` comes from a success envelope;
|
||||
/// `AdoptFailed` is post-reply).
|
||||
pub fn call_error(&self) -> Option<&CallError> {
|
||||
match self {
|
||||
ChannelOpenError::CallFailed { error } => Some(error),
|
||||
ChannelOpenError::MissingChannelId | ChannelOpenError::AdoptFailed(_) => None,
|
||||
}
|
||||
}
|
||||
|
||||
/// The establishment-failure reason code — `details.reason` of a
|
||||
/// `channel:open_failed` error (`dial_failed`, `unknown_resource`,
|
||||
/// `resource_shortage`, `handler_error`, `timeout`). `None` for
|
||||
/// any other failure shape.
|
||||
pub fn establishment_reason(&self) -> Option<&str> {
|
||||
match self {
|
||||
ChannelOpenError::CallFailed { error }
|
||||
if error.code == super::operations::CHANNEL_OPEN_FAILED =>
|
||||
{
|
||||
error
|
||||
.details
|
||||
.as_ref()
|
||||
.and_then(|d| d.get("reason"))
|
||||
.and_then(|r| r.as_str())
|
||||
}
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// The client-side handle for a channels connection. Constructed via
|
||||
/// [`ChannelClient::from_connection`] from an established
|
||||
/// `Connection`. Holds the `ChannelManager` (for the demux/mux state)
|
||||
@@ -224,26 +280,33 @@ impl ChannelClient {
|
||||
///
|
||||
/// `alpn` is the data-plane ALPN (e.g. `alk/tty`), used for
|
||||
/// observability in the local manager.
|
||||
///
|
||||
/// The error is typed (ADR-049 §4 — review 006 N-1): a failed open
|
||||
/// resolves [`ChannelOpenError::CallFailed`] carrying the
|
||||
/// wire `CallError` verbatim — branch on
|
||||
/// [`ChannelOpenError::establishment_reason`] for
|
||||
/// `channel:open_failed`'s reason code (`dial_failed`,
|
||||
/// `unknown_resource`, `resource_shortage`, `handler_error`,
|
||||
/// `timeout`).
|
||||
pub async fn open_channel(
|
||||
&self,
|
||||
operation_id: &str,
|
||||
input: Value,
|
||||
alpn: &str,
|
||||
) -> Result<(u32, MpscSendStream, MpscRecvStream), String> {
|
||||
) -> Result<(u32, MpscSendStream, MpscRecvStream), ChannelOpenError> {
|
||||
let response = self.call_open_op(operation_id, input).await;
|
||||
let out = response
|
||||
.result
|
||||
.map_err(|e| format!("open op failed: {e:?}"))?;
|
||||
.map_err(|error| ChannelOpenError::CallFailed { error })?;
|
||||
let channel_id = out
|
||||
.get("channel_id")
|
||||
.and_then(|v| v.as_u64())
|
||||
.ok_or_else(|| "open op response missing channel_id".to_string())?
|
||||
as u32;
|
||||
.ok_or(ChannelOpenError::MissingChannelId)? as u32;
|
||||
|
||||
self.manager
|
||||
.adopt_channel(channel_id, alpn, None)
|
||||
.await
|
||||
.map_err(|e| format!("adopt_channel failed: {e}"))
|
||||
.map_err(ChannelOpenError::AdoptFailed)
|
||||
.map(|(send, recv)| (channel_id, send, recv))
|
||||
}
|
||||
|
||||
@@ -1256,6 +1319,260 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
/// ADR-049 Unit 1 acceptance gate (review 006 E-01/N-1) — the
|
||||
/// establisher fails after channel allocation: the consumer's
|
||||
/// `open_channel` resolves `Err` carrying the wire `CallError`
|
||||
/// (`channel:open_failed` + `details.reason`), no `channel_id` was
|
||||
/// ever returned (the SSH "channel never exists opener-side"
|
||||
/// property, consumer-visible), and the accept side's manager has
|
||||
/// no channel afterward (teardown + ledger un-increment).
|
||||
#[tokio::test]
|
||||
async fn establisher_failure_resolves_typed_open_failed_on_the_client() {
|
||||
use crate::channels::operations::{
|
||||
ChannelCore, EstablishmentError, OpenEstablisher, OpenHandler,
|
||||
};
|
||||
use crate::channels::policy::PerIdentityChannelPolicy;
|
||||
use crate::registry::spec::ChannelOpenSpec;
|
||||
|
||||
let policy: Arc<PerIdentityChannelPolicy> = Arc::new(PerIdentityChannelPolicy::new(256));
|
||||
let policy_for_hook: Arc<dyn ChannelLifecyclePolicy> =
|
||||
Arc::clone(&policy) as Arc<dyn ChannelLifecyclePolicy>;
|
||||
let policy_for_assert = Arc::clone(&policy);
|
||||
|
||||
let open_handler: OpenHandler =
|
||||
Arc::new(|_input, _channel_conn, _auth| tokio::spawn(async {}));
|
||||
let establisher: OpenEstablisher = Arc::new(|_input, _auth| {
|
||||
Box::pin(async {
|
||||
Err(EstablishmentError::DialFailed {
|
||||
message: "target refused the connection".to_string(),
|
||||
})
|
||||
})
|
||||
});
|
||||
|
||||
let install_hook: crate::channels::adapter::InstallChannelZero =
|
||||
Arc::new(move |manager, channel0_conn, auth| {
|
||||
let open_handler = Arc::clone(&open_handler);
|
||||
let establisher = Arc::clone(&establisher);
|
||||
let policy = Arc::clone(&policy_for_hook);
|
||||
tokio::spawn(async move {
|
||||
let channel0_bidi = match channel0_conn.accept_bi().await {
|
||||
Ok(s) => s,
|
||||
Err(_) => return,
|
||||
};
|
||||
let (writer, reader) = split_single_stream(channel0_bidi);
|
||||
let core = ChannelCore::new(manager, policy);
|
||||
let registry = crate::registry::registration::OperationRegistry::new();
|
||||
let spec = OperationSpec::new(
|
||||
"channels/tty/sub",
|
||||
OperationType::Sub,
|
||||
Visibility::External,
|
||||
serde_json::json!({}),
|
||||
serde_json::json!({
|
||||
"type": "object",
|
||||
"properties": { "channel_id": { "type": "integer" } }
|
||||
}),
|
||||
vec![],
|
||||
AccessControl::default(),
|
||||
None,
|
||||
)
|
||||
.with_channel_open(ChannelOpenSpec::new("alk/tty"));
|
||||
core.register_openable_with_establisher(
|
||||
spec,
|
||||
Some(establisher),
|
||||
open_handler,
|
||||
®istry,
|
||||
auth.clone(),
|
||||
None,
|
||||
)
|
||||
.expect("register_openable_with_establisher");
|
||||
let registry = Arc::new(registry);
|
||||
let provider: Arc<dyn IdentityProvider> = Arc::new(NoopIdProvider);
|
||||
let call_connection = Arc::new(CallConnection::new_single_stream(
|
||||
channel0_conn,
|
||||
Arc::clone(&writer),
|
||||
));
|
||||
let dp = Dispatcher::new(registry, provider);
|
||||
dp.run_loop_single_stream(call_connection, reader, writer)
|
||||
.await;
|
||||
})
|
||||
});
|
||||
|
||||
let (client_end, server_end) = tokio::io::duplex(64 * 1024);
|
||||
let client_conn =
|
||||
Connection::from_bidi(client_end, b"alk/channels".to_vec(), Some(TEST_ADDR));
|
||||
let server_conn =
|
||||
Connection::from_bidi(server_end, b"alk/channels".to_vec(), Some(TEST_ADDR));
|
||||
|
||||
let adapter = ChannelsAdapter::new(install_hook, Arc::new(NoCap));
|
||||
let auth = AuthContext::anonymous(b"alk/channels");
|
||||
let _server_handle = tokio::spawn(async move {
|
||||
let _ = crate::core::types::ProtocolHandler::handle(&adapter, server_conn, &auth).await;
|
||||
});
|
||||
|
||||
let client = ChannelClient::from_connection(client_conn)
|
||||
.await
|
||||
.expect("channel client init");
|
||||
|
||||
let err = tokio::time::timeout(
|
||||
std::time::Duration::from_secs(10),
|
||||
client.open_channel(
|
||||
"channels/tty/sub",
|
||||
serde_json::json!({ "container": "abc" }),
|
||||
"alk/tty",
|
||||
),
|
||||
)
|
||||
.await
|
||||
.expect("open_channel timed out");
|
||||
let err = match err {
|
||||
Err(e) => e,
|
||||
Ok(_) => panic!("establisher failure must resolve Err through the client"),
|
||||
};
|
||||
|
||||
// N-1 gate: the typed error carries the wire CallError verbatim.
|
||||
let call_error = err.call_error().expect("CallFailed carries the CallError");
|
||||
assert_eq!(call_error.code, "channel:open_failed");
|
||||
assert_eq!(
|
||||
err.establishment_reason(),
|
||||
Some("dial_failed"),
|
||||
"consumer branches on the typed reason code"
|
||||
);
|
||||
|
||||
// The SSH property, consumer-visible: no channel_id was ever
|
||||
// returned, and the accept side holds no channel for the failed
|
||||
// open (channel 0 is the pre-negotiated call channel — always
|
||||
// present, never the failed open's).
|
||||
assert!(
|
||||
client.manager().channel_ids().into_iter().all(|id| id == 0),
|
||||
"no data channel exists for the failed open"
|
||||
);
|
||||
let anonymous = crate::core::auth::Identity {
|
||||
id: "anonymous".to_string(),
|
||||
scopes: vec![],
|
||||
resources: Default::default(),
|
||||
};
|
||||
assert_eq!(
|
||||
policy_for_assert.count_for(&anonymous),
|
||||
0,
|
||||
"ledger un-incremented on the accept side after establishment failure"
|
||||
);
|
||||
}
|
||||
|
||||
/// ADR-049 Unit 1 acceptance gate (companion to the E-01 gate) —
|
||||
/// establisher success through a real channels connection: the
|
||||
/// client's `open_channel` adopts and the pumps flow (the happy
|
||||
/// path is unchanged by the establishment phase).
|
||||
#[tokio::test]
|
||||
async fn establisher_success_opens_and_pumps_data_end_to_end() {
|
||||
use crate::channels::operations::{
|
||||
ChannelCore, Establishment, OpenEstablisher, OpenHandler,
|
||||
};
|
||||
use crate::channels::policy::NoCap;
|
||||
use crate::registry::spec::ChannelOpenSpec;
|
||||
use tokio::io::AsyncReadExt;
|
||||
|
||||
let (data_tx, mut data_rx) = tokio::sync::mpsc::channel::<Vec<u8>>(1);
|
||||
|
||||
let open_handler: OpenHandler = Arc::new(move |_input, channel_conn, _auth| {
|
||||
let data_tx = data_tx.clone();
|
||||
tokio::spawn(async move {
|
||||
let mut bidi = channel_conn.accept_bi().await.expect("accept_bi");
|
||||
let mut buf = [0u8; 4];
|
||||
bidi.read_exact(&mut buf).await.expect("read");
|
||||
data_tx.send(buf.to_vec()).await.expect("send to channel");
|
||||
})
|
||||
});
|
||||
let establisher: OpenEstablisher =
|
||||
Arc::new(|_input, _auth| Box::pin(async { Ok(Establishment::default()) }));
|
||||
|
||||
let establisher_for_hook = Arc::clone(&establisher);
|
||||
let open_handler_for_hook = Arc::clone(&open_handler);
|
||||
let install_hook: crate::channels::adapter::InstallChannelZero =
|
||||
Arc::new(move |manager, channel0_conn, auth| {
|
||||
let establisher = Arc::clone(&establisher_for_hook);
|
||||
let open_handler = Arc::clone(&open_handler_for_hook);
|
||||
tokio::spawn(async move {
|
||||
let channel0_bidi = match channel0_conn.accept_bi().await {
|
||||
Ok(s) => s,
|
||||
Err(_) => return,
|
||||
};
|
||||
let (writer, reader) = split_single_stream(channel0_bidi);
|
||||
let core = ChannelCore::new(manager, Arc::new(NoCap));
|
||||
let registry = crate::registry::registration::OperationRegistry::new();
|
||||
let spec = OperationSpec::new(
|
||||
"channels/tty/sub",
|
||||
OperationType::Sub,
|
||||
Visibility::External,
|
||||
serde_json::json!({
|
||||
"type": "object",
|
||||
"properties": { "container": { "type": "string" } },
|
||||
"required": ["container"]
|
||||
}),
|
||||
serde_json::json!({
|
||||
"type": "object",
|
||||
"properties": { "channel_id": { "type": "integer" } }
|
||||
}),
|
||||
vec![],
|
||||
AccessControl::default(),
|
||||
None,
|
||||
)
|
||||
.with_channel_open(ChannelOpenSpec::new("alk/tty"));
|
||||
core.register_openable_with_establisher(
|
||||
spec,
|
||||
Some(establisher),
|
||||
open_handler,
|
||||
®istry,
|
||||
auth.clone(),
|
||||
None,
|
||||
)
|
||||
.expect("register_openable_with_establisher");
|
||||
let registry = Arc::new(registry);
|
||||
let provider: Arc<dyn IdentityProvider> = Arc::new(NoopIdProvider);
|
||||
let call_connection = Arc::new(CallConnection::new_single_stream(
|
||||
channel0_conn,
|
||||
Arc::clone(&writer),
|
||||
));
|
||||
let dp = Dispatcher::new(registry, provider);
|
||||
dp.run_loop_single_stream(call_connection, reader, writer)
|
||||
.await;
|
||||
})
|
||||
});
|
||||
|
||||
let (client_end, server_end) = tokio::io::duplex(64 * 1024);
|
||||
let client_conn =
|
||||
Connection::from_bidi(client_end, b"alk/channels".to_vec(), Some(TEST_ADDR));
|
||||
let server_conn =
|
||||
Connection::from_bidi(server_end, b"alk/channels".to_vec(), Some(TEST_ADDR));
|
||||
|
||||
let adapter = ChannelsAdapter::new(install_hook, Arc::new(NoCap));
|
||||
let auth = AuthContext::anonymous(b"alk/channels");
|
||||
let _server_handle = tokio::spawn(async move {
|
||||
let _ = crate::core::types::ProtocolHandler::handle(&adapter, server_conn, &auth).await;
|
||||
});
|
||||
|
||||
let client = ChannelClient::from_connection(client_conn)
|
||||
.await
|
||||
.expect("channel client init");
|
||||
|
||||
let (channel_id, mut send, _recv) = client
|
||||
.open_channel(
|
||||
"channels/tty/sub",
|
||||
serde_json::json!({ "container": "abc" }),
|
||||
"alk/tty",
|
||||
)
|
||||
.await
|
||||
.expect("open_channel with establisher success");
|
||||
|
||||
assert!(channel_id > 0, "channel_id should be non-zero");
|
||||
send.write_all(b"ping").await.expect("write ping");
|
||||
drop(send);
|
||||
|
||||
let data = tokio::time::timeout(std::time::Duration::from_secs(5), data_rx.recv())
|
||||
.await
|
||||
.expect("timed out waiting for handler data")
|
||||
.expect("handler should receive data");
|
||||
assert_eq!(&data, b"ping", "handler received ping post-establishment");
|
||||
}
|
||||
|
||||
// --- review 004 Unit 3 acceptance gates (F-04 serving half) -----------
|
||||
|
||||
/// F-04 gate 1: hub→consumer call over an existing `ChannelClient`
|
||||
|
||||
+642
-15
@@ -14,7 +14,9 @@
|
||||
//! See `docs/architecture/channel-operations.md` for the spec.
|
||||
|
||||
use std::sync::Arc;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
use futures::future::BoxFuture;
|
||||
use serde_json::{json, Value};
|
||||
|
||||
use crate::core::auth::{AuthContext, Identity};
|
||||
@@ -295,7 +297,10 @@ fn make_resources_subscribe_handler(manager: ChannelManager) -> StreamingHandler
|
||||
}
|
||||
|
||||
/// `ChannelCore` — the channel machinery the ALPN crate's open-op
|
||||
/// wrapper uses (ADR-047 §3). Provides [`ChannelCore::register_openable`],
|
||||
/// wrapper uses (ADR-047 §3, as amended by ADR-049 — the wrapper
|
||||
/// gains an awaited, bounded establishment phase via
|
||||
/// [`ChannelCore::register_openable_with_establisher`]). Provides
|
||||
/// [`ChannelCore::register_openable`],
|
||||
/// which wraps the ALPN's open handler with channel-id allocation,
|
||||
/// `ChannelManager` integration, opener-ledger recording,
|
||||
/// `ChannelLifecyclePolicy` consultation, and teardown hooks.
|
||||
@@ -333,6 +338,85 @@ pub struct ChannelCore {
|
||||
pub type OpenHandler =
|
||||
Arc<dyn Fn(Value, Connection, AuthContext) -> tokio::task::JoinHandle<()> + Send + Sync>;
|
||||
|
||||
/// The establishment-phase result (ADR-049 §1). Reserved for a channel
|
||||
/// plan — today the wrapper consults only success/failure, so `()`
|
||||
/// carries no data.
|
||||
#[derive(Debug, Clone, Default)]
|
||||
pub struct Establishment {}
|
||||
|
||||
/// The establishment-failure reason code the wrapper puts in the open
|
||||
/// reply's `details.reason` (ADR-049 §3). The vocabulary maps 1:1 onto
|
||||
/// what an establisher can actually produce (per the SSH four;
|
||||
/// review 006's survey) plus the wrapper's own deadline.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
|
||||
pub enum EstablishmentError {
|
||||
/// The backend/target could not be reached or refused the
|
||||
/// connection.
|
||||
#[error("dial failed: {message}")]
|
||||
DialFailed { message: String },
|
||||
/// The requested resource does not exist.
|
||||
#[error("unknown resource: {message}")]
|
||||
UnknownResource { message: String },
|
||||
/// The backend is out of capacity (ports, fds, slots).
|
||||
#[error("resource shortage: {message}")]
|
||||
ResourceShortage { message: String },
|
||||
/// Establisher-internal failure not covered above.
|
||||
#[error("handler error: {message}")]
|
||||
HandlerError { message: String },
|
||||
}
|
||||
|
||||
impl EstablishmentError {
|
||||
/// The wire-visible reason code (`details.reason` in
|
||||
/// `channel:open_failed`).
|
||||
pub fn reason(&self) -> &'static str {
|
||||
match self {
|
||||
EstablishmentError::DialFailed { .. } => "dial_failed",
|
||||
EstablishmentError::UnknownResource { .. } => "unknown_resource",
|
||||
EstablishmentError::ResourceShortage { .. } => "resource_shortage",
|
||||
EstablishmentError::HandlerError { .. } => "handler_error",
|
||||
}
|
||||
}
|
||||
|
||||
/// The human-readable detail message.
|
||||
pub fn message(&self) -> &str {
|
||||
match self {
|
||||
EstablishmentError::DialFailed { message }
|
||||
| EstablishmentError::UnknownResource { message }
|
||||
| EstablishmentError::ResourceShortage { message }
|
||||
| EstablishmentError::HandlerError { message } => message,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// The establisher hook (ADR-049 §1): the awaited establishment phase
|
||||
/// of a channel open — validate params semantically, consult ownership,
|
||||
/// prepare/dial the backend. **Awaited by the open-op wrapper,
|
||||
/// bounded** (dispatch deadline or `ESTABLISHMENT_TIMEOUT`, whichever
|
||||
/// is earlier) — before the reply is written, before the pump handler
|
||||
/// is spawned. On failure the wrapper tears down the just-allocated
|
||||
/// channel and replies `channel:open_failed` with `details.reason`.
|
||||
///
|
||||
/// The establisher takes the open op's `input` (params) and the peer's
|
||||
/// `AuthContext`. It deliberately does **not** receive the channel's
|
||||
/// `Connection`: the channel `BiStream` is yield-once
|
||||
/// (`ChannelBidiStreamSource`) and belongs exclusively to the pump
|
||||
/// handler (the `OpenHandler`); the establisher is pre-data-plane (its
|
||||
/// dial targets the backend, not the channel). Amendment to the
|
||||
/// ADR-049 §1 signature sketch.
|
||||
pub type OpenEstablisher = Arc<
|
||||
dyn Fn(Value, AuthContext) -> BoxFuture<'static, Result<Establishment, EstablishmentError>>
|
||||
+ Send
|
||||
+ Sync,
|
||||
>;
|
||||
|
||||
/// The default bound on the establisher await (ADR-049 §2) when the
|
||||
/// `OperationContext` carries no deadline and the registration set no
|
||||
/// override.
|
||||
pub const ESTABLISHMENT_TIMEOUT: Duration = Duration::from_secs(10);
|
||||
|
||||
/// The wire-visible error code for establishment failure (ADR-049 §3).
|
||||
pub const CHANNEL_OPEN_FAILED: &str = "channel:open_failed";
|
||||
|
||||
/// The error code prefix for channel-open failures (ADR-047 §3). The
|
||||
/// wrapper maps `ChannelError` / `ManagerError` to `CallError` with
|
||||
/// codes prefixed by `channel:` so the initiator can branch on the
|
||||
@@ -374,10 +458,17 @@ impl ChannelCore {
|
||||
/// spawn the protocol on the channel's `BiStream`). This wrapper
|
||||
/// does the channel machinery: ACL check (run by the registry's
|
||||
/// `invoke`/`invoke_streaming`/`invoke_sink` before the wrapper) →
|
||||
/// `check_open(identity)` → `manager.open_channel(alpn, opener)`
|
||||
/// → spawn the ALPN handler on the channel's `Connection` →
|
||||
/// `check_open(identity)` → `manager.open_channel(alpn, opener)` →
|
||||
/// spawn the ALPN handler on the channel's `Connection` →
|
||||
/// respond with `{ "channel_id": <id> }`.
|
||||
///
|
||||
/// **No establisher**: the open op cannot fail after allocation —
|
||||
/// establishment work inside the handler is invisible to the open
|
||||
/// reply (review 006 E-01's subject). Use
|
||||
/// [`ChannelCore::register_openable_with_establisher`] when the
|
||||
/// establishment phase can fail and the failure belongs in the
|
||||
/// open reply (ADR-049).
|
||||
///
|
||||
/// The op is registered on the given `registry` (the connection
|
||||
/// overlay registry, Layer 2 per ADR-019 — this is the
|
||||
/// per-connection registration the §4 amendment blesses). The
|
||||
@@ -402,6 +493,38 @@ impl ChannelCore {
|
||||
open_handler: OpenHandler,
|
||||
registry: &OperationRegistry,
|
||||
auth: AuthContext,
|
||||
) -> Result<(), String> {
|
||||
self.register_openable_with_establisher(spec, None, open_handler, registry, auth, None)
|
||||
}
|
||||
|
||||
/// Register a per-ALPN open op with an establishment phase
|
||||
/// (ADR-049 §1/§5). The `establisher` is **awaited by the wrapper,
|
||||
/// bounded** — after allocation, before the reply and before the
|
||||
/// pump handler is spawned. On establisher failure (error or
|
||||
/// deadline) the wrapper tears down the just-allocated channel and
|
||||
/// replies `channel:open_failed` with
|
||||
/// `details: { reason, message }`; on success the pump handler is
|
||||
/// spawned exactly as [`ChannelCore::register_openable`] does.
|
||||
///
|
||||
/// `timeout` bounds the establisher when the dispatch carries no
|
||||
/// deadline (`None` = [`ESTABLISHMENT_TIMEOUT`], 10s). When the
|
||||
/// `OperationContext` carries a deadline, the effective bound is
|
||||
/// the earlier of the two (ADR-049 §2). The bound applies only to
|
||||
/// the establisher — the spawned pump's lifetime is governed by
|
||||
/// the existing teardown machinery, unchanged.
|
||||
///
|
||||
/// Backward compatible by construction: `establisher: None` behaves
|
||||
/// exactly like [`ChannelCore::register_openable`] (the
|
||||
/// no-establisher shape), and existing `OpenHandler`s compile
|
||||
/// unchanged.
|
||||
pub fn register_openable_with_establisher(
|
||||
&self,
|
||||
spec: OperationSpec,
|
||||
establisher: Option<OpenEstablisher>,
|
||||
open_handler: OpenHandler,
|
||||
registry: &OperationRegistry,
|
||||
auth: AuthContext,
|
||||
timeout: Option<Duration>,
|
||||
) -> Result<(), String> {
|
||||
let op_type = spec.op_type;
|
||||
let channel_open = spec.channel_open.clone();
|
||||
@@ -417,17 +540,22 @@ impl ChannelCore {
|
||||
.alpn
|
||||
.to_string();
|
||||
|
||||
let establisher = establisher.map(|establisher| EstablisherHook {
|
||||
establisher,
|
||||
timeout,
|
||||
});
|
||||
let manager = self.manager.clone();
|
||||
let policy = Arc::clone(&self.policy);
|
||||
let open_handler = Arc::clone(&open_handler);
|
||||
|
||||
let handler_kind = match op_type {
|
||||
OperationType::Query | OperationType::Mutation => HandlerKind::Once(
|
||||
make_open_handler_once(manager, policy, open_handler, auth, alpn),
|
||||
make_open_handler_once(manager, policy, establisher, open_handler, auth, alpn),
|
||||
),
|
||||
OperationType::Sub => HandlerKind::Stream(make_open_handler_stream(
|
||||
manager,
|
||||
policy,
|
||||
establisher,
|
||||
open_handler,
|
||||
auth,
|
||||
alpn,
|
||||
@@ -435,6 +563,7 @@ impl ChannelCore {
|
||||
OperationType::Pub => HandlerKind::Sink(make_open_handler_sink(
|
||||
manager,
|
||||
policy,
|
||||
establisher,
|
||||
open_handler,
|
||||
auth,
|
||||
alpn,
|
||||
@@ -467,22 +596,103 @@ fn opener_identity_from_context(ctx: &OperationContext) -> Identity {
|
||||
})
|
||||
}
|
||||
|
||||
/// The establisher hook + its per-registration bound override,
|
||||
/// threaded from `ChannelCore::register_openable_with_establisher`
|
||||
/// into the wrapper factories. `timeout: None` = `ESTABLISHMENT_TIMEOUT`.
|
||||
#[derive(Clone)]
|
||||
struct EstablisherHook {
|
||||
establisher: OpenEstablisher,
|
||||
timeout: Option<Duration>,
|
||||
}
|
||||
|
||||
/// The effective bound on the establisher await (ADR-049 §2): the
|
||||
/// earlier of the dispatch deadline (when the `OperationContext`
|
||||
/// carries one) and the establishment timeout (the per-registration
|
||||
/// override, else `ESTABLISHMENT_TIMEOUT`). A deadline already in the
|
||||
/// past yields a zero bound — the timeout fires immediately.
|
||||
fn establishment_bound(timeout_override: Option<Duration>, deadline: Option<Instant>) -> Duration {
|
||||
let timeout = timeout_override.unwrap_or(ESTABLISHMENT_TIMEOUT);
|
||||
deadline
|
||||
.map(|d| d.saturating_duration_since(Instant::now()).min(timeout))
|
||||
.unwrap_or(timeout)
|
||||
}
|
||||
|
||||
/// Tear down a just-allocated channel after establishment failure
|
||||
/// (ADR-049 §3): drop the demux sender (EOF to nothing — no handler
|
||||
/// was ever spawned), take the opener-ledger entry, and un-reserve the
|
||||
/// per-identity count (`policy.on_close` — the same un-increment path
|
||||
/// the allocation-failure arms run; the ledger `take` is the atomic
|
||||
/// gate, ADR-047 §7).
|
||||
fn teardown_failed_channel(
|
||||
manager: &ChannelManager,
|
||||
policy: &Arc<dyn ChannelLifecyclePolicy>,
|
||||
channel_id: u32,
|
||||
opener: &Identity,
|
||||
) {
|
||||
if let Err(e) = manager.teardown_channel(channel_id) {
|
||||
tracing::debug!(
|
||||
channel_id,
|
||||
error = %e,
|
||||
"open wrapper: establisher-failure teardown found no channel to remove"
|
||||
);
|
||||
}
|
||||
if manager.opener_ledger().take(channel_id).is_some() {
|
||||
policy.on_close(opener);
|
||||
}
|
||||
}
|
||||
|
||||
/// The `channel:open_failed` reply for an establisher error
|
||||
/// (ADR-049 §3): `details` carries the reason code + message.
|
||||
fn establishment_error_to_call_error(err: &EstablishmentError) -> CallError {
|
||||
CallError::new(
|
||||
CHANNEL_OPEN_FAILED,
|
||||
format!("channel establishment failed: {}", err.message()),
|
||||
false,
|
||||
)
|
||||
.with_details(json!({ "reason": err.reason(), "message": err.message() }))
|
||||
}
|
||||
|
||||
/// The `channel:open_failed` reply for a deadline expiry — same
|
||||
/// `details` shape, reason `timeout`.
|
||||
fn establishment_timeout_call_error() -> CallError {
|
||||
CallError::new(
|
||||
CHANNEL_OPEN_FAILED,
|
||||
"channel establishment exceeded the deadline",
|
||||
false,
|
||||
)
|
||||
.with_details(json!({
|
||||
"reason": "timeout",
|
||||
"message": "establishment exceeded the deadline",
|
||||
}))
|
||||
}
|
||||
|
||||
/// Run the open-op wrapper's shared post-ACL steps: `check_open` →
|
||||
/// `open_channel` → build the channel `Connection` → spawn the ALPN
|
||||
/// handler → return `{ channel_id }` on success, or a `CallError` on
|
||||
/// any failure. The ACL check has already run (by the registry's
|
||||
/// `invoke`/`invoke_streaming`/`invoke_sink` before the wrapper); this
|
||||
/// function is the `check_open` → `allocate` → `ledger` → `spawn` →
|
||||
/// `respond` tail of the ADR-047 §3 flow.
|
||||
/// `open_channel` → **await the establisher bounded** (ADR-049 §1/§2 —
|
||||
/// no-op when no establisher is registered) → build the channel
|
||||
/// `Connection` → spawn the ALPN handler → return `{ channel_id }` on
|
||||
/// success, or a `CallError` on any failure. The ACL check has already
|
||||
/// run (by the registry's `invoke`/`invoke_streaming`/`invoke_sink`
|
||||
/// before the wrapper); this function is the `check_open` →
|
||||
/// `allocate` → `establish` → `ledger` → `spawn` → `respond` tail of
|
||||
/// the ADR-047 §3 flow (as amended by ADR-049).
|
||||
///
|
||||
/// On establishment failure (error or deadline) the wrapper tears down
|
||||
/// the just-allocated channel (demux sender, opener-ledger entry,
|
||||
/// `policy.on_close` un-increment) and replies `channel:open_failed`
|
||||
/// with `details: { reason, message }` — the SSH "channel never exists
|
||||
/// opener-side" property: the consumer's open resolves `Err` and no
|
||||
/// `channel_id` was ever returned.
|
||||
///
|
||||
/// `input` is the open op's input (params), passed through to the
|
||||
/// ALPN's `OpenHandler` so it can validate params and prepare the
|
||||
/// backend. `opener_id` is the `PeerId` recorded in the opener ledger
|
||||
/// (the direct caller's identity id).
|
||||
/// establisher (when registered) and the ALPN's `OpenHandler`.
|
||||
/// `opener_id` is the `PeerId` recorded in the opener ledger (the
|
||||
/// direct caller's identity id). `opener_identity` is the identity the
|
||||
/// cap was reserved against (un-reserved on failure paths).
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
async fn run_open_wrapper(
|
||||
manager: &ChannelManager,
|
||||
policy: &Arc<dyn ChannelLifecyclePolicy>,
|
||||
establisher: Option<&EstablisherHook>,
|
||||
open_handler: &OpenHandler,
|
||||
auth: &AuthContext,
|
||||
alpn: &str,
|
||||
@@ -490,6 +700,7 @@ async fn run_open_wrapper(
|
||||
opener_id: String,
|
||||
opener_identity: Identity,
|
||||
request_id: String,
|
||||
deadline: Option<Instant>,
|
||||
) -> ResponseEnvelope {
|
||||
if let Err(channel_err) = policy.check_open(&opener_identity) {
|
||||
return ResponseEnvelope::error(request_id, map_channel_error_to_call_error(&channel_err));
|
||||
@@ -497,6 +708,39 @@ async fn run_open_wrapper(
|
||||
|
||||
let channel_id = match manager.open_channel(alpn, opener_id, None).await {
|
||||
Ok((id, send, recv)) => {
|
||||
// The establishment phase (ADR-049 §1): awaited bounded,
|
||||
// before the reply and before the pump handler is spawned.
|
||||
// No establisher registered = an always-OK establisher
|
||||
// (the pre-ADR-049 shape, unchanged).
|
||||
if let Some(hook) = establisher {
|
||||
let bound = establishment_bound(hook.timeout, deadline);
|
||||
let establishment =
|
||||
tokio::time::timeout(bound, (hook.establisher)(input.clone(), auth.clone()))
|
||||
.await;
|
||||
match establishment {
|
||||
Err(_elapsed) => {
|
||||
teardown_failed_channel(manager, policy, id, &opener_identity);
|
||||
return ResponseEnvelope::error(
|
||||
request_id,
|
||||
establishment_timeout_call_error(),
|
||||
);
|
||||
}
|
||||
Ok(Err(e)) => {
|
||||
tracing::debug!(
|
||||
channel_id = id,
|
||||
reason = e.reason(),
|
||||
"open wrapper: establisher failed; tearing down channel"
|
||||
);
|
||||
teardown_failed_channel(manager, policy, id, &opener_identity);
|
||||
return ResponseEnvelope::error(
|
||||
request_id,
|
||||
establishment_error_to_call_error(&e),
|
||||
);
|
||||
}
|
||||
Ok(Ok(_establishment)) => {}
|
||||
}
|
||||
}
|
||||
|
||||
let remote_addr = manager.remote_addr();
|
||||
let source = super::source::channel_source(recv, send, remote_addr);
|
||||
let channel_conn = Connection::from_source(source, alpn.as_bytes().to_vec());
|
||||
@@ -584,6 +828,7 @@ fn map_channel_error_to_call_error(err: &ChannelError) -> CallError {
|
||||
fn make_open_handler_once(
|
||||
manager: ChannelManager,
|
||||
policy: Arc<dyn ChannelLifecyclePolicy>,
|
||||
establisher: Option<EstablisherHook>,
|
||||
open_handler: OpenHandler,
|
||||
auth: AuthContext,
|
||||
alpn: String,
|
||||
@@ -591,6 +836,7 @@ fn make_open_handler_once(
|
||||
Arc::new(move |input: Value, ctx: OperationContext| {
|
||||
let manager = manager.clone();
|
||||
let policy = Arc::clone(&policy);
|
||||
let establisher = establisher.clone();
|
||||
let open_handler = Arc::clone(&open_handler);
|
||||
let auth = auth.clone();
|
||||
let alpn = alpn.clone();
|
||||
@@ -601,6 +847,7 @@ fn make_open_handler_once(
|
||||
run_open_wrapper(
|
||||
&manager,
|
||||
&policy,
|
||||
establisher.as_ref(),
|
||||
&open_handler,
|
||||
&auth,
|
||||
&alpn,
|
||||
@@ -608,6 +855,7 @@ fn make_open_handler_once(
|
||||
opener_id,
|
||||
opener_identity,
|
||||
request_id,
|
||||
ctx.deadline,
|
||||
)
|
||||
.await
|
||||
})
|
||||
@@ -625,6 +873,7 @@ fn make_open_handler_once(
|
||||
fn make_open_handler_stream(
|
||||
manager: ChannelManager,
|
||||
policy: Arc<dyn ChannelLifecyclePolicy>,
|
||||
establisher: Option<EstablisherHook>,
|
||||
open_handler: OpenHandler,
|
||||
auth: AuthContext,
|
||||
alpn: String,
|
||||
@@ -632,6 +881,7 @@ fn make_open_handler_stream(
|
||||
Arc::new(move |input: Value, ctx: OperationContext| {
|
||||
let manager = manager.clone();
|
||||
let policy = Arc::clone(&policy);
|
||||
let establisher = establisher.clone();
|
||||
let open_handler = Arc::clone(&open_handler);
|
||||
let auth = auth.clone();
|
||||
let alpn = alpn.clone();
|
||||
@@ -642,6 +892,7 @@ fn make_open_handler_stream(
|
||||
run_open_wrapper(
|
||||
&manager,
|
||||
&policy,
|
||||
establisher.as_ref(),
|
||||
&open_handler,
|
||||
&auth,
|
||||
&alpn,
|
||||
@@ -649,6 +900,7 @@ fn make_open_handler_stream(
|
||||
opener_id,
|
||||
opener_identity,
|
||||
request_id,
|
||||
ctx.deadline,
|
||||
)
|
||||
.await
|
||||
}))
|
||||
@@ -675,6 +927,7 @@ fn make_open_handler_stream(
|
||||
fn make_open_handler_sink(
|
||||
_manager: ChannelManager,
|
||||
_policy: Arc<dyn ChannelLifecyclePolicy>,
|
||||
_establisher: Option<EstablisherHook>,
|
||||
_open_handler: OpenHandler,
|
||||
_auth: AuthContext,
|
||||
_alpn: String,
|
||||
@@ -1039,6 +1292,7 @@ mod tests {
|
||||
let handler = make_open_handler_once(
|
||||
manager.clone(),
|
||||
policy,
|
||||
None,
|
||||
open_handler,
|
||||
auth,
|
||||
"alk/tty".to_string(),
|
||||
@@ -1069,6 +1323,7 @@ mod tests {
|
||||
let handler = make_open_handler_stream(
|
||||
manager.clone(),
|
||||
policy,
|
||||
None,
|
||||
open_handler,
|
||||
auth,
|
||||
"alk/tty".to_string(),
|
||||
@@ -1096,8 +1351,14 @@ mod tests {
|
||||
let policy = super::super::policy::default_policy();
|
||||
let open_handler: OpenHandler = Arc::new(|_input, _conn, _auth| tokio::spawn(async {}));
|
||||
let auth = AuthContext::anonymous(b"alk/call");
|
||||
let handler =
|
||||
make_open_handler_sink(manager, policy, open_handler, auth, "alk/tty".to_string());
|
||||
let handler = make_open_handler_sink(
|
||||
manager,
|
||||
policy,
|
||||
None,
|
||||
open_handler,
|
||||
auth,
|
||||
"alk/tty".to_string(),
|
||||
);
|
||||
let publish_stream: crate::registry::registration::PublishStream =
|
||||
Box::pin(futures::stream::empty());
|
||||
let env = handler(json!({}), test_context("open-sink-1"), publish_stream).await;
|
||||
@@ -1265,6 +1526,7 @@ mod tests {
|
||||
let env = run_open_wrapper(
|
||||
&manager,
|
||||
&policy,
|
||||
None,
|
||||
&open_handler,
|
||||
&auth,
|
||||
"alk/tty",
|
||||
@@ -1272,6 +1534,7 @@ mod tests {
|
||||
opener_id.id.clone(),
|
||||
opener_id,
|
||||
"req-cap-deny".to_string(),
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
match env.result {
|
||||
@@ -1279,4 +1542,368 @@ mod tests {
|
||||
Ok(_) => panic!("cap 0 should deny open"),
|
||||
}
|
||||
}
|
||||
|
||||
// --- ADR-049 Unit 1: the establishment phase -------------------------
|
||||
|
||||
fn failing_establisher() -> OpenEstablisher {
|
||||
Arc::new(|_input, _auth| {
|
||||
Box::pin(async {
|
||||
Err(EstablishmentError::DialFailed {
|
||||
message: "connection refused".to_string(),
|
||||
})
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
fn ok_establisher() -> OpenEstablisher {
|
||||
Arc::new(|_input, _auth| Box::pin(async { Ok(Establishment::default()) }))
|
||||
}
|
||||
|
||||
fn hook(establisher: OpenEstablisher, timeout: Option<Duration>) -> Option<EstablisherHook> {
|
||||
Some(EstablisherHook {
|
||||
establisher,
|
||||
timeout,
|
||||
})
|
||||
}
|
||||
|
||||
fn identity(id: &str) -> Identity {
|
||||
Identity {
|
||||
id: id.to_string(),
|
||||
scopes: vec![],
|
||||
resources: HashMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn run_open_wrapper_establisher_failure_tears_down_and_replies_open_failed() {
|
||||
let manager = make_manager().await;
|
||||
let concrete_policy = Arc::new(super::super::policy::PerIdentityChannelPolicy::new(8));
|
||||
let policy: Arc<dyn ChannelLifecyclePolicy> =
|
||||
Arc::clone(&concrete_policy) as Arc<dyn ChannelLifecyclePolicy>;
|
||||
let spawned = Arc::new(AtomicBool::new(false));
|
||||
let spawned_clone = Arc::clone(&spawned);
|
||||
let open_handler: OpenHandler = Arc::new(move |_input, _conn, _auth| {
|
||||
spawned_clone.store(true, Ordering::SeqCst);
|
||||
tokio::spawn(async {})
|
||||
});
|
||||
let auth = AuthContext::anonymous(b"alk/call");
|
||||
let env = run_open_wrapper(
|
||||
&manager,
|
||||
&policy,
|
||||
hook(failing_establisher(), None).as_ref(),
|
||||
&open_handler,
|
||||
&auth,
|
||||
"alk/tty",
|
||||
json!({}),
|
||||
"alice".to_string(),
|
||||
identity("alice"),
|
||||
"req-est-fail".to_string(),
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
match env.result {
|
||||
Err(e) => {
|
||||
assert_eq!(e.code, "channel:open_failed");
|
||||
assert!(!e.retryable);
|
||||
let details = e.details.expect("details carry the reason");
|
||||
assert_eq!(details["reason"], "dial_failed");
|
||||
assert_eq!(details["message"], "connection refused");
|
||||
}
|
||||
Ok(_) => panic!("establisher failure must fail the open"),
|
||||
}
|
||||
// The SSH "channel never exists opener-side" property: the
|
||||
// channel was torn down and no `channel_id` was ever returned.
|
||||
assert!(
|
||||
manager.channel_ids().is_empty(),
|
||||
"no channel survives a failed establishment"
|
||||
);
|
||||
assert_eq!(manager.open_count(), 0);
|
||||
// The per-identity cap reservation was un-incremented.
|
||||
assert_eq!(
|
||||
concrete_policy.count_for(&identity("alice")),
|
||||
0,
|
||||
"ledger take + policy.on_close restored the cap count"
|
||||
);
|
||||
// The pump handler was never spawned.
|
||||
assert!(
|
||||
!spawned.load(Ordering::SeqCst),
|
||||
"pump handler must not spawn on establishment failure"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn run_open_wrapper_establisher_timeout_replies_timeout_reason() {
|
||||
let manager = make_manager().await;
|
||||
let concrete_policy = Arc::new(super::super::policy::PerIdentityChannelPolicy::new(8));
|
||||
let policy: Arc<dyn ChannelLifecyclePolicy> =
|
||||
Arc::clone(&concrete_policy) as Arc<dyn ChannelLifecyclePolicy>;
|
||||
let open_handler: OpenHandler = Arc::new(|_input, _conn, _auth| tokio::spawn(async {}));
|
||||
let auth = AuthContext::anonymous(b"alk/call");
|
||||
let hanging: OpenEstablisher = Arc::new(|_input, _auth| {
|
||||
Box::pin(async {
|
||||
tokio::time::sleep(Duration::from_secs(5)).await;
|
||||
Ok(Establishment::default())
|
||||
})
|
||||
});
|
||||
let env = run_open_wrapper(
|
||||
&manager,
|
||||
&policy,
|
||||
hook(hanging, Some(Duration::from_millis(50))).as_ref(),
|
||||
&open_handler,
|
||||
&auth,
|
||||
"alk/tty",
|
||||
json!({}),
|
||||
"alice".to_string(),
|
||||
identity("alice"),
|
||||
"req-est-timeout".to_string(),
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
match env.result {
|
||||
Err(e) => {
|
||||
assert_eq!(e.code, "channel:open_failed");
|
||||
let details = e.details.expect("details carry the reason");
|
||||
assert_eq!(details["reason"], "timeout");
|
||||
}
|
||||
Ok(_) => panic!("a hung establisher must fail the open within the bound"),
|
||||
}
|
||||
assert!(manager.channel_ids().is_empty(), "channel torn down");
|
||||
assert_eq!(
|
||||
concrete_policy.count_for(&identity("alice")),
|
||||
0,
|
||||
"ledger decremented on timeout"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn run_open_wrapper_dispatch_deadline_bounds_establisher() {
|
||||
let manager = make_manager().await;
|
||||
let policy = super::super::policy::default_policy();
|
||||
let open_handler: OpenHandler = Arc::new(|_input, _conn, _auth| tokio::spawn(async {}));
|
||||
let auth = AuthContext::anonymous(b"alk/call");
|
||||
let hanging: OpenEstablisher = Arc::new(|_input, _auth| {
|
||||
Box::pin(async {
|
||||
tokio::time::sleep(Duration::from_secs(5)).await;
|
||||
Ok(Establishment::default())
|
||||
})
|
||||
});
|
||||
// No per-registration override; the dispatch deadline (50ms) is
|
||||
// the bound — the earlier of the two per ADR-049 §2.
|
||||
let env = run_open_wrapper(
|
||||
&manager,
|
||||
&policy,
|
||||
hook(hanging, None).as_ref(),
|
||||
&open_handler,
|
||||
&auth,
|
||||
"alk/tty",
|
||||
json!({}),
|
||||
"alice".to_string(),
|
||||
identity("alice"),
|
||||
"req-est-deadline".to_string(),
|
||||
Some(Instant::now() + Duration::from_millis(50)),
|
||||
)
|
||||
.await;
|
||||
match env.result {
|
||||
Err(e) => {
|
||||
assert_eq!(e.code, "channel:open_failed");
|
||||
assert_eq!(e.details.as_ref().expect("details")["reason"], "timeout");
|
||||
}
|
||||
Ok(_) => panic!("deadline must bound the establisher"),
|
||||
}
|
||||
assert!(manager.channel_ids().is_empty());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn run_open_wrapper_establisher_success_spawns_pumps_and_replies_channel_id() {
|
||||
let manager = make_manager().await;
|
||||
let policy = super::super::policy::default_policy();
|
||||
let spawned = Arc::new(AtomicBool::new(false));
|
||||
let spawned_clone = Arc::clone(&spawned);
|
||||
let open_handler: OpenHandler = Arc::new(move |_input, _conn, _auth| {
|
||||
spawned_clone.store(true, Ordering::SeqCst);
|
||||
tokio::spawn(async {})
|
||||
});
|
||||
let auth = AuthContext::anonymous(b"alk/call");
|
||||
let env = run_open_wrapper(
|
||||
&manager,
|
||||
&policy,
|
||||
hook(ok_establisher(), None).as_ref(),
|
||||
&open_handler,
|
||||
&auth,
|
||||
"alk/tty",
|
||||
json!({}),
|
||||
"alice".to_string(),
|
||||
identity("alice"),
|
||||
"req-est-ok".to_string(),
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
match env.result {
|
||||
Ok(v) => {
|
||||
let channel_id = v["channel_id"].as_u64().expect("channel_id");
|
||||
assert!(channel_id > 0);
|
||||
assert!(manager.has_channel(channel_id as u32));
|
||||
}
|
||||
Err(e) => panic!("successful establishment should open, got: {e:?}"),
|
||||
}
|
||||
assert!(spawned.load(Ordering::SeqCst));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn establishment_error_reason_mapping_covers_the_vocabulary() {
|
||||
let cases: Vec<(EstablishmentError, &str)> = vec![
|
||||
(
|
||||
EstablishmentError::DialFailed {
|
||||
message: "refused".into(),
|
||||
},
|
||||
"dial_failed",
|
||||
),
|
||||
(
|
||||
EstablishmentError::UnknownResource {
|
||||
message: "gone".into(),
|
||||
},
|
||||
"unknown_resource",
|
||||
),
|
||||
(
|
||||
EstablishmentError::ResourceShortage {
|
||||
message: "no fds".into(),
|
||||
},
|
||||
"resource_shortage",
|
||||
),
|
||||
(
|
||||
EstablishmentError::HandlerError {
|
||||
message: "boom".into(),
|
||||
},
|
||||
"handler_error",
|
||||
),
|
||||
];
|
||||
for (err, reason) in cases {
|
||||
let call_err = establishment_error_to_call_error(&err);
|
||||
assert_eq!(call_err.code, "channel:open_failed");
|
||||
let details = call_err.details.expect("details");
|
||||
assert_eq!(details["reason"], reason);
|
||||
assert_eq!(details["message"], err.message());
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn establishment_timeout_error_carries_timeout_reason() {
|
||||
let err = establishment_timeout_call_error();
|
||||
assert_eq!(err.code, "channel:open_failed");
|
||||
assert_eq!(err.details.as_ref().expect("details")["reason"], "timeout");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn establishment_bound_is_the_earlier_of_deadline_and_timeout() {
|
||||
// No deadline → the override (or the 10s default).
|
||||
assert_eq!(establishment_bound(None, None), ESTABLISHMENT_TIMEOUT);
|
||||
assert_eq!(
|
||||
establishment_bound(Some(Duration::from_secs(3)), None),
|
||||
Duration::from_secs(3)
|
||||
);
|
||||
// Deadline earlier than the timeout → the deadline.
|
||||
let soon = Instant::now() + Duration::from_millis(100);
|
||||
let bound = establishment_bound(None, Some(soon));
|
||||
assert!(bound <= Duration::from_millis(100));
|
||||
// Deadline later than the timeout → the timeout wins.
|
||||
let late = Instant::now() + Duration::from_secs(60);
|
||||
assert_eq!(
|
||||
establishment_bound(Some(Duration::from_secs(3)), Some(late)),
|
||||
Duration::from_secs(3)
|
||||
);
|
||||
// Deadline already past → zero bound (timeout fires immediately).
|
||||
let past = Instant::now() - Duration::from_secs(1);
|
||||
assert_eq!(establishment_bound(None, Some(past)), Duration::ZERO);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn register_openable_with_establisher_fails_via_invoke_streaming() {
|
||||
let manager = make_manager().await;
|
||||
let concrete_policy = Arc::new(super::super::policy::PerIdentityChannelPolicy::new(8));
|
||||
let policy: Arc<dyn ChannelLifecyclePolicy> =
|
||||
Arc::clone(&concrete_policy) as Arc<dyn ChannelLifecyclePolicy>;
|
||||
let core = ChannelCore::new(manager.clone(), policy);
|
||||
let spec = OperationSpec::new(
|
||||
"channels/tty/sub",
|
||||
OperationType::Sub,
|
||||
Visibility::External,
|
||||
json!({}),
|
||||
json!({ "type": "object", "properties": { "channel_id": { "type": "integer" } } }),
|
||||
vec![],
|
||||
AccessControl::default(),
|
||||
None,
|
||||
)
|
||||
.with_channel_open(ChannelOpenSpec::new("alk/tty"));
|
||||
let open_handler: OpenHandler = Arc::new(|_input, _conn, _auth| tokio::spawn(async {}));
|
||||
let registry = OperationRegistry::new();
|
||||
core.register_openable_with_establisher(
|
||||
spec,
|
||||
Some(failing_establisher()),
|
||||
open_handler,
|
||||
®istry,
|
||||
AuthContext::anonymous(b"alk/call"),
|
||||
None,
|
||||
)
|
||||
.expect("register");
|
||||
let mut stream =
|
||||
registry.invoke_streaming("channels/tty/sub", json!({}), test_context("reg-est-fail"));
|
||||
let env = stream.next().await.expect("one envelope");
|
||||
match env.result {
|
||||
Err(e) => {
|
||||
assert_eq!(e.code, "channel:open_failed");
|
||||
assert_eq!(
|
||||
e.details.as_ref().expect("details")["reason"],
|
||||
"dial_failed"
|
||||
);
|
||||
}
|
||||
Ok(_) => panic!("establisher failure must surface as a call error"),
|
||||
}
|
||||
assert!(manager.channel_ids().is_empty());
|
||||
assert_eq!(
|
||||
concrete_policy.count_for(&identity("alice")),
|
||||
0,
|
||||
"no-establisher compat: registry-level failure un-reserves the ledger"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn register_openable_no_establisher_behaves_as_today() {
|
||||
let manager = make_manager().await;
|
||||
let policy = super::super::policy::default_policy();
|
||||
let core = ChannelCore::new(manager.clone(), policy);
|
||||
let spec = OperationSpec::new(
|
||||
"channels/tty/sub",
|
||||
OperationType::Query,
|
||||
Visibility::External,
|
||||
json!({}),
|
||||
json!({ "type": "object", "properties": { "channel_id": { "type": "integer" } } }),
|
||||
vec![],
|
||||
AccessControl::default(),
|
||||
None,
|
||||
)
|
||||
.with_channel_open(ChannelOpenSpec::new("alk/tty"));
|
||||
let spawned = Arc::new(AtomicBool::new(false));
|
||||
let spawned_clone = Arc::clone(&spawned);
|
||||
let open_handler: OpenHandler = Arc::new(move |_input, _conn, _auth| {
|
||||
spawned_clone.store(true, Ordering::SeqCst);
|
||||
tokio::spawn(async {})
|
||||
});
|
||||
let registry = OperationRegistry::new();
|
||||
core.register_openable(
|
||||
spec,
|
||||
open_handler,
|
||||
®istry,
|
||||
AuthContext::anonymous(b"alk/call"),
|
||||
)
|
||||
.expect("register");
|
||||
let env = registry
|
||||
.invoke("channels/tty/sub", json!({}), test_context("reg-compat"))
|
||||
.await;
|
||||
match env.result {
|
||||
Ok(v) => assert!(v["channel_id"].as_u64().is_some()),
|
||||
Err(e) => panic!("no-establisher open should succeed, got: {e:?}"),
|
||||
}
|
||||
assert!(spawned.load(Ordering::SeqCst));
|
||||
assert_eq!(manager.open_count(), 1, "channel stays open (compat)");
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user