diff --git a/CHANGELOG.md b/CHANGELOG.md index e9986ce..41bfe39 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,44 @@ 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. + +### Changed + +- **`Establishment` carries the channel plan (review 007 R-01 — + breaking at 0.6.0).** The reserved field is filled: + `Establishment { plan: Option }` with + `ChannelPlan = Arc` — 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, 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 + diff --git a/Cargo.lock b/Cargo.lock index f961f22..6d8d58c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -27,7 +27,7 @@ dependencies = [ [[package]] name = "alkcall" -version = "0.5.0" +version = "0.6.0" dependencies = [ "async-trait", "bytes", diff --git a/Cargo.toml b/Cargo.toml index 082887c..42aeacd 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "alkcall" -version = "0.5.0" +version = "0.6.0" edition = "2021" rust-version = "1.85" license = "MIT OR Apache-2.0" diff --git a/docs/architecture/decisions/049-channel-open-establishment-phase.md b/docs/architecture/decisions/049-channel-open-establishment-phase.md index 7064fd2..4add3bb 100644 --- a/docs/architecture/decisions/049-channel-open-establishment-phase.md +++ b/docs/architecture/decisions/049-channel-open-establishment-phase.md @@ -353,4 +353,56 @@ CallError }` / `MissingChannelId` / `AdoptFailed`, with 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). \ No newline at end of file +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 }` +with `ChannelPlan = Arc` — **typed-opaque, not +`serde_json::Value`**. The review's `Option` 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, 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. \ No newline at end of file diff --git a/src/channels/client.rs b/src/channels/client.rs index bb335ad..9547141 100644 --- a/src/channels/client.rs +++ b/src/channels/client.rs @@ -897,7 +897,7 @@ mod tests { Arc::clone(&policy) as Arc; 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 @@ -1220,7 +1220,7 @@ mod tests { let (data_tx, mut data_rx) = tokio::sync::mpsc::channel::>(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"); @@ -1340,7 +1340,7 @@ mod tests { let policy_for_assert = Arc::clone(&policy); let open_handler: OpenHandler = - Arc::new(|_input, _channel_conn, _auth| tokio::spawn(async {})); + Arc::new(|_input, _plan, _channel_conn, _auth| tokio::spawn(async {})); let establisher: OpenEstablisher = Arc::new(|_input, _auth| { Box::pin(async { Err(EstablishmentError::DialFailed { @@ -1472,7 +1472,7 @@ mod tests { let (data_tx, mut data_rx) = tokio::sync::mpsc::channel::>(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"); @@ -2333,7 +2333,7 @@ mod tests { let policy_for_hook: Arc = Arc::clone(&policy) as Arc; - 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; }) diff --git a/src/channels/operations.rs b/src/channels/operations.rs index cbbcc23..04a042b 100644 --- a/src/channels/operations.rs +++ b/src/channels/operations.rs @@ -319,30 +319,75 @@ pub struct ChannelCore { policy: Arc, } -/// The ALPN-specific open handler (ADR-047 §3). The ALPN crate -/// provides this; [`ChannelCore::register_openable`] wraps it with the -/// channel machinery. The handler receives the open op's `input` -/// (params — e.g. which container for tty, which target for tunnel), -/// the channel's [`Connection`] (carrying the data-plane ALPN, built -/// by the wrapper from the channel's reassembled read half + mux write +/// The ALPN-specific open handler (ADR-047 §3, as amended by ADR-049 +/// amendment 2 — the establishment `plan` flows to the handler). The +/// ALPN crate provides this; [`ChannelCore::register_openable`] wraps +/// it with the channel machinery. The handler receives the open op's +/// `input` (params — e.g. which container for tty, which target for +/// tunnel), the establishment `plan` (the typed-opaque payload the +/// establisher returned — e.g. the dialed `SubstrateHandle`; `None` +/// when no establisher is registered or it returned no plan), the +/// channel's [`Connection`] (carrying the data-plane ALPN, built by +/// the wrapper from the channel's reassembled read half + mux write /// half), and the peer's [`AuthContext`], spawns its protocol on the /// channel's `BiStream`, and returns the `JoinHandle` so the wrapper /// can record it for teardown (abort on `channel/close` / connection /// drop). /// +/// **The returned `JoinHandle` must track the data-plane lifetime** +/// (R-02): 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 the pumps inline inside the returned +/// task, never spawn-and-forget them and return early. +/// /// This mirrors the `InstallChannelZero` hook shape (the call adapter /// spawns `run_loop_single_stream` on channel 0's `Connection`): the /// ALPN handler does the same for its data-plane protocol on the /// allocated channel's `Connection`. The handler owns its protocol's /// sub-stream multiplexing on the `BiStream` it receives (ADR-035). -pub type OpenHandler = - Arc tokio::task::JoinHandle<()> + Send + Sync>; +pub type OpenHandler = Arc< + dyn Fn(Value, Option, 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. +/// The establishment plan payload (ADR-049 §1, as filled by ADR-049 +/// amendment 2 — R-01): the typed-opaque data the establisher hands to +/// the pump handler. Opaque to alkcall by construction — the +/// establisher and the handler agree on the concrete type; the +/// downcast (`Arc::clone` + `Any::downcast_ref`) happens in the ALPN +/// crate. +/// +/// The plan is **process-local**: establisher → wrapper → handler on +/// the producing side. Nothing crosses the transport that isn't +/// already the open op's input — the wire surface is unchanged. +pub type ChannelPlan = Arc; + +/// The establishment-phase result (ADR-049 §1, as filled by ADR-049 +/// amendment 2 — R-01). Carries the channel plan the wrapper threads +/// to the pump handler's `plan` parameter ([`OpenHandler`]). The +/// plan is optional: an establisher that only validates (TTY's +/// lookup/ownership shape) returns [`Establishment::default`]. +/// +/// `#[non_exhaustive]` so a future carrier change is not another +/// breaking release. Construct with [`Establishment::new`] (plan) or +/// [`Establishment::default`] (no plan). #[derive(Debug, Clone, Default)] -pub struct Establishment {} +#[non_exhaustive] +pub struct Establishment { + /// The channel plan — the payload the pump handler receives. The + /// establisher and the handler agree on the concrete type + /// ([`ChannelPlan`]); alkcall never inspects it. + pub plan: Option, +} + +impl Establishment { + /// A result carrying `plan` to the pump handler. + pub fn new(plan: ChannelPlan) -> Self { + Self { plan: Some(plan) } + } +} /// The establishment-failure reason code the wrapper puts in the open /// reply's `details.reason` (ADR-049 §3). The vocabulary maps 1:1 onto @@ -469,6 +514,12 @@ impl ChannelCore { /// establishment phase can fail and the failure belongs in the /// open reply (ADR-049). /// + /// **Lifetime contract (R-02):** the returned `OpenHandler` + /// `JoinHandle` must track the data-plane lifetime — see the + /// [`OpenHandler`] type docs. The wrapper awaits it and its + /// completion triggers channel teardown; early return = teardown + /// at birth. + /// /// 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 @@ -517,6 +568,13 @@ impl ChannelCore { /// exactly like [`ChannelCore::register_openable`] (the /// no-establisher shape), and existing `OpenHandler`s compile /// unchanged. + /// + /// **Lifetime contract (R-02):** the returned `OpenHandler` + /// `JoinHandle` must track the data-plane lifetime — see the + /// [`OpenHandler`] type docs. The wrapper awaits it and its + /// completion triggers channel teardown; early return = teardown + /// at birth. On success the establisher's plan flows to the + /// handler's `plan` parameter (ADR-049 amendment 2). pub fn register_openable_with_establisher( &self, spec: OperationSpec, @@ -711,40 +769,45 @@ async fn run_open_wrapper( // 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(), - ); + // (the pre-ADR-049 shape, unchanged) — its plan is `None`. + let plan = match establisher { + Some(hook) => { + 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)) => establishment.plan, } - 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)) => {} } - } + None => None, + }; 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()); - let raw_task = open_handler(input, channel_conn, auth.clone()); + let raw_task = open_handler(input, plan, channel_conn, auth.clone()); let teardown_manager = manager.clone(); let teardown_policy = Arc::clone(policy); @@ -976,7 +1039,7 @@ mod tests { use crate::registry::spec::ChannelOpenSpec; use futures::stream::StreamExt; use std::collections::HashMap; - use std::sync::atomic::{AtomicBool, Ordering}; + use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::Arc; use tokio::io::duplex; @@ -1300,7 +1363,7 @@ mod tests { 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| { + let open_handler: OpenHandler = Arc::new(move |_input, _plan, _conn, _auth| { spawned_clone.store(true, Ordering::SeqCst); tokio::spawn(async {}) }); @@ -1331,7 +1394,7 @@ mod tests { 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| { + let open_handler: OpenHandler = Arc::new(move |_input, _plan, _conn, _auth| { spawned_clone.store(true, Ordering::SeqCst); tokio::spawn(async {}) }); @@ -1365,7 +1428,8 @@ mod tests { async fn make_open_handler_sink_returns_not_implemented() { 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 open_handler: OpenHandler = + Arc::new(|_input, _plan, _conn, _auth| tokio::spawn(async {})); let auth = AuthContext::anonymous(b"alk/call"); let handler = make_open_handler_sink( manager, @@ -1402,7 +1466,7 @@ mod tests { .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| { + let open_handler: OpenHandler = Arc::new(move |_input, _plan, _conn, _auth| { spawned_clone.store(true, Ordering::SeqCst); tokio::spawn(async {}) }); @@ -1445,7 +1509,7 @@ mod tests { .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| { + let open_handler: OpenHandler = Arc::new(move |_input, _plan, _conn, _auth| { spawned_clone.store(true, Ordering::SeqCst); tokio::spawn(async {}) }); @@ -1485,7 +1549,8 @@ mod tests { None, ) .with_channel_open(ChannelOpenSpec::new("alk/tty")); - let open_handler: OpenHandler = Arc::new(|_input, _conn, _auth| tokio::spawn(async {})); + let open_handler: OpenHandler = + Arc::new(|_input, _plan, _conn, _auth| tokio::spawn(async {})); let registry = OperationRegistry::new(); core.register_openable( spec, @@ -1515,7 +1580,8 @@ mod tests { AccessControl::default(), None, ); - let open_handler: OpenHandler = Arc::new(|_input, _conn, _auth| tokio::spawn(async {})); + let open_handler: OpenHandler = + Arc::new(|_input, _plan, _conn, _auth| tokio::spawn(async {})); let registry = OperationRegistry::new(); let result = core.register_openable( spec, @@ -1532,7 +1598,8 @@ mod tests { let manager = make_manager().await; let policy: Arc = Arc::new(super::super::policy::PerIdentityChannelPolicy::new(0)); - let open_handler: OpenHandler = Arc::new(|_input, _conn, _auth| tokio::spawn(async {})); + let open_handler: OpenHandler = + Arc::new(|_input, _plan, _conn, _auth| tokio::spawn(async {})); let auth = AuthContext::anonymous(b"alk/call"); let opener_id = Identity { id: "alice".to_string(), @@ -1598,7 +1665,7 @@ mod tests { Arc::clone(&concrete_policy) as Arc; let spawned = Arc::new(AtomicBool::new(false)); let spawned_clone = Arc::clone(&spawned); - let open_handler: OpenHandler = Arc::new(move |_input, _conn, _auth| { + let open_handler: OpenHandler = Arc::new(move |_input, _plan, _conn, _auth| { spawned_clone.store(true, Ordering::SeqCst); tokio::spawn(async {}) }); @@ -1653,7 +1720,8 @@ mod tests { let concrete_policy = Arc::new(super::super::policy::PerIdentityChannelPolicy::new(8)); let policy: Arc = Arc::clone(&concrete_policy) as Arc; - let open_handler: OpenHandler = Arc::new(|_input, _conn, _auth| tokio::spawn(async {})); + let open_handler: OpenHandler = + Arc::new(|_input, _plan, _conn, _auth| tokio::spawn(async {})); let auth = AuthContext::anonymous(b"alk/call"); let hanging: OpenEstablisher = Arc::new(|_input, _auth| { Box::pin(async { @@ -1695,7 +1763,8 @@ mod tests { 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 open_handler: OpenHandler = + Arc::new(|_input, _plan, _conn, _auth| tokio::spawn(async {})); let auth = AuthContext::anonymous(b"alk/call"); let hanging: OpenEstablisher = Arc::new(|_input, _auth| { Box::pin(async { @@ -1735,7 +1804,7 @@ mod tests { 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| { + let open_handler: OpenHandler = Arc::new(move |_input, _plan, _conn, _auth| { spawned_clone.store(true, Ordering::SeqCst); tokio::spawn(async {}) }); @@ -1850,7 +1919,8 @@ mod tests { None, ) .with_channel_open(ChannelOpenSpec::new("alk/tty")); - let open_handler: OpenHandler = Arc::new(|_input, _conn, _auth| tokio::spawn(async {})); + let open_handler: OpenHandler = + Arc::new(|_input, _plan, _conn, _auth| tokio::spawn(async {})); let registry = OperationRegistry::new(); core.register_openable_with_establisher( spec, @@ -1900,7 +1970,7 @@ mod tests { .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| { + let open_handler: OpenHandler = Arc::new(move |_input, _plan, _conn, _auth| { spawned_clone.store(true, Ordering::SeqCst); tokio::spawn(async {}) }); @@ -1922,4 +1992,166 @@ mod tests { assert!(spawned.load(Ordering::SeqCst)); assert_eq!(manager.open_count(), 1, "channel stays open (compat)"); } + + // --- ADR-049 amendment 2 (R-01): the plan payload --------------------- + + /// The test-side stand-in for a live substrate handle: the + /// establisher's dialed resource, handed to the pump handler via + /// the plan. + #[derive(Debug, PartialEq, Eq)] + struct TestHandle(u64); + + #[tokio::test] + async fn run_open_wrapper_plan_flows_from_establisher_to_handler() { + let manager = make_manager().await; + let policy = super::super::policy::default_policy(); + let establishing: OpenEstablisher = Arc::new(|_input, _auth| { + Box::pin(async { Ok(Establishment::new(Arc::new(TestHandle(42)) as ChannelPlan)) }) + }); + let received: Arc>> = Arc::new(std::sync::Mutex::new(None)); + let received_clone = Arc::clone(&received); + let open_handler: OpenHandler = Arc::new(move |_input, plan, _conn, _auth| { + let handle = plan + .as_ref() + .and_then(|p| p.downcast_ref::()) + .map(|h| h.0); + *received_clone.lock().unwrap() = handle; + tokio::spawn(async {}) + }); + let auth = AuthContext::anonymous(b"alk/call"); + let env = run_open_wrapper( + &manager, + &policy, + hook(establishing, None).as_ref(), + &open_handler, + &auth, + "alk/tty", + json!({}), + "alice".to_string(), + identity("alice"), + "req-plan-flow".to_string(), + None, + ) + .await; + assert!(env.result.is_ok(), "open with plan should succeed"); + assert_eq!( + *received.lock().unwrap(), + Some(42), + "the establisher's plan reaches the pump handler typed-opaque" + ); + } + + #[tokio::test] + async fn concurrent_same_resource_opens_receive_their_own_plans() { + let manager = make_manager().await; + let policy = super::super::policy::default_policy(); + // The establisher dials per open (here: a distinct TestHandle + // per invocation) and delivers via the plan — no shared + // resource-keyed slot, so the POC's same-resource race is + // unreachable by construction. + let next_handle = Arc::new(AtomicU64::new(1)); + let establishing: OpenEstablisher = { + let next_handle = Arc::clone(&next_handle); + Arc::new(move |_input, _auth| { + let handle = TestHandle(next_handle.fetch_add(1, Ordering::SeqCst)); + Box::pin(async move { Ok(Establishment::new(Arc::new(handle) as ChannelPlan)) }) + }) + }; + let observed: Arc>> = Arc::new(std::sync::Mutex::new(vec![])); + let observed_clone = Arc::clone(&observed); + let open_handler: OpenHandler = Arc::new(move |_input, plan, _conn, _auth| { + let got = plan + .as_ref() + .and_then(|p| p.downcast_ref::()) + .map(|h| h.0); + observed_clone.lock().unwrap().push(got.expect("plan")); + tokio::spawn(async {}) + }); + let auth = AuthContext::anonymous(b"alk/call"); + let hook_a = hook(establishing.clone(), None); + let hook_b = hook(establishing, None); + // Two concurrent opens of the SAME resource key: both succeed, + // each handler sees its own establisher's handle. + let (env_a, env_b) = tokio::join!( + run_open_wrapper( + &manager, + &policy, + hook_a.as_ref(), + &open_handler, + &auth, + "alk/tty", + json!({ "resource": "same" }), + "alice".to_string(), + identity("alice"), + "req-race-a".to_string(), + None, + ), + run_open_wrapper( + &manager, + &policy, + hook_b.as_ref(), + &open_handler, + &auth, + "alk/tty", + json!({ "resource": "same" }), + "alice".to_string(), + identity("alice"), + "req-race-b".to_string(), + None, + ) + ); + assert!(env_a.result.is_ok(), "open A should succeed"); + assert!(env_b.result.is_ok(), "open B should succeed"); + let mut got = observed.lock().unwrap().clone(); + got.sort_unstable(); + assert_eq!( + got, + vec![1, 2], + "each concurrent same-resource open got its own handle — the handoff-slot race is unreachable" + ); + assert_eq!(manager.open_count(), 2, "both channels stay open"); + } + + #[tokio::test] + async fn no_establisher_passes_none_plan_to_handler() { + let manager = make_manager().await; + let policy = super::super::policy::default_policy(); + let received: Arc>> = Arc::new(std::sync::Mutex::new(None)); + let received_clone = Arc::clone(&received); + let open_handler: OpenHandler = Arc::new(move |_input, plan, _conn, _auth| { + *received_clone.lock().unwrap() = Some(plan.is_none()); + tokio::spawn(async {}) + }); + let auth = AuthContext::anonymous(b"alk/call"); + let env = run_open_wrapper( + &manager, + &policy, + None, + &open_handler, + &auth, + "alk/tty", + json!({}), + "alice".to_string(), + identity("alice"), + "req-no-plan".to_string(), + None, + ) + .await; + assert!(env.result.is_ok(), "no-establisher open should succeed"); + assert_eq!( + *received.lock().unwrap(), + Some(true), + "no establisher = plan is None" + ); + } + + #[test] + fn establishment_default_has_no_plan() { + let e = Establishment::default(); + assert!(e.plan.is_none()); + let e = Establishment::new(Arc::new(TestHandle(7)) as ChannelPlan); + let plan = e.plan.expect("plan"); + let h = plan.downcast_ref::().expect("typed"); + assert_eq!(h.0, 7); + } }