From a2159f61b9d5630d6a43bfb142a7bb0fc49af916 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Thu, 20 Aug 2026 05:02:14 -0700 Subject: [PATCH 1/6] test(network): require bounded BiDi WebSocket opening response read --- ...bdriver_bidi_websocket_opening_response.rs | 182 ++++++++++++++++++ 1 file changed, 182 insertions(+) create mode 100644 crates/originweave-network/tests/webdriver_bidi_websocket_opening_response.rs diff --git a/crates/originweave-network/tests/webdriver_bidi_websocket_opening_response.rs b/crates/originweave-network/tests/webdriver_bidi_websocket_opening_response.rs new file mode 100644 index 000000000..17a66ff9e --- /dev/null +++ b/crates/originweave-network/tests/webdriver_bidi_websocket_opening_response.rs @@ -0,0 +1,182 @@ +use std::{ + io::{self, Read, Write}, + net::TcpListener, + thread, + time::Duration, +}; + +use originweave_core::WebDriverBiDiWebSocketEndpoint; +use originweave_network::{ + MAX_WEBSOCKET_OPENING_READ_TIMEOUT, WebDriverBiDiTcpConnectionPlan, + WebDriverBiDiWebSocketClientKey, WebDriverBiDiWebSocketHandshakePlan, + WebDriverBiDiWebSocketOpeningReadError, +}; + +const SESSION_ID: &str = "01234567-89ab-cdef-0123-456789abcdef"; +const RFC6455_SAMPLE_KEY: &str = "dGhlIHNhbXBsZSBub25jZQ=="; +const SERVER_OPENING_RESPONSE: &[u8] = b"HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=\r\n\r\n"; + +fn connect(endpoint: &str) -> originweave_network::WebDriverBiDiTcpConnection { + let admitted = WebDriverBiDiWebSocketEndpoint::new(endpoint); + assert!(admitted.is_ok(), "{admitted:?}"); + let Ok(admitted) = admitted else { + unreachable!("asserted valid endpoint") + }; + let correlated = admitted.correlate_session_id(SESSION_ID); + assert!(correlated.is_ok(), "{correlated:?}"); + let Ok(correlated) = correlated else { + unreachable!("asserted correlated endpoint") + }; + let target = correlated.into_explicit_connect_target(); + assert!(target.is_ok(), "{target:?}"); + let Ok(target) = target else { + unreachable!("asserted explicit target") + }; + let plan = WebDriverBiDiTcpConnectionPlan::new(target, Duration::from_secs(1), 1); + assert!(plan.is_ok(), "{plan:?}"); + let Ok(plan) = plan else { + unreachable!("asserted connection plan") + }; + let connection = plan.connect(); + assert!(connection.is_ok(), "{connection:?}"); + let Ok(connection) = connection else { + unreachable!("asserted loopback connection") + }; + connection +} + +fn read_opening_request(mut stream: &std::net::TcpStream) -> io::Result> { + stream.set_read_timeout(Some(Duration::from_secs(2)))?; + let mut request = Vec::new(); + let mut byte = [0_u8; 1]; + while !request.ends_with(b"\r\n\r\n") { + let count = stream.read(&mut byte)?; + if count == 0 { + break; + } + request.push(byte[0]); + } + stream.set_read_timeout(None)?; + Ok(request) +} + +fn opening_request_sent( + local_addr: std::net::SocketAddr, +) -> originweave_network::WebDriverBiDiWebSocketOpeningRequestSent { + let endpoint = format!("ws://{local_addr}/session/{SESSION_ID}"); + let connection = connect(&endpoint); + let key = WebDriverBiDiWebSocketClientKey::new(RFC6455_SAMPLE_KEY); + assert!(key.is_ok(), "{key:?}"); + let Ok(key) = key else { + unreachable!("asserted canonical sample key") + }; + let plan = WebDriverBiDiWebSocketHandshakePlan::new(connection, key); + assert!(plan.is_ok(), "{plan:?}"); + let Ok(plan) = plan else { + unreachable!("asserted plain WebSocket handshake plan") + }; + let sent = plan.write_opening_request(Duration::from_millis(500)); + assert!(sent.is_ok(), "{sent:?}"); + let Ok(sent) = sent else { + unreachable!("asserted bounded opening write") + }; + sent +} + +#[test] +fn bounded_opening_response_read_preserves_exact_transport_and_request_evidence() { + let listener = TcpListener::bind(("127.0.0.1", 0)); + assert!(listener.is_ok(), "{listener:?}"); + let Ok(listener) = listener else { + return; + }; + let local_addr = listener.local_addr(); + assert!(local_addr.is_ok(), "{local_addr:?}"); + let Ok(local_addr) = local_addr else { + return; + }; + let server = thread::spawn(move || -> io::Result> { + let accepted = listener.accept()?; + let mut stream = accepted.0; + let request = read_opening_request(&stream)?; + stream.write_all(SERVER_OPENING_RESPONSE)?; + Ok(request) + }); + + let sent = opening_request_sent(local_addr); + let request_byte_count = sent.request_byte_count(); + let write_timeout = sent.write_timeout(); + let read_timeout = Duration::from_millis(500); + let response = sent.read_opening_response(read_timeout); + assert!(response.is_ok(), "{response:?}"); + let Ok(response) = response else { + return; + }; + + assert_eq!(response.header_bytes(), SERVER_OPENING_RESPONSE); + assert_eq!(response.read_timeout(), read_timeout); + assert_eq!(response.request_byte_count(), request_byte_count); + assert_eq!(response.write_timeout(), write_timeout); + assert_eq!(response.client_key().as_str(), RFC6455_SAMPLE_KEY); + assert_eq!( + response.transport_evidence().verified_peer().socket_addr(), + local_addr + ); + assert_eq!( + response.transport_evidence().verified_peer().session_id(), + SESSION_ID + ); + assert_eq!(response.transport_evidence().attempt_number(), 1); + let debug = format!("{response:?}"); + assert!(debug.contains("WebDriverBiDiWebSocketOpeningResponseRead")); + assert!(!debug.contains(RFC6455_SAMPLE_KEY)); + + let server_result = server.join(); + assert!(server_result.is_ok(), "{server_result:?}"); + if let Ok(received) = server_result { + assert!(received.is_ok(), "{received:?}"); + if let Ok(received) = received { + assert!(received.starts_with(b"GET /session/")); + assert!(received.ends_with(b"\r\n\r\n")); + } + } +} + +#[test] +fn opening_response_read_rejects_zero_and_excessive_deadlines_before_reading() { + for timeout in [ + Duration::ZERO, + MAX_WEBSOCKET_OPENING_READ_TIMEOUT + Duration::from_nanos(1), + ] { + let listener = TcpListener::bind(("127.0.0.1", 0)); + assert!(listener.is_ok(), "{listener:?}"); + let Ok(listener) = listener else { + continue; + }; + let local_addr = listener.local_addr(); + assert!(local_addr.is_ok(), "{local_addr:?}"); + let Ok(local_addr) = local_addr else { + continue; + }; + let server = thread::spawn(move || -> io::Result> { + let accepted = listener.accept()?; + read_opening_request(&accepted.0) + }); + + let sent = opening_request_sent(local_addr); + let result = sent.read_opening_response(timeout); + assert!(matches!( + result, + Err(WebDriverBiDiWebSocketOpeningReadError::InvalidReadTimeout { + read_timeout, + maximum_timeout: MAX_WEBSOCKET_OPENING_READ_TIMEOUT, + }) if read_timeout == timeout + )); + + let server_result = server.join(); + assert!(server_result.is_ok(), "{server_result:?}"); + if let Ok(received) = server_result { + assert!(received.is_ok(), "{received:?}"); + } + } +} From 6fe37b62df4503662cf49e5b2f077266202cb576 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Thu, 20 Aug 2026 05:05:04 -0700 Subject: [PATCH 2/6] feat(network): bound WebSocket opening response reads --- ...bdriver_bidi_websocket_opening_response.rs | 608 ++++++++++++++++++ 1 file changed, 608 insertions(+) create mode 100644 crates/originweave-network/src/webdriver_bidi_websocket_opening_response.rs diff --git a/crates/originweave-network/src/webdriver_bidi_websocket_opening_response.rs b/crates/originweave-network/src/webdriver_bidi_websocket_opening_response.rs new file mode 100644 index 000000000..152461184 --- /dev/null +++ b/crates/originweave-network/src/webdriver_bidi_websocket_opening_response.rs @@ -0,0 +1,608 @@ +use std::{ + error::Error, + fmt, + io::{self, Read}, + net::TcpStream, + time::{Duration, Instant}, +}; + +use crate::{ + WebDriverBiDiTcpConnectionEvidence, WebDriverBiDiWebSocketClientKey, + WebDriverBiDiWebSocketOpeningRequestSent, +}; + +const OPENING_RESPONSE_TERMINATOR: &[u8; 4] = b"\r\n\r\n"; + +/// Maximum wall-clock budget accepted for reading one WebSocket opening-response header section. +/// +/// This is an OriginWeave resource-safety ceiling, not an RFC 6455 protocol limit. Callers may +/// choose any smaller nonzero deadline. +pub const MAX_WEBSOCKET_OPENING_READ_TIMEOUT: Duration = Duration::from_secs(5); + +/// Maximum number of bytes accepted before the opening-response header terminator is observed. +/// +/// This is an OriginWeave resource-safety budget, not an HTTP or RFC 6455 maximum. The terminating +/// `CRLF CRLF` bytes are included in the budget. +pub const MAX_WEBSOCKET_OPENING_RESPONSE_HEADER_BYTES: usize = 16 * 1024; + +/// Fail-closed errors while reading one bounded WebDriver BiDi WebSocket opening response. +#[derive(Debug)] +pub enum WebDriverBiDiWebSocketOpeningReadError { + /// The requested total read deadline was zero or above the reviewed resource ceiling. + InvalidReadTimeout { + /// Rejected caller-supplied deadline. + read_timeout: Duration, + /// Maximum reviewed deadline accepted by this boundary. + maximum_timeout: Duration, + }, + /// The monotonic total read deadline elapsed before the complete header section was received. + ReadDeadlineExceeded { + /// Number of response bytes received before the deadline elapsed. + bytes_read: usize, + }, + /// Applying the remaining operating-system read timeout failed. + ReadTimeoutConfigurationFailed { + /// Number of response bytes already received before configuration failed. + bytes_read: usize, + /// Underlying operating-system error. + source: io::Error, + }, + /// A bounded socket read reported timeout or would-block before the header section completed. + ReadTimedOut { + /// Number of response bytes received before the timed-out operation. + bytes_read: usize, + /// Underlying operating-system error. + source: io::Error, + }, + /// The peer closed the stream before the opening-response header section completed. + UnexpectedEof { + /// Number of response bytes received before EOF. + bytes_read: usize, + }, + /// The response exceeded the reviewed header-byte budget before its terminator was observed. + ResponseHeaderTooLarge { + /// Maximum reviewed header bytes, including the terminator. + maximum_bytes: usize, + }, + /// A non-recoverable socket read failed before the complete header section was received. + ReadFailed { + /// Number of response bytes received before the failure. + bytes_read: usize, + /// Underlying operating-system error. + source: io::Error, + }, + /// Clearing the operation-local socket read timeout failed after the header was received. + ReadTimeoutCleanupFailed { + /// Number of response bytes already received before cleanup failed. + bytes_read: usize, + /// Underlying operating-system error. + source: io::Error, + }, +} + +impl fmt::Display for WebDriverBiDiWebSocketOpeningReadError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::InvalidReadTimeout { .. } => formatter.write_str( + "WebDriver BiDi WebSocket opening read timeout is outside the reviewed bound", + ), + Self::ReadDeadlineExceeded { .. } => formatter.write_str( + "WebDriver BiDi WebSocket opening read exceeded its monotonic deadline", + ), + Self::ReadTimeoutConfigurationFailed { .. } => formatter.write_str( + "failed to configure the bounded WebDriver BiDi WebSocket opening read timeout", + ), + Self::ReadTimedOut { .. } => formatter.write_str( + "WebDriver BiDi WebSocket opening read timed out before the response header completed", + ), + Self::UnexpectedEof { .. } => formatter.write_str( + "WebDriver BiDi WebSocket peer closed before the opening response header completed", + ), + Self::ResponseHeaderTooLarge { .. } => formatter.write_str( + "WebDriver BiDi WebSocket opening response header exceeded the reviewed byte bound", + ), + Self::ReadFailed { .. } => formatter.write_str( + "WebDriver BiDi WebSocket opening response read failed before the header completed", + ), + Self::ReadTimeoutCleanupFailed { .. } => formatter.write_str( + "failed to clear the WebDriver BiDi WebSocket opening read timeout before handoff", + ), + } + } +} + +impl Error for WebDriverBiDiWebSocketOpeningReadError { + fn source(&self) -> Option<&(dyn Error + 'static)> { + match self { + Self::ReadTimeoutConfigurationFailed { source, .. } + | Self::ReadTimedOut { source, .. } + | Self::ReadFailed { source, .. } + | Self::ReadTimeoutCleanupFailed { source, .. } => Some(source), + Self::InvalidReadTimeout { .. } + | Self::ReadDeadlineExceeded { .. } + | Self::UnexpectedEof { .. } + | Self::ResponseHeaderTooLarge { .. } => None, + } + } +} + +/// A live verified stream after one bounded opening-response header section has been received. +/// +/// This state preserves the exact client key and transport/request evidence from the preceding +/// opening write together with the exact `CRLF CRLF`-terminated response header bytes. It does not +/// parse or accept the server handshake: `101 Switching Protocols`, `Upgrade`, `Connection`, and +/// `Sec-WebSocket-Accept` remain untrusted bytes until a later fail-closed validator consumes this +/// value. No WebSocket OPEN state, browser identity, policy decision, browser action, or Agent +/// authority is established here. +pub struct WebDriverBiDiWebSocketOpeningResponseRead { + sent: WebDriverBiDiWebSocketOpeningRequestSent, + header: Vec, + read_timeout: Duration, +} + +impl fmt::Debug for WebDriverBiDiWebSocketOpeningResponseRead { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("WebDriverBiDiWebSocketOpeningResponseRead") + .field("transport_evidence", self.sent.transport_evidence()) + .field( + "client_key", + &"", + ) + .field("request_byte_count", &self.sent.request_byte_count()) + .field("write_timeout", &self.sent.write_timeout()) + .field("header_byte_count", &self.header.len()) + .field("read_timeout", &self.read_timeout) + .finish() + } +} + +impl WebDriverBiDiWebSocketOpeningResponseRead { + /// Borrow the exact unparsed opening-response header bytes, including the final `CRLF CRLF`. + #[must_use] + pub fn header_bytes(&self) -> &[u8] { + &self.header + } + + /// Return the total read deadline configured for this opening response. + #[must_use] + pub const fn read_timeout(&self) -> Duration { + self.read_timeout + } + + /// Borrow the exact verified transport evidence retained from the opening request. + #[must_use] + pub const fn transport_evidence(&self) -> &WebDriverBiDiTcpConnectionEvidence { + self.sent.transport_evidence() + } + + /// Borrow the exact client key required by a later `Sec-WebSocket-Accept` validator. + #[must_use] + pub const fn client_key(&self) -> &WebDriverBiDiWebSocketClientKey { + self.sent.client_key() + } + + /// Return the exact number of opening-request bytes written before this read boundary. + #[must_use] + pub const fn request_byte_count(&self) -> usize { + self.sent.request_byte_count() + } + + /// Return the total write deadline used by the preceding opening request. + #[must_use] + pub const fn write_timeout(&self) -> Duration { + self.sent.write_timeout() + } + + pub(crate) fn into_parts(self) -> (WebDriverBiDiWebSocketOpeningRequestSent, Vec) { + (self.sent, self.header) + } +} + +impl WebDriverBiDiWebSocketOpeningRequestSent { + /// Read one complete bounded RFC 6455 server opening-response header section. + /// + /// The already-verified stream is consumed into the result on success. Zero and over-ceiling + /// deadlines fail closed before reading. Reads use one total monotonic deadline and retry only an + /// interrupted system call. The reader consumes exactly one byte at a time until `CRLF CRLF`, so + /// it never discards bytes belonging to the first WebSocket frame. A response that closes early, + /// times out, exceeds the product header budget, or raises any other I/O failure yields no success + /// evidence. Before success, the operation-local socket read timeout is cleared so the next stage + /// cannot inherit stale timeout authority. + pub fn read_opening_response( + mut self, + read_timeout: Duration, + ) -> Result + { + if read_timeout.is_zero() || read_timeout > MAX_WEBSOCKET_OPENING_READ_TIMEOUT { + return Err(WebDriverBiDiWebSocketOpeningReadError::InvalidReadTimeout { + read_timeout, + maximum_timeout: MAX_WEBSOCKET_OPENING_READ_TIMEOUT, + }); + } + + let mut now = Instant::now; + let header = read_response_header_with_clock(&mut self.stream, read_timeout, &mut now)?; + Ok(WebDriverBiDiWebSocketOpeningResponseRead { + sent: self, + header, + read_timeout, + }) + } +} + +trait OpeningResponseReader { + fn set_read_timeout(&self, timeout: Duration) -> io::Result<()>; + fn clear_read_timeout(&self) -> io::Result<()>; + fn read_response_byte(&mut self, byte: &mut [u8; 1]) -> io::Result; +} + +impl OpeningResponseReader for TcpStream { + fn set_read_timeout(&self, timeout: Duration) -> io::Result<()> { + TcpStream::set_read_timeout(self, Some(timeout)) + } + + fn clear_read_timeout(&self) -> io::Result<()> { + TcpStream::set_read_timeout(self, None) + } + + fn read_response_byte(&mut self, byte: &mut [u8; 1]) -> io::Result { + self.read(byte) + } +} + +fn read_response_header_with_clock( + reader: &mut dyn OpeningResponseReader, + read_timeout: Duration, + now: &mut dyn FnMut() -> Instant, +) -> Result, WebDriverBiDiWebSocketOpeningReadError> { + let deadline = now() + read_timeout; + let mut response = Vec::with_capacity(512); + let mut byte = [0_u8; 1]; + + loop { + let remaining = deadline.saturating_duration_since(now()); + if remaining.is_zero() { + return Err(WebDriverBiDiWebSocketOpeningReadError::ReadDeadlineExceeded { + bytes_read: response.len(), + }); + } + reader.set_read_timeout(remaining).map_err(|source| { + WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutConfigurationFailed { + bytes_read: response.len(), + source, + } + })?; + + match reader.read_response_byte(&mut byte) { + Ok(0) => { + return Err(WebDriverBiDiWebSocketOpeningReadError::UnexpectedEof { + bytes_read: response.len(), + }); + } + Ok(_) => { + response.push(byte[0]); + if deadline.saturating_duration_since(now()).is_zero() { + return Err(WebDriverBiDiWebSocketOpeningReadError::ReadDeadlineExceeded { + bytes_read: response.len(), + }); + } + if response.ends_with(OPENING_RESPONSE_TERMINATOR) { + reader.clear_read_timeout().map_err(|source| { + WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutCleanupFailed { + bytes_read: response.len(), + source, + } + })?; + return Ok(response); + } + if response.len() >= MAX_WEBSOCKET_OPENING_RESPONSE_HEADER_BYTES { + return Err( + WebDriverBiDiWebSocketOpeningReadError::ResponseHeaderTooLarge { + maximum_bytes: MAX_WEBSOCKET_OPENING_RESPONSE_HEADER_BYTES, + }, + ); + } + } + Err(source) if source.kind() == io::ErrorKind::Interrupted => {} + Err(source) + if matches!( + source.kind(), + io::ErrorKind::TimedOut | io::ErrorKind::WouldBlock + ) => + { + return Err(WebDriverBiDiWebSocketOpeningReadError::ReadTimedOut { + bytes_read: response.len(), + source, + }); + } + Err(source) => { + return Err(WebDriverBiDiWebSocketOpeningReadError::ReadFailed { + bytes_read: response.len(), + source, + }); + } + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::collections::VecDeque; + + #[derive(Debug)] + enum ReadStep { + Byte(u8), + Eof, + Error(io::ErrorKind), + } + + struct FakeReader { + steps: VecDeque, + set_timeout_error: Option, + clear_timeout_error: Option, + } + + impl FakeReader { + fn from_steps(steps: impl IntoIterator) -> Self { + Self { + steps: steps.into_iter().collect(), + set_timeout_error: None, + clear_timeout_error: None, + } + } + } + + impl OpeningResponseReader for FakeReader { + fn set_read_timeout(&self, _timeout: Duration) -> io::Result<()> { + match self.set_timeout_error { + Some(kind) => Err(io::Error::from(kind)), + None => Ok(()), + } + } + + fn clear_read_timeout(&self) -> io::Result<()> { + match self.clear_timeout_error { + Some(kind) => Err(io::Error::from(kind)), + None => Ok(()), + } + } + + fn read_response_byte(&mut self, byte: &mut [u8; 1]) -> io::Result { + match self.steps.pop_front() { + Some(ReadStep::Byte(value)) => { + byte[0] = value; + Ok(1) + } + Some(ReadStep::Eof) | None => Ok(0), + Some(ReadStep::Error(kind)) => Err(io::Error::from(kind)), + } + } + } + + fn fixed_clock(instants: impl IntoIterator) -> impl FnMut() -> Instant { + let mut instants: VecDeque<_> = instants.into_iter().collect(); + let fallback = Instant::now(); + move || instants.pop_front().unwrap_or(fallback) + } + + fn assert_source_kind( + error: &WebDriverBiDiWebSocketOpeningReadError, + expected: io::ErrorKind, + ) { + let source = error.source(); + assert!(source.is_some()); + if let Some(source) = source { + assert_eq!(source.downcast_ref::().map(io::Error::kind), Some(expected)); + } + } + + #[test] + fn bounded_reader_retries_interrupted_and_returns_exact_header() { + let mut reader = FakeReader::from_steps([ + ReadStep::Error(io::ErrorKind::Interrupted), + ReadStep::Byte(b'X'), + ReadStep::Byte(b'\r'), + ReadStep::Byte(b'\n'), + ReadStep::Byte(b'\r'), + ReadStep::Byte(b'\n'), + ]); + let start = Instant::now(); + let mut now = move || start; + let result = read_response_header_with_clock(&mut reader, Duration::from_secs(1), &mut now); + assert_eq!(result.as_deref(), Ok(&b"X\r\n\r\n"[..])); + } + + #[test] + fn bounded_reader_rejects_deadline_before_and_after_read() { + let start = Instant::now(); + let mut before_reader = FakeReader::from_steps([]); + let mut before_clock = fixed_clock([start, start + Duration::from_secs(1)]); + let before = read_response_header_with_clock( + &mut before_reader, + Duration::from_secs(1), + &mut before_clock, + ); + assert!(matches!( + before, + Err(WebDriverBiDiWebSocketOpeningReadError::ReadDeadlineExceeded { bytes_read: 0 }) + )); + + let mut after_reader = FakeReader::from_steps([ReadStep::Byte(b'X')]); + let mut after_clock = fixed_clock([ + start, + start, + start + Duration::from_secs(1), + ]); + let after = read_response_header_with_clock( + &mut after_reader, + Duration::from_secs(1), + &mut after_clock, + ); + assert!(matches!( + after, + Err(WebDriverBiDiWebSocketOpeningReadError::ReadDeadlineExceeded { bytes_read: 1 }) + )); + } + + #[test] + fn bounded_reader_classifies_configuration_timeout_eof_and_io_failures() { + let start = Instant::now(); + + let mut configuration_reader = FakeReader::from_steps([]); + configuration_reader.set_timeout_error = Some(io::ErrorKind::InvalidInput); + let mut configuration_clock = move || start; + let configuration = read_response_header_with_clock( + &mut configuration_reader, + Duration::from_secs(1), + &mut configuration_clock, + ); + assert!(matches!( + configuration, + Err(WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutConfigurationFailed { + bytes_read: 0, + .. + }) + )); + if let Err(error) = configuration { + assert_source_kind(&error, io::ErrorKind::InvalidInput); + } + + let mut timeout_reader = FakeReader::from_steps([ReadStep::Error(io::ErrorKind::TimedOut)]); + let mut timeout_clock = move || start; + let timeout = read_response_header_with_clock( + &mut timeout_reader, + Duration::from_secs(1), + &mut timeout_clock, + ); + assert!(matches!( + timeout, + Err(WebDriverBiDiWebSocketOpeningReadError::ReadTimedOut { bytes_read: 0, .. }) + )); + if let Err(error) = timeout { + assert_source_kind(&error, io::ErrorKind::TimedOut); + } + + let mut eof_reader = FakeReader::from_steps([ReadStep::Eof]); + let mut eof_clock = move || start; + let eof = read_response_header_with_clock( + &mut eof_reader, + Duration::from_secs(1), + &mut eof_clock, + ); + assert!(matches!( + eof, + Err(WebDriverBiDiWebSocketOpeningReadError::UnexpectedEof { bytes_read: 0 }) + )); + + let mut failure_reader = + FakeReader::from_steps([ReadStep::Error(io::ErrorKind::ConnectionReset)]); + let mut failure_clock = move || start; + let failure = read_response_header_with_clock( + &mut failure_reader, + Duration::from_secs(1), + &mut failure_clock, + ); + assert!(matches!( + failure, + Err(WebDriverBiDiWebSocketOpeningReadError::ReadFailed { bytes_read: 0, .. }) + )); + if let Err(error) = failure { + assert_source_kind(&error, io::ErrorKind::ConnectionReset); + } + } + + #[test] + fn bounded_reader_rejects_oversize_and_cleanup_failure() { + let start = Instant::now(); + let mut oversize_reader = FakeReader::from_steps( + std::iter::repeat_n(ReadStep::Byte(b'X'), MAX_WEBSOCKET_OPENING_RESPONSE_HEADER_BYTES), + ); + let mut oversize_clock = move || start; + let oversize = read_response_header_with_clock( + &mut oversize_reader, + Duration::from_secs(1), + &mut oversize_clock, + ); + assert!(matches!( + oversize, + Err(WebDriverBiDiWebSocketOpeningReadError::ResponseHeaderTooLarge { + maximum_bytes: MAX_WEBSOCKET_OPENING_RESPONSE_HEADER_BYTES, + }) + )); + + let mut cleanup_reader = FakeReader::from_steps([ + ReadStep::Byte(b'\r'), + ReadStep::Byte(b'\n'), + ReadStep::Byte(b'\r'), + ReadStep::Byte(b'\n'), + ]); + cleanup_reader.clear_timeout_error = Some(io::ErrorKind::InvalidInput); + let mut cleanup_clock = move || start; + let cleanup = read_response_header_with_clock( + &mut cleanup_reader, + Duration::from_secs(1), + &mut cleanup_clock, + ); + assert!(matches!( + cleanup, + Err(WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutCleanupFailed { + bytes_read: 4, + .. + }) + )); + if let Err(error) = cleanup { + assert_source_kind(&error, io::ErrorKind::InvalidInput); + } + } + + #[test] + fn read_error_display_and_source_contracts_are_complete() { + let errors = [ + WebDriverBiDiWebSocketOpeningReadError::InvalidReadTimeout { + read_timeout: Duration::ZERO, + maximum_timeout: MAX_WEBSOCKET_OPENING_READ_TIMEOUT, + }, + WebDriverBiDiWebSocketOpeningReadError::ReadDeadlineExceeded { bytes_read: 1 }, + WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutConfigurationFailed { + bytes_read: 1, + source: io::Error::from(io::ErrorKind::InvalidInput), + }, + WebDriverBiDiWebSocketOpeningReadError::ReadTimedOut { + bytes_read: 1, + source: io::Error::from(io::ErrorKind::TimedOut), + }, + WebDriverBiDiWebSocketOpeningReadError::UnexpectedEof { bytes_read: 1 }, + WebDriverBiDiWebSocketOpeningReadError::ResponseHeaderTooLarge { + maximum_bytes: MAX_WEBSOCKET_OPENING_RESPONSE_HEADER_BYTES, + }, + WebDriverBiDiWebSocketOpeningReadError::ReadFailed { + bytes_read: 1, + source: io::Error::from(io::ErrorKind::ConnectionReset), + }, + WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutCleanupFailed { + bytes_read: 1, + source: io::Error::from(io::ErrorKind::InvalidInput), + }, + ]; + + for error in errors { + assert!(!error.to_string().is_empty()); + match error { + WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutConfigurationFailed { .. } + | WebDriverBiDiWebSocketOpeningReadError::ReadTimedOut { .. } + | WebDriverBiDiWebSocketOpeningReadError::ReadFailed { .. } + | WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutCleanupFailed { .. } => { + assert!(error.source().is_some()); + } + WebDriverBiDiWebSocketOpeningReadError::InvalidReadTimeout { .. } + | WebDriverBiDiWebSocketOpeningReadError::ReadDeadlineExceeded { .. } + | WebDriverBiDiWebSocketOpeningReadError::UnexpectedEof { .. } + | WebDriverBiDiWebSocketOpeningReadError::ResponseHeaderTooLarge { .. } => { + assert!(error.source().is_none()); + } + } + } + } +} From d7b11b083f6ee869ddb1efbd0d6b651cf11fbffb Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Thu, 20 Aug 2026 05:05:20 -0700 Subject: [PATCH 3/6] feat(network): expose bounded opening response read --- crates/originweave-network/src/lib.rs | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 deletions(-) diff --git a/crates/originweave-network/src/lib.rs b/crates/originweave-network/src/lib.rs index fc72c341d..fb2ab526a 100644 --- a/crates/originweave-network/src/lib.rs +++ b/crates/originweave-network/src/lib.rs @@ -5,9 +5,10 @@ //! peers before exposing transport I/O, and emits credential-free evidence. //! It also bridges a session-correlated WebDriver BiDi loopback target from //! `originweave-core` into one bounded exact TCP connection, binds an RFC 6455 -//! opening request to that verified plain stream, and can write that exact request -//! under one bounded deadline without claiming a completed WebSocket handshake or -//! granting browser, WebSocket, TLS, policy, or Agent authority. +//! opening request to that verified plain stream, writes that exact request under +//! one bounded deadline, and can receive one bounded opening-response header +//! section without claiming a completed WebSocket handshake or granting browser, +//! WebSocket, TLS, policy, or Agent authority. #![forbid(unsafe_code)] #![deny(missing_docs)] @@ -15,6 +16,7 @@ mod connection; mod webdriver_bidi_connection; mod webdriver_bidi_websocket_handshake; +mod webdriver_bidi_websocket_opening_response; pub use connection::{ ConnectionPlan, DirectTcpConnection, MAX_CONNECT_TIMEOUT, MAX_CONNECTION_ATTEMPTS, @@ -29,3 +31,7 @@ pub use webdriver_bidi_websocket_handshake::{ WebDriverBiDiWebSocketHandshakeError, WebDriverBiDiWebSocketHandshakePlan, WebDriverBiDiWebSocketOpeningRequestSent, WebDriverBiDiWebSocketOpeningWriteError, }; +pub use webdriver_bidi_websocket_opening_response::{ + MAX_WEBSOCKET_OPENING_READ_TIMEOUT, MAX_WEBSOCKET_OPENING_RESPONSE_HEADER_BYTES, + WebDriverBiDiWebSocketOpeningReadError, WebDriverBiDiWebSocketOpeningResponseRead, +}; From bc7016c4b652a6775df408e66403327444769407 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Thu, 20 Aug 2026 05:11:35 -0700 Subject: [PATCH 4/6] fix(network): harden bounded opening response read --- ...bdriver_bidi_websocket_opening_response.rs | 106 +++++++++--------- 1 file changed, 56 insertions(+), 50 deletions(-) diff --git a/crates/originweave-network/src/webdriver_bidi_websocket_opening_response.rs b/crates/originweave-network/src/webdriver_bidi_websocket_opening_response.rs index 152461184..f917dfa09 100644 --- a/crates/originweave-network/src/webdriver_bidi_websocket_opening_response.rs +++ b/crates/originweave-network/src/webdriver_bidi_websocket_opening_response.rs @@ -193,10 +193,6 @@ impl WebDriverBiDiWebSocketOpeningResponseRead { pub const fn write_timeout(&self) -> Duration { self.sent.write_timeout() } - - pub(crate) fn into_parts(self) -> (WebDriverBiDiWebSocketOpeningRequestSent, Vec) { - (self.sent, self.header) - } } impl WebDriverBiDiWebSocketOpeningRequestSent { @@ -263,9 +259,11 @@ fn read_response_header_with_clock( loop { let remaining = deadline.saturating_duration_since(now()); if remaining.is_zero() { - return Err(WebDriverBiDiWebSocketOpeningReadError::ReadDeadlineExceeded { - bytes_read: response.len(), - }); + return Err( + WebDriverBiDiWebSocketOpeningReadError::ReadDeadlineExceeded { + bytes_read: response.len(), + }, + ); } reader.set_read_timeout(remaining).map_err(|source| { WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutConfigurationFailed { @@ -283,9 +281,11 @@ fn read_response_header_with_clock( Ok(_) => { response.push(byte[0]); if deadline.saturating_duration_since(now()).is_zero() { - return Err(WebDriverBiDiWebSocketOpeningReadError::ReadDeadlineExceeded { - bytes_read: response.len(), - }); + return Err( + WebDriverBiDiWebSocketOpeningReadError::ReadDeadlineExceeded { + bytes_read: response.len(), + }, + ); } if response.ends_with(OPENING_RESPONSE_TERMINATOR) { reader.clear_read_timeout().map_err(|source| { @@ -387,14 +387,14 @@ mod tests { move || instants.pop_front().unwrap_or(fallback) } - fn assert_source_kind( - error: &WebDriverBiDiWebSocketOpeningReadError, - expected: io::ErrorKind, - ) { + fn assert_source_kind(error: &WebDriverBiDiWebSocketOpeningReadError, expected: io::ErrorKind) { let source = error.source(); assert!(source.is_some()); if let Some(source) = source { - assert_eq!(source.downcast_ref::().map(io::Error::kind), Some(expected)); + assert_eq!( + source.downcast_ref::().map(io::Error::kind), + Some(expected) + ); } } @@ -430,11 +430,7 @@ mod tests { )); let mut after_reader = FakeReader::from_steps([ReadStep::Byte(b'X')]); - let mut after_clock = fixed_clock([ - start, - start, - start + Duration::from_secs(1), - ]); + let mut after_clock = fixed_clock([start, start, start + Duration::from_secs(1)]); let after = read_response_header_with_clock( &mut after_reader, Duration::from_secs(1), @@ -459,29 +455,33 @@ mod tests { &mut configuration_clock, ); assert!(matches!( - configuration, - Err(WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutConfigurationFailed { - bytes_read: 0, - .. - }) + &configuration, + Err( + WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutConfigurationFailed { + bytes_read: 0, + .. + } + ) )); if let Err(error) = configuration { assert_source_kind(&error, io::ErrorKind::InvalidInput); } - let mut timeout_reader = FakeReader::from_steps([ReadStep::Error(io::ErrorKind::TimedOut)]); - let mut timeout_clock = move || start; - let timeout = read_response_header_with_clock( - &mut timeout_reader, - Duration::from_secs(1), - &mut timeout_clock, - ); - assert!(matches!( - timeout, - Err(WebDriverBiDiWebSocketOpeningReadError::ReadTimedOut { bytes_read: 0, .. }) - )); - if let Err(error) = timeout { - assert_source_kind(&error, io::ErrorKind::TimedOut); + for kind in [io::ErrorKind::TimedOut, io::ErrorKind::WouldBlock] { + let mut timeout_reader = FakeReader::from_steps([ReadStep::Error(kind)]); + let mut timeout_clock = move || start; + let timeout = read_response_header_with_clock( + &mut timeout_reader, + Duration::from_secs(1), + &mut timeout_clock, + ); + assert!(matches!( + &timeout, + Err(WebDriverBiDiWebSocketOpeningReadError::ReadTimedOut { bytes_read: 0, .. }) + )); + if let Err(error) = timeout { + assert_source_kind(&error, kind); + } } let mut eof_reader = FakeReader::from_steps([ReadStep::Eof]); @@ -505,7 +505,7 @@ mod tests { &mut failure_clock, ); assert!(matches!( - failure, + &failure, Err(WebDriverBiDiWebSocketOpeningReadError::ReadFailed { bytes_read: 0, .. }) )); if let Err(error) = failure { @@ -517,7 +517,7 @@ mod tests { fn bounded_reader_rejects_oversize_and_cleanup_failure() { let start = Instant::now(); let mut oversize_reader = FakeReader::from_steps( - std::iter::repeat_n(ReadStep::Byte(b'X'), MAX_WEBSOCKET_OPENING_RESPONSE_HEADER_BYTES), + (0..MAX_WEBSOCKET_OPENING_RESPONSE_HEADER_BYTES).map(|_| ReadStep::Byte(b'X')), ); let mut oversize_clock = move || start; let oversize = read_response_header_with_clock( @@ -527,9 +527,11 @@ mod tests { ); assert!(matches!( oversize, - Err(WebDriverBiDiWebSocketOpeningReadError::ResponseHeaderTooLarge { - maximum_bytes: MAX_WEBSOCKET_OPENING_RESPONSE_HEADER_BYTES, - }) + Err( + WebDriverBiDiWebSocketOpeningReadError::ResponseHeaderTooLarge { + maximum_bytes: MAX_WEBSOCKET_OPENING_RESPONSE_HEADER_BYTES, + } + ) )); let mut cleanup_reader = FakeReader::from_steps([ @@ -546,11 +548,13 @@ mod tests { &mut cleanup_clock, ); assert!(matches!( - cleanup, - Err(WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutCleanupFailed { - bytes_read: 4, - .. - }) + &cleanup, + Err( + WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutCleanupFailed { + bytes_read: 4, + .. + } + ) )); if let Err(error) = cleanup { assert_source_kind(&error, io::ErrorKind::InvalidInput); @@ -589,8 +593,10 @@ mod tests { for error in errors { assert!(!error.to_string().is_empty()); - match error { - WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutConfigurationFailed { .. } + match &error { + WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutConfigurationFailed { + .. + } | WebDriverBiDiWebSocketOpeningReadError::ReadTimedOut { .. } | WebDriverBiDiWebSocketOpeningReadError::ReadFailed { .. } | WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutCleanupFailed { .. } => { From d3a0e4b6b462640e41692d72c93d803ed9cbe3c3 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Thu, 20 Aug 2026 05:15:47 -0700 Subject: [PATCH 5/6] test(network): avoid error equality in opening read assertion --- .../src/webdriver_bidi_websocket_opening_response.rs | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/crates/originweave-network/src/webdriver_bidi_websocket_opening_response.rs b/crates/originweave-network/src/webdriver_bidi_websocket_opening_response.rs index f917dfa09..c0a9a5911 100644 --- a/crates/originweave-network/src/webdriver_bidi_websocket_opening_response.rs +++ b/crates/originweave-network/src/webdriver_bidi_websocket_opening_response.rs @@ -411,7 +411,10 @@ mod tests { let start = Instant::now(); let mut now = move || start; let result = read_response_header_with_clock(&mut reader, Duration::from_secs(1), &mut now); - assert_eq!(result.as_deref(), Ok(&b"X\r\n\r\n"[..])); + assert!(result.is_ok(), "{result:?}"); + if let Ok(header) = result { + assert_eq!(header, b"X\r\n\r\n"); + } } #[test] @@ -594,9 +597,7 @@ mod tests { for error in errors { assert!(!error.to_string().is_empty()); match &error { - WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutConfigurationFailed { - .. - } + WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutConfigurationFailed { .. } | WebDriverBiDiWebSocketOpeningReadError::ReadTimedOut { .. } | WebDriverBiDiWebSocketOpeningReadError::ReadFailed { .. } | WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutCleanupFailed { .. } => { From 60eb9d43b4a1e413d3f5f4145892c26d6e55b0bf Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Thu, 20 Aug 2026 06:17:46 -0700 Subject: [PATCH 6/6] style(network): restore canonical rustfmt --- .../src/webdriver_bidi_websocket_opening_response.rs | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/crates/originweave-network/src/webdriver_bidi_websocket_opening_response.rs b/crates/originweave-network/src/webdriver_bidi_websocket_opening_response.rs index c0a9a5911..252f2b41c 100644 --- a/crates/originweave-network/src/webdriver_bidi_websocket_opening_response.rs +++ b/crates/originweave-network/src/webdriver_bidi_websocket_opening_response.rs @@ -597,7 +597,9 @@ mod tests { for error in errors { assert!(!error.to_string().is_empty()); match &error { - WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutConfigurationFailed { .. } + WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutConfigurationFailed { + .. + } | WebDriverBiDiWebSocketOpeningReadError::ReadTimedOut { .. } | WebDriverBiDiWebSocketOpeningReadError::ReadFailed { .. } | WebDriverBiDiWebSocketOpeningReadError::ReadTimeoutCleanupFailed { .. } => {