//! The end-to-end suite — every validated POC behavior, re-expressed //! against the crate's public API (tests/common/harness.rs, the two //! POC harnesses' topologies unified). The spec's executable form: the //! quality gate before the implementation review. //! //! Suites: //! 1. **Forward (`-L`)** — open → establisher dials → pump → round- //! trip; 1 MiB backpressure; typed errors; per-call opener identity //! (CF-006); establishment deadline → `timeout` (ADR-049 §2). //! 2. **Reverse (`-R`)** — `open_reverse_channel` + `adopt` + //! `pump_against`; odd-ID allocation (ADR-047 §5); the identity //! precedence chain (CF-005); half-close (W4); out-of-band close + //! self-reaping (W3); same-resource concurrency (R-01). //! 3. **Listen producer** — the accept-queue flow + typed errors. //! 4. **Datagram (`udp`)** — forward + reverse datagram sessions via //! the codec; round-trip, empty datagram (F-2), chunk-split //! survival, TCP+UDP concurrent on one connection. //! 5. **Teardown matrix** (ADR-005) — close/join/Drop; `channel_ids()` //! asserts no leaks on either side. //! 6. Spec-conformance assertions — no phantom channel on a failed //! open, no leaked session on a failed adopt, R-02's //! JoinHandle-tracks-the-data-plane shape. //! //! (`local`-feature real-socket suites live in tests/local_halves.rs //! — the duplex suites here need no sockets.) mod harness; use std::sync::Arc; use alkcall::channels::client::ChannelOpenError; use alktunnels::params::{establishment_reason, Substrate, TunnelParams, TUNNEL_OPEN_SCOPE}; use alktunnels::producer::{AcceptQueue, ResourceRegistry, TargetHandle, TunnelEstablishError}; use alktunnels::{open_reverse_channel, TunnelSession}; use harness::{ echo_dial, failing_dial, framed_udp_echo_dial, wire, wire_forward, wire_listen, wire_serving_identity, wire_with, ForwardTopology, RegistrationMode, Topology, TEST_AUTH_TOKEN, }; fn params(resource: &str, substrate: Substrate) -> TunnelParams { TunnelParams { resource: resource.to_string(), substrate, } } /// Only channel 0 (the pre-negotiated call channel) may remain — /// asserted on both sides of both topologies after every teardown. trait LeakCheck { fn no_data_channels(&self) -> bool; } impl LeakCheck for Topology { fn no_data_channels(&self) -> bool { self.consumer_manager .channel_ids() .iter() .all(|&id| id == 0) && self .producer .manager() .channel_ids() .iter() .all(|&id| id == 0) } } impl LeakCheck for ForwardTopology { fn no_data_channels(&self) -> bool { self.consumer .manager() .channel_ids() .iter() .all(|&id| id == 0) && self .producer_manager .channel_ids() .iter() .all(|&id| id == 0) } } fn only_channel_0(topo: &T) -> bool { topo.no_data_channels() } fn reason_of_open_failure(err: &TunnelSessionOpenError) -> Option { establishment_reason(err.open_ref()).map(str::to_string) } type TunnelSessionOpenError = alktunnels::TunnelOpenError; fn call_error_reason(err: &alktunnels::ReverseOpenError) -> Option { match err { alktunnels::ReverseOpenError::Call(call_err) => call_err .details .as_ref() .and_then(|d| d.get("reason")) .and_then(|r| r.as_str()) .map(str::to_string), _ => None, } } // ===================================================================== // Suite 1 — Forward (`-L`): open → establisher → pump → round-trip // ===================================================================== /// A 1 MiB patterned payload round-trips through the two-pump shape /// (the bounded-buffer backpressure path; the POCs' sizing probe). /// The consumer writes via taken halves and reads the echo back on /// the same halves. #[tokio::test] async fn forward_large_payload_backpressure_round_trip() { let registry = ResourceRegistry::new(); registry .register("echo", Substrate::Tcp, "in-process") .await; let topo = wire_forward(registry, echo_dial()).await; let session = TunnelSession::open(&topo.consumer, params("echo", Substrate::Tcp)) .await .expect("open"); let Taken { mut read, mut write, session, } = split_taken(session.take_halves()); let payload: Vec = (0..1024 * 1024u32).map(|i| (i % 251) as u8).collect(); let expected = payload.clone(); let writer = tokio::spawn(async move { use tokio::io::AsyncWriteExt; write.write_all(&payload).await.expect("write all"); write.flush().await.expect("flush"); write }); use tokio::io::AsyncReadExt; let mut received = vec![0u8; expected.len()]; let mut got = 0; while got < received.len() { let n = tokio::time::timeout( std::time::Duration::from_secs(10), read.read(&mut received[got..]), ) .await .expect("read timed out (backpressure deadlock?)") .expect("read"); assert!(n > 0, "EOF before the full echo"); got += n; } writer.await.expect("writer task"); assert_eq!(got, expected.len(), "1 MiB round trip intact"); assert_eq!(&received[..], &expected[..]); drop(read); let (_, _, reaped) = session.join().await; assert!(reaped); tokio::time::sleep(std::time::Duration::from_millis(100)).await; assert!(only_channel_0(&topo)); } /// A hanging establisher is bounded by the per-registration timeout /// (ADR-049 §2's override): the open fails `channel:open_failed` with /// reason `timeout`, and no channel survives on either side (the /// no-phantom property, deadline-flavored). #[tokio::test] async fn forward_establishment_deadline_maps_to_timeout_typed() { let registry = ResourceRegistry::new(); registry .register("echo", Substrate::Tcp, "in-process") .await; let producer_op_registry = Arc::new(alkcall::registry::registration::OperationRegistry::new()); // The timeout topology: the default wire() shape with the hanging // establisher replacing the dial establisher — wired by hand here // via wire_with's Timeout mode on the reverse topology, then the // forward direction asserted through the session API on a fresh // forward wiring below. let topo = wire_with( registry.clone(), RegistrationMode::Dial(echo_dial()), Some(harness::consumer_identity()), None, Arc::new(alkcall::core::auth::NoopIdentityProvider), ) .await; // Reverse-side open against a hanging establisher with a 300ms // per-registration bound (the deadline-expiry probe). let queue = AcceptQueue::new(); queue.close().await; let timeout_topo = harness::wire_with( ResourceRegistry::new(), RegistrationMode::Timeout(std::time::Duration::from_millis(300)), Some(harness::consumer_identity()), None, Arc::new(alkcall::core::auth::NoopIdentityProvider), ) .await; let started = std::time::Instant::now(); match open_reverse_channel( &timeout_topo.consumer_call, ¶ms("echo", Substrate::Tcp), None, ) .await { Ok(_) => panic!("hanging establisher must time out"), Err(e) => { let err = match e { alktunnels::ReverseOpenError::Call(c) => c, other => panic!("expected a call error, got {other:?}"), }; assert_eq!(err.code, "channel:open_failed"); let reason = err .details .as_ref() .and_then(|d| d.get("reason")) .and_then(|r| r.as_str()) .expect("reason in details"); assert_eq!(reason, "timeout"); } } assert!( started.elapsed() < std::time::Duration::from_secs(5), "the bound fired, not the harness backstop: {:?}", started.elapsed() ); tokio::time::sleep(std::time::Duration::from_millis(100)).await; assert!(only_channel_0(&timeout_topo), "no phantom channel"); // The generic channel ops on the same registry still work (the // timeout is per-registration, ADR-049 §2). let response = topo .producer .call_open_op("consumer/serves/nothing", serde_json::json!({})) .await; assert!( response.result.is_err(), "the point is the reply RESOLVED (per-op isolation)" ); let _ = queue; let _ = producer_op_registry; } /// The identity-witness end-to-end (CF-006): the establisher's /// per-call auth carried the END CALLER's identity (the consumer), on /// both the transport-identity path and the token path. #[tokio::test] async fn forward_opener_identity_witnessed_on_both_paths() { // Transport identity path (the mTLS posture — no token). let registry = ResourceRegistry::new(); registry .register("echo", Substrate::Tcp, "in-process") .await; let topo = wire_forward(registry, echo_dial()).await; let session = TunnelSession::open(&topo.consumer, params("echo", Substrate::Tcp)) .await .expect("transport-identity open"); let witness = topo.identity_witness.lock().await.clone(); assert_eq!( witness.as_deref(), Some("consumer"), "establisher saw the per-call opener identity (CF-006)" ); session.close().await; // Token path (the hub-forwarding posture): the token resolves to // the consumer identity and the establisher sees it end-to-end. let registry = ResourceRegistry::new(); registry .register("echo", Substrate::Tcp, "in-process") .await; let topo = wire(registry, echo_dial()).await; let channel_id = open_reverse_channel( &topo.consumer_call, ¶ms("echo", Substrate::Tcp), Some(TEST_AUTH_TOKEN), ) .await .expect("token open authorized"); let witness = topo.identity_witness.lock().await.clone(); assert_eq!(witness.as_deref(), Some("consumer")); assert_eq!(channel_id % 2, 1, "the connect side allocated an odd id"); let session = TunnelSession::adopt( &topo.consumer_manager, channel_id, Substrate::Tcp, alktunnels::TUNNEL_ALPN, ) .await .expect("adopt"); session.close().await; tokio::time::sleep(std::time::Duration::from_millis(100)).await; assert!(only_channel_0(&topo)); } // ===================================================================== // Suite 2 — Reverse (`-R`): open_reverse_channel + adopt + pump_against // ===================================================================== /// Odd-ID allocation under ADR-047 §5: the serving side (producer, /// `ChannelSide::Connect` here) allocates odd ids; the initiator /// adopts. Asserted across several opens on one connection. #[tokio::test] async fn reverse_channel_ids_are_odd_and_distinct() { let registry = ResourceRegistry::new(); registry .register("echo", Substrate::Tcp, "in-process") .await; let topo = wire(registry, echo_dial()).await; let mut ids = Vec::new(); for _ in 0..3 { let id = open_reverse_channel(&topo.consumer_call, ¶ms("echo", Substrate::Tcp), None) .await .expect("reverse open"); assert_eq!(id % 2, 1, "ADR-047 §5: the serving side allocates odd ids"); ids.push(id); } assert_eq!( ids.len(), ids.iter().collect::>().len() ); for id in ids { // The reverse-open'd channel was never adopted: the manager // parked its early arrivals. `adopt_channel` drains the park // and installs the entry — adopt, then close (the session's // teardown; teardown_channel alone is UnknownChannel for a // parked-but-never-adopted id). let session = TunnelSession::adopt( &topo.consumer_manager, id, Substrate::Tcp, alktunnels::TUNNEL_ALPN, ) .await .expect("adopt parked channel"); assert!(session.close().await, "close reaps the entry"); } tokio::time::sleep(std::time::Duration::from_millis(100)).await; assert!(only_channel_0(&topo)); } /// The identity precedence chain (CF-005), asserted through the /// establisher's per-call identity: token > ServingConfig.identity > /// transport; identity-less everywhere → fail-closed FORBIDDEN. #[tokio::test] async fn reverse_identity_precedence_chain() { // (a) ServingConfig.identity overrides the transport identity. let registry = ResourceRegistry::new(); registry .register("echo", Substrate::Tcp, "in-process") .await; let override_identity = alkcall::core::auth::Identity { id: "worker-effective".to_string(), scopes: vec![TUNNEL_OPEN_SCOPE.to_string()], resources: Default::default(), }; let topo = wire_serving_identity(registry, echo_dial(), Some(override_identity)).await; let id = open_reverse_channel(&topo.consumer_call, ¶ms("echo", Substrate::Tcp), None) .await .expect("override identity authorizes"); let witness = topo.identity_witness.lock().await.clone(); assert_eq!( witness.as_deref(), Some("worker-effective"), "ServingConfig.identity won over the transport identity" ); let session = TunnelSession::adopt( &topo.consumer_manager, id, Substrate::Tcp, alktunnels::TUNNEL_ALPN, ) .await .expect("adopt parked channel"); session.close().await; tokio::time::sleep(std::time::Duration::from_millis(100)).await; // (b) The token wins over the transport identity: the provider // resolves the scoped token to a DISTINCT id — the witness proves // which one the establisher saw. struct DistinctTokenProvider; impl alkcall::core::auth::IdentityProvider for DistinctTokenProvider { fn resolve_from_fingerprint(&self, _: &str) -> Option { None } fn resolve_from_token( &self, token: &alkcall::core::auth::AuthToken, ) -> Option { if token.raw == TEST_AUTH_TOKEN.as_bytes() { Some(alkcall::core::auth::Identity { id: "token-caller".to_string(), scopes: vec![TUNNEL_OPEN_SCOPE.to_string()], resources: Default::default(), }) } else { None } } } let registry = ResourceRegistry::new(); registry .register("echo", Substrate::Tcp, "in-process") .await; let topo = harness::wire_with( registry, RegistrationMode::Dial(echo_dial()), Some(harness::consumer_identity()), None, Arc::new(DistinctTokenProvider), ) .await; let id = open_reverse_channel( &topo.consumer_call, ¶ms("echo", Substrate::Tcp), Some(TEST_AUTH_TOKEN), ) .await .expect("token open authorized"); let witness = topo.identity_witness.lock().await.clone(); assert_eq!( witness.as_deref(), Some("token-caller"), "the payload token won over the transport identity (CF-005 precedence)" ); let session = TunnelSession::adopt( &topo.consumer_manager, id, Substrate::Tcp, alktunnels::TUNNEL_ALPN, ) .await .expect("adopt parked channel"); session.close().await; tokio::time::sleep(std::time::Duration::from_millis(100)).await; assert!(only_channel_0(&topo)); // (c) Identity-less everywhere → fail-closed. let registry = ResourceRegistry::new(); registry .register("echo", Substrate::Tcp, "in-process") .await; let topo = wire_with( registry, RegistrationMode::Dial(echo_dial()), None, None, Arc::new(alkcall::core::auth::NoopIdentityProvider), ) .await; let err = open_reverse_channel(&topo.consumer_call, ¶ms("echo", Substrate::Tcp), None) .await .expect_err("identity-less open must be denied"); match err { alktunnels::ReverseOpenError::Call(call_err) => { assert_eq!(call_err.code, "FORBIDDEN"); } other => panic!("expected a call error, got {other:?}"), } assert!(only_channel_0(&topo)); } /// Concurrent opens of the SAME resource (R-01): each open's plan /// carries its own dialed handle; each round-trips its own marker; no /// handoff race. #[tokio::test] async fn reverse_concurrent_same_resource_no_handoff_race() { let registry = ResourceRegistry::new(); registry .register("echo", Substrate::Tcp, "in-process") .await; let topo = wire(registry, echo_dial()).await; use tokio::io::{AsyncReadExt, AsyncWriteExt}; let mut sessions = Vec::new(); for i in 0..3 { let channel_id = open_reverse_channel(&topo.consumer_call, ¶ms("echo", Substrate::Tcp), None) .await .expect("reverse open"); assert_eq!(channel_id % 2, 1); let session = TunnelSession::adopt( &topo.consumer_manager, channel_id, Substrate::Tcp, alktunnels::TUNNEL_ALPN, ) .await .expect("adopt"); let (accepted_end, local_end) = tokio::io::duplex(64 * 1024); let session = session.pump_against(accepted_end).await; let (mut l_read, mut l_write) = tokio::io::split(local_end); let msg = format!("racer-{i}-{}", channel_id); l_write.write_all(msg.as_bytes()).await.expect("write"); l_write.flush().await.expect("flush"); let mut buf = vec![0u8; msg.len()]; tokio::time::timeout( std::time::Duration::from_secs(5), l_read.read_exact(&mut buf), ) .await .expect("round trip") .expect("read"); assert_eq!(buf, msg.as_bytes()); sessions.push(session); } let ids: std::collections::HashSet = sessions.iter().map(|s| s.channel_id).collect(); assert_eq!(ids.len(), sessions.len(), "distinct odd ids"); assert!(ids.iter().all(|&id| id % 2 == 1), "all odd (ADR-047 §5)"); for session in sessions { let (_, _, reaped) = session.join().await; assert!(reaped); } assert!(only_channel_0(&topo)); } // ===================================================================== // Suite 3 — Listen producer (shape 2, the accept-queue flow) // ===================================================================== /// The listen-queue flow: FIFO always-before-take across two opens; /// the pump handler is the same one the dial shape registers. #[tokio::test] async fn listen_fifo_ordering_two_opens_two_handles() { let queue = AcceptQueue::new(); queue.push(echo_handle("first")).await.expect("push 1"); queue.push(echo_handle("second")).await.expect("push 2"); let registry = ResourceRegistry::new(); registry .register("listener", Substrate::Tcp, "assembly-owned") .await; let topo = wire_listen(registry, accept_from_queue(queue.clone())).await; let id1 = open_reverse_channel( &topo.consumer_call, ¶ms("listener", Substrate::Tcp), None, ) .await .expect("open 1"); let id2 = open_reverse_channel( &topo.consumer_call, ¶ms("listener", Substrate::Tcp), None, ) .await .expect("open 2"); assert_ne!(id1, id2); assert_eq!(id1 % 2, 1); assert_eq!(id2 % 2, 1); for (channel_id, marker) in [(id1, "first"), (id2, "second")] { let mut session = TunnelSession::adopt( &topo.consumer_manager, channel_id, Substrate::Tcp, alktunnels::TUNNEL_ALPN, ) .await .expect("adopt"); use tokio::io::{AsyncReadExt, AsyncWriteExt}; let (read, write) = session.stream_halves().expect("halves"); write.write_all(marker.as_bytes()).await.expect("write"); write.flush().await.expect("flush"); let mut buf = vec![0u8; marker.len()]; tokio::time::timeout(std::time::Duration::from_secs(5), read.read_exact(&mut buf)) .await .expect("round trip") .expect("read"); assert_eq!(buf, marker.as_bytes()); session.close().await; } tokio::time::sleep(std::time::Duration::from_millis(150)).await; assert!(only_channel_0(&topo)); } /// The listen-shape typed-error table (producer.md's mapping): /// unknown resource, empty-queue → resource_shortage, closed listener /// → dial_failed — each leaves no channel (the no-phantom property). #[tokio::test] async fn listen_typed_error_table_leaves_no_channel() { // Empty queue (the assembly's fail-fast posture: a bounded accept // budget maps empty-pop to the typed error). let empty: alktunnels::producer::AcceptFn = Arc::new(|| { Box::pin(async { Err(TunnelEstablishError::ResourceShortage("drained".into())) }) }); let registry = ResourceRegistry::new(); registry .register("listener", Substrate::Tcp, "assembly-owned") .await; let topo = wire_listen(registry, empty).await; let err = open_reverse_channel( &topo.consumer_call, ¶ms("listener", Substrate::Tcp), None, ) .await .expect_err("empty-queue open must fail"); assert_eq!( call_error_reason(&err).as_deref(), Some("resource_shortage") ); tokio::time::sleep(std::time::Duration::from_millis(100)).await; assert!(only_channel_0(&topo)); // Closed listener → dial_failed. let queue = AcceptQueue::new(); queue.close().await; let registry = ResourceRegistry::new(); registry .register("listener", Substrate::Tcp, "assembly-owned") .await; let topo = wire_listen(registry, accept_from_queue(queue.clone())).await; let err = open_reverse_channel( &topo.consumer_call, ¶ms("listener", Substrate::Tcp), None, ) .await .expect_err("closed-listener open must fail"); assert_eq!(call_error_reason(&err).as_deref(), Some("dial_failed")); tokio::time::sleep(std::time::Duration::from_millis(100)).await; assert!(only_channel_0(&topo)); // Unknown resource → unknown_resource (the namespace gate). let queue = AcceptQueue::new(); queue.push(echo_handle("unused")).await.expect("push"); let registry = ResourceRegistry::new(); registry .register("listener", Substrate::Tcp, "assembly-owned") .await; let topo = wire_listen(registry, accept_from_queue(queue.clone())).await; let err = open_reverse_channel(&topo.consumer_call, ¶ms("ghost", Substrate::Tcp), None) .await .expect_err("unknown resource must fail"); assert_eq!(call_error_reason(&err).as_deref(), Some("unknown_resource")); tokio::time::sleep(std::time::Duration::from_millis(100)).await; assert!(only_channel_0(&topo)); } // ===================================================================== // Suite 4 — Datagram (`udp`): forward + reverse via the codec // ===================================================================== /// Many datagrams in each direction, boundaries preserved, batched /// frames survive (one chunk may carry several datagrams — the /// pending-queue cover). #[tokio::test] async fn datagram_boundaries_preserved_both_directions() { let registry = ResourceRegistry::new(); registry.register("dns", Substrate::Udp, "in-process").await; let topo = wire_forward(registry, framed_udp_echo_dial()).await; let mut session = TunnelSession::open(&topo.consumer, params("dns", Substrate::Udp)) .await .expect("udp open"); // Burst 40 datagrams (varied sizes) — they may batch into shared // chunks; each must come back as exactly one datagram. let mut payloads = Vec::new(); for i in 0..40u16 { payloads.push(vec![i as u8; (i as usize * 137 % 500) + 1]); } for p in &payloads { session.send_datagram(p).await.expect("send"); } let mut received = Vec::new(); while received.len() < payloads.len() { let dg = tokio::time::timeout(std::time::Duration::from_secs(5), session.recv_datagram()) .await .expect("recv timed out") .expect("io ok") .expect("not eof"); received.push(dg.to_vec()); } assert_eq!(received, payloads, "boundaries preserved under batching"); // The reverse direction: the target echo only repeats — drive the // other way by sending again after draining; order holds. session.send_datagram(b"final").await.expect("send"); let dg = tokio::time::timeout(std::time::Duration::from_secs(5), session.recv_datagram()) .await .expect("recv timed out") .expect("io ok") .expect("not eof"); assert_eq!(&dg[..], b"final"); let reaped = session.close().await; assert!(reaped); tokio::time::sleep(std::time::Duration::from_millis(200)).await; assert!(only_channel_0(&topo)); } /// The codec survives maximal chunk splitting (the forward POC's /// chunk-split test, driven through a live session): datagrams /// chopped into awkward chunks still decode exactly. #[tokio::test] async fn datagrams_survive_chunk_splitting_codec_level() { // Codec-level probe (the mux may split anywhere): 100 datagrams // framed, chopped into 7-byte chunks, decoded incrementally. let mut stream = Vec::new(); for i in 0..100u16 { let payload = vec![i as u8; (i as usize % 300) + 1]; stream.extend_from_slice(&alktunnels::wire::frame_datagram(&payload).expect("frame")); } let mut reader = alktunnels::wire::DatagramReader::new(); let mut recovered = Vec::new(); for chunk in stream.chunks(7) { for dg in reader.feed(chunk).expect("decode") { recovered.push(dg.to_vec()); } } assert_eq!(recovered.len(), 100); for (i, dg) in recovered.iter().enumerate() { assert_eq!(dg.len(), (i % 300) + 1, "size preserved"); assert!(dg.iter().all(|&b| b == i as u8), "content preserved"); } } /// A reverse datagram session: the consumer pumps a UDP "accepted" /// socket against the codec-framed channel; datagrams round-trip via /// the pump (the reverse POC's UDP shape, codec-flavored). #[tokio::test] async fn reverse_udp_datagram_session_round_trips() { let registry = ResourceRegistry::new(); registry.register("dns", Substrate::Udp, "in-process").await; let topo = wire(registry, framed_udp_echo_dial()).await; let channel_id = open_reverse_channel(&topo.consumer_call, ¶ms("dns", Substrate::Udp), None) .await .expect("reverse udp open"); assert_eq!(channel_id % 2, 1); let mut session = TunnelSession::adopt( &topo.consumer_manager, channel_id, Substrate::Udp, alktunnels::TUNNEL_ALPN, ) .await .expect("adopt"); // No pump_against: the session's codec data plane is driven // directly (the pump-less datagram shape — pump_bidi copies raw // bytes, so the pump is the consumer's job here only when it // holds local halves; the direct drive is the session API's). session.send_datagram(b"query").await.expect("send"); let dg = tokio::time::timeout(std::time::Duration::from_secs(5), session.recv_datagram()) .await .expect("recv timed out") .expect("io ok") .expect("datagram"); assert_eq!(&dg[..], b"query"); let (_, _, reaped) = session.join().await; assert!(reaped, "pump-less join reaps (0, 0, reaped)"); tokio::time::sleep(std::time::Duration::from_millis(150)).await; assert!(only_channel_0(&topo)); } /// TCP + UDP sessions concurrently on ONE connection — the mux is /// transparent across substrates (both POCs' concurrent shape). #[tokio::test] async fn tcp_and_udp_concurrent_on_one_connection() { let registry = ResourceRegistry::new(); registry .register("echo", Substrate::Tcp, "in-process") .await; registry.register("dns", Substrate::Udp, "in-process").await; let topo = wire_forward(registry, harness::all_substrate_dial()).await; let mut tcp_session = TunnelSession::open(&topo.consumer, params("echo", Substrate::Tcp)) .await .expect("tcp open"); let mut udp_session = TunnelSession::open(&topo.consumer, params("dns", Substrate::Udp)) .await .expect("udp open"); assert_ne!(tcp_session.channel_id, udp_session.channel_id); // Interleaved traffic on both planes. use tokio::io::{AsyncReadExt, AsyncWriteExt}; let (tcp_read, tcp_write) = tcp_session.stream_halves().expect("tcp halves"); tcp_write.write_all(b"tcp-mixed").await.expect("tcp write"); tcp_write.flush().await.expect("tcp flush"); udp_session .send_datagram(b"udp-mixed") .await .expect("udp send"); let mut tbuf = vec![0u8; 9]; tokio::time::timeout( std::time::Duration::from_secs(5), tcp_read.read_exact(&mut tbuf), ) .await .expect("tcp round trip") .expect("tcp read"); assert_eq!(&tbuf, b"tcp-mixed"); let dg = tokio::time::timeout( std::time::Duration::from_secs(5), udp_session.recv_datagram(), ) .await .expect("udp recv timed out") .expect("udp io ok") .expect("udp datagram"); assert_eq!(&dg[..], b"udp-mixed"); let reaped = tcp_session.close().await; assert!(reaped); let reaped = udp_session.close().await; assert!(reaped); tokio::time::sleep(std::time::Duration::from_millis(200)).await; assert!(only_channel_0(&topo)); } // ===================================================================== // Suite 5 — Teardown matrix (ADR-005): close / join / Drop // ===================================================================== /// close on a session WITH a spawned pump: aborts the pump, reaps the /// entry on the consumer side; the producer side reaps via the /// wrapper. Both managers end on channel 0 only. #[tokio::test] async fn teardown_close_with_pump_reaps_both_sides() { let registry = ResourceRegistry::new(); registry .register("echo", Substrate::Tcp, "in-process") .await; let topo = wire(registry, echo_dial()).await; let channel_id = open_reverse_channel(&topo.consumer_call, ¶ms("echo", Substrate::Tcp), None) .await .expect("reverse open"); let session = TunnelSession::adopt( &topo.consumer_manager, channel_id, Substrate::Tcp, alktunnels::TUNNEL_ALPN, ) .await .expect("adopt"); let session = session.pump_against(tokio::io::duplex(64 * 1024).0).await; assert!(topo.consumer_manager.has_channel(channel_id)); let reaped = session.close().await; assert!(reaped, "close reaps the adopted entry"); tokio::time::sleep(std::time::Duration::from_millis(150)).await; assert!(only_channel_0(&topo)); } /// join on a session WITH a pump whose local ends drop: both pumps /// complete, copy counts are reported, both sides reaped (the /// two-pump contract observed from the session). #[tokio::test] async fn teardown_join_reports_copy_counts_and_reaps() { let registry = ResourceRegistry::new(); registry .register("echo", Substrate::Tcp, "in-process") .await; let topo = wire(registry, echo_dial()).await; let channel_id = open_reverse_channel(&topo.consumer_call, ¶ms("echo", Substrate::Tcp), None) .await .expect("reverse open"); let session = TunnelSession::adopt( &topo.consumer_manager, channel_id, Substrate::Tcp, alktunnels::TUNNEL_ALPN, ) .await .expect("adopt"); let (accepted_end, local_end) = tokio::io::duplex(64 * 1024); let session = session.pump_against(accepted_end).await; use tokio::io::{AsyncReadExt, AsyncWriteExt}; let (mut l_read, mut l_write) = tokio::io::split(local_end); l_write.write_all(b"joinme").await.expect("write"); l_write.flush().await.expect("flush"); let mut buf = vec![0u8; 6]; tokio::time::timeout( std::time::Duration::from_secs(5), l_read.read_exact(&mut buf), ) .await .expect("round trip") .expect("read"); assert_eq!(&buf, b"joinme"); drop(l_write); drop(l_read); let (c2p, p2c, reaped) = tokio::time::timeout(std::time::Duration::from_secs(5), session.join()) .await .expect("join timed out"); assert!(c2p > 0 && p2c > 0, "both pumps moved bytes: {c2p}, {p2c}"); assert!(reaped); tokio::time::sleep(std::time::Duration::from_millis(100)).await; assert!(only_channel_0(&topo)); } /// Drop after pump_against: the pump aborts AND the entry reaps — /// the full teardown, not just the pump (the drop matrix's /// pump-carrying arm). #[tokio::test] async fn teardown_drop_with_pump_reaps_entry() { let registry = ResourceRegistry::new(); registry .register("echo", Substrate::Tcp, "in-process") .await; let topo = wire(registry, echo_dial()).await; let channel_id = open_reverse_channel(&topo.consumer_call, ¶ms("echo", Substrate::Tcp), None) .await .expect("reverse open"); let session = TunnelSession::adopt( &topo.consumer_manager, channel_id, Substrate::Tcp, alktunnels::TUNNEL_ALPN, ) .await .expect("adopt"); let session = session.pump_against(tokio::io::duplex(64 * 1024).0).await; drop(session); assert!( !topo.consumer_manager.has_channel(channel_id), "Drop reaped the entry alongside the aborted pump" ); tokio::time::sleep(std::time::Duration::from_millis(150)).await; assert!(only_channel_0(&topo)); } /// A failed adopt leaks nothing: the manager's parked entry is /// released (the failed-adopt arm of the matrix) — a bogus-ID adopt /// SUCCEEDS by design (the manager parks any unseen id), so the /// collision path is the honest probe; after the probe the queue is /// clean on both sides. #[tokio::test] async fn teardown_failed_adopt_leaves_no_leak() { let registry = ResourceRegistry::new(); registry .register("echo", Substrate::Tcp, "in-process") .await; let topo = wire(registry, echo_dial()).await; let channel_id = open_reverse_channel(&topo.consumer_call, ¶ms("echo", Substrate::Tcp), None) .await .expect("reverse open"); let first = TunnelSession::adopt( &topo.consumer_manager, channel_id, Substrate::Tcp, alktunnels::TUNNEL_ALPN, ) .await .expect("first adopt"); let second = TunnelSession::adopt( &topo.consumer_manager, channel_id, Substrate::Tcp, alktunnels::TUNNEL_ALPN, ) .await; match second .err() .expect("adopting a live id twice must fail") .open_ref() { ChannelOpenError::AdoptFailed(_) => {} other => panic!("expected AdoptFailed, got {other:?}"), } first.close().await; tokio::time::sleep(std::time::Duration::from_millis(150)).await; assert!(only_channel_0(&topo)); } // ===================================================================== // Suite 6 — Spec-conformance assertions (kept visible) // ===================================================================== /// A failed open NEVER returns a channel id (the no-phantom-channel /// property, dial shape): unknown_resource + dial_failed both leave /// both sides' managers at channel 0 only, and the consumer's typed /// error exposes the establishment reason. #[tokio::test] async fn failed_open_never_yields_channel_id() { let registry = ResourceRegistry::new(); registry .register("echo", Substrate::Tcp, "in-process") .await; let topo = wire_forward(registry, echo_dial()).await; let err = match TunnelSession::open(&topo.consumer, params("ghost", Substrate::Tcp)).await { Ok(_) => panic!("unknown resource must not open"), Err(e) => e, }; assert_eq!( reason_of_open_failure(&err).as_deref(), Some("unknown_resource") ); assert!(only_channel_0(&topo)); let registry = ResourceRegistry::new(); registry .register("dead", Substrate::Tcp, "unreachable") .await; let topo = wire_forward(registry, failing_dial("no route")).await; let err = match TunnelSession::open(&topo.consumer, params("dead", Substrate::Tcp)).await { Ok(_) => panic!("failing dial must not open"), Err(e) => e, }; assert_eq!(reason_of_open_failure(&err).as_deref(), Some("dial_failed")); assert!(only_channel_0(&topo)); } /// The pump handler's JoinHandle tracks the data plane (R-02): the /// session-side pump completes when the local ends close (not /// teardown-at-birth), and copy counts prove bytes moved AFTER the /// open resolved. #[tokio::test] async fn pump_handle_tracks_the_data_plane_r02() { let registry = ResourceRegistry::new(); registry .register("echo", Substrate::Tcp, "in-process") .await; let topo = wire(registry, echo_dial()).await; let channel_id = open_reverse_channel(&topo.consumer_call, ¶ms("echo", Substrate::Tcp), None) .await .expect("reverse open"); let session = TunnelSession::adopt( &topo.consumer_manager, channel_id, Substrate::Tcp, alktunnels::TUNNEL_ALPN, ) .await .expect("adopt"); let (accepted_end, local_end) = tokio::io::duplex(64 * 1024); let session = session.pump_against(accepted_end).await; // The channel is alive and pumpable long after the open resolved // (the early-return shape would have torn it down at birth). use tokio::io::{AsyncReadExt, AsyncWriteExt}; let (mut l_read, mut l_write) = tokio::io::split(local_end); for round in 0..3 { let msg = format!("alive-{round}"); l_write.write_all(msg.as_bytes()).await.expect("write"); l_write.flush().await.expect("flush"); let mut buf = vec![0u8; msg.len()]; tokio::time::timeout( std::time::Duration::from_secs(5), l_read.read_exact(&mut buf), ) .await .expect("round trip") .expect("read"); assert_eq!(buf, msg.as_bytes()); } // Pump completion arrives only when BOTH directions end. drop(l_write); drop(l_read); let (c2p, p2c, reaped) = tokio::time::timeout(std::time::Duration::from_secs(5), session.join()) .await .expect("join timed out (teardown-at-birth would show here)"); assert!(c2p > 0 && p2c > 0 && reaped); assert!(only_channel_0(&topo)); } // ===================================================================== // Harness helpers (the accepted-handle stand-ins reused by suite 3) // ===================================================================== fn split_taken(taken: alktunnels::TakenHalves) -> Taken { Taken { read: taken.read, write: taken.write, session: taken.session, } } struct Taken { read: Box, write: Box, session: TunnelSession, } /// An in-process accepted handle: a duplex pair whose far end echoes /// (the same shape the dial harness echoes with — a listen producer's /// "accepted connection"). fn echo_handle(_label: &str) -> TargetHandle { let (consumer_side, target_side) = tokio::io::duplex(64 * 1024); let (mut t_read, mut t_write) = tokio::io::split(target_side); tokio::spawn(async move { use tokio::io::{AsyncReadExt, AsyncWriteExt}; let mut buf = vec![0u8; 8192]; loop { match t_read.read(&mut buf).await { Ok(0) | Err(_) => return, Ok(n) => { if t_write.write_all(&buf[..n]).await.is_err() { return; } } } } }); let (c_read, c_write) = tokio::io::split(consumer_side); TargetHandle { read: Box::new(c_read), write: Box::new(c_write), } } /// An accept closure popping from the queue (the assembly-side /// contract; a closed queue maps to `dial_failed`). fn accept_from_queue(queue: AcceptQueue) -> alktunnels::producer::AcceptFn { Arc::new(move || { let queue = queue.clone(); Box::pin(async move { queue .pop() .await .ok_or_else(|| TunnelEstablishError::DialFailed("listener closed".to_string())) }) }) }