diff --git a/CHANGELOG.md b/CHANGELOG.md index c41fa3705..26192cdeb 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,7 @@ All notable changes to OriginWeave are documented in this file. The format follo ### Added +- Bounded RFC 6455 WebDriver BiDi opening-response validation on the exact peer-verified stream: it admits only HTTP/1.1 `101`, case-insensitive `Upgrade`/`Connection` tokens, and the client-key-correlated `Sec-WebSocket-Accept` value within monotonic time and header-size ceilings; it restores blocking mode and still does not implement WebSocket frames or grant browser/Agent authority. - Bounded WebDriver BiDi loopback TCP transport that consumes one exact no-DNS connect target, retries only explicitly recoverable local transport failures within repository timeout and attempt ceilings, exposes the stream only after operating-system peer inspection and exact peer verification, supports a consuming handoff of the original stream with typed credential-free peer/session/TLS and bounded-attempt evidence, preserves typed causal errors, and performs no DNS, proxy/PAC, process authentication, TLS, WebSocket, BiDi message, browser-action, or Agent-authority step. - Exact WebDriver BiDi socket-peer verification that consumes an approved no-DNS connect target, requires the observed IP address and port to match exactly, preserves the TLS requirement and exact correlated session id, and remains inert metadata that does not authenticate an OS process, does not negotiate TLS, perform a WebSocket handshake, or grant Agent authority. - Explicit no-DNS WebDriver BiDi loopback connection targets that derive exact IPv4/IPv6 loopback `SocketAddr` metadata from a session-correlated endpoint, reject `localhost` as requiring separately trusted name resolution, preserve the TLS requirement and exact session id, perform no socket I/O, and grant no Agent authority. @@ -64,7 +65,9 @@ All notable changes to OriginWeave are documented in this file. The format follo - Moved autonomous-agent Cargo targets and Python bytecode caches outside the proposed source tree and prefetched locked Cargo dependencies for offline verification. - Kept the real loopback WebDriver BiDi opening-write regression test fail-fast with test-only diagnostics, while explicitly covering successful and panicked server-thread handoffs so strict all-target Clippy and exact coverage remain clean. - Kept the loopback peer alive until opening-write timeout cleanup completes, removing a macOS close race that could report `EINVAL` after a successful request write without weakening production cleanup failures. +- Kept the invalid opening-response deadline fixture's accepted peer alive through opening-write cleanup, removing the same macOS `EINVAL` race from the integration coverage path. - Kept the revoked-stream fixture peer alive until local shutdown and fail-closed write classification complete, removing a macOS `ENOTCONN` race from the coverage path. +- Made opening-exchange tests wait for the complete client request and retain the peer until each client assertion finishes, avoiding premature connection closure in both successful and rejected handshakes without changing production error handling. - Updated research doctoring to pin Chromium canonicalizer evidence to an immutable revision, add RFC 9293, RFC 5280, RFC 8446, RFC 9525, rustls 0.23.42, and Rust `TcpStream` evidence, distinguish the April 2026 Fugu beta from the June 2026 release, and treat vendor benchmark claims as first-party evidence rather than independent validation. ### Security diff --git a/Cargo.lock b/Cargo.lock index e2ada3c4e..90b2ed7c5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -285,8 +285,10 @@ dependencies = [ name = "originweave-network" version = "0.1.0" dependencies = [ + "base64", "originweave-core", "originweave-destination", + "sha1", ] [[package]] @@ -448,6 +450,17 @@ dependencies = [ "syn 3.0.3", ] +[[package]] +name = "sha1" +version = "0.10.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a978451301f4db1d02937a4ab3ccce137717b81826e79b7d49ffe3244a13c3b8" +dependencies = [ + "cfg-if", + "cpufeatures", + "digest", +] + [[package]] name = "sha2" version = "0.10.9" diff --git a/crates/originweave-network/Cargo.toml b/crates/originweave-network/Cargo.toml index 3d800d8de..68ac4fc8c 100644 --- a/crates/originweave-network/Cargo.toml +++ b/crates/originweave-network/Cargo.toml @@ -11,8 +11,10 @@ homepage.workspace = true publish = false [dependencies] +base64 = "0.22.1" originweave-core = { path = "../originweave-core" } originweave-destination = { path = "../originweave-destination" } +sha1 = "0.10.6" [lints] workspace = true diff --git a/crates/originweave-network/src/lib.rs b/crates/originweave-network/src/lib.rs index bf8c21ca9..e20c86d1f 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, can write that exact request +//! under one bounded deadline, and validates its bounded RFC 6455 opening response +//! without implementing WebSocket framing or granting browser, WebSocket, TLS, +//! policy, or Agent authority. #![forbid(unsafe_code)] #![deny(missing_docs)] @@ -26,8 +27,10 @@ pub use webdriver_bidi_connection::{ WebDriverBiDiTcpConnectionEvidence, WebDriverBiDiTcpConnectionPlan, }; pub use webdriver_bidi_websocket_handshake::{ + MAX_WEBSOCKET_OPENING_RESPONSE_SIZE, MAX_WEBSOCKET_OPENING_RESPONSE_TIMEOUT, MAX_WEBSOCKET_OPENING_WRITE_TIMEOUT, WebDriverBiDiWebSocketClientKey, - WebDriverBiDiWebSocketHandshakeError, WebDriverBiDiWebSocketHandshakePlan, + WebDriverBiDiWebSocketEstablished, WebDriverBiDiWebSocketHandshakeError, + WebDriverBiDiWebSocketHandshakePlan, WebDriverBiDiWebSocketHandshakeResponseError, WebDriverBiDiWebSocketOpeningRequestSent, WebDriverBiDiWebSocketOpeningWriteError, }; pub use webdriver_bidi_websocket_opening_recovery::WebDriverBiDiWebSocketOpeningWriteRecoveryDisposition; diff --git a/crates/originweave-network/src/webdriver_bidi_websocket_handshake.rs b/crates/originweave-network/src/webdriver_bidi_websocket_handshake.rs index 822f4154c..f05443f99 100644 --- a/crates/originweave-network/src/webdriver_bidi_websocket_handshake.rs +++ b/crates/originweave-network/src/webdriver_bidi_websocket_handshake.rs @@ -1,16 +1,21 @@ use std::{ error::Error, fmt, - io::{self, Write}, + io::{self, Read, Write}, net::TcpStream, + thread, time::{Duration, Instant}, }; +use base64::{Engine, engine::general_purpose::STANDARD}; use originweave_core::VerifiedWebDriverBiDiSocketPeer; +use sha1::{Digest, Sha1}; use crate::{WebDriverBiDiTcpConnection, WebDriverBiDiTcpConnectionEvidence}; const WEBSOCKET_CLIENT_KEY_LENGTH: usize = 24; +const RFC6455_WEBSOCKET_GUID: &[u8] = b"258EAFA5-E914-47DA-95CA-C5AB0DC85B11"; +const MAX_WEBSOCKET_OPENING_RESPONSE_BYTES: usize = 16 * 1024; const REDACTED_WEBSOCKET_CLIENT_NONCE: &str = ""; /// Maximum wall-clock budget accepted for writing one bounded WebSocket opening request. @@ -19,6 +24,18 @@ const REDACTED_WEBSOCKET_CLIENT_NONCE: &str = " /// already bounded before this budget is applied. Callers may choose any smaller nonzero deadline. pub const MAX_WEBSOCKET_OPENING_WRITE_TIMEOUT: Duration = Duration::from_secs(5); +/// Maximum wall-clock budget accepted for reading one bounded WebSocket opening response. +/// +/// 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_RESPONSE_TIMEOUT: Duration = Duration::from_secs(5); + +/// Maximum bytes admitted while reading one WebSocket HTTP opening response. +/// +/// The response is consumed only through its terminating `CRLF CRLF`; WebSocket frames are not +/// read or interpreted by this boundary. +pub const MAX_WEBSOCKET_OPENING_RESPONSE_SIZE: usize = MAX_WEBSOCKET_OPENING_RESPONSE_BYTES; + fn is_base64_data_byte(byte: u8) -> bool { byte.is_ascii_alphanumeric() || matches!(byte, b'+' | b'/') } @@ -266,6 +283,515 @@ impl WebDriverBiDiWebSocketOpeningRequestSent { pub const fn write_timeout(&self) -> Duration { self.write_timeout } + + /// Read and validate the bounded RFC 6455 server opening response on this exact stream. + /// + /// Success proves only an HTTP/1.1 `101 Switching Protocols` response with the required + /// `Upgrade`, `Connection`, and client-key-correlated `Sec-WebSocket-Accept` headers. The + /// response body, WebSocket frames, browser process identity, TLS, and browser/Agent authority + /// remain separate boundaries. + pub fn read_opening_response( + self, + response_timeout: Duration, + ) -> Result + { + if response_timeout.is_zero() || response_timeout > MAX_WEBSOCKET_OPENING_RESPONSE_TIMEOUT { + return Err( + WebDriverBiDiWebSocketHandshakeResponseError::InvalidResponseTimeout { + response_timeout, + maximum_timeout: MAX_WEBSOCKET_OPENING_RESPONSE_TIMEOUT, + }, + ); + } + + let Self { + mut stream, + transport_evidence, + client_key, + request_byte_count, + write_timeout, + } = self; + let mut now = Instant::now; + let (response_status, response_byte_count) = + read_opening_response_with_clock(&mut stream, &client_key, response_timeout, &mut now)?; + + Ok(WebDriverBiDiWebSocketEstablished { + stream, + transport_evidence, + client_key, + response_status, + response_byte_count, + response_timeout, + request_byte_count, + write_timeout, + }) + } +} + +/// A live verified stream after both RFC 6455 opening messages were validated. +/// +/// This state does not implement WebSocket framing or grant browser, page, policy, or Agent +/// authority. It retains the exact transport evidence and client key so later protocol stages can +/// remain correlated with the verified peer and opening handshake. +pub struct WebDriverBiDiWebSocketEstablished { + pub(crate) stream: TcpStream, + transport_evidence: WebDriverBiDiTcpConnectionEvidence, + client_key: WebDriverBiDiWebSocketClientKey, + response_status: u16, + response_byte_count: usize, + response_timeout: Duration, + request_byte_count: usize, + write_timeout: Duration, +} + +impl fmt::Debug for WebDriverBiDiWebSocketEstablished { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("WebDriverBiDiWebSocketEstablished") + .field("stream_local_addr", &self.stream.local_addr().ok()) + .field("transport_evidence", &self.transport_evidence) + .field( + "client_key", + &"", + ) + .field("response_status", &self.response_status) + .field("response_byte_count", &self.response_byte_count) + .field("response_timeout", &self.response_timeout) + .field("request_byte_count", &self.request_byte_count) + .field("write_timeout", &self.write_timeout) + .finish() + } +} + +impl WebDriverBiDiWebSocketEstablished { + /// Borrow the exact verified transport evidence retained with this live stream. + #[must_use] + pub const fn transport_evidence(&self) -> &WebDriverBiDiTcpConnectionEvidence { + &self.transport_evidence + } + + /// Borrow the exact client key correlated with the validated server accept value. + #[must_use] + pub const fn client_key(&self) -> &WebDriverBiDiWebSocketClientKey { + &self.client_key + } + + /// Return the validated HTTP status code, currently always `101` on success. + #[must_use] + pub const fn response_status(&self) -> u16 { + self.response_status + } + + /// Return the number of HTTP opening-response bytes consumed through its header terminator. + #[must_use] + pub const fn response_byte_count(&self) -> usize { + self.response_byte_count + } + + /// Return the total response deadline configured for this opening response. + #[must_use] + pub const fn response_timeout(&self) -> Duration { + self.response_timeout + } + + /// Return the number of request bytes written before the response was read. + #[must_use] + pub const fn request_byte_count(&self) -> usize { + self.request_byte_count + } + + /// Return the total write deadline configured for the preceding opening request. + #[must_use] + pub const fn write_timeout(&self) -> Duration { + self.write_timeout + } +} + +/// Fail-closed errors while reading one bounded WebDriver BiDi WebSocket opening response. +#[derive(Debug)] +pub enum WebDriverBiDiWebSocketHandshakeResponseError { + /// The requested total response deadline was zero or above the reviewed resource ceiling. + InvalidResponseTimeout { + /// Rejected caller-supplied deadline. + response_timeout: Duration, + /// Maximum reviewed deadline accepted by this boundary. + maximum_timeout: Duration, + }, + /// The monotonic total response deadline elapsed before validation completed. + ResponseDeadlineExceeded { + /// Number of response bytes consumed before the deadline elapsed. + bytes_read: usize, + }, + /// The response exceeded the reviewed header-size ceiling before its terminator was found. + ResponseTooLarge { + /// Number of response bytes consumed before rejection. + bytes_read: usize, + /// Maximum response bytes admitted by this boundary. + maximum_bytes: usize, + }, + /// Applying the operation-local nonblocking read mode failed. + ResponseReadModeConfigurationFailed { + /// Number of response bytes consumed before configuration failed. + bytes_read: usize, + /// Underlying operating-system error. + source: io::Error, + }, + /// A bounded socket read timed out before the opening response was complete. + ResponseReadTimedOut { + /// Number of response bytes consumed before the timed-out operation. + bytes_read: usize, + /// Underlying operating-system error. + source: io::Error, + }, + /// A non-recoverable socket read failed before the opening response was complete. + ResponseReadFailed { + /// Number of response bytes consumed before the failure. + bytes_read: usize, + /// Underlying operating-system error. + source: io::Error, + }, + /// The peer closed the stream before sending a complete HTTP header block. + ResponseEndedBeforeHeaders { + /// Number of response bytes consumed before the peer closed the stream. + bytes_read: usize, + }, + /// The HTTP response was not a valid, required WebSocket opening response. + MalformedResponse { + /// Stable, non-secret reason for the rejected response shape. + reason: &'static str, + }, + /// The response's `Sec-WebSocket-Accept` did not correlate with the sent client key. + AcceptMismatch, + /// Restoring blocking mode failed after validation. + ReadModeCleanupFailed { + /// Underlying operating-system error. + source: io::Error, + }, +} + +impl fmt::Display for WebDriverBiDiWebSocketHandshakeResponseError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::InvalidResponseTimeout { .. } => formatter.write_str( + "WebDriver BiDi WebSocket opening response timeout is outside the reviewed bound", + ), + Self::ResponseDeadlineExceeded { .. } => formatter.write_str( + "WebDriver BiDi WebSocket opening response exceeded its monotonic deadline", + ), + Self::ResponseTooLarge { .. } => formatter.write_str( + "WebDriver BiDi WebSocket opening response exceeded its bounded header size", + ), + Self::ResponseReadModeConfigurationFailed { .. } => formatter.write_str( + "failed to configure bounded nonblocking WebDriver BiDi WebSocket response reads", + ), + Self::ResponseReadTimedOut { .. } => formatter.write_str( + "WebDriver BiDi WebSocket opening response timed out before completion", + ), + Self::ResponseReadFailed { .. } => formatter.write_str( + "WebDriver BiDi WebSocket opening response read failed before completion", + ), + Self::ResponseEndedBeforeHeaders { .. } => formatter.write_str( + "WebDriver BiDi WebSocket peer ended the stream before completing response headers", + ), + Self::MalformedResponse { .. } => formatter.write_str( + "WebDriver BiDi WebSocket opening response was malformed or missing a required header", + ), + Self::AcceptMismatch => formatter.write_str( + "WebDriver BiDi WebSocket opening response accept value did not match the client key", + ), + Self::ReadModeCleanupFailed { .. } => formatter.write_str( + "failed to restore blocking WebDriver BiDi WebSocket response reads before handoff", + ), + } + } +} + +impl Error for WebDriverBiDiWebSocketHandshakeResponseError { + fn source(&self) -> Option<&(dyn Error + 'static)> { + match self { + Self::ResponseReadModeConfigurationFailed { source, .. } + | Self::ResponseReadTimedOut { source, .. } + | Self::ResponseReadFailed { source, .. } + | Self::ReadModeCleanupFailed { source } => Some(source), + Self::InvalidResponseTimeout { .. } + | Self::ResponseDeadlineExceeded { .. } + | Self::ResponseTooLarge { .. } + | Self::ResponseEndedBeforeHeaders { .. } + | Self::MalformedResponse { .. } + | Self::AcceptMismatch => None, + } + } +} + +struct ParsedOpeningResponse { + status_code: u16, + byte_count: usize, +} + +fn expected_accept_value(client_key: &WebDriverBiDiWebSocketClientKey) -> String { + let mut digest = Sha1::new(); + digest.update(client_key.as_str().as_bytes()); + digest.update(RFC6455_WEBSOCKET_GUID); + STANDARD.encode(digest.finalize()) +} + +fn is_http_token_byte(byte: u8) -> bool { + byte.is_ascii_alphanumeric() + || matches!( + byte, + b'!' | b'#' + | b'$' + | b'%' + | b'&' + | b'\'' + | b'*' + | b'+' + | b'-' + | b'.' + | b'^' + | b'_' + | b'`' + | b'|' + | b'~' + ) +} + +fn has_header_token(value: &str, expected: &str) -> bool { + value + .split(',') + .map(str::trim) + .any(|token| token.eq_ignore_ascii_case(expected)) +} + +#[allow(clippy::collapsible_if)] +fn parse_opening_response( + response: &[u8], + client_key: &WebDriverBiDiWebSocketClientKey, +) -> Result { + if !response.ends_with(b"\r\n\r\n") { + return Err( + WebDriverBiDiWebSocketHandshakeResponseError::MalformedResponse { + reason: "response is missing its CRLF header terminator", + }, + ); + } + let response_text = String::from_utf8_lossy(response); + let header_text = &response_text[..response_text.len() - 4]; + let (status_line, header_lines) = header_text + .split_once("\r\n") + .map_or((header_text, ""), |(line, rest)| (line, rest)); + if status_line.bytes().any(|byte| byte < 0x20 || byte == 0x7f) { + return Err( + WebDriverBiDiWebSocketHandshakeResponseError::MalformedResponse { + reason: "status line contains a control byte", + }, + ); + } + if !status_line.starts_with("HTTP/1.1 101 ") { + return Err( + WebDriverBiDiWebSocketHandshakeResponseError::MalformedResponse { + reason: "status line is not HTTP/1.1 101", + }, + ); + } + + let mut upgrade_has_websocket = false; + let mut connection_has_upgrade = false; + let mut accept = None; + for line in header_lines.split("\r\n") { + if line.is_empty() + || line + .as_bytes() + .first() + .is_some_and(|byte| matches!(byte, b' ' | b'\t')) + { + return Err( + WebDriverBiDiWebSocketHandshakeResponseError::MalformedResponse { + reason: "header line is empty or folded", + }, + ); + } + let (name, value) = line.split_once(':').ok_or( + WebDriverBiDiWebSocketHandshakeResponseError::MalformedResponse { + reason: "header line has no colon", + }, + )?; + if name.is_empty() || !name.bytes().all(is_http_token_byte) { + return Err( + WebDriverBiDiWebSocketHandshakeResponseError::MalformedResponse { + reason: "header name is not an HTTP token", + }, + ); + } + let value = value.trim_matches([' ', '\t']); + if value.bytes().any(|byte| byte < 0x20 || byte == 0x7f) { + return Err( + WebDriverBiDiWebSocketHandshakeResponseError::MalformedResponse { + reason: "header value contains a control byte", + }, + ); + } + if name.eq_ignore_ascii_case("upgrade") { + upgrade_has_websocket |= has_header_token(value, "websocket"); + } else if name.eq_ignore_ascii_case("connection") { + connection_has_upgrade |= has_header_token(value, "upgrade"); + } else if name.eq_ignore_ascii_case("sec-websocket-accept") { + if accept.is_some() { + return Err( + WebDriverBiDiWebSocketHandshakeResponseError::MalformedResponse { + reason: "response repeats the Sec-WebSocket-Accept header", + }, + ); + } + accept = Some(value); + } + } + + if !upgrade_has_websocket { + return Err( + WebDriverBiDiWebSocketHandshakeResponseError::MalformedResponse { + reason: "Upgrade header does not contain websocket", + }, + ); + } + if !connection_has_upgrade { + return Err( + WebDriverBiDiWebSocketHandshakeResponseError::MalformedResponse { + reason: "Connection header does not contain Upgrade", + }, + ); + } + let Some(accept) = accept else { + return Err( + WebDriverBiDiWebSocketHandshakeResponseError::MalformedResponse { + reason: "response has no Sec-WebSocket-Accept header", + }, + ); + }; + if accept != expected_accept_value(client_key) { + return Err(WebDriverBiDiWebSocketHandshakeResponseError::AcceptMismatch); + } + + Ok(ParsedOpeningResponse { + status_code: 101, + byte_count: response.len(), + }) +} + +trait OpeningResponseReader { + fn set_nonblocking(&self, nonblocking: bool) -> io::Result<()>; + fn read_response_bytes(&mut self, bytes: &mut [u8]) -> io::Result; +} + +impl OpeningResponseReader for TcpStream { + fn set_nonblocking(&self, nonblocking: bool) -> io::Result<()> { + TcpStream::set_nonblocking(self, nonblocking) + } + + fn read_response_bytes(&mut self, bytes: &mut [u8]) -> io::Result { + self.read(bytes) + } +} + +fn read_opening_response_with_clock( + reader: &mut dyn OpeningResponseReader, + client_key: &WebDriverBiDiWebSocketClientKey, + response_timeout: Duration, + now: &mut dyn FnMut() -> Instant, +) -> Result<(u16, usize), WebDriverBiDiWebSocketHandshakeResponseError> { + let deadline = now() + response_timeout; + let mut response = Vec::new(); + + reader.set_nonblocking(true).map_err(|source| { + WebDriverBiDiWebSocketHandshakeResponseError::ResponseReadModeConfigurationFailed { + bytes_read: 0, + source, + } + })?; + + loop { + let remaining = deadline.saturating_duration_since(now()); + if remaining.is_zero() { + return Err( + WebDriverBiDiWebSocketHandshakeResponseError::ResponseDeadlineExceeded { + bytes_read: response.len(), + }, + ); + } + if response.len() >= MAX_WEBSOCKET_OPENING_RESPONSE_BYTES { + return Err( + WebDriverBiDiWebSocketHandshakeResponseError::ResponseTooLarge { + bytes_read: response.len(), + maximum_bytes: MAX_WEBSOCKET_OPENING_RESPONSE_BYTES, + }, + ); + } + let mut byte = [0_u8; 1]; + match reader.read_response_bytes(&mut byte) { + Ok(0) => { + return Err( + WebDriverBiDiWebSocketHandshakeResponseError::ResponseEndedBeforeHeaders { + bytes_read: response.len(), + }, + ); + } + Ok(1) => { + response.push(byte[0]); + if response.ends_with(b"\r\n\r\n") { + if deadline.saturating_duration_since(now()).is_zero() { + return Err( + WebDriverBiDiWebSocketHandshakeResponseError::ResponseDeadlineExceeded { + bytes_read: response.len(), + }, + ); + } + let parsed = parse_opening_response(&response, client_key)?; + reader.set_nonblocking(false).map_err(|source| { + WebDriverBiDiWebSocketHandshakeResponseError::ReadModeCleanupFailed { + source, + } + })?; + return Ok((parsed.status_code, parsed.byte_count)); + } + } + Ok(_) => { + return Err( + WebDriverBiDiWebSocketHandshakeResponseError::ResponseReadFailed { + bytes_read: response.len(), + source: io::Error::new( + io::ErrorKind::InvalidData, + "response reader returned more bytes than requested", + ), + }, + ); + } + Err(source) if source.kind() == io::ErrorKind::Interrupted => {} + Err(source) + if matches!( + source.kind(), + io::ErrorKind::TimedOut | io::ErrorKind::WouldBlock + ) => + { + if deadline.saturating_duration_since(now()).is_zero() { + return Err( + WebDriverBiDiWebSocketHandshakeResponseError::ResponseReadTimedOut { + bytes_read: response.len(), + source, + }, + ); + } + thread::sleep(Duration::from_millis(1)); + } + Err(source) => { + return Err( + WebDriverBiDiWebSocketHandshakeResponseError::ResponseReadFailed { + bytes_read: response.len(), + source, + }, + ); + } + } + } } /// Fail-closed errors while writing one bounded WebDriver BiDi WebSocket opening request. @@ -417,19 +943,19 @@ fn write_request_with_clock( ); } } - Err(source) if source.kind() == io::ErrorKind::Interrupted => {} - Err(source) + Err(source) => { + if source.kind() == io::ErrorKind::Interrupted { + continue; + } if matches!( source.kind(), io::ErrorKind::TimedOut | io::ErrorKind::WouldBlock - ) => - { - return Err(WebDriverBiDiWebSocketOpeningWriteError::WriteTimedOut { - bytes_written, - source, - }); - } - Err(source) => { + ) { + return Err(WebDriverBiDiWebSocketOpeningWriteError::WriteTimedOut { + bytes_written, + source, + }); + } return Err(WebDriverBiDiWebSocketOpeningWriteError::WriteFailed { bytes_written, source, @@ -504,6 +1030,239 @@ mod opening_write_tests { } } + #[derive(Clone, Debug)] + enum ReadAction { + Byte(u8), + Count(usize), + End, + Error(io::ErrorKind), + } + + #[derive(Debug)] + struct FakeReader { + actions: VecDeque, + mode_error: Option, + cleanup_error: Option, + } + + impl FakeReader { + fn new(actions: impl IntoIterator) -> Self { + Self { + actions: actions.into_iter().collect(), + mode_error: None, + cleanup_error: None, + } + } + } + + impl OpeningResponseReader for FakeReader { + fn set_nonblocking(&self, nonblocking: bool) -> io::Result<()> { + let error = if nonblocking { + self.mode_error + } else { + self.cleanup_error + }; + error.map_or(Ok(()), |kind| Err(io::Error::from(kind))) + } + + fn read_response_bytes(&mut self, bytes: &mut [u8]) -> io::Result { + match self.actions.pop_front().unwrap_or(ReadAction::End) { + ReadAction::Byte(byte) => { + bytes[0] = byte; + Ok(1) + } + ReadAction::Count(count) => Ok(count), + ReadAction::End => Ok(0), + ReadAction::Error(kind) => Err(io::Error::from(kind)), + } + } + } + + fn client_key() -> WebDriverBiDiWebSocketClientKey { + WebDriverBiDiWebSocketClientKey::new("dGhlIHNhbXBsZSBub25jZQ==") + .expect("test client key must be valid") + } + + fn valid_response() -> Vec { + b"HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=\r\n\r\n".to_vec() + } + + fn byte_actions(bytes: &[u8]) -> Vec { + bytes.iter().copied().map(ReadAction::Byte).collect() + } + + fn is_malformed_response(response: &[u8], key: &WebDriverBiDiWebSocketClientKey) -> bool { + matches!( + parse_opening_response(response, key), + Err(WebDriverBiDiWebSocketHandshakeResponseError::MalformedResponse { .. }) + ) + } + + fn read_with_fake( + reader: &mut FakeReader, + now_values: impl IntoIterator, + ) -> Result<(u16, usize), WebDriverBiDiWebSocketHandshakeResponseError> { + let key = client_key(); + let fallback = Instant::now(); + let mut now_values = now_values.into_iter(); + let mut now = || now_values.next().unwrap_or(fallback); + read_opening_response_with_clock(reader, &key, Duration::from_secs(1), &mut now) + } + + #[test] + fn parser_accepts_case_insensitive_upgrade_tokens_and_rejects_malformed_headers() { + let key = client_key(); + let response = b"HTTP/1.1 101 Switching Protocols\r\nUpGrAdE: WebSocket\r\nConnection: keep-alive, Upgrade\r\nSec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=\r\nX-Test: retained\r\n\r\n"; + let parsed = parse_opening_response(response, &key).expect("valid response"); + assert_eq!(parsed.status_code, 101); + assert_eq!(parsed.byte_count, response.len()); + assert!(!is_malformed_response(response, &key)); + let same_length_mismatch = String::from_utf8(response.to_vec()) + .expect("valid response fixture") + .replace( + "s3pPLMBiTxaQ9kYGzzhZRbK+xOo=", + "s3pPLMBiTxaQ9kYGzzhZRbK+xOoX", + ); + assert!(parse_opening_response(same_length_mismatch.as_bytes(), &key).is_err()); + + let malformed_responses = [ + b"HTTP/1.1 101".to_vec(), + vec![0xff, b'\r', b'\n', b'\r', b'\n'], + b"HTTP/1.1 101\0 Switching Protocols\r\n\r\n".to_vec(), + b"HTTP/1.1 200 OK\r\n\r\n".to_vec(), + b"HTTP/1.1 101 Switching Protocols\r\n\r\n".to_vec(), + b"HTTP/1.1 101 Switching Protocols\r\n Upgrade: websocket\r\n\r\n".to_vec(), + b"HTTP/1.1 101 Switching Protocols\r\nUpgrade\r\n\r\n".to_vec(), + b"HTTP/1.1 101 Switching Protocols\r\nBad Header: value\r\n\r\n".to_vec(), + b"HTTP/1.1 101 Switching Protocols\r\n: value\r\n\r\n".to_vec(), + b"HTTP/1.1 101 Switching Protocols\r\nUpgrade: web\x01socket\r\n\r\n".to_vec(), + b"HTTP/1.1 101 Switching Protocols\r\nUpgrade: web\x80socket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=\r\n\r\n".to_vec(), + b"HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nUpgrade: websocket\r\n\r\n".to_vec(), + b"HTTP/1.1 101 Switching Protocols\r\nConnection: Upgrade\r\nConnection: Upgrade\r\n\r\n".to_vec(), + b"HTTP/1.1 101 Switching Protocols\r\nSec-WebSocket-Accept: one\r\nSec-WebSocket-Accept: two\r\n\r\n".to_vec(), + b"HTTP/1.1 101 Switching Protocols\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=\r\n\r\n".to_vec(), + b"HTTP/1.1 101 Switching Protocols\r\nUpgrade: h2c\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=\r\n\r\n".to_vec(), + b"HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nSec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=\r\n\r\n".to_vec(), + b"HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: keep-alive\r\nSec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=\r\n\r\n".to_vec(), + b"HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n\r\n".to_vec(), + ]; + for response in malformed_responses { + assert!(is_malformed_response(&response, &key)); + } + } + + #[test] + fn bounded_response_reader_covers_deadlines_size_io_and_cleanup() { + let start = Instant::now(); + + let mut valid_reader = FakeReader::new(byte_actions(&valid_response())); + let valid = read_with_fake(&mut valid_reader, [start]); + assert!(valid.is_ok()); + + let mut malformed_reader = FakeReader::new(byte_actions(b"HTTP/1.1 200 OK\r\n\r\n")); + assert!(read_with_fake(&mut malformed_reader, [start]).is_err()); + + let mut interrupted_reader = FakeReader::new( + std::iter::once(ReadAction::Error(io::ErrorKind::Interrupted)) + .chain(byte_actions(&valid_response())), + ); + assert!(read_with_fake(&mut interrupted_reader, [start]).is_ok()); + + let mut mode_error_reader = FakeReader::new([]); + mode_error_reader.mode_error = Some(io::ErrorKind::InvalidInput); + assert!(read_with_fake(&mut mode_error_reader, [start]).is_err()); + + let mut ended_reader = FakeReader::new([ReadAction::End]); + assert!(read_with_fake(&mut ended_reader, [start]).is_err()); + + let mut count_reader = FakeReader::new([ReadAction::Count(2)]); + assert!(read_with_fake(&mut count_reader, [start]).is_err()); + + let mut failed_reader = FakeReader::new([ReadAction::Error(io::ErrorKind::BrokenPipe)]); + assert!(read_with_fake(&mut failed_reader, [start]).is_err()); + + let mut retrying_reader = FakeReader::new( + std::iter::once(ReadAction::Error(io::ErrorKind::WouldBlock)) + .chain(byte_actions(&valid_response())), + ); + assert!(read_with_fake(&mut retrying_reader, [start]).is_ok()); + + let mut timed_out_reader = FakeReader::new([ReadAction::Error(io::ErrorKind::TimedOut)]); + assert!( + read_with_fake( + &mut timed_out_reader, + [start, start, start + Duration::from_secs(1)] + ) + .is_err() + ); + + let mut deadline_reader = FakeReader::new([ReadAction::End]); + assert!( + read_with_fake( + &mut deadline_reader, + [start, start + Duration::from_secs(1)] + ) + .is_err() + ); + + let mut late_response_reader = FakeReader::new(byte_actions(&valid_response())); + let mut late_response_times = vec![start; valid_response().len() + 1]; + late_response_times.push(start + Duration::from_secs(1)); + assert!(read_with_fake(&mut late_response_reader, late_response_times).is_err()); + + let mut cleanup_reader = FakeReader::new(byte_actions(&valid_response())); + cleanup_reader.cleanup_error = Some(io::ErrorKind::InvalidInput); + assert!(read_with_fake(&mut cleanup_reader, [start]).is_err()); + + let mut too_large_reader = FakeReader::new(std::iter::repeat_n( + ReadAction::Byte(b'a'), + MAX_WEBSOCKET_OPENING_RESPONSE_BYTES, + )); + assert!(read_with_fake(&mut too_large_reader, [start]).is_err()); + } + + #[test] + fn response_errors_have_deterministic_messages_and_sources() { + let source = io::Error::from(io::ErrorKind::InvalidInput); + let errors = [ + WebDriverBiDiWebSocketHandshakeResponseError::InvalidResponseTimeout { + response_timeout: Duration::ZERO, + maximum_timeout: MAX_WEBSOCKET_OPENING_RESPONSE_TIMEOUT, + }, + WebDriverBiDiWebSocketHandshakeResponseError::ResponseDeadlineExceeded { + bytes_read: 1, + }, + WebDriverBiDiWebSocketHandshakeResponseError::ResponseTooLarge { + bytes_read: 1, + maximum_bytes: 1, + }, + WebDriverBiDiWebSocketHandshakeResponseError::ResponseReadModeConfigurationFailed { + bytes_read: 1, + source: io::Error::from(io::ErrorKind::InvalidInput), + }, + WebDriverBiDiWebSocketHandshakeResponseError::ResponseReadTimedOut { + bytes_read: 1, + source: io::Error::from(io::ErrorKind::TimedOut), + }, + WebDriverBiDiWebSocketHandshakeResponseError::ResponseReadFailed { + bytes_read: 1, + source: io::Error::from(io::ErrorKind::BrokenPipe), + }, + WebDriverBiDiWebSocketHandshakeResponseError::ResponseEndedBeforeHeaders { + bytes_read: 1, + }, + WebDriverBiDiWebSocketHandshakeResponseError::MalformedResponse { reason: "test" }, + WebDriverBiDiWebSocketHandshakeResponseError::AcceptMismatch, + WebDriverBiDiWebSocketHandshakeResponseError::ReadModeCleanupFailed { source }, + ]; + for (error, has_source) in errors.iter().zip([ + false, false, false, true, true, true, false, false, false, true, + ]) { + assert!(!error.to_string().is_empty()); + assert_eq!(error.source().is_some(), has_source); + } + } + #[test] fn bounded_writer_completes_partial_and_interrupted_writes() { let mut writer = FakeWriter::new([ diff --git a/crates/originweave-network/tests/webdriver_bidi_websocket_opening_write.rs b/crates/originweave-network/tests/webdriver_bidi_websocket_opening_write.rs index 433774dd4..eb6448c5f 100644 --- a/crates/originweave-network/tests/webdriver_bidi_websocket_opening_write.rs +++ b/crates/originweave-network/tests/webdriver_bidi_websocket_opening_write.rs @@ -1,6 +1,7 @@ use std::{ - io::{self, Read}, + io::{self, Read, Write}, net::TcpListener, + sync::mpsc, thread, time::Duration, }; @@ -9,7 +10,7 @@ use originweave_core::WebDriverBiDiWebSocketEndpoint; use originweave_network::{ MAX_WEBSOCKET_OPENING_WRITE_TIMEOUT, WebDriverBiDiTcpConnectionPlan, WebDriverBiDiWebSocketClientKey, WebDriverBiDiWebSocketHandshakePlan, - WebDriverBiDiWebSocketOpeningWriteError, + WebDriverBiDiWebSocketHandshakeResponseError, WebDriverBiDiWebSocketOpeningWriteError, }; const SESSION_ID: &str = "01234567-89ab-cdef-0123-456789abcdef"; @@ -44,7 +45,7 @@ fn connect(endpoint: &str) -> originweave_network::WebDriverBiDiTcpConnection { connection } -fn read_opening_request(mut stream: std::net::TcpStream) -> io::Result> { +fn read_opening_request(stream: &mut std::net::TcpStream) -> io::Result> { stream.set_read_timeout(Some(Duration::from_secs(2)))?; let mut request = Vec::new(); let mut buffer = [0_u8; 512]; @@ -58,6 +59,24 @@ fn read_opening_request(mut stream: std::net::TcpStream) -> io::Result> Ok(request) } +fn serve_opening_exchange( + listener: TcpListener, + response: &'static [u8], +) -> ( + mpsc::SyncSender<()>, + thread::JoinHandle>>, +) { + let (release_server, await_release) = mpsc::sync_channel(0); + let server = thread::spawn(move || { + let (mut stream, _) = listener.accept()?; + let request = read_opening_request(&mut stream)?; + stream.write_all(response)?; + await_release.recv().map_err(io::Error::other)?; + Ok(request) + }); + (release_server, server) +} + #[test] fn bounded_opening_write_sends_exact_request_and_preserves_transport_evidence() { let listener = TcpListener::bind(("127.0.0.1", 0)); @@ -70,10 +89,7 @@ fn bounded_opening_write_sends_exact_request_and_preserves_transport_evidence() let Ok(local_addr) = local_addr else { return; }; - let server = thread::spawn(move || { - let accepted = listener.accept()?; - read_opening_request(accepted.0) - }); + let (release_server, server) = serve_opening_exchange(listener, b""); let endpoint = format!("ws://{local_addr}/session/{SESSION_ID}"); let connection = connect(&endpoint); @@ -114,6 +130,7 @@ fn bounded_opening_write_sends_exact_request_and_preserves_transport_evidence() assert!(debug.contains("WebDriverBiDiWebSocketOpeningRequestSent")); assert!(!debug.contains(RFC6455_SAMPLE_KEY)); + assert!(release_server.send(()).is_ok()); let server_result = server.join(); assert!(server_result.is_ok(), "{server_result:?}"); if let Ok(received) = server_result { @@ -171,3 +188,160 @@ fn opening_write_rejects_zero_and_excessive_deadlines_before_success_evidence() } } } + +#[test] +fn opening_response_requires_rfc6455_switching_protocols_and_matching_accept() { + 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 (release_server, server) = serve_opening_exchange( + listener, + b"HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=\r\n\r\n", + ); + + let endpoint = format!("ws://{local_addr}/session/{SESSION_ID}"); + let key = WebDriverBiDiWebSocketClientKey::new(RFC6455_SAMPLE_KEY); + assert!(key.is_ok(), "{key:?}"); + let Ok(key) = key else { + return; + }; + let plan = WebDriverBiDiWebSocketHandshakePlan::new(connect(&endpoint), key); + assert!(plan.is_ok(), "{plan:?}"); + let Ok(plan) = plan else { + return; + }; + let written = plan.write_opening_request(Duration::from_millis(500)); + assert!(written.is_ok(), "{written:?}"); + let Ok(written) = written else { + return; + }; + + let established = written.read_opening_response(Duration::from_millis(500)); + assert!(established.is_ok(), "{established:?}"); + let Ok(established) = established else { + return; + }; + assert_eq!(established.response_status(), 101); + assert!(established.response_byte_count() > 0); + assert!(established.request_byte_count() > 0); + assert_eq!(established.response_timeout(), Duration::from_millis(500)); + assert_eq!(established.write_timeout(), Duration::from_millis(500)); + assert_eq!(established.client_key().as_str(), RFC6455_SAMPLE_KEY); + assert_eq!( + established + .transport_evidence() + .verified_peer() + .socket_addr(), + local_addr + ); + let debug = format!("{established:?}"); + assert!(debug.contains("WebDriverBiDiWebSocketEstablished")); + assert!(!debug.contains(RFC6455_SAMPLE_KEY)); + drop(established); + assert!(release_server.send(()).is_ok()); + assert!(server.join().is_ok()); +} + +#[test] +fn opening_response_rejects_a_mismatched_accept_value() { + 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 (release_server, server) = serve_opening_exchange( + listener, + b"HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: invalid\r\n\r\n", + ); + + let endpoint = format!("ws://{local_addr}/session/{SESSION_ID}"); + let key = WebDriverBiDiWebSocketClientKey::new(RFC6455_SAMPLE_KEY); + assert!(key.is_ok(), "{key:?}"); + let Ok(key) = key else { + return; + }; + let plan = WebDriverBiDiWebSocketHandshakePlan::new(connect(&endpoint), key); + assert!(plan.is_ok(), "{plan:?}"); + let Ok(plan) = plan else { + return; + }; + let written = plan.write_opening_request(Duration::from_millis(500)); + assert!(written.is_ok(), "{written:?}"); + let Ok(written) = written else { + return; + }; + + assert!(matches!( + written.read_opening_response(Duration::from_millis(500)), + Err(WebDriverBiDiWebSocketHandshakeResponseError::AcceptMismatch) + )); + assert!(release_server.send(()).is_ok()); + assert!(server.join().is_ok()); +} + +#[test] +fn opening_response_rejects_zero_and_excessive_deadlines_before_socket_mode_change() { + for timeout in [ + Duration::ZERO, + originweave_network::MAX_WEBSOCKET_OPENING_RESPONSE_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 (release_server, await_release) = mpsc::sync_channel(0); + let server = thread::spawn(move || { + let accepted = listener.accept()?; + await_release.recv().map_err(io::Error::other)?; + drop(accepted); + Ok::<(), io::Error>(()) + }); + + let endpoint = format!("ws://{local_addr}/session/{SESSION_ID}"); + let key = WebDriverBiDiWebSocketClientKey::new(RFC6455_SAMPLE_KEY); + assert!(key.is_ok(), "{key:?}"); + let Ok(key) = key else { + continue; + }; + let plan = WebDriverBiDiWebSocketHandshakePlan::new(connect(&endpoint), key); + assert!(plan.is_ok(), "{plan:?}"); + let Ok(plan) = plan else { + continue; + }; + let written = plan.write_opening_request(Duration::from_millis(500)); + assert!(written.is_ok(), "{written:?}"); + let Ok(written) = written else { + continue; + }; + + assert!(matches!( + written.read_opening_response(timeout), + Err(WebDriverBiDiWebSocketHandshakeResponseError::InvalidResponseTimeout { + response_timeout, + maximum_timeout, + }) if response_timeout == timeout + && maximum_timeout + == originweave_network::MAX_WEBSOCKET_OPENING_RESPONSE_TIMEOUT + )); + assert!(release_server.send(()).is_ok()); + assert!(server.join().is_ok()); + } +} diff --git a/crates/originweave-network/tests/webdriver_bidi_websocket_repeated_headers.rs b/crates/originweave-network/tests/webdriver_bidi_websocket_repeated_headers.rs new file mode 100644 index 000000000..ffea9c8e7 --- /dev/null +++ b/crates/originweave-network/tests/webdriver_bidi_websocket_repeated_headers.rs @@ -0,0 +1,135 @@ +use std::{ + error::Error, + io::{self, Read, Write}, + net::TcpListener, + thread, + time::Duration, +}; + +use originweave_core::WebDriverBiDiWebSocketEndpoint; +use originweave_network::{ + WebDriverBiDiTcpConnectionPlan, WebDriverBiDiWebSocketClientKey, + WebDriverBiDiWebSocketHandshakePlan, WebDriverBiDiWebSocketHandshakeResponseError, +}; + +const SESSION_ID: &str = "01234567-89ab-cdef-0123-456789abcdef"; +const RFC6455_SAMPLE_KEY: &str = "dGhlIHNhbXBsZSBub25jZQ=="; +const RFC6455_SAMPLE_ACCEPT: &str = "s3pPLMBiTxaQ9kYGzzhZRbK+xOo="; + +type TestResult = Result>; + +fn connect(endpoint: &str) -> TestResult { + let endpoint = WebDriverBiDiWebSocketEndpoint::new(endpoint)?; + let correlated = endpoint.correlate_session_id(SESSION_ID)?; + let target = correlated.into_explicit_connect_target()?; + let connection = + WebDriverBiDiTcpConnectionPlan::new(target, Duration::from_secs(1), 1)?.connect()?; + Ok(connection) +} + +fn read_opening_request(stream: &mut std::net::TcpStream) -> io::Result<()> { + stream.set_read_timeout(Some(Duration::from_secs(2)))?; + let mut request = Vec::new(); + let mut buffer = [0_u8; 256]; + while !request.ends_with(b"\r\n\r\n") { + let count = stream.read(&mut buffer)?; + if count == 0 { + break; + } + request.extend_from_slice(&buffer[..count]); + } + if !request.ends_with(b"\r\n\r\n") { + return Err(io::Error::new( + io::ErrorKind::UnexpectedEof, + "opening request ended before headers", + )); + } + Ok(()) +} + +fn exercise_response( + response: Vec, +) -> TestResult> { + let listener = TcpListener::bind(("127.0.0.1", 0))?; + let local_addr = listener.local_addr()?; + let server = thread::spawn(move || -> io::Result<()> { + let (mut stream, _) = listener.accept()?; + read_opening_request(&mut stream)?; + stream.write_all(&response)?; + Ok(()) + }); + + let endpoint = format!("ws://{local_addr}/session/{SESSION_ID}"); + let key = WebDriverBiDiWebSocketClientKey::new(RFC6455_SAMPLE_KEY)?; + let plan = WebDriverBiDiWebSocketHandshakePlan::new(connect(&endpoint)?, key)?; + let written = plan.write_opening_request(Duration::from_millis(500))?; + let result = written + .read_opening_response(Duration::from_millis(500)) + .map(|established| { + drop(established); + }); + match server.join() { + Ok(server_result) => server_result?, + Err(_) => return Err(io::Error::other("loopback fixture thread panicked").into()), + } + Ok(result) +} + +#[test] +fn repeated_list_valued_upgrade_and_connection_lines_are_combined_semantically() -> TestResult<()> { + let response = format!( + "HTTP/1.1 101 Switching Protocols\r\n\ +Upgrade: h2c\r\n\ +Upgrade: websocket\r\n\ +Connection: keep-alive\r\n\ +Connection: Upgrade\r\n\ +Sec-WebSocket-Accept: {RFC6455_SAMPLE_ACCEPT}\r\n\r\n" + ) + .into_bytes(); + + let result = exercise_response(response)?; + assert!( + result.is_ok(), + "RFC 9110 list-valued fields must combine: {result:?}" + ); + Ok(()) +} + +#[test] +fn unknown_header_obs_text_is_ignored_without_weakening_required_fields() -> TestResult<()> { + let mut response = b"HTTP/1.1 101 Switching Protocols\r\nX-OriginWeave-Ignored: ".to_vec(); + response.push(0x80); + response.extend_from_slice( + format!( + "\r\nUpgrade: websocket\r\n\ +Connection: Upgrade\r\n\ +Sec-WebSocket-Accept: {RFC6455_SAMPLE_ACCEPT}\r\n\r\n" + ) + .as_bytes(), + ); + + let result = exercise_response(response)?; + assert!( + result.is_ok(), + "RFC 9110 obs-text in an ignored extension field must not invalidate the opening response: {result:?}" + ); + Ok(()) +} + +#[test] +fn repeated_sec_websocket_accept_remains_fail_closed() -> TestResult<()> { + let response = format!( + "HTTP/1.1 101 Switching Protocols\r\n\ +Upgrade: websocket\r\n\ +Connection: Upgrade\r\n\ +Sec-WebSocket-Accept: {RFC6455_SAMPLE_ACCEPT}\r\n\ +Sec-WebSocket-Accept: {RFC6455_SAMPLE_ACCEPT}\r\n\r\n" + ) + .into_bytes(); + + assert!(matches!( + exercise_response(response)?, + Err(WebDriverBiDiWebSocketHandshakeResponseError::MalformedResponse { .. }) + )); + Ok(()) +} diff --git a/crates/originweave-network/tests/webdriver_bidi_websocket_status_line_strict.rs b/crates/originweave-network/tests/webdriver_bidi_websocket_status_line_strict.rs new file mode 100644 index 000000000..03b3e69ca --- /dev/null +++ b/crates/originweave-network/tests/webdriver_bidi_websocket_status_line_strict.rs @@ -0,0 +1,77 @@ +use std::{ + error::Error, + io::{self, Read, Write}, + net::TcpListener, + thread, + time::Duration, +}; + +use originweave_core::WebDriverBiDiWebSocketEndpoint; +use originweave_network::{ + WebDriverBiDiTcpConnectionPlan, WebDriverBiDiWebSocketClientKey, + WebDriverBiDiWebSocketHandshakePlan, WebDriverBiDiWebSocketHandshakeResponseError, +}; + +const SESSION_ID: &str = "01234567-89ab-cdef-0123-456789abcdef"; +const RFC6455_SAMPLE_KEY: &str = "dGhlIHNhbXBsZSBub25jZQ=="; +const RFC6455_SAMPLE_ACCEPT: &str = "s3pPLMBiTxaQ9kYGzzhZRbK+xOo="; + +type TestResult = Result>; + +fn read_opening_request(stream: &mut std::net::TcpStream) -> io::Result<()> { + stream.set_read_timeout(Some(Duration::from_secs(2)))?; + let mut request = Vec::new(); + let mut buffer = [0_u8; 256]; + while !request.ends_with(b"\r\n\r\n") { + let count = stream.read(&mut buffer)?; + if count == 0 { + break; + } + request.extend_from_slice(&buffer[..count]); + } + if request.ends_with(b"\r\n\r\n") { + Ok(()) + } else { + Err(io::Error::new( + io::ErrorKind::UnexpectedEof, + "opening request ended before headers", + )) + } +} + +#[test] +fn response_without_mandatory_space_after_status_code_fails_closed() -> TestResult<()> { + let listener = TcpListener::bind(("127.0.0.1", 0))?; + let local_addr = listener.local_addr()?; + let server = thread::spawn(move || -> io::Result<()> { + let (mut stream, _) = listener.accept()?; + read_opening_request(&mut stream)?; + write!( + stream, + "HTTP/1.1 101\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: {RFC6455_SAMPLE_ACCEPT}\r\n\r\n" + )?; + Ok(()) + }); + + let endpoint_url = format!("ws://{local_addr}/session/{SESSION_ID}"); + let endpoint = WebDriverBiDiWebSocketEndpoint::new(&endpoint_url)?; + let correlated = endpoint.correlate_session_id(SESSION_ID)?; + let target = correlated.into_explicit_connect_target()?; + let connection = + WebDriverBiDiTcpConnectionPlan::new(target, Duration::from_secs(1), 1)?.connect()?; + let key = WebDriverBiDiWebSocketClientKey::new(RFC6455_SAMPLE_KEY)?; + let written = WebDriverBiDiWebSocketHandshakePlan::new(connection, key)? + .write_opening_request(Duration::from_millis(500))?; + + let result = written.read_opening_response(Duration::from_millis(500)); + match server.join() { + Ok(server_result) => server_result?, + Err(_) => return Err(io::Error::other("loopback fixture thread panicked").into()), + } + + assert!(matches!( + result, + Err(WebDriverBiDiWebSocketHandshakeResponseError::MalformedResponse { .. }) + )); + Ok(()) +} diff --git a/docs/doctoring.md b/docs/doctoring.md index 866d8766f..86d1f3ee2 100644 --- a/docs/doctoring.md +++ b/docs/doctoring.md @@ -14,6 +14,8 @@ WAI-ARIA 1.2 defines host-language `role` values as a token list: user agents sp UTS #39 Revision 32 is the current Unicode security-mechanisms standard and marks Default_Ignorable and bidirectional format characters as restricted in identifier profiles. UAX #9 defines the bidirectional format controls that can reorder displayed protocol text. UTR #36 Revision 15 remains a stabilized historical security-considerations report; its identifier recommendations are superseded by UTS #39 rather than cited as current normative profile rules. OriginWeave therefore rejects the reviewed format-character set in roles, shared identifiers, and registry external identifiers, and rejects those same characters inside accessible names while still allowing ordinary U+0020 spaces. +RFC 6455 carries the WebSocket opening handshake over HTTP/1.1, and RFC 9110 permits `obs-text` octets (`%x80-FF`) in field values while retaining ASCII field-name and delimiter syntax. RFC 6455 also specifies that unknown opening-handshake header fields are ignored. OriginWeave therefore treats unknown extension-field values as opaque compatibility data rather than requiring the entire opening response to be UTF-8, while keeping the authority-bearing `Upgrade`, `Connection`, and `Sec-WebSocket-Accept` checks fail closed: opaque replacement material cannot satisfy the reviewed ASCII token or exact accept-value contracts. Ignoring an unknown field never grants browser, network, secret, approval, or Agent authority. + ### Browser origin equivalence The WHATWG URL host parser and Chromium canonicalizer classify shortened decimal, integer, hexadecimal, legacy octal-looking, and mixed-component numeric hosts as IPv4 or broken IPv4 candidates rather than ordinary DNS names. Chromium's regression suite includes values such as `192`, `0xC0a80001`, `030052000001`, and mixed hexadecimal components. A non-final empty `0x` component can participate in Chromium's multi-part IPv4 truncation behavior, but a final `0x` label does not produce an IPv4 number because stripping its prefix leaves no digits; it remains a domain label. Chromium also warns that broken IP-like hosts must not be connected because another resolver could accept them. OriginWeave therefore admits only canonical dotted-decimal IPv4 into its policy origin type, rejects browser-special numeric spellings before DNS validation, and preserves final non-numeric DNS labels such as `0x`. @@ -96,6 +98,12 @@ TRINITY uses a compact learned coordinator to select models and assign Thinker, These results motivate explicit OriginWeave configuration for model routing, workflow stage, decomposition, recursion depth, permitted access, role assignment, and role-specific reasoning effort. They do not justify always using multiple agents. OriginWeave must compare bounded single-model, routed-model, and deeper multi-agent configurations through task-success, safety, variance, token, and compute ablations. No learned coordinator may expand browser capabilities, origins, destinations, approvals, secrets, or deterministic policy. +### Opening-exchange fixture lifetime + +On 5 September 2026, a complete Rust run after integrating PR #242 head `55fef0c3fae1724eddada53e52c4a0311f509aa3` into #243 reproduced `WriteTimeoutCleanupFailed` with macOS `EINVAL` after 198 request bytes in `opening_response_rejects_a_mismatched_accept_value`. That fixture returned its invalid response and closed immediately, before the client could finish opening-write cleanup. The successful-handshake fixture's one-byte close probe could consume the first request byte rather than observe closure, and the request-only fixture also closed immediately after reading the request. All three paths therefore shared a premature peer-lifetime assumption; the previously repaired invalid-deadline fixture did not cover them. + +The test-only `serve_opening_exchange` helper reuses the bounded request reader, reads the complete request before sending the configured response, and retains the accepted stream until the client explicitly releases it after its assertions. Successful, mismatched-accept, and request-only tests all use that helper. No sleep, retry-based acceptance, production timeout change, ignored cleanup error, dependency, or coverage exclusion is introduced. Existing real-socket assertions remain the regression checks, including the requirement that an invalid accept value reaches `AcceptMismatch` rather than an earlier fixture-induced transport failure. Descendant stacks must adopt the owner fix and rerun their own gates; the reproduced failure and this fixture repair do not imply protected-main or browser-runtime delivery. + ## References Amazon Web Services. (n.d.). *Set up the Amazon EKS Pod Identity Agent*. Retrieved August 6, 2026, from https://docs.aws.amazon.com/eks/latest/userguide/pod-id-agent-setup.html @@ -118,6 +126,8 @@ Eddy, W. M. (Ed.). (2022). *Transmission Control Protocol (TCP)* (RFC 9293). Int Evtimov, I., Zharmagambetov, A., Grattafiori, A., Guo, C., & Chaudhuri, K. (2025). *WASP: Benchmarking web agent security against prompt injection attacks*. arXiv. https://doi.org/10.48550/arXiv.2504.18575 +Fette, I., & Melnikov, A. (2011). *The WebSocket Protocol* (RFC 6455). Internet Engineering Task Force. https://doi.org/10.17487/RFC6455 + Fielding, R., Nottingham, M., & Reschke, J. (2022). *HTTP semantics* (RFC 9110). Internet Engineering Task Force. https://doi.org/10.17487/RFC9110 Fugu Team, Sakana AI. (2026). *Sakana Fugu technical report* [Technical report]. arXiv. https://doi.org/10.48550/arXiv.2606.21228