From e937b1cfe89ccd782ee251b4fb30a8e4eaa6c147 Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Fri, 31 Jul 2026 16:50:04 -0600 Subject: [PATCH] Add U4 split tunnel pipe --- PROGRESS.md | 10 + core/crates/solstone-core-spl/src/lib.rs | 4 +- .../solstone-core-spl/src/loopback_pipe.rs | 298 +++++++++++++++++- 3 files changed, 304 insertions(+), 8 deletions(-) diff --git a/PROGRESS.md b/PROGRESS.md index 875de1a1a..b397ea75b 100644 --- a/PROGRESS.md +++ b/PROGRESS.md @@ -57,6 +57,11 @@ HPKE sealing is randomized; the reverse remains the supervisor's live differential. The short/long header rows are parser vectors, never socket-wire vectors. U6 may use the corpus when U1–U5 are accepted, but remains blocked on those orchestration units. +- The delegated U4 split WS-to-loopback pipe is accepted after fresh-eyes review. It replays + `drain_buffer()` before forwarding, uses the transport halves concurrently, limits TCP→WS + frames to 64 KiB, writes local EOF after a WS close, and cancels the opposite direction on + first completion. Its 93-SPL/17-HPKE crate gate is green; the complete relay client, + dispatch handoff, and live seam tests remain unaccepted. - Checkpoint gates after the correction-driven units: `cargo fmt --all -- --check`, strict combined clippy, and combined locked tests are green (85 SPL + 12 HPKE); both HPKE and SPL libraries pass the explicit `aarch64-apple-ios` check without an exclusion; `cargo deny` @@ -124,6 +129,11 @@ the seam's PKCS#8/SPKI boundary. This was a brief defect, not code rework. read and write halves. `BufferedWsReader` retains its accepted read-only ownership and U2 receives a sibling `WsByteSink` for `READY`/`ACK`/close. This was a contract defect, not a delegate rejection. +- **U3 source-close taxonomy:** `WsByteSource` exposes one `WsClosed` outcome for both clean + EOF and a source-side read failure. U4's split pipe can therefore preserve the shared + close-to-local-EOF behavior but cannot classify those two causes without changing an accepted + U3 seam. This is logged rather than silently expanded; no owner-visible behavior requires a + distinction today. ### Contract rework diff --git a/core/crates/solstone-core-spl/src/lib.rs b/core/crates/solstone-core-spl/src/lib.rs index 1bb4acb74..4959174af 100644 --- a/core/crates/solstone-core-spl/src/lib.rs +++ b/core/crates/solstone-core-spl/src/lib.rs @@ -55,7 +55,9 @@ pub use link_state_files::{ LinkServiceToken, LinkServiceTokenRead, LinkState, LinkStateRead, load_link_service_token, load_link_state, }; -pub use loopback_pipe::pipe_loopback; +pub use loopback_pipe::{ + TCP_TO_WS_READ_MAX, TunnelPipeError, TunnelPipeProgress, pipe_loopback, pipe_tunnel, +}; pub use posture_gate::{ PostureGate, PostureInput, RelayBlocked, RelayDecision, RelayPermit, ServiceToken, TokenInput, }; diff --git a/core/crates/solstone-core-spl/src/loopback_pipe.rs b/core/crates/solstone-core-spl/src/loopback_pipe.rs index a73cb1bf3..1dc17a416 100644 --- a/core/crates/solstone-core-spl/src/loopback_pipe.rs +++ b/core/crates/solstone-core-spl/src/loopback_pipe.rs @@ -5,14 +5,180 @@ use std::io; -use tokio::io::{AsyncRead, AsyncWrite, AsyncWriteExt}; +use bytes::Bytes; +use thiserror::Error; +use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt}; -/// Replay `initial_prefix` to the loopback listener, then forward both streams. +use crate::{BufferedWsReader, WsBufferError, WsByteSink, WsByteSource}; + +/// The largest local TCP read sent in one WebSocket message. +pub const TCP_TO_WS_READ_MAX: usize = 64 * 1024; + +/// Class-only failures at the relay-tunnel-to-loopback seam. +/// +/// The variants deliberately contain no transport details: a relay URL can +/// carry a service token, so callers must surface only this taxonomy rather +/// than formatting an underlying WebSocket or I/O error. +#[derive(Debug, Error, Clone, Copy, PartialEq, Eq)] +pub enum TunnelPipeError { + #[error("websocket read failed")] + WebSocketRead, + #[error("websocket write failed")] + WebSocketWrite, + #[error("loopback read failed")] + LoopbackRead, + #[error("loopback write failed")] + LoopbackWrite, + #[error("forwarded byte count overflowed")] + ByteCountOverflow, +} + +/// Directional byte counts observed before either task completed. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct TunnelPipeProgress { + /// Bytes copied from the WebSocket reader to the loopback TCP writer. + pub websocket_to_tcp: u64, + /// Bytes copied from the loopback TCP reader to the WebSocket sink. + pub tcp_to_websocket: u64, +} + +/// Replay buffered TLS bytes and forward a split WebSocket tunnel to TCP. +/// +/// The relay dispatcher may have peeked its TLS prefix already. That prefix +/// remains in `reader`, so draining the buffer before either forwarding task +/// starts is what preserves the local byte stream. The two directions then +/// race exactly as Python's `_pipe_tunnel` does: the first clean end or +/// genuine failure cancels the other direction. A clean WebSocket end writes +/// EOF to the local TCP writer before returning. +/// +/// # Errors +/// +/// Returns the class-only transport failure that ended the tunnel first. +pub async fn pipe_tunnel( + reader: &mut BufferedWsReader, + sink: &mut W, + loopback: Loopback, +) -> Result +where + R: WsByteSource, + W: WsByteSink, + Loopback: AsyncRead + AsyncWrite + Unpin, +{ + let initial = reader.drain_buffer(); + let (mut tcp_reader, mut tcp_writer) = tokio::io::split(loopback); + + write_loopback(&mut tcp_writer, &initial).await?; + + tokio::select! { + result = forward_websocket_to_tcp(reader, &mut tcp_writer) => { + let websocket_to_tcp = result?; + Ok(TunnelPipeProgress { + websocket_to_tcp, + tcp_to_websocket: 0, + }) + } + result = forward_tcp_to_websocket(&mut tcp_reader, sink) => { + let tcp_to_websocket = result?; + Ok(TunnelPipeProgress { + websocket_to_tcp: 0, + tcp_to_websocket, + }) + } + } +} + +async fn forward_websocket_to_tcp( + reader: &mut BufferedWsReader, + tcp_writer: &mut Writer, +) -> Result +where + R: WsByteSource, + Writer: AsyncWrite + Unpin, +{ + let mut transferred = 0_u64; + + loop { + let first = match reader.read_exactly(1).await { + Ok(bytes) => bytes, + Err(WsBufferError::Closed) => { + tcp_writer + .shutdown() + .await + .map_err(|_| TunnelPipeError::LoopbackWrite)?; + return Ok(transferred); + } + Err(WsBufferError::ReadTimeout | WsBufferError::ProgressTimeout) => { + return Err(TunnelPipeError::WebSocketRead); + } + }; + let buffered = reader.drain_buffer(); + + write_loopback(tcp_writer, &first).await?; + write_loopback(tcp_writer, &buffered).await?; + transferred = add_bytes(transferred, first.len())?; + transferred = add_bytes(transferred, buffered.len())?; + } +} + +async fn forward_tcp_to_websocket( + tcp_reader: &mut Reader, + sink: &mut W, +) -> Result +where + Reader: AsyncRead + Unpin, + W: WsByteSink, +{ + let mut buffer = vec![0_u8; TCP_TO_WS_READ_MAX].into_boxed_slice(); + let mut transferred = 0_u64; + + loop { + let count = tcp_reader + .read(&mut buffer) + .await + .map_err(|_| TunnelPipeError::LoopbackRead)?; + if count == 0 { + return Ok(transferred); + } + let data = buffer.get(..count).ok_or(TunnelPipeError::LoopbackRead)?; + sink.send(Bytes::copy_from_slice(data)) + .await + .map_err(|_| TunnelPipeError::WebSocketWrite)?; + transferred = add_bytes(transferred, count)?; + } +} + +async fn write_loopback(writer: &mut Writer, bytes: &[u8]) -> Result<(), TunnelPipeError> +where + Writer: AsyncWrite + Unpin, +{ + if bytes.is_empty() { + return Ok(()); + } + writer + .write_all(bytes) + .await + .map_err(|_| TunnelPipeError::LoopbackWrite)?; + writer + .flush() + .await + .map_err(|_| TunnelPipeError::LoopbackWrite) +} + +fn add_bytes(total: u64, count: usize) -> Result { + let count = u64::try_from(count).map_err(|_| TunnelPipeError::ByteCountOverflow)?; + total + .checked_add(count) + .ok_or(TunnelPipeError::ByteCountOverflow) +} + +/// Compatibility helper for unframed duplex streams. +/// +/// This remains useful for callers whose tunnel is not a WebSocket split. It +/// has the same prefix-before-forwarding guarantee as [`pipe_tunnel`]. /// -/// The prefix has already been consumed by the relay-side protocol classifier. -/// Replaying it before either forwarding direction starts preserves the original -/// local byte stream exactly. `copy_bidirectional` shuts down the opposite -/// writer after an EOF and returns the underlying I/O error unchanged. +/// # Errors +/// +/// Returns the first underlying local or tunnel stream I/O failure. pub async fn pipe_loopback( mut tunnel: Tunnel, mut loopback: Loopback, @@ -29,11 +195,129 @@ where #[cfg(test)] mod tests { + use std::future::Future; use std::io; + use bytes::Bytes; use tokio::io::{AsyncReadExt, AsyncWriteExt}; + use tokio::sync::mpsc; + + use crate::{BufferedWsReader, WsByteSink, WsByteSource, WsClosed}; + + use super::{TunnelPipeProgress, pipe_loopback, pipe_tunnel}; + + struct Frames { + frames: mpsc::UnboundedReceiver, + } + + impl Frames { + fn from_receiver(frames: mpsc::UnboundedReceiver) -> Self { + Self { frames } + } + + fn fixed(frames: &[&[u8]]) -> Self { + let (sender, receiver) = mpsc::unbounded_channel(); + for frame in frames { + let _ = sender.send(Bytes::copy_from_slice(frame)); + } + drop(sender); + Self::from_receiver(receiver) + } + } + + impl WsByteSource for Frames { + async fn next_message(&mut self) -> Result, WsClosed> { + Ok(self.frames.recv().await) + } + } + + struct RecordingSink { + sent: mpsc::UnboundedSender, + } + + impl WsByteSink for RecordingSink { + fn send(&mut self, bytes: Bytes) -> impl Future> + Send { + let result = self.sent.send(bytes).map_err(|_| WsClosed); + std::future::ready(result) + } + + fn close(&mut self) -> impl Future> + Send { + std::future::ready(Ok(())) + } + } + + #[tokio::test] + async fn replays_buffered_tls_bytes_first_and_forwards_both_directions() -> io::Result<()> { + let (source_send, source_receive) = mpsc::unbounded_channel(); + let mut reader = BufferedWsReader::new(Frames::from_receiver(source_receive)); + source_send + .send(Bytes::from_static(b"TLS-client-hello")) + .map_err(|_| io::Error::other("test source receiver missing"))?; + assert_eq!(reader.peek(1).await, Ok(Bytes::from_static(b"T"))); + + let (sink_send, mut sink_receive) = mpsc::unbounded_channel(); + let mut sink = RecordingSink { sent: sink_send }; + let (loopback, mut loopback_peer) = tokio::io::duplex(128); + + let pipe = tokio::spawn(async move { pipe_tunnel(&mut reader, &mut sink, loopback).await }); + + let mut initial = [0_u8; 16]; + loopback_peer.read_exact(&mut initial).await?; + assert_eq!(&initial, b"TLS-client-hello"); + + source_send + .send(Bytes::from_static(b"-continued")) + .map_err(|_| io::Error::other("test source receiver missing"))?; + let mut continued = [0_u8; 10]; + loopback_peer.read_exact(&mut continued).await?; + assert_eq!(&continued, b"-continued"); - use super::pipe_loopback; + loopback_peer.write_all(b"server-response").await?; + let sent = sink_receive + .recv() + .await + .ok_or_else(|| io::Error::other("test sink did not receive response"))?; + assert_eq!(sent, Bytes::from_static(b"server-response")); + + drop(source_send); + let progress = pipe + .await + .map_err(|error| io::Error::other(error.to_string()))? + .map_err(|error| io::Error::other(error.to_string()))?; + assert_eq!( + progress, + TunnelPipeProgress { + websocket_to_tcp: 10, + tcp_to_websocket: 0, + } + ); + Ok(()) + } + + #[tokio::test] + async fn websocket_eof_writes_tcp_eof_and_cancels_reverse_read() -> io::Result<()> { + let mut reader = BufferedWsReader::new(Frames::fixed(&[b"TLS"])); + assert_eq!(reader.peek(1).await, Ok(Bytes::from_static(b"T"))); + + let (sink_send, _sink_receive) = mpsc::unbounded_channel(); + let mut sink = RecordingSink { sent: sink_send }; + let (loopback, mut loopback_peer) = tokio::io::duplex(32); + + let pipe = tokio::spawn(async move { pipe_tunnel(&mut reader, &mut sink, loopback).await }); + + let mut prefix = [0_u8; 3]; + loopback_peer.read_exact(&mut prefix).await?; + assert_eq!(&prefix, b"TLS"); + let mut eof = [0_u8; 1]; + assert_eq!(loopback_peer.read(&mut eof).await?, 0); + + let progress = pipe + .await + .map_err(|error| io::Error::other(error.to_string()))? + .map_err(|error| io::Error::other(error.to_string()))?; + assert_eq!(progress, TunnelPipeProgress::default()); + Ok(()) + } #[tokio::test] async fn replays_initial_prefix_once_before_tunnel_payload() -> io::Result<()> { -- 2.51.2