diff --git a/AGENTS.md b/AGENTS.md index 6f747c38e..ad94ed713 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -113,6 +113,27 @@ A skipped security, GPU, browser, TLS, or statistical test is not passing eviden - Run reasoning-effort and orchestration-depth ablations before claiming an LLM path is superior. - Scheduled agents may create bounded reviewed PRs but may not merge, tag, publish, alter workflows, add secrets, or weaken checks. +### Coverage diagnosis lesson + +Keep generic wire-code validation separate from endpoint-role validation. For +WebSocket Close, reject server-sent client-only 1010 in the client closure state +machine before echoing it, and pair that RED case with valid server 1011 so a +blanket rejection cannot satisfy the regression. + +Repeated control-frame support must preserve the total deadline and per-write masking authority. Supply a separate caller-owned key per Ping, consume none for unsolicited Pong, and keep the Close key independent. Test exact-budget success, budget overflow, exhausted keys, and reused keys with a peer that verifies no rejected reply. The 64-control budget is local resource policy, not an RFC limit. + +For Close-code validation, check the current IANA registry as well as RFC 6455: the protocol-reserved range is not an allowance for unassigned values. Keep application/private ranges separate, test assigned and reserved boundaries, and verify rejected peer codes produce neither an echo nor closure evidence. Record the registry date; future assignments need a reviewed update, not ambient network lookup during frame parsing. + +For a multi-step socket deadline, reproduce a sequence whose individual waits fit the limit but whose sum does not. Carry one monotonic expiry through every step and recheck before admitting final evidence. Pair real delayed-peer tests with a controlled clock at each read/write/evidence transition; a pre-I/O check alone cannot reject late completion. This is not a hard real-time host scheduling guarantee. + +Aggregate LLVM code regions by source coordinates across function instantiations before identifying a missing path. An invalid-input test at a public entry point may stop at an earlier guard; it does not prove a later private writer's error return executed. Exercise that writer directly, assert the exact error, and verify the peer received no bytes. Do not rewrite production predicates based only on a file-level coverage deficit. + +Coordinate union alone does not reproduce LLVM's region summary: `RegionCoverageInfo::merge` takes maximum covered/total counts across instantiations. Complementary unit-test and integration-test executions can therefore leave a deficit. Exercise successful status-bearing and empty Close writes, invalid deadlines, and adjacent masking-key rejection in the same unit-test binary; compare literal wire bytes and join the peer. Reference: LLVM Project. (n.d.). *CoverageSummaryInfo.h* [Source code]. https://github.com/llvm/llvm-project/blob/main/llvm/tools/llvm-cov/CoverageSummaryInfo.h + +Inspect uncovered coordinates inside test assertions too: guarded `matches!` expressions can contribute never-taken failure branches to the file summary. Preserve exact variant and field checks rather than broadening the accepted error to make coverage pass. + +Mask-reuse socket tests must consume the preceding text/Pong frame before sending the frame that triggers rejection. Assert the exact reuse error, literal preceding bytes, and EOF with no rejected response; propagate peer thread errors. A broad transport-error assertion plus an ignored join can pass because the peer rejected the fixture's own wrong opcode. An intervening fresh-key Pong also changes the adjacent-key history, so an older text key does not test adjacent Close-key reuse. + ## Release contract A release requires all current-head checks, complete coverage and docs, updated `CHANGELOG.md`, SBOM and provenance, reproducible artifacts, compatibility evidence, security review, and an explicit version decision. Pre-alpha commits are not releases. diff --git a/CHANGELOG.md b/CHANGELOG.md index 6127f0374..ccf2f4eb1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,12 +4,34 @@ All notable changes to OriginWeave are documented in this file. The format follo ## [Unreleased] +### Changed + +- Browser transport shutdown rejects the client-only Close status 1010 when received from a server, before sending a reply or recording completion; server status 1011 remains supported. + +- Browser transport shutdown now handles repeated keepalive traffic within one time budget, using a separate supplied key for each reply and rejecting exhausted keys or excess traffic without reporting completion. + +- Browser transport shutdown rejects unassigned protocol close codes before replying or recording completion; application and private-use code ranges remain supported. + +- Browser transport shutdown now shares one time budget across control replies and final connection closure; late closure cannot become successful completion evidence. + +- Restored the simple frame-timeout validation after correcting coverage diagnosis; direct Close-writer tests verify invalid deadlines send no bytes and return the expected error. +- Kept exact invalid-deadline error checks without compound test-only guards in the coverage measurement. +- Verified literal masked Close bytes with and without a status code, and that a reused masking key emits no Close bytes after the preceding text frame. +- Corrected three masking-rejection test peers so unrelated socket errors cannot masquerade as key-reuse protection; each now verifies the preceding frame, the exact rejection, and no subsequent bytes. + +- Removed an unused private correlated-response accessor while retaining connection-generation validation at the receiving-message boundary, and corrected the Rust `AtomicU64` standard-library reference to its canonical type-alias page. +- Integrated the current teardown prerequisites into transport-closure observation, including the previously uncollected release-record check, while retaining the unresolved connection-provenance finding and its downstream repair ownership. +- Integrated the verified opening-exchange and closure prerequisites into the connection-bound response repair, preserving its sender, receiver and closure provenance checks while restoring the inherited executable release contract; process and profile cleanup remain unproven. + ### Added - Session-ending replies from a replacement connection can no longer complete the original pending request. The original reply remains usable, and a protocol acknowledgment still does not prove browser shutdown or cleanup. - The session-ending command stack now retains the status-reply protections from its current parent. A reply from a replacement connection is rejected while the original pending status request remains recoverable; sending the end command still does not prove that the browser session ended. - Typed outbound WebDriver BiDi `session.end` over the bounded client WebSocket stream: it serializes only the standards-defined method with empty params, rejects invalid frame deadlines before correlation registration, retires only the just-registered id when frame preflight proves no command bytes were emitted, preserves exact command-kind correlation across ambiguous writes, and does not treat frame-write success as proof that the browser session ended. - Typed `session.end` response admission that consumes only the exact outstanding command-kind correlation after complete envelope validation, preserves remote protocol errors as failures, and does not claim browser-process exit or resource cleanup from a protocol acknowledgment. +- Fail-closed `session.end` teardown assessment that binds only the typed observation produced by consuming the exact transport, keeps browser-process-exit and task-profile-removal evidence unavailable until their runtime owners exist, and therefore cannot report operational completion from caller-supplied booleans. +- Bounded WebDriver BiDi transport-closure observation that consumes the established stream, accepts only a validated peer Close frame or clean pre-frame EOF, permits at most one unsolicited Pong, and keeps transport closure separate from process-exit and profile-cleanup claims. + - Regression checks now exercise fragmented browser replies, interleaved control messages, and rejected replies without losing a pending request. These checks do not establish browser readiness or release acceptance. - The typed browser-status response stack now includes its verified command and opening-exchange prerequisites, including the release-record check that previously did not execute; parsing remains bounded and does not grant browser authority or prove operational readiness. - 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. @@ -62,6 +84,7 @@ All notable changes to OriginWeave are documented in this file. The format follo ### Changed +- Carried current response prerequisites and the executable release-record check into the teardown-assessment stack; caller-supplied cleanup claims remain unverified and cannot establish operational acceptance. - Carried verified command prerequisites and the executable release-record check into session-end response validation without changing response admission or treating an acknowledgment as proof of resource cleanup. - Carried the verified status-response prerequisites into the session-end sender, preserving its command behavior and making the inherited release-record check execute in the existing test suite. - Kept the `session.status` frame-failure coverage contract focused on observable correlation state, avoiding assertion-internal uncovered branches without weakening preflight retirement or ambiguous-write retention checks. diff --git a/crates/originweave-network/src/lib.rs b/crates/originweave-network/src/lib.rs index b2cc56cb1..9ebfa3588 100644 --- a/crates/originweave-network/src/lib.rs +++ b/crates/originweave-network/src/lib.rs @@ -7,10 +7,14 @@ //! `originweave-core` into one bounded exact TCP connection, binds and validates //! the RFC 6455 opening exchange, provides bounded masked client writes and //! unmasked server-frame reads, assembles bounded WebDriver BiDi text messages, -//! classifies complete local-end JSON envelopes, tracks bounded command-response -//! correlation, sends narrowly typed `session.status` and `session.end` commands, -//! and admits typed correlated status and end responses without exposing generic -//! JSON bodies or granting browser, TLS, policy, secret, or Agent authority. +//! binds received fragmented text to one exact verified connection, classifies +//! complete local-end JSON envelopes, tracks bounded command-response correlation, +//! sends narrowly typed `session.status` and `session.end` commands, admits typed +//! correlated status and end responses, binds `session.end` ACK and closure evidence +//! to one private process-local connection generation, observes bounded peer Close +//! or clean-EOF transport cessation, and keeps protocol/transport evidence separate +//! from explicit operational teardown observations without exposing generic JSON +//! bodies or granting browser, TLS, policy, secret, process, profile, or Agent authority. #![forbid(unsafe_code)] #![deny(missing_docs)] @@ -24,14 +28,24 @@ mod webdriver_bidi_session_end_command; mod webdriver_bidi_session_end_response; mod webdriver_bidi_session_status_command; mod webdriver_bidi_session_status_response; +mod webdriver_bidi_session_teardown; mod webdriver_bidi_websocket_frame; mod webdriver_bidi_websocket_handshake; mod webdriver_bidi_websocket_message; mod webdriver_bidi_websocket_opening_recovery; +mod webdriver_bidi_websocket_transport_closure; #[cfg(test)] mod webdriver_bidi_json_envelope_public_boundary_tests; +// LLVM coverage keeps the crate unit-test instantiation separate from integration-test binaries, +// so compile the same realistic 1010/1011 loopback contract here instead of maintaining a copy. +#[cfg(test)] +extern crate self as originweave_network; +#[cfg(test)] +#[path = "../tests/webdriver_bidi_transport_close_role_validation.rs"] +mod webdriver_bidi_transport_close_role_validation_unit; + pub use connection::{ ConnectionPlan, DirectTcpConnection, MAX_CONNECT_TIMEOUT, MAX_CONNECTION_ATTEMPTS, NetworkError, SocketConnectionEvidence, @@ -67,6 +81,10 @@ pub use webdriver_bidi_session_status_response::{ MAX_WEBDRIVER_BIDI_SESSION_STATUS_MESSAGE_SIZE, WebDriverBiDiSessionStatusResponseError, WebDriverBiDiSessionStatusResult, }; +pub use webdriver_bidi_session_teardown::{ + WebDriverBiDiSessionTeardownAssessment, WebDriverBiDiSessionTeardownAssessmentError, + WebDriverBiDiSessionTeardownDisposition, WebDriverBiDiSessionTeardownObservations, +}; pub use webdriver_bidi_websocket_frame::{ MAX_WEBSOCKET_FRAME_PAYLOAD_SIZE, MAX_WEBSOCKET_FRAME_TIMEOUT, WebDriverBiDiWebSocketEstablished, WebDriverBiDiWebSocketFrame, @@ -86,3 +104,7 @@ pub use webdriver_bidi_websocket_message::{ WebDriverBiDiWebSocketTextMessage, }; pub use webdriver_bidi_websocket_opening_recovery::WebDriverBiDiWebSocketOpeningWriteRecoveryDisposition; +pub use webdriver_bidi_websocket_transport_closure::{ + WebDriverBiDiWebSocketTransportClosureError, WebDriverBiDiWebSocketTransportClosureKind, + WebDriverBiDiWebSocketTransportClosureObservation, +}; diff --git a/crates/originweave-network/src/webdriver_bidi_command_correlation.rs b/crates/originweave-network/src/webdriver_bidi_command_correlation.rs index a3dd0d53b..9cd3ebc16 100644 --- a/crates/originweave-network/src/webdriver_bidi_command_correlation.rs +++ b/crates/originweave-network/src/webdriver_bidi_command_correlation.rs @@ -32,6 +32,25 @@ struct OutstandingCommand { connection_generation: Option, } +fn response_route( + envelope: &WebDriverBiDiJsonEnvelope, +) -> Result<(u64, WebDriverBiDiCorrelatedResponseOutcome), WebDriverBiDiCommandCorrelationError> { + match envelope.routing() { + WebDriverBiDiJsonEnvelopeRouting::Event => { + Err(WebDriverBiDiCommandCorrelationError::EventIsNotResponse) + } + WebDriverBiDiJsonEnvelopeRouting::CommandError { command_id: None } => { + Err(WebDriverBiDiCommandCorrelationError::UncorrelatableErrorResponse) + } + WebDriverBiDiJsonEnvelopeRouting::CommandError { + command_id: Some(command_id), + } => Ok((command_id, WebDriverBiDiCorrelatedResponseOutcome::Error)), + WebDriverBiDiJsonEnvelopeRouting::CommandSuccess { command_id } => { + Ok((command_id, WebDriverBiDiCorrelatedResponseOutcome::Success)) + } + } +} + /// Outcome of a response after it has consumed the matching outstanding command identifier. #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub enum WebDriverBiDiCorrelatedResponseOutcome { @@ -239,8 +258,26 @@ impl WebDriverBiDiCommandCorrelation { envelope: &WebDriverBiDiJsonEnvelope, expected_kind: WebDriverBiDiCommandKind, ) -> Result { - let (command_id, outcome) = response_route(envelope)?; - self.complete(command_id, expected_kind, outcome) + match envelope.routing() { + WebDriverBiDiJsonEnvelopeRouting::Event => { + Err(WebDriverBiDiCommandCorrelationError::EventIsNotResponse) + } + WebDriverBiDiJsonEnvelopeRouting::CommandError { command_id: None } => { + Err(WebDriverBiDiCommandCorrelationError::UncorrelatableErrorResponse) + } + WebDriverBiDiJsonEnvelopeRouting::CommandError { + command_id: Some(command_id), + } => self.complete( + command_id, + expected_kind, + WebDriverBiDiCorrelatedResponseOutcome::Error, + ), + WebDriverBiDiJsonEnvelopeRouting::CommandSuccess { command_id } => self.complete( + command_id, + expected_kind, + WebDriverBiDiCorrelatedResponseOutcome::Success, + ), + } } pub(crate) fn correlate_response_for_connection( @@ -317,25 +354,6 @@ impl WebDriverBiDiCommandCorrelation { } } -fn response_route( - envelope: &WebDriverBiDiJsonEnvelope, -) -> Result<(u64, WebDriverBiDiCorrelatedResponseOutcome), WebDriverBiDiCommandCorrelationError> { - match envelope.routing() { - WebDriverBiDiJsonEnvelopeRouting::Event => { - Err(WebDriverBiDiCommandCorrelationError::EventIsNotResponse) - } - WebDriverBiDiJsonEnvelopeRouting::CommandError { command_id: None } => { - Err(WebDriverBiDiCommandCorrelationError::UncorrelatableErrorResponse) - } - WebDriverBiDiJsonEnvelopeRouting::CommandError { - command_id: Some(command_id), - } => Ok((command_id, WebDriverBiDiCorrelatedResponseOutcome::Error)), - WebDriverBiDiJsonEnvelopeRouting::CommandSuccess { command_id } => { - Ok((command_id, WebDriverBiDiCorrelatedResponseOutcome::Success)) - } - } -} - #[cfg(test)] mod tests { use super::{WebDriverBiDiCommandCorrelationError, WebDriverBiDiCommandKind}; diff --git a/crates/originweave-network/src/webdriver_bidi_session_end_command.rs b/crates/originweave-network/src/webdriver_bidi_session_end_command.rs index e9d3f7174..82dc95cf7 100644 --- a/crates/originweave-network/src/webdriver_bidi_session_end_command.rs +++ b/crates/originweave-network/src/webdriver_bidi_session_end_command.rs @@ -41,12 +41,12 @@ impl WebDriverBiDiSessionEndCommand { /// Register and write this exact command on an already established verified BiDi stream. /// /// Locally invalid frame deadlines fail before correlation registration and before any remote - /// side effect. Correlation then binds the command to this connection before the first possible - /// frame write. Only a reply received on this same connection can complete that registration. - /// A frame-owner preflight rejection that proves no write began retires this exact command - /// again. Once frame emission can have begun, a later failure leaves the identifier outstanding - /// because partial or full emission is ambiguous. A successful write also leaves the identifier - /// outstanding until a later correlated response proves completion. + /// side effect. Correlation then binds the command id and command family to the private + /// process-local generation of this exact established connection before the first possible + /// frame write. A frame-owner preflight rejection that proves no write began retires this exact + /// command again. Once frame emission can have begun, a later failure leaves the identifier + /// outstanding because partial or full emission is ambiguous. A successful write also leaves + /// the identifier outstanding until a later correlated response proves completion. pub fn send( self, established: WebDriverBiDiWebSocketEstablished, @@ -62,11 +62,12 @@ impl WebDriverBiDiSessionEndCommand { }, }); } + let connection_generation = established.transport_evidence().connection_generation(); correlation .register_command_for_connection( self.command_id, WebDriverBiDiCommandKind::SessionEnd, - established.transport_evidence().connection_generation(), + connection_generation, ) .map_err(|source| WebDriverBiDiSessionEndCommandError::Correlation { source })?; let message = self.serialized(); diff --git a/crates/originweave-network/src/webdriver_bidi_session_end_response.rs b/crates/originweave-network/src/webdriver_bidi_session_end_response.rs index cd284f134..003bee300 100644 --- a/crates/originweave-network/src/webdriver_bidi_session_end_response.rs +++ b/crates/originweave-network/src/webdriver_bidi_session_end_response.rs @@ -4,6 +4,7 @@ use crate::{ WebDriverBiDiCommandCorrelation, WebDriverBiDiCommandCorrelationError, WebDriverBiDiCommandKind, WebDriverBiDiCorrelatedResponseOutcome, WebDriverBiDiJsonEnvelope, WebDriverBiDiJsonEnvelopeError, WebDriverBiDiReceivedTextMessage, + webdriver_bidi_connection::WebDriverBiDiConnectionGeneration, }; /// Typed protocol acknowledgment for one correlated WebDriver BiDi `session.end` command. @@ -11,24 +12,29 @@ use crate::{ /// WebDriver BiDi defines `session.EndResult` as the extensible `EmptyResult` object. The common /// local-end envelope parser already validates the complete JSON document and requires a success /// `result` object, so this command-specific boundary intentionally retains no generic result body -/// and accepts extension members. This value proves only that the remote end returned a correlated -/// protocol success; it does not prove Chromium process exit, profile deletion, resource release, -/// or any other OriginWeave operational teardown postcondition. +/// and accepts extension members. The result retains the private process-local generation of the +/// exact connection on which the command was registered before I/O, and response admission requires +/// the complete received message to carry the same non-forgeable generation. This prevents another +/// verified connection from acknowledging the command even when both reuse the same WebDriver +/// session and command id. This value proves only a correlated protocol success; it does not prove +/// Chromium process exit, profile deletion, resource release, or any other OriginWeave operational +/// teardown postcondition. #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub struct WebDriverBiDiSessionEndResult { command_id: u64, + connection_generation: WebDriverBiDiConnectionGeneration, } impl WebDriverBiDiSessionEndResult { - /// Parse one bounded local-end message and consume its exact outstanding command on response. + /// Parse one connection-bound local-end message and consume its exact outstanding command. /// /// Complete JSON and common WebDriver BiDi envelope validation occur before correlation state - /// can be consumed. Successful responses retain only the matched command id. A correlatable - /// protocol-error response consumes its matching id and returns a typed remote failure, while - /// events, null-id errors, malformed envelopes, unknown ids, and command-kind mismatches fail - /// closed without consuming unrelated outstanding state. Only a sealed reply from the same - /// connection that registered the command can consume it; a replacement connection cannot - /// complete the request even when its session and command identifiers match. + /// can be consumed. The response must have been assembled by the connection-bound message reader + /// on the same private transport generation registered by the typed `session.end` sender. + /// Connection mismatch or missing command provenance fails without consuming the outstanding id. + /// A correlatable protocol-error response on the correct connection consumes its matching id and + /// returns a typed remote failure, while events, null-id errors, malformed envelopes, unknown ids, + /// and command-kind mismatches leave unrelated outstanding state untouched. pub fn parse_and_correlate( message: &WebDriverBiDiReceivedTextMessage, correlation: &mut WebDriverBiDiCommandCorrelation, @@ -41,11 +47,22 @@ impl WebDriverBiDiSessionEndResult { WebDriverBiDiCommandKind::SessionEnd, message.connection_generation(), ) - .map_err(|source| WebDriverBiDiSessionEndResponseError::Correlation { source })?; + .map_err(|source| match source { + WebDriverBiDiCommandCorrelationError::CommandConnectionProvenanceMissing { + command_id, + } => { + WebDriverBiDiSessionEndResponseError::MissingConnectionProvenance { command_id } + } + WebDriverBiDiCommandCorrelationError::ResponseConnectionMismatch { command_id } => { + WebDriverBiDiSessionEndResponseError::TransportConnectionMismatch { command_id } + } + source => WebDriverBiDiSessionEndResponseError::Correlation { source }, + })?; match completed.outcome() { WebDriverBiDiCorrelatedResponseOutcome::Success => Ok(Self { command_id: completed.command_id(), + connection_generation: message.connection_generation(), }), WebDriverBiDiCorrelatedResponseOutcome::Error => { Err(WebDriverBiDiSessionEndResponseError::RemoteProtocolError { @@ -60,6 +77,10 @@ impl WebDriverBiDiSessionEndResult { pub const fn command_id(&self) -> u64 { self.command_id } + + pub(crate) const fn connection_generation(&self) -> WebDriverBiDiConnectionGeneration { + self.connection_generation + } } /// Fail-closed failures while admitting one typed WebDriver BiDi `session.end` response. @@ -75,6 +96,16 @@ pub enum WebDriverBiDiSessionEndResponseError { /// Exact typed correlation failure. source: WebDriverBiDiCommandCorrelationError, }, + /// The outstanding `session.end` command lacked connection provenance required for admission. + MissingConnectionProvenance { + /// Exact local command identifier left outstanding after rejection. + command_id: u64, + }, + /// The response was received on a different verified transport from the outstanding command. + TransportConnectionMismatch { + /// Exact local command identifier left outstanding after rejection. + command_id: u64, + }, /// The remote end returned a correlatable WebDriver BiDi protocol error for this command. RemoteProtocolError { /// Exact local command identifier consumed by the protocol-error response. @@ -91,6 +122,10 @@ impl fmt::Display for WebDriverBiDiSessionEndResponseError { Self::Correlation { .. } => { formatter.write_str("WebDriver BiDi session.end response correlation failed") } + Self::MissingConnectionProvenance { .. } => formatter + .write_str("WebDriver BiDi session.end response lacks connection provenance"), + Self::TransportConnectionMismatch { .. } => formatter + .write_str("WebDriver BiDi session.end response arrived on a different connection"), Self::RemoteProtocolError { .. } => { formatter.write_str("WebDriver BiDi session.end returned a protocol error") } @@ -103,7 +138,9 @@ impl Error for WebDriverBiDiSessionEndResponseError { match self { Self::Envelope { source } => Some(source), Self::Correlation { source } => Some(source), - Self::RemoteProtocolError { .. } => None, + Self::MissingConnectionProvenance { .. } + | Self::TransportConnectionMismatch { .. } + | Self::RemoteProtocolError { .. } => None, } } } @@ -132,6 +169,22 @@ mod tests { ); assert!(correlation.source().is_some()); + let missing = + WebDriverBiDiSessionEndResponseError::MissingConnectionProvenance { command_id: 7 }; + assert_eq!( + missing.to_string(), + "WebDriver BiDi session.end response lacks connection provenance" + ); + assert!(missing.source().is_none()); + + let mismatch = + WebDriverBiDiSessionEndResponseError::TransportConnectionMismatch { command_id: 7 }; + assert_eq!( + mismatch.to_string(), + "WebDriver BiDi session.end response arrived on a different connection" + ); + assert!(mismatch.source().is_none()); + let remote = WebDriverBiDiSessionEndResponseError::RemoteProtocolError { command_id: 7 }; assert_eq!( remote.to_string(), diff --git a/crates/originweave-network/src/webdriver_bidi_session_teardown.rs b/crates/originweave-network/src/webdriver_bidi_session_teardown.rs new file mode 100644 index 000000000..d83d947db --- /dev/null +++ b/crates/originweave-network/src/webdriver_bidi_session_teardown.rs @@ -0,0 +1,147 @@ +use std::{error::Error, fmt}; + +use crate::{WebDriverBiDiSessionEndResult, WebDriverBiDiWebSocketTransportClosureObservation}; + +/// Fail-closed operational disposition derived from teardown observations. +/// +/// The current assessment can report only `OperationalTeardownPending`. Transport closure already +/// has a typed owner, while browser-process exit and task-profile removal are not represented until +/// their own typed runtime owners are connected. Missing owner evidence cannot establish operational +/// completion or grant process, profile, browser, network, policy, secret, or Agent authority. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum WebDriverBiDiSessionTeardownDisposition { + /// Operational teardown is not yet proven by typed observations from every owning boundary. + OperationalTeardownPending, +} + +/// Explicit operational observations retained after a correlated WebDriver BiDi `session.end` ack. +/// +/// Transport closure is represented by the typed observation produced only by consuming the exact +/// established WebSocket at its bounded closure-observation boundary. Browser-process exit and task +/// profile removal are deliberately absent until their own typed runtime owners can supply +/// non-forgeable evidence. +#[derive(Debug, Eq, PartialEq)] +pub struct WebDriverBiDiSessionTeardownObservations { + transport_closure_observation: Option, +} + +impl WebDriverBiDiSessionTeardownObservations { + /// Construct the typed transport observation retained by this boundary. + /// + /// Transport closure cannot be asserted with a caller-supplied boolean. `Some` requires the + /// typed closure observation returned by the bounded transport owner; `None` keeps teardown + /// pending. Process and profile state cannot be supplied as placeholders and remain unproven + /// until their typed runtime owners are connected. + #[must_use] + pub const fn new( + transport_closure_observation: Option, + ) -> Self { + Self { + transport_closure_observation, + } + } + + /// Return whether typed closure evidence for the exact acknowledged transport was supplied. + /// + /// A teardown assessment is constructed only after connection-generation equality has been + /// checked, so a retained observation belongs to the same exact connection as its protocol ack. + #[must_use] + pub const fn transport_closed_observed(&self) -> bool { + self.transport_closure_observation.is_some() + } + + /// Borrow the typed transport-closure observation when one was supplied. + #[must_use] + pub const fn transport_closure_observation( + &self, + ) -> Option<&WebDriverBiDiWebSocketTransportClosureObservation> { + self.transport_closure_observation.as_ref() + } +} + +/// Fail-closed failures while binding protocol acknowledgment to operational teardown evidence. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum WebDriverBiDiSessionTeardownAssessmentError { + /// The closure observation belongs to a different process-local connection generation. + TransportConnectionMismatch, +} + +impl fmt::Display for WebDriverBiDiSessionTeardownAssessmentError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::TransportConnectionMismatch => formatter.write_str( + "WebDriver BiDi transport closure does not match the acknowledged connection", + ), + } + } +} + +impl Error for WebDriverBiDiSessionTeardownAssessmentError {} + +/// One correlated `session.end` acknowledgment kept separate from operational teardown evidence. +/// +/// A protocol acknowledgment alone is never operational completion. Typed transport closure is +/// retained only when its private process-local connection generation matches the generation bound +/// to the acknowledged `session.end` command before I/O. Process/profile evidence remains absent +/// until its owning runtime boundaries can provide typed observations. Consequently this assessment +/// remains fail closed in the current dependency-ordered slice. +#[derive(Debug, Eq, PartialEq)] +pub struct WebDriverBiDiSessionTeardownAssessment { + protocol_ack: WebDriverBiDiSessionEndResult, + observations: WebDriverBiDiSessionTeardownObservations, +} + +impl WebDriverBiDiSessionTeardownAssessment { + /// Bind one correlated protocol acknowledgment to separately supplied operational observations. + /// + /// If typed transport closure is present, its private connection generation must equal the one + /// retained by the protocol acknowledgment. A closure from another socket is rejected even when + /// both transports reuse the same WebDriver session identifier and command id. + pub fn from_protocol_ack( + protocol_ack: WebDriverBiDiSessionEndResult, + observations: WebDriverBiDiSessionTeardownObservations, + ) -> Result { + if observations + .transport_closure_observation() + .is_some_and(|transport_closure| { + transport_closure.connection_generation() != protocol_ack.connection_generation() + }) + { + return Err(WebDriverBiDiSessionTeardownAssessmentError::TransportConnectionMismatch); + } + Ok(Self { + protocol_ack, + observations, + }) + } + + /// Return the exact command id proven by the correlated protocol acknowledgment. + #[must_use] + pub const fn command_id(&self) -> u64 { + self.protocol_ack.command_id() + } + + /// Borrow the explicit operational observations bound to this assessment. + #[must_use] + pub const fn observations(&self) -> &WebDriverBiDiSessionTeardownObservations { + &self.observations + } + + /// Return whether typed evidence proves every required operational teardown boundary. + /// + /// This is always `false` while typed browser-process-exit and task-profile-removal evidence is + /// absent from this boundary. + #[must_use] + pub const fn is_operationally_complete(&self) -> bool { + false + } + + /// Return the fail-closed disposition for the currently supplied observations. + /// + /// The current boundary cannot emit an operational-completion disposition because process and + /// profile cleanup still lack typed owner evidence. + #[must_use] + pub const fn disposition(&self) -> WebDriverBiDiSessionTeardownDisposition { + WebDriverBiDiSessionTeardownDisposition::OperationalTeardownPending + } +} diff --git a/crates/originweave-network/src/webdriver_bidi_websocket_frame.rs b/crates/originweave-network/src/webdriver_bidi_websocket_frame.rs index 4a680daed..8a0ec303a 100644 --- a/crates/originweave-network/src/webdriver_bidi_websocket_frame.rs +++ b/crates/originweave-network/src/webdriver_bidi_websocket_frame.rs @@ -284,6 +284,29 @@ impl WebDriverBiDiWebSocketEstablished { write_frame_with_clock(&mut self.raw.stream, &frame, frame_timeout, &mut now).map(|_| self) } + /// Write one final masked RFC 6455 Close response on this verified stream. + /// + /// This crate-private operation is used only after the frame reader has validated a peer Close. + /// It echoes only the validated status code when one was present and deliberately does not replay + /// arbitrary peer reason text. The caller supplies fresh masking entropy; the same adjacent-key + /// guard used by all client frame writers remains in force. + pub(crate) fn write_close_frame( + mut self, + peer_close_status_code: Option, + masking_key: WebDriverBiDiWebSocketMaskKey, + frame_timeout: Duration, + ) -> Result { + validate_frame_timeout(frame_timeout)?; + self.client_mask_keys.reserve(masking_key)?; + let status_bytes = peer_close_status_code.map(u16::to_be_bytes); + let payload = status_bytes + .as_ref() + .map_or(&[][..], |bytes| bytes.as_slice()); + let frame = serialize_client_frame(0x8, payload, masking_key); + let mut now = Instant::now; + write_frame_with_clock(&mut self.raw.stream, &frame, frame_timeout, &mut now).map(|_| self) + } + /// Read one bounded RFC 6455 frame from this verified stream. /// /// Server frames must be unmasked. Reserved bits/opcodes, non-minimal lengths, oversized @@ -469,7 +492,9 @@ impl Error for WebDriverBiDiWebSocketFrameError { } } -fn validate_frame_timeout(frame_timeout: Duration) -> Result<(), WebDriverBiDiWebSocketFrameError> { +pub(crate) fn validate_frame_timeout( + frame_timeout: Duration, +) -> Result<(), WebDriverBiDiWebSocketFrameError> { if frame_timeout.is_zero() || frame_timeout > MAX_WEBSOCKET_FRAME_TIMEOUT { return Err(WebDriverBiDiWebSocketFrameError::InvalidFrameTimeout { frame_timeout, @@ -788,7 +813,7 @@ fn validate_close_frame( }); } let status_code = u16::from_be_bytes([frame.payload()[0], frame.payload()[1]]); - if !(1000..=4999).contains(&status_code) || matches!(status_code, 1004 | 1005 | 1006 | 1015) { + if !(1000..=4999).contains(&status_code) || matches!(status_code, 1004..=1006 | 1015..=2999) { return Err(WebDriverBiDiWebSocketFrameError::MalformedFrame { reason: "Close frame status code is not valid on the wire", }); @@ -803,6 +828,128 @@ mod tests { use super::*; + #[test] + fn close_writer_checks_deadlines_masks_and_exact_wire_bytes() { + use crate::WebDriverBiDiTcpConnectionPlan; + use originweave_core::WebDriverBiDiWebSocketEndpoint; + use std::net::{TcpListener, TcpStream}; + use std::thread; + + let excessive = MAX_WEBSOCKET_FRAME_TIMEOUT + Duration::from_nanos(1); + for (timeout, status, seed_text, expected_bytes, expected_result) in [ + ( + Duration::ZERO, + Some(1000), + false, + vec![], + Err(WebDriverBiDiWebSocketFrameError::InvalidFrameTimeout { + frame_timeout: Duration::ZERO, + maximum_timeout: MAX_WEBSOCKET_FRAME_TIMEOUT, + }), + ), + ( + excessive, + Some(1000), + false, + vec![], + Err(WebDriverBiDiWebSocketFrameError::InvalidFrameTimeout { + frame_timeout: excessive, + maximum_timeout: MAX_WEBSOCKET_FRAME_TIMEOUT, + }), + ), + ( + Duration::from_secs(1), + Some(1000), + false, + vec![0x88, 0x82, 1, 2, 3, 4, 2, 0xea], + Ok(()), + ), + ( + Duration::from_secs(1), + None, + false, + vec![0x88, 0x80, 1, 2, 3, 4], + Ok(()), + ), + ( + Duration::from_secs(1), + Some(1000), + true, + vec![0x81, 0x82, 1, 2, 3, 4, 0x7a, 0x7f], + Err(WebDriverBiDiWebSocketFrameError::MalformedFrame { + reason: REUSED_CLIENT_MASK_KEY_REASON, + }), + ), + ] { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind peer"); + let address = listener.local_addr().expect("peer address"); + let peer = thread::spawn(move || { + let (mut stream, _) = listener.accept().expect("accept client"); + stream + .set_read_timeout(Some(Duration::from_secs(2))) + .expect("bound peer read"); + let mut request = Vec::new(); + while !request.ends_with(b"\r\n\r\n") { + let mut byte = [0]; + stream.read_exact(&mut byte).expect("opening request"); + request.push(byte[0]); + } + stream.write_all(b"HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=\r\n\r\n").expect("opening response"); + let mut received = vec![0; expected_bytes.len()]; + stream + .read_exact(&mut received) + .expect("expected client frame"); + assert_eq!(received, expected_bytes); + let mut byte = [0]; + assert_eq!( + TcpStream::read(&mut stream, &mut byte).expect("client EOF"), + 0 + ); + }); + let session = "01234567-89ab-cdef-0123-456789abcdef"; + let endpoint = + WebDriverBiDiWebSocketEndpoint::new(&format!("ws://{address}/session/{session}")) + .expect("endpoint"); + let target = endpoint + .correlate_session_id(session) + .expect("session") + .into_explicit_connect_target() + .expect("target"); + let connection = WebDriverBiDiTcpConnectionPlan::new(target, Duration::from_secs(1), 1) + .expect("plan") + .connect() + .expect("connect"); + let key = crate::WebDriverBiDiWebSocketClientKey::new("dGhlIHNhbXBsZSBub25jZQ==") + .expect("key"); + let established = WebDriverBiDiWebSocketHandshakePlan::new(connection, key) + .expect("handshake") + .write_opening_request(Duration::from_secs(1)) + .expect("write opening") + .read_opening_response(Duration::from_secs(1)) + .expect("read opening"); + let established = if seed_text { + established + .write_text_frame( + "{}", + WebDriverBiDiWebSocketMaskKey::new([1, 2, 3, 4]), + timeout, + ) + .expect("seed text frame") + } else { + established + }; + let result = established + .write_close_frame( + status, + WebDriverBiDiWebSocketMaskKey::new([1, 2, 3, 4]), + timeout, + ) + .map(drop); + assert_eq!(format!("{result:?}"), format!("{expected_result:?}")); + peer.join().expect("peer completed"); + } + } + #[derive(Clone, Debug)] enum ReadAction { Bytes(Vec), @@ -1310,16 +1457,30 @@ mod tests { payload: vec![0x03, 0xe8, 0xff], }; assert!(validate_close_frame(&invalid_utf8).is_err()); - for status in [999_u16, 1004, 1005, 1006, 1015, 5000] { + for status in [ + 0_u16, + 999, + 1004, + 1005, + 1006, + 1015, + 1016, + 2000, + 2999, + 5000, + u16::MAX, + ] { let payload = status.to_be_bytes().to_vec(); let frame = WebDriverBiDiWebSocketFrame { fin: true, opcode: 8, payload, }; - assert!(validate_close_frame(&frame).is_err()); + assert!(validate_close_frame(&frame).is_err(), "status {status}"); } - for status in [1000_u16, 3000, 4000] { + for status in [ + 1000_u16, 1003, 1007, 1011, 1012, 1013, 1014, 3000, 3999, 4000, 4999, + ] { let mut payload = status.to_be_bytes().to_vec(); payload.extend_from_slice(b"ok"); let frame = WebDriverBiDiWebSocketFrame { diff --git a/crates/originweave-network/src/webdriver_bidi_websocket_transport_closure.rs b/crates/originweave-network/src/webdriver_bidi_websocket_transport_closure.rs new file mode 100644 index 000000000..0269f3241 --- /dev/null +++ b/crates/originweave-network/src/webdriver_bidi_websocket_transport_closure.rs @@ -0,0 +1,626 @@ +use std::{ + error::Error, + fmt, + time::{Duration, Instant}, +}; + +use crate::{ + WebDriverBiDiWebSocketEstablished, WebDriverBiDiWebSocketFrameError, + WebDriverBiDiWebSocketMaskKey, webdriver_bidi_connection::WebDriverBiDiConnectionGeneration, +}; + +const MAX_PRE_CLOSE_CONTROL_FRAMES: usize = 64; + +struct ClosureDeadline<'a> { + expires_at: Instant, + now: &'a mut dyn FnMut() -> Instant, +} + +impl ClosureDeadline<'_> { + fn remaining(&mut self) -> Result { + let remaining = self.expires_at.saturating_duration_since((self.now)()); + if remaining.is_zero() { + return Err(WebDriverBiDiWebSocketTransportClosureError::DeadlineExpired); + } + Ok(remaining) + } +} + +/// Bounded transport-closure condition observed on one consumed WebDriver BiDi WebSocket. +/// +/// Every successful variant proves clean TCP EOF on the exact established connection. A preceding +/// RFC 6455 Close exchange is retained separately from EOF-only cessation so a consumer cannot +/// confuse entering CLOSING with the transport actually becoming closed. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum WebDriverBiDiWebSocketTransportClosureKind { + /// The peer ended the TCP byte stream cleanly before any new WebSocket frame byte was read. + PeerEof, + /// A validated peer Close was answered with a masked client Close before clean TCP EOF. + PeerCloseThenEof, +} + +/// Credential-free observation that one established WebDriver BiDi transport actually closed. +/// +/// Construction consumes the established WebSocket, so this value cannot be used to regain the +/// underlying connection. It retains the private process-local generation of that exact connection +/// so a later teardown consumer can reject closure observed on another socket even when session and +/// command identifiers are reused. Peer Close enters CLOSING only: OriginWeave answers it with a +/// caller-keyed masked Close and emits this final observation only after subsequent clean TCP EOF. +/// Pre-Close Ping is answered with an equally caller-keyed masked Pong carrying identical payload. +/// No masking entropy, browser authority, process state, profile cleanup, policy, secret, reconnect, +/// retry, or Agent authority is invented by this transport adapter. +#[derive(Debug, Eq, PartialEq)] +pub struct WebDriverBiDiWebSocketTransportClosureObservation { + kind: WebDriverBiDiWebSocketTransportClosureKind, + peer_close_status_code: Option, + connection_generation: WebDriverBiDiConnectionGeneration, +} + +impl WebDriverBiDiWebSocketTransportClosureObservation { + /// Consume one established connection and prove bounded transport closure. + /// + /// The caller supplies independent fresh masking keys for each pre-Close Ping response and + /// the required Close response. Up to 64 pre-Close Ping/Pong frames are admitted as a local + /// resource limit. Pings consume keys in slice order; unsolicited Pongs consume no key. + /// Exhausted keys, excess control traffic, application data, timeout, partial EOF, malformed + /// frames, write failures, a server-sent client-only status 1010, or any post-Close frame fail + /// closed. Peer Close alone is never success: after the masked response this boundary requires + /// zero-byte TCP EOF within one operation-wide deadline covering all reads and writes. Evidence + /// arriving after that deadline is rejected. + pub fn observe( + established: WebDriverBiDiWebSocketEstablished, + pong_masking_keys: &[WebDriverBiDiWebSocketMaskKey], + close_masking_key: WebDriverBiDiWebSocketMaskKey, + frame_timeout: Duration, + ) -> Result { + Self::observe_with_clock( + established, + pong_masking_keys, + close_masking_key, + frame_timeout, + &mut Instant::now, + ) + } + + fn observe_with_clock( + established: WebDriverBiDiWebSocketEstablished, + pong_masking_keys: &[WebDriverBiDiWebSocketMaskKey], + close_masking_key: WebDriverBiDiWebSocketMaskKey, + frame_timeout: Duration, + now: &mut dyn FnMut() -> Instant, + ) -> Result { + crate::webdriver_bidi_websocket_frame::validate_frame_timeout(frame_timeout) + .map_err(|source| WebDriverBiDiWebSocketTransportClosureError::Frame { source })?; + let mut deadline = ClosureDeadline { + expires_at: now() + frame_timeout, + now, + }; + let connection_generation = established.transport_evidence().connection_generation(); + Self::observe_pre_close( + established, + pong_masking_keys, + close_masking_key, + &mut deadline, + MAX_PRE_CLOSE_CONTROL_FRAMES, + connection_generation, + ) + } + + fn observe_pre_close( + established: WebDriverBiDiWebSocketEstablished, + pong_masking_keys: &[WebDriverBiDiWebSocketMaskKey], + close_masking_key: WebDriverBiDiWebSocketMaskKey, + deadline: &mut ClosureDeadline<'_>, + remaining_control_frames: usize, + connection_generation: WebDriverBiDiConnectionGeneration, + ) -> Result { + match established.read_frame(deadline.remaining()?) { + Ok((_established, frame)) + if matches!(frame.opcode(), 0x9 | 0xa) && remaining_control_frames == 0 => + { + Err(WebDriverBiDiWebSocketTransportClosureError::ControlFrameLimitExceeded) + } + Ok((established, frame)) if frame.opcode() == 0xa => Self::observe_pre_close( + established, + pong_masking_keys, + close_masking_key, + deadline, + remaining_control_frames - 1, + connection_generation, + ), + Ok((established, frame)) if frame.opcode() == 0x9 => { + let (pong_masking_key, remaining_keys) = pong_masking_keys + .split_first() + .ok_or(WebDriverBiDiWebSocketTransportClosureError::PongMaskingKeysExhausted)?; + let established = established + .write_pong_frame(frame.payload(), *pong_masking_key, deadline.remaining()?) + .map_err( + |source| WebDriverBiDiWebSocketTransportClosureError::Frame { source }, + )?; + Self::observe_pre_close( + established, + remaining_keys, + close_masking_key, + deadline, + remaining_control_frames - 1, + connection_generation, + ) + } + Ok((established, frame)) if frame.opcode() == 0x8 => { + let peer_close_status_code = frame + .payload() + .get(..2) + .map(|bytes| u16::from_be_bytes([bytes[0], bytes[1]])); + if peer_close_status_code == Some(1010) { + return Err( + WebDriverBiDiWebSocketTransportClosureError::PeerCloseStatusNotAllowed { + status_code: 1010, + }, + ); + } + let established = established + .write_close_frame( + peer_close_status_code, + close_masking_key, + deadline.remaining()?, + ) + .map_err( + |source| WebDriverBiDiWebSocketTransportClosureError::Frame { source }, + )?; + Self::observe_eof_after_close( + established, + deadline, + peer_close_status_code, + connection_generation, + ) + } + Ok((_established, frame)) => Err( + WebDriverBiDiWebSocketTransportClosureError::UnexpectedFrame { + opcode: frame.opcode(), + }, + ), + Err(WebDriverBiDiWebSocketFrameError::FrameEnded { bytes_read: 0 }) => { + deadline.remaining()?; + Ok(Self { + kind: WebDriverBiDiWebSocketTransportClosureKind::PeerEof, + peer_close_status_code: None, + connection_generation, + }) + } + Err(source) => Err(WebDriverBiDiWebSocketTransportClosureError::Frame { source }), + } + } + + fn observe_eof_after_close( + established: WebDriverBiDiWebSocketEstablished, + deadline: &mut ClosureDeadline<'_>, + peer_close_status_code: Option, + connection_generation: WebDriverBiDiConnectionGeneration, + ) -> Result { + match established.read_frame(deadline.remaining()?) { + Err(WebDriverBiDiWebSocketFrameError::FrameEnded { bytes_read: 0 }) => { + deadline.remaining()?; + Ok(Self { + kind: WebDriverBiDiWebSocketTransportClosureKind::PeerCloseThenEof, + peer_close_status_code, + connection_generation, + }) + } + Ok((_established, frame)) => Err( + WebDriverBiDiWebSocketTransportClosureError::UnexpectedFrame { + opcode: frame.opcode(), + }, + ), + Err(source) => Err(WebDriverBiDiWebSocketTransportClosureError::Frame { source }), + } + } + + /// Return the exact transport-closure condition that produced this observation. + #[must_use] + pub const fn kind(&self) -> WebDriverBiDiWebSocketTransportClosureKind { + self.kind + } + + /// Return the validated peer Close status when the completed closing exchange carried one. + #[must_use] + pub const fn peer_close_status_code(&self) -> Option { + self.peer_close_status_code + } + + pub(crate) const fn connection_generation(&self) -> WebDriverBiDiConnectionGeneration { + self.connection_generation + } +} + +/// Fail-closed errors while converting one established BiDi transport into closure evidence. +#[derive(Debug)] +pub enum WebDriverBiDiWebSocketTransportClosureError { + /// More than 64 pre-Close Ping/Pong frames exceeded the local resource budget. + ControlFrameLimitExceeded, + /// A Ping required another caller-supplied masking key; no response was emitted for it. + PongMaskingKeysExhausted, + /// The operation-wide deadline expired before transport-closure evidence was admitted. + DeadlineExpired, + /// The server sent a wire-valid Close status whose meaning is reserved for clients. + PeerCloseStatusNotAllowed { + /// Exact peer-supplied Close status rejected by the known client/server role boundary. + status_code: u16, + }, + /// The peer sent a valid WebSocket frame outside the bounded closing state machine. + UnexpectedFrame { + /// Exact validated RFC 6455 opcode observed instead of admissible bounded closing traffic. + opcode: u8, + }, + /// The existing bounded WebSocket frame reader or writer failed before closure was proven. + Frame { + /// Original typed frame failure retained as the causal source. + source: WebDriverBiDiWebSocketFrameError, + }, +} + +impl fmt::Display for WebDriverBiDiWebSocketTransportClosureError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::ControlFrameLimitExceeded => { + formatter.write_str("WebDriver BiDi transport closure control-frame limit exceeded") + } + Self::PongMaskingKeysExhausted => formatter + .write_str("WebDriver BiDi transport closure requires another Pong masking key"), + Self::DeadlineExpired => { + formatter.write_str("WebDriver BiDi transport closure deadline expired") + } + Self::PeerCloseStatusNotAllowed { .. } => formatter + .write_str("WebDriver BiDi server sent a Close status reserved for clients"), + Self::UnexpectedFrame { .. } => formatter + .write_str("WebDriver BiDi peer sent non-closure traffic instead of closing"), + Self::Frame { .. } => { + formatter.write_str("WebDriver BiDi transport closure could not be observed safely") + } + } + } +} + +impl Error for WebDriverBiDiWebSocketTransportClosureError { + fn source(&self) -> Option<&(dyn Error + 'static)> { + match self { + Self::UnexpectedFrame { .. } + | Self::DeadlineExpired + | Self::ControlFrameLimitExceeded + | Self::PongMaskingKeysExhausted + | Self::PeerCloseStatusNotAllowed { .. } => None, + Self::Frame { source } => Some(source), + } + } +} + +#[cfg(test)] +#[allow(clippy::expect_used)] +mod tests { + use super::*; + use crate::{ + WebDriverBiDiTcpConnectionPlan, WebDriverBiDiWebSocketClientKey, + WebDriverBiDiWebSocketHandshakePlan, + }; + use originweave_core::WebDriverBiDiWebSocketEndpoint; + use std::{ + io::{Read, Write}, + net::{Shutdown, TcpListener, TcpStream}, + thread, + }; + + fn established_peer(frames: &[u8]) -> (WebDriverBiDiWebSocketEstablished, TcpStream) { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind peer"); + let address = listener.local_addr().expect("peer address"); + let peer = thread::spawn(move || { + let (mut stream, _) = listener.accept().expect("accept client"); + stream + .set_read_timeout(Some(Duration::from_secs(2))) + .expect("bound peer"); + let mut request = Vec::new(); + while !request.ends_with(b"\r\n\r\n") { + let mut byte = [0]; + stream.read_exact(&mut byte).expect("opening request"); + request.push(byte[0]); + } + stream.write_all(b"HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=\r\n\r\n").expect("opening response"); + stream + }); + let session = "01234567-89ab-cdef-0123-456789abcdef"; + let endpoint = + WebDriverBiDiWebSocketEndpoint::new(&format!("ws://{address}/session/{session}")) + .expect("endpoint"); + let target = endpoint + .correlate_session_id(session) + .expect("session") + .into_explicit_connect_target() + .expect("target"); + let connection = WebDriverBiDiTcpConnectionPlan::new(target, Duration::from_secs(1), 1) + .expect("plan") + .connect() + .expect("connect"); + let key = WebDriverBiDiWebSocketClientKey::new("dGhlIHNhbXBsZSBub25jZQ==").expect("key"); + let established = WebDriverBiDiWebSocketHandshakePlan::new(connection, key) + .expect("handshake") + .write_opening_request(Duration::from_secs(1)) + .expect("opening write") + .read_opening_response(Duration::from_secs(1)) + .expect("opening read"); + let mut peer = peer.join().expect("peer handshake completed"); + peer.write_all(frames).expect("peer frames"); + peer.shutdown(Shutdown::Write).expect("peer write EOF"); + (established, peer) + } + + #[test] + fn control_budget_and_key_exhaustion_preserve_exact_failures() { + use WebDriverBiDiWebSocketTransportClosureError::{ + ControlFrameLimitExceeded, PongMaskingKeysExhausted, + }; + for (frames, expected, message) in [ + ( + vec![0x89, 0], + PongMaskingKeysExhausted, + "WebDriver BiDi transport closure requires another Pong masking key", + ), + ( + [[0x8a, 0].repeat(64), vec![0x89, 0]].concat(), + ControlFrameLimitExceeded, + "WebDriver BiDi transport closure control-frame limit exceeded", + ), + ( + [0x8a, 0].repeat(65), + ControlFrameLimitExceeded, + "WebDriver BiDi transport closure control-frame limit exceeded", + ), + ] { + let (established, mut peer) = established_peer(&frames); + let now = Instant::now(); + let error = WebDriverBiDiWebSocketTransportClosureObservation::observe_with_clock( + established, + &[], + WebDriverBiDiWebSocketMaskKey::new([5, 6, 7, 8]), + Duration::from_secs(1), + &mut || now, + ) + .expect_err("control traffic must stay bounded"); + assert_eq!(format!("{error:?}"), format!("{expected:?}")); + assert_eq!(error.to_string(), message); + assert!(error.source().is_none()); + let mut reply = [0]; + use std::io::Read; + assert_eq!(peer.read(&mut reply).expect("peer EOF"), 0); + } + } + + #[test] + fn exact_control_budget_allows_close_with_fresh_ping_keys() { + for opcode in [0x9, 0xa] { + let frames = [[0x80 | opcode, 0].repeat(64), vec![0x88, 0]].concat(); + let (established, _peer) = established_peer(&frames); + let keys: Vec<_> = (0..64) + .map(|index| WebDriverBiDiWebSocketMaskKey::new([index, 1, 2, 3])) + .collect(); + let now = Instant::now(); + let observation = + WebDriverBiDiWebSocketTransportClosureObservation::observe_with_clock( + established, + &keys, + WebDriverBiDiWebSocketMaskKey::new([5, 6, 7, 8]), + Duration::from_secs(1), + &mut || now, + ) + .expect("Close after exactly 64 controls"); + assert_eq!( + observation.kind(), + WebDriverBiDiWebSocketTransportClosureKind::PeerCloseThenEof + ); + } + } + + #[test] + fn fixed_clock_accepts_complete_closure_before_deadline() { + for (frames, kind, status) in [ + ( + &[][..], + WebDriverBiDiWebSocketTransportClosureKind::PeerEof, + None, + ), + ( + &[0x88, 0][..], + WebDriverBiDiWebSocketTransportClosureKind::PeerCloseThenEof, + None, + ), + ( + &[0x88, 2, 3, 0xe8][..], + WebDriverBiDiWebSocketTransportClosureKind::PeerCloseThenEof, + Some(1000), + ), + ( + &[0x89, 0, 0x88, 0][..], + WebDriverBiDiWebSocketTransportClosureKind::PeerCloseThenEof, + None, + ), + ( + &[0x8a, 0, 0x88, 0][..], + WebDriverBiDiWebSocketTransportClosureKind::PeerCloseThenEof, + None, + ), + ] { + let (established, _peer) = established_peer(frames); + let now = Instant::now(); + let observation = + WebDriverBiDiWebSocketTransportClosureObservation::observe_with_clock( + established, + &[WebDriverBiDiWebSocketMaskKey::new([1, 2, 3, 4])], + WebDriverBiDiWebSocketMaskKey::new([5, 6, 7, 8]), + Duration::from_secs(1), + &mut || now, + ) + .expect("closure before deadline"); + assert_eq!(observation.kind(), kind); + assert_eq!(observation.peer_close_status_code(), status); + } + } + + #[test] + fn invalid_deadline_is_rejected_before_clock_use() { + for (timeout, expected_calls) in [ + (Duration::ZERO, 0), + ( + crate::MAX_WEBSOCKET_FRAME_TIMEOUT + Duration::from_nanos(1), + 0, + ), + (Duration::from_secs(1), 3), + ] { + let (established, _peer) = established_peer(&[]); + let mut clock_calls = 0; + let result = WebDriverBiDiWebSocketTransportClosureObservation::observe_with_clock( + established, + &[WebDriverBiDiWebSocketMaskKey::new([1, 2, 3, 4])], + WebDriverBiDiWebSocketMaskKey::new([5, 6, 7, 8]), + timeout, + &mut || { + clock_calls += 1; + Instant::now() + }, + ); + let expected = if expected_calls == 0 { + Err(WebDriverBiDiWebSocketTransportClosureError::Frame { + source: WebDriverBiDiWebSocketFrameError::InvalidFrameTimeout { + frame_timeout: timeout, + maximum_timeout: crate::MAX_WEBSOCKET_FRAME_TIMEOUT, + }, + }) + } else { + Ok(()) + }; + assert_eq!(format!("{:?}", result.map(|_| ())), format!("{expected:?}")); + assert_eq!(clock_calls, expected_calls); + } + } + + #[test] + fn deadline_expiry_is_checked_at_every_closure_transition() { + for (frames, expire_tick) in [ + (&[][..], 2), + (&[0x88, 0][..], 1), + (&[0x88, 0][..], 2), + (&[0x88, 0][..], 3), + (&[0x88, 0][..], 4), + (&[0x89, 0, 0x88, 0][..], 2), + (&[0x89, 0, 0x88, 0][..], 3), + (&[0x89, 0, 0x88, 0][..], 4), + (&[0x89, 0, 0x88, 0][..], 5), + (&[0x89, 0, 0x88, 0][..], 6), + (&[0x8a, 0, 0x88, 0][..], 2), + ] { + let (established, _peer) = established_peer(frames); + let start = Instant::now(); + let mut tick = 0; + let mut now = || { + let elapsed = Duration::from_secs(u64::from(tick >= expire_tick)); + tick += 1; + start + elapsed + }; + let error = WebDriverBiDiWebSocketTransportClosureObservation::observe_with_clock( + established, + &[WebDriverBiDiWebSocketMaskKey::new([1, 2, 3, 4])], + WebDriverBiDiWebSocketMaskKey::new([5, 6, 7, 8]), + Duration::from_secs(1), + &mut now, + ) + .expect_err("expired closure must not produce evidence"); + assert_eq!(format!("{error:?}"), "DeadlineExpired"); + assert_eq!( + error.to_string(), + "WebDriver BiDi transport closure deadline expired" + ); + assert!(error.source().is_none()); + } + } + + #[test] + fn fixed_deadline_preserves_protocol_and_masking_failures() { + use WebDriverBiDiWebSocketFrameError::{FrameEnded, MalformedFrame}; + use WebDriverBiDiWebSocketTransportClosureError::{Frame, UnexpectedFrame}; + let reused = + "client masking key was reused for consecutive frames on this established WebSocket"; + for (frames, seed_text, pong_key, close_key, expected) in [ + ( + &[0x81, 0][..], + false, + [1, 2, 3, 4], + [5, 6, 7, 8], + UnexpectedFrame { opcode: 1 }, + ), + ( + &[0x88][..], + false, + [1, 2, 3, 4], + [5, 6, 7, 8], + Frame { + source: FrameEnded { bytes_read: 1 }, + }, + ), + ( + &[0x88, 0, 0x8a, 0][..], + false, + [1, 2, 3, 4], + [5, 6, 7, 8], + UnexpectedFrame { opcode: 0xa }, + ), + ( + &[0x88, 0, 0x88][..], + false, + [1, 2, 3, 4], + [5, 6, 7, 8], + Frame { + source: FrameEnded { bytes_read: 1 }, + }, + ), + ( + &[0x89, 0][..], + true, + [1, 2, 3, 4], + [5, 6, 7, 8], + Frame { + source: MalformedFrame { reason: reused }, + }, + ), + ( + &[0x88, 0][..], + true, + [5, 6, 7, 8], + [1, 2, 3, 4], + Frame { + source: MalformedFrame { reason: reused }, + }, + ), + ] { + let (established, _peer) = established_peer(frames); + let established = if seed_text { + established + .write_text_frame( + "{}", + WebDriverBiDiWebSocketMaskKey::new([1, 2, 3, 4]), + Duration::from_secs(1), + ) + .expect("seed masking history") + } else { + established + }; + let now = Instant::now(); + let error = WebDriverBiDiWebSocketTransportClosureObservation::observe_with_clock( + established, + &[WebDriverBiDiWebSocketMaskKey::new(pong_key)], + WebDriverBiDiWebSocketMaskKey::new(close_key), + Duration::from_secs(1), + &mut || now, + ) + .expect_err("protocol failure must survive the deadline wrapper"); + assert_eq!(format!("{error:?}"), format!("{expected:?}")); + assert_eq!(error.to_string(), expected.to_string()); + assert_eq!(error.source().is_some(), expected.source().is_some()); + } + } +} diff --git a/crates/originweave-network/tests/webdriver_bidi_session_end_response.rs b/crates/originweave-network/tests/webdriver_bidi_session_end_response.rs index 610526556..2d243d33b 100644 --- a/crates/originweave-network/tests/webdriver_bidi_session_end_response.rs +++ b/crates/originweave-network/tests/webdriver_bidi_session_end_response.rs @@ -9,7 +9,7 @@ use std::{ use originweave_core::WebDriverBiDiWebSocketEndpoint; use originweave_network::{ WebDriverBiDiCommandCorrelation, WebDriverBiDiCommandCorrelationError, - WebDriverBiDiConnectionMessageRead, WebDriverBiDiReceivedTextMessage, + WebDriverBiDiCommandKind, WebDriverBiDiConnectionMessageRead, WebDriverBiDiReceivedTextMessage, WebDriverBiDiSessionEndCommand, WebDriverBiDiSessionEndResponseError, WebDriverBiDiSessionEndResult, WebDriverBiDiTcpConnectionPlan, WebDriverBiDiWebSocketClientKey, WebDriverBiDiWebSocketHandshakePlan, WebDriverBiDiWebSocketMaskKey, @@ -23,76 +23,14 @@ const END_SUCCESS_RESPONSE: &[u8] = br#"{"type":"success","id":7,"result":{"vendorExtension":{"clean":true}}}"#; const END_REMOTE_ERROR_RESPONSE: &[u8] = br#"{"type":"error","id":7,"error":"unknown error","message":"remote refused"}"#; +const END_NULL_ID_ERROR_RESPONSE: &[u8] = + br#"{"type":"error","id":null,"error":"unknown error","message":"unattributed"}"#; +const END_EVENT_RESPONSE: &[u8] = + br#"{"type":"event","method":"log.entryAdded","params":{"text":"not a response"}}"#; const END_UNKNOWN_ID_RESPONSE: &[u8] = br#"{"type":"success","id":8,"result":{"vendorExtension":true}}"#; const END_MALFORMED_RESPONSE: &[u8] = br#"{"type":"success","id":7}"#; -#[test] -fn unbound_end_command_cannot_consume_a_connection_bound_reply() -> Result<(), Box> { - use originweave_network::WebDriverBiDiCommandKind; - - let (message, mut original) = send_end_and_read_response(END_SUCCESS_RESPONSE)?; - let mut unbound = WebDriverBiDiCommandCorrelation::new(); - unbound.register_command_for(7, WebDriverBiDiCommandKind::SessionEnd)?; - assert!(matches!( - WebDriverBiDiSessionEndResult::parse_and_correlate(&message, &mut unbound), - Err(WebDriverBiDiSessionEndResponseError::Correlation { - source: WebDriverBiDiCommandCorrelationError::CommandConnectionProvenanceMissing { - command_id: 7, - }, - }) - )); - assert_eq!(unbound.outstanding_count(), 1); - let result = WebDriverBiDiSessionEndResult::parse_and_correlate(&message, &mut original)?; - assert_eq!(result.command_id(), 7); - assert_eq!(original.outstanding_count(), 0); - Ok(()) -} - -#[test] -fn event_and_null_id_error_preserve_the_sent_end_command() -> Result<(), Box> { - for (document, expected) in [ - ( - br#"{"type":"event","method":"log.entryAdded","params":{}}"#.as_slice(), - WebDriverBiDiCommandCorrelationError::EventIsNotResponse, - ), - ( - br#"{"type":"error","id":null,"error":"unknown error","message":"remote"}"#.as_slice(), - WebDriverBiDiCommandCorrelationError::UncorrelatableErrorResponse, - ), - ] { - let (message, mut correlation) = send_end_and_read_response(document)?; - assert!(matches!( - WebDriverBiDiSessionEndResult::parse_and_correlate(&message, &mut correlation), - Err(WebDriverBiDiSessionEndResponseError::Correlation { source }) if source == expected - )); - assert_eq!(correlation.outstanding_count(), 1); - } - Ok(()) -} - -#[test] -fn replacement_end_replies_preserve_original_pending_request_and_recovery() --> Result<(), Box> { - for response in [END_SUCCESS_RESPONSE, END_REMOTE_ERROR_RESPONSE] { - let (original, mut pending) = send_end_and_read_response(END_SUCCESS_RESPONSE)?; - let (replacement, _) = send_end_and_read_response(response)?; - assert!(matches!( - WebDriverBiDiSessionEndResult::parse_and_correlate(&replacement, &mut pending), - Err(WebDriverBiDiSessionEndResponseError::Correlation { - source: WebDriverBiDiCommandCorrelationError::ResponseConnectionMismatch { - command_id: 7 - } - }) - )); - assert_eq!(pending.outstanding_count(), 1); - let result = WebDriverBiDiSessionEndResult::parse_and_correlate(&original, &mut pending)?; - assert_eq!(result.command_id(), 7); - assert_eq!(pending.outstanding_count(), 0); - } - Ok(()) -} - fn read_opening_request(stream: &mut TcpStream) -> io::Result<()> { stream.set_read_timeout(Some(Duration::from_secs(2)))?; let mut request = Vec::new(); @@ -181,15 +119,17 @@ fn send_end_and_read_response( Duration::from_millis(500), )?; - let text = match WebDriverBiDiWebSocketMessageReader::new(established) - .read_next(Duration::from_millis(500))? - { - WebDriverBiDiConnectionMessageRead::Text { message, .. } => message, - other => { - return Err(io::Error::other(format!( - "session.end response produced unexpected assembly state: {other:?}" - )) - .into()); + let mut reader = WebDriverBiDiWebSocketMessageReader::new(established); + let text = loop { + match reader.read_next(Duration::from_millis(500))? { + WebDriverBiDiConnectionMessageRead::Pending(next) => reader = next, + WebDriverBiDiConnectionMessageRead::Text { message, .. } => break message, + WebDriverBiDiConnectionMessageRead::Control { message, .. } => { + return Err(io::Error::other(format!( + "session.end response produced unexpected control message: {message:?}" + )) + .into()); + } } }; @@ -200,8 +140,7 @@ fn send_end_and_read_response( } #[test] -fn session_end_success_accepts_extensible_empty_result_and_consumes_exact_correlation() --> Result<(), Box> { +fn session_end_success_consumes_exact_connection_correlation() -> Result<(), Box> { let (text, mut correlation) = send_end_and_read_response(END_SUCCESS_RESPONSE)?; assert_eq!(correlation.outstanding_count(), 1); @@ -211,6 +150,32 @@ fn session_end_success_accepts_extensible_empty_result_and_consumes_exact_correl Ok(()) } +#[test] +fn session_end_rejects_unbound_connection_provenance() -> Result<(), Box> { + let (text, _connection_bound_correlation) = send_end_and_read_response(END_SUCCESS_RESPONSE)?; + let mut unbound = WebDriverBiDiCommandCorrelation::new(); + unbound.register_command_for(7, WebDriverBiDiCommandKind::SessionEnd)?; + + let parsed = WebDriverBiDiSessionEndResult::parse_and_correlate(&text, &mut unbound); + let error = parsed + .err() + .ok_or_else(|| io::Error::other("unbound session.end correlation was accepted"))?; + assert!(matches!( + error, + WebDriverBiDiSessionEndResponseError::MissingConnectionProvenance { command_id: 7 } + )); + assert_eq!( + error.to_string(), + "WebDriver BiDi session.end response lacks connection provenance" + ); + assert_eq!( + unbound.outstanding_count(), + 1, + "missing provenance must not consume the outstanding command" + ); + Ok(()) +} + #[test] fn session_end_remote_error_consumes_only_the_correlated_command() -> Result<(), Box> { let (text, mut correlation) = send_end_and_read_response(END_REMOTE_ERROR_RESPONSE)?; @@ -238,8 +203,43 @@ fn session_end_remote_error_consumes_only_the_correlated_command() -> Result<(), } #[test] -fn malformed_session_end_envelope_fails_before_consuming_correlation() -> Result<(), Box> -{ +fn null_id_and_event_preserve_bound_correlation() -> Result<(), Box> { + for (response, expected) in [ + ( + END_NULL_ID_ERROR_RESPONSE, + WebDriverBiDiCommandCorrelationError::UncorrelatableErrorResponse, + ), + ( + END_EVENT_RESPONSE, + WebDriverBiDiCommandCorrelationError::EventIsNotResponse, + ), + ] { + let (text, mut correlation) = send_end_and_read_response(response)?; + let error = WebDriverBiDiSessionEndResult::parse_and_correlate(&text, &mut correlation) + .err() + .ok_or_else(|| io::Error::other("uncorrelatable message was accepted"))?; + match error { + WebDriverBiDiSessionEndResponseError::Correlation { source } => { + assert_eq!(source, expected); + } + other => { + return Err(io::Error::other(format!( + "uncorrelatable message produced unexpected error: {other}" + )) + .into()); + } + } + assert_eq!( + correlation.outstanding_count(), + 1, + "uncorrelatable message must leave session.end outstanding" + ); + } + Ok(()) +} + +#[test] +fn malformed_envelope_preserves_correlation() -> Result<(), Box> { let (text, mut correlation) = send_end_and_read_response(END_MALFORMED_RESPONSE)?; let parsed = WebDriverBiDiSessionEndResult::parse_and_correlate(&text, &mut correlation); let error = match parsed { @@ -263,8 +263,7 @@ fn malformed_session_end_envelope_fails_before_consuming_correlation() -> Result } #[test] -fn unknown_session_end_response_id_does_not_consume_the_outstanding_command() --> Result<(), Box> { +fn unknown_response_id_preserves_outstanding_command() -> Result<(), Box> { let (text, mut correlation) = send_end_and_read_response(END_UNKNOWN_ID_RESPONSE)?; let parsed = WebDriverBiDiSessionEndResult::parse_and_correlate(&text, &mut correlation); let error = match parsed { diff --git a/crates/originweave-network/tests/webdriver_bidi_session_end_response_connection_provenance.rs b/crates/originweave-network/tests/webdriver_bidi_session_end_response_connection_provenance.rs new file mode 100644 index 000000000..2f3c69fa0 --- /dev/null +++ b/crates/originweave-network/tests/webdriver_bidi_session_end_response_connection_provenance.rs @@ -0,0 +1,160 @@ +use std::{ + error::Error, + io::{self, Read, Write}, + net::{SocketAddr, TcpListener, TcpStream}, + thread, + time::Duration, +}; + +use originweave_core::WebDriverBiDiWebSocketEndpoint; +use originweave_network::{ + WebDriverBiDiCommandCorrelation, WebDriverBiDiConnectionMessageRead, + WebDriverBiDiReceivedTextMessage, WebDriverBiDiSessionEndCommand, + WebDriverBiDiSessionEndResponseError, WebDriverBiDiSessionEndResult, + WebDriverBiDiTcpConnectionPlan, WebDriverBiDiWebSocketClientKey, + WebDriverBiDiWebSocketHandshakePlan, WebDriverBiDiWebSocketMaskKey, + WebDriverBiDiWebSocketMessageReader, +}; + +const SESSION_ID: &str = "01234567-89ab-cdef-0123-456789abcdef"; +const RFC6455_SAMPLE_KEY: &str = "dGhlIHNhbXBsZSBub25jZQ=="; +const 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"; +const END_SUCCESS_RESPONSE: &[u8] = br#"{"type":"success","id":7,"result":{}}"#; + +fn read_opening_request(stream: &mut TcpStream) -> io::Result<()> { + stream.set_read_timeout(Some(Duration::from_secs(2)))?; + let mut request = Vec::new(); + let mut buffer = [0_u8; 512]; + while !request.ends_with(b"\r\n\r\n") { + let count = stream.read(&mut buffer)?; + if count == 0 { + return Err(io::Error::new( + io::ErrorKind::UnexpectedEof, + "client opening request ended before the header terminator", + )); + } + request.extend_from_slice(&buffer[..count]); + } + Ok(()) +} + +fn read_masked_text_frame(stream: &mut TcpStream) -> io::Result> { + let mut header = [0_u8; 2]; + stream.read_exact(&mut header)?; + if header[0] != 0x81 || header[1] & 0x80 == 0 { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "expected one final masked client text frame", + )); + } + let length = usize::from(header[1] & 0x7f); + if length > 125 { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "session.end command unexpectedly required extended framing", + )); + } + let mut mask = [0_u8; 4]; + stream.read_exact(&mut mask)?; + let mut payload = vec![0_u8; length]; + stream.read_exact(&mut payload)?; + for (index, byte) in payload.iter_mut().enumerate() { + *byte ^= mask[index % mask.len()]; + } + Ok(payload) +} + +fn establish( + local_addr: SocketAddr, +) -> Result> { + let endpoint = format!("ws://{local_addr}/session/{SESSION_ID}"); + let target = WebDriverBiDiWebSocketEndpoint::new(&endpoint)? + .correlate_session_id(SESSION_ID)? + .into_explicit_connect_target()?; + let connection = + WebDriverBiDiTcpConnectionPlan::new(target, Duration::from_secs(1), 1)?.connect()?; + let key = WebDriverBiDiWebSocketClientKey::new(RFC6455_SAMPLE_KEY)?; + Ok(WebDriverBiDiWebSocketHandshakePlan::new(connection, key)? + .write_opening_request(Duration::from_millis(500))? + .read_opening_response(Duration::from_millis(500))?) +} + +fn read_one_connection_bound_text( + established: originweave_network::WebDriverBiDiWebSocketEstablished, +) -> Result> { + let mut reader = WebDriverBiDiWebSocketMessageReader::new(established); + loop { + match reader.read_next(Duration::from_millis(500))? { + WebDriverBiDiConnectionMessageRead::Pending(next) => reader = next, + WebDriverBiDiConnectionMessageRead::Text { message, .. } => return Ok(message), + WebDriverBiDiConnectionMessageRead::Control { message, .. } => { + return Err(io::Error::other(format!( + "foreign response produced unexpected control message: {message:?}" + )) + .into()); + } + } + } +} + +#[test] +fn reconnected_response_cannot_consume_prior_connection_command() -> Result<(), Box> { + let listener = TcpListener::bind(("127.0.0.1", 0))?; + let local_addr = listener.local_addr()?; + let server = thread::spawn(move || -> io::Result<()> { + let (mut first, _) = listener.accept()?; + read_opening_request(&mut first)?; + first.write_all(OPENING_RESPONSE)?; + let command = read_masked_text_frame(&mut first)?; + if command != br#"{"id":7,"method":"session.end","params":{}}"# { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "unexpected session.end command on first connection", + )); + } + + let (mut second, _) = listener.accept()?; + read_opening_request(&mut second)?; + second.write_all(OPENING_RESPONSE)?; + second.write_all(&[0x81, END_SUCCESS_RESPONSE.len() as u8])?; + second.write_all(END_SUCCESS_RESPONSE)?; + Ok(()) + }); + + let first = establish(local_addr)?; + let mut correlation = WebDriverBiDiCommandCorrelation::new(); + let _first = WebDriverBiDiSessionEndCommand::new(7)?.send( + first, + &mut correlation, + WebDriverBiDiWebSocketMaskKey::new([1, 2, 3, 4]), + Duration::from_millis(500), + )?; + assert_eq!(correlation.outstanding_count(), 1); + + let second = establish(local_addr)?; + let foreign_response = read_one_connection_bound_text(second)?; + let parsed = + WebDriverBiDiSessionEndResult::parse_and_correlate(&foreign_response, &mut correlation); + let error = parsed + .err() + .ok_or_else(|| io::Error::other("foreign response acknowledged prior connection"))?; + + assert!(matches!( + error, + WebDriverBiDiSessionEndResponseError::TransportConnectionMismatch { command_id: 7 } + )); + assert_eq!( + error.to_string(), + "WebDriver BiDi session.end response arrived on a different connection" + ); + assert_eq!( + correlation.outstanding_count(), + 1, + "foreign response must not consume connection A correlation" + ); + + server + .join() + .map_err(|_| io::Error::other("connection-provenance test server panicked"))??; + Ok(()) +} diff --git a/crates/originweave-network/tests/webdriver_bidi_session_status_response.rs b/crates/originweave-network/tests/webdriver_bidi_session_status_response.rs index 6779da6c1..19a4b2219 100644 --- a/crates/originweave-network/tests/webdriver_bidi_session_status_response.rs +++ b/crates/originweave-network/tests/webdriver_bidi_session_status_response.rs @@ -25,6 +25,22 @@ const STATUS_RESPONSE_MISSING_READY: &[u8] = br#"{"type":"success","id":7,"result":{"message":"capacity available"}}"#; const STATUS_RESPONSE_EMPTY_RESULT: &[u8] = br#"{"type":"success","id":7,"result":{}}"#; +#[test] +fn teardown_stack_rejects_replacement_status_reply_and_preserves_original_request() +-> Result<(), Box> { + let (original, mut pending) = send_status_and_read_response(STATUS_RESPONSE)?; + let (replacement, _) = send_status_and_read_response(STATUS_RESPONSE)?; + assert!( + WebDriverBiDiSessionStatusResult::parse_and_correlate(&replacement, &mut pending).is_err(), + "a replacement status reply must not complete the original request" + ); + assert_eq!(pending.outstanding_count(), 1); + let result = WebDriverBiDiSessionStatusResult::parse_and_correlate(&original, &mut pending)?; + assert_eq!(result.command_id(), 7); + assert_eq!(pending.outstanding_count(), 0); + Ok(()) +} + #[test] fn unbound_command_cannot_consume_a_connection_bound_reply() -> Result<(), Box> { use originweave_network::{WebDriverBiDiCommandCorrelationError, WebDriverBiDiCommandKind}; diff --git a/crates/originweave-network/tests/webdriver_bidi_session_teardown_assessment.rs b/crates/originweave-network/tests/webdriver_bidi_session_teardown_assessment.rs new file mode 100644 index 000000000..545de3d8a --- /dev/null +++ b/crates/originweave-network/tests/webdriver_bidi_session_teardown_assessment.rs @@ -0,0 +1,327 @@ +use std::{ + error::Error, + io::{self, Read, Write}, + net::{TcpListener, TcpStream}, + thread, + time::Duration, +}; + +use originweave_core::WebDriverBiDiWebSocketEndpoint; +use originweave_network::{ + WebDriverBiDiCommandCorrelation, WebDriverBiDiConnectionMessageRead, + WebDriverBiDiSessionEndCommand, WebDriverBiDiSessionEndResult, + WebDriverBiDiSessionTeardownAssessment, WebDriverBiDiSessionTeardownAssessmentError, + WebDriverBiDiSessionTeardownDisposition, WebDriverBiDiSessionTeardownObservations, + WebDriverBiDiTcpConnectionPlan, WebDriverBiDiWebSocketClientKey, + WebDriverBiDiWebSocketEstablished, WebDriverBiDiWebSocketHandshakePlan, + WebDriverBiDiWebSocketMaskKey, WebDriverBiDiWebSocketMessageReader, + WebDriverBiDiWebSocketTransportClosureKind, WebDriverBiDiWebSocketTransportClosureObservation, +}; + +const SESSION_ID: &str = "01234567-89ab-cdef-0123-456789abcdef"; +const SECOND_SESSION_ID: &str = "fedcba98-7654-3210-fedc-ba9876543210"; +const RFC6455_SAMPLE_KEY: &str = "dGhlIHNhbXBsZSBub25jZQ=="; +const 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"; +const END_SUCCESS_RESPONSE: &[u8] = br#"{"type":"success","id":7,"result":{}}"#; +const NORMAL_CLOSE_FRAME: &[u8] = &[0x88, 0x02, 0x03, 0xe8]; +const PONG_MASK_KEY: [u8; 4] = [5, 6, 7, 8]; +const CLOSE_MASK_KEY: [u8; 4] = [9, 10, 11, 12]; + +type SessionEndFixture = ( + WebDriverBiDiSessionEndResult, + WebDriverBiDiWebSocketEstablished, + thread::JoinHandle>, +); + +fn read_opening_request(stream: &mut TcpStream) -> io::Result<()> { + stream.set_read_timeout(Some(Duration::from_secs(2)))?; + let mut request = Vec::new(); + let mut buffer = [0_u8; 512]; + while !request.ends_with(b"\r\n\r\n") { + let count = stream.read(&mut buffer)?; + if count == 0 { + return Err(io::Error::new( + io::ErrorKind::UnexpectedEof, + "client opening request ended before the header terminator", + )); + } + request.extend_from_slice(&buffer[..count]); + } + Ok(()) +} + +fn read_masked_text_frame(stream: &mut TcpStream) -> io::Result> { + let mut header = [0_u8; 2]; + stream.read_exact(&mut header)?; + if header[0] != 0x81 || header[1] & 0x80 == 0 { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "expected one final masked client text frame", + )); + } + let length = usize::from(header[1] & 0x7f); + if length > 125 { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "session.end command unexpectedly required extended framing", + )); + } + let mut mask = [0_u8; 4]; + stream.read_exact(&mut mask)?; + let mut payload = vec![0_u8; length]; + stream.read_exact(&mut payload)?; + for (index, byte) in payload.iter_mut().enumerate() { + *byte ^= mask[index % mask.len()]; + } + Ok(payload) +} + +fn read_masked_close_frame(stream: &mut TcpStream) -> io::Result> { + let mut header = [0_u8; 2]; + stream.read_exact(&mut header)?; + if header[0] != 0x88 || header[1] & 0x80 == 0 { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "expected one final masked client Close frame", + )); + } + let length = usize::from(header[1] & 0x7f); + if length > 125 { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "client Close response exceeded the RFC 6455 control-frame bound", + )); + } + let mut mask = [0_u8; 4]; + stream.read_exact(&mut mask)?; + let mut payload = vec![0_u8; length]; + stream.read_exact(&mut payload)?; + for (index, byte) in payload.iter_mut().enumerate() { + *byte ^= mask[index % mask.len()]; + } + Ok(payload) +} + +fn correlated_session_end_ack_and_transport_on( + listener: &TcpListener, + session_id: &str, + expect_close_response: bool, +) -> Result> { + let local_addr = listener.local_addr()?; + let server_listener = listener.try_clone()?; + let server = thread::spawn(move || -> io::Result<()> { + let (mut stream, _) = server_listener.accept()?; + read_opening_request(&mut stream)?; + stream.write_all(OPENING_RESPONSE)?; + let command = read_masked_text_frame(&mut stream)?; + if command != br#"{"id":7,"method":"session.end","params":{}}"# { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "unexpected session.end command", + )); + } + stream.write_all(&[0x81, END_SUCCESS_RESPONSE.len() as u8])?; + stream.write_all(END_SUCCESS_RESPONSE)?; + stream.write_all(NORMAL_CLOSE_FRAME)?; + if expect_close_response { + let close_reply = read_masked_close_frame(&mut stream)?; + if close_reply != [0x03, 0xe8] { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "session teardown did not return Close(1000)", + )); + } + } + Ok(()) + }); + + let endpoint = format!("ws://{local_addr}/session/{session_id}"); + let target = WebDriverBiDiWebSocketEndpoint::new(&endpoint)? + .correlate_session_id(session_id)? + .into_explicit_connect_target()?; + let connection = + WebDriverBiDiTcpConnectionPlan::new(target, Duration::from_secs(1), 1)?.connect()?; + let key = WebDriverBiDiWebSocketClientKey::new(RFC6455_SAMPLE_KEY)?; + let established = WebDriverBiDiWebSocketHandshakePlan::new(connection, key)? + .write_opening_request(Duration::from_millis(500))? + .read_opening_response(Duration::from_millis(500))?; + + let mut correlation = WebDriverBiDiCommandCorrelation::new(); + let established = WebDriverBiDiSessionEndCommand::new(7)?.send( + established, + &mut correlation, + WebDriverBiDiWebSocketMaskKey::new([1, 2, 3, 4]), + Duration::from_millis(500), + )?; + let mut reader = WebDriverBiDiWebSocketMessageReader::new(established); + let (established, text) = loop { + match reader.read_next(Duration::from_millis(500))? { + WebDriverBiDiConnectionMessageRead::Pending(next) => reader = next, + WebDriverBiDiConnectionMessageRead::Text { + established, + message, + } => break (established, message), + WebDriverBiDiConnectionMessageRead::Control { message, .. } => { + return Err(io::Error::other(format!( + "session.end response produced unexpected control message: {message:?}" + )) + .into()); + } + } + }; + let acknowledged = WebDriverBiDiSessionEndResult::parse_and_correlate(&text, &mut correlation)?; + Ok((acknowledged, established, server)) +} + +fn correlated_session_end_ack_and_transport_for( + session_id: &str, + expect_close_response: bool, +) -> Result> { + let listener = TcpListener::bind(("127.0.0.1", 0))?; + correlated_session_end_ack_and_transport_on(&listener, session_id, expect_close_response) +} + +fn correlated_session_end_ack_and_transport( + expect_close_response: bool, +) -> Result> { + correlated_session_end_ack_and_transport_for(SESSION_ID, expect_close_response) +} + +fn join_server(server: thread::JoinHandle>) -> Result<(), Box> { + server + .join() + .map_err(|_| io::Error::other("session.end teardown test server panicked"))??; + Ok(()) +} + +fn observed_transport_closure( + established: WebDriverBiDiWebSocketEstablished, + server: thread::JoinHandle>, +) -> Result> { + let observation = WebDriverBiDiWebSocketTransportClosureObservation::observe( + established, + &[WebDriverBiDiWebSocketMaskKey::new(PONG_MASK_KEY)], + WebDriverBiDiWebSocketMaskKey::new(CLOSE_MASK_KEY), + Duration::from_millis(500), + )?; + join_server(server)?; + assert_eq!( + observation.kind(), + WebDriverBiDiWebSocketTransportClosureKind::PeerCloseThenEof + ); + assert_eq!(observation.peer_close_status_code(), Some(1000)); + Ok(observation) +} + +#[test] +fn missing_typed_transport_observation_keeps_teardown_pending() -> Result<(), Box> { + for transport_closed in [false, true] { + let (acknowledged, established, server) = + correlated_session_end_ack_and_transport(transport_closed)?; + let transport_closure = if transport_closed { + Some(observed_transport_closure(established, server)?) + } else { + drop(established); + join_server(server)?; + None + }; + let assessment = WebDriverBiDiSessionTeardownAssessment::from_protocol_ack( + acknowledged, + WebDriverBiDiSessionTeardownObservations::new(transport_closure), + )?; + + assert_eq!(assessment.command_id(), 7); + assert_eq!( + assessment.observations().transport_closed_observed(), + transport_closed + ); + assert_eq!( + assessment + .observations() + .transport_closure_observation() + .map(WebDriverBiDiWebSocketTransportClosureObservation::kind), + transport_closed + .then_some(WebDriverBiDiWebSocketTransportClosureKind::PeerCloseThenEof) + ); + assert!(!assessment.is_operationally_complete()); + assert_eq!( + assessment.disposition(), + WebDriverBiDiSessionTeardownDisposition::OperationalTeardownPending + ); + } + Ok(()) +} + +#[test] +fn typed_closure_still_requires_process_and_profile_evidence() -> Result<(), Box> { + let (acknowledged, established, server) = correlated_session_end_ack_and_transport(true)?; + let transport_closure = observed_transport_closure(established, server)?; + let assessment = WebDriverBiDiSessionTeardownAssessment::from_protocol_ack( + acknowledged, + WebDriverBiDiSessionTeardownObservations::new(Some(transport_closure)), + )?; + + assert!(!assessment.is_operationally_complete()); + assert_eq!( + assessment.disposition(), + WebDriverBiDiSessionTeardownDisposition::OperationalTeardownPending + ); + assert_eq!(assessment.command_id(), 7); + assert_eq!( + assessment + .observations() + .transport_closure_observation() + .map(WebDriverBiDiWebSocketTransportClosureObservation::peer_close_status_code), + Some(Some(1000)) + ); + Ok(()) +} + +fn assert_cross_connection_closure_rejected( + acknowledged: WebDriverBiDiSessionEndResult, + foreign_established: WebDriverBiDiWebSocketEstablished, + foreign_server: thread::JoinHandle>, +) -> Result<(), Box> { + let foreign_closure = observed_transport_closure(foreign_established, foreign_server)?; + let result = WebDriverBiDiSessionTeardownAssessment::from_protocol_ack( + acknowledged, + WebDriverBiDiSessionTeardownObservations::new(Some(foreign_closure)), + ); + + let error = result + .err() + .ok_or_else(|| io::Error::other("closure from another connection was accepted"))?; + assert_eq!( + error, + WebDriverBiDiSessionTeardownAssessmentError::TransportConnectionMismatch + ); + assert_eq!( + error.to_string(), + "WebDriver BiDi transport closure does not match the acknowledged connection" + ); + assert!(error.source().is_none()); + Ok(()) +} + +#[test] +fn same_session_reconnect_cannot_supply_prior_closure() -> Result<(), Box> { + let listener = TcpListener::bind(("127.0.0.1", 0))?; + let (acknowledged_a, established_a, server_a) = + correlated_session_end_ack_and_transport_on(&listener, SESSION_ID, false)?; + drop(established_a); + join_server(server_a)?; + let (_acknowledged_b, established_b, server_b) = + correlated_session_end_ack_and_transport_on(&listener, SESSION_ID, true)?; + assert_cross_connection_closure_rejected(acknowledged_a, established_b, server_b) +} + +#[test] +fn another_session_cannot_supply_acknowledged_transport_closure() -> Result<(), Box> { + let (acknowledged_a, established_a, server_a) = + correlated_session_end_ack_and_transport_for(SESSION_ID, false)?; + drop(established_a); + join_server(server_a)?; + let (_acknowledged_b, established_b, server_b) = + correlated_session_end_ack_and_transport_for(SECOND_SESSION_ID, true)?; + assert_cross_connection_closure_rejected(acknowledged_a, established_b, server_b) +} diff --git a/crates/originweave-network/tests/webdriver_bidi_transport_close_evidence.rs b/crates/originweave-network/tests/webdriver_bidi_transport_close_evidence.rs new file mode 100644 index 000000000..9850efb3b --- /dev/null +++ b/crates/originweave-network/tests/webdriver_bidi_transport_close_evidence.rs @@ -0,0 +1,757 @@ +use std::{ + error::Error, + io::{self, Read, Write}, + net::{TcpListener, TcpStream}, + thread, + time::Duration, +}; + +use originweave_core::WebDriverBiDiWebSocketEndpoint; +use originweave_network::{ + WebDriverBiDiTcpConnectionPlan, WebDriverBiDiWebSocketClientKey, + WebDriverBiDiWebSocketFrameError, WebDriverBiDiWebSocketHandshakePlan, + WebDriverBiDiWebSocketMaskKey, WebDriverBiDiWebSocketTransportClosureError, + WebDriverBiDiWebSocketTransportClosureKind, WebDriverBiDiWebSocketTransportClosureObservation, +}; + +const SESSION_ID: &str = "01234567-89ab-cdef-0123-456789abcdef"; +const RFC6455_SAMPLE_KEY: &str = "dGhlIHNhbXBsZSBub25jZQ=="; +const 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"; +const PONG_MASK_KEY: [u8; 4] = [5, 6, 7, 8]; +const CLOSE_MASK_KEY: [u8; 4] = [9, 10, 11, 12]; + +type EstablishedWithServer = ( + originweave_network::WebDriverBiDiWebSocketEstablished, + thread::JoinHandle>, +); + +fn read_opening_request(stream: &mut TcpStream) -> io::Result<()> { + stream.set_read_timeout(Some(Duration::from_secs(2)))?; + let mut request = Vec::new(); + let mut buffer = [0_u8; 512]; + while !request.ends_with(b"\r\n\r\n") { + let count = stream.read(&mut buffer)?; + if count == 0 { + return Err(io::Error::new( + io::ErrorKind::UnexpectedEof, + "client opening request ended before the header terminator", + )); + } + request.extend_from_slice(&buffer[..count]); + } + Ok(()) +} + +fn read_masked_control_frame(stream: &mut TcpStream, expected_opcode: u8) -> io::Result> { + let mut header = [0_u8; 2]; + stream.read_exact(&mut header)?; + if header[0] != (0x80 | expected_opcode) || header[1] & 0x80 == 0 { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "expected one final masked client control frame", + )); + } + let payload_len = usize::from(header[1] & 0x7f); + if payload_len > 125 { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "client control-frame payload exceeded RFC 6455 bound", + )); + } + let mut mask = [0_u8; 4]; + stream.read_exact(&mut mask)?; + let mut payload = vec![0_u8; payload_len]; + stream.read_exact(&mut payload)?; + for (index, byte) in payload.iter_mut().enumerate() { + *byte ^= mask[index % mask.len()]; + } + Ok(payload) +} + +fn established_with_server_frame( + frame: Option<&'static [u8]>, +) -> Result> { + established_with_peer_script(move |stream| { + if let Some(frame) = frame { + stream.write_all(frame)?; + } + Ok(()) + }) +} + +fn established_with_peer_script( + script: impl FnOnce(&mut TcpStream) -> io::Result<()> + Send + 'static, +) -> Result> { + 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(OPENING_RESPONSE)?; + script(&mut stream) + }); + + let endpoint = format!("ws://{local_addr}/session/{SESSION_ID}"); + let target = WebDriverBiDiWebSocketEndpoint::new(&endpoint)? + .correlate_session_id(SESSION_ID)? + .into_explicit_connect_target()?; + let connection = + WebDriverBiDiTcpConnectionPlan::new(target, Duration::from_secs(1), 1)?.connect()?; + let key = WebDriverBiDiWebSocketClientKey::new(RFC6455_SAMPLE_KEY)?; + let established = WebDriverBiDiWebSocketHandshakePlan::new(connection, key)? + .write_opening_request(Duration::from_millis(500))? + .read_opening_response(Duration::from_millis(500))?; + Ok((established, server)) +} + +fn established_with_peer_close_then_eof( + close_frame: &'static [u8], + expected_reply_payload: &'static [u8], +) -> Result> { + 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(OPENING_RESPONSE)?; + stream.write_all(close_frame)?; + let close_reply = read_masked_control_frame(&mut stream, 0x08)?; + if close_reply != expected_reply_payload { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "Close response did not preserve the validated peer status contract", + )); + } + Ok(()) + }); + + let endpoint = format!("ws://{local_addr}/session/{SESSION_ID}"); + let target = WebDriverBiDiWebSocketEndpoint::new(&endpoint)? + .correlate_session_id(SESSION_ID)? + .into_explicit_connect_target()?; + let connection = + WebDriverBiDiTcpConnectionPlan::new(target, Duration::from_secs(1), 1)?.connect()?; + let key = WebDriverBiDiWebSocketClientKey::new(RFC6455_SAMPLE_KEY)?; + let established = WebDriverBiDiWebSocketHandshakePlan::new(connection, key)? + .write_opening_request(Duration::from_millis(500))? + .read_opening_response(Duration::from_millis(500))?; + Ok((established, server)) +} + +fn established_with_held_open_peer_close() -> Result> { + 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(OPENING_RESPONSE)?; + stream.write_all(&[0x88, 0x02, 0x03, 0xe8])?; + let close_reply = read_masked_control_frame(&mut stream, 0x08)?; + if close_reply != [0x03, 0xe8] { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "held-open peer did not receive the required Close response", + )); + } + thread::sleep(Duration::from_millis(250)); + Ok(()) + }); + + let endpoint = format!("ws://{local_addr}/session/{SESSION_ID}"); + let target = WebDriverBiDiWebSocketEndpoint::new(&endpoint)? + .correlate_session_id(SESSION_ID)? + .into_explicit_connect_target()?; + let connection = + WebDriverBiDiTcpConnectionPlan::new(target, Duration::from_secs(1), 1)?.connect()?; + let key = WebDriverBiDiWebSocketClientKey::new(RFC6455_SAMPLE_KEY)?; + let established = WebDriverBiDiWebSocketHandshakePlan::new(connection, key)? + .write_opening_request(Duration::from_millis(500))? + .read_opening_response(Duration::from_millis(500))?; + Ok((established, server)) +} + +fn established_with_ping_then_held_open_peer_close() -> Result> +{ + 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(OPENING_RESPONSE)?; + stream.write_all(&[0x89, 0x05, b'p', b'r', b'o', b'b', b'e'])?; + let pong_payload = read_masked_control_frame(&mut stream, 0x0a)?; + if pong_payload != b"probe" { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "Pong did not preserve the Ping application data", + )); + } + stream.write_all(&[0x88, 0x02, 0x03, 0xe8])?; + let close_reply = read_masked_control_frame(&mut stream, 0x08)?; + if close_reply != [0x03, 0xe8] { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "Ping/Close peer did not receive the required Close response", + )); + } + thread::sleep(Duration::from_millis(250)); + Ok(()) + }); + + let endpoint = format!("ws://{local_addr}/session/{SESSION_ID}"); + let target = WebDriverBiDiWebSocketEndpoint::new(&endpoint)? + .correlate_session_id(SESSION_ID)? + .into_explicit_connect_target()?; + let connection = + WebDriverBiDiTcpConnectionPlan::new(target, Duration::from_secs(1), 1)?.connect()?; + let key = WebDriverBiDiWebSocketClientKey::new(RFC6455_SAMPLE_KEY)?; + let established = WebDriverBiDiWebSocketHandshakePlan::new(connection, key)? + .write_opening_request(Duration::from_millis(500))? + .read_opening_response(Duration::from_millis(500))?; + Ok((established, server)) +} + +fn established_with_pong_then_peer_close() -> Result> { + 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(OPENING_RESPONSE)?; + stream.write_all(&[0x8a, 0x00, 0x88, 0x02, 0x03, 0xe8])?; + let close_reply = read_masked_control_frame(&mut stream, 0x08)?; + if close_reply != [0x03, 0xe8] { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "Pong/Close peer did not receive the required Close response", + )); + } + Ok(()) + }); + + let endpoint = format!("ws://{local_addr}/session/{SESSION_ID}"); + let target = WebDriverBiDiWebSocketEndpoint::new(&endpoint)? + .correlate_session_id(SESSION_ID)? + .into_explicit_connect_target()?; + let connection = + WebDriverBiDiTcpConnectionPlan::new(target, Duration::from_secs(1), 1)?.connect()?; + let key = WebDriverBiDiWebSocketClientKey::new(RFC6455_SAMPLE_KEY)?; + let established = WebDriverBiDiWebSocketHandshakePlan::new(connection, key)? + .write_opening_request(Duration::from_millis(500))? + .read_opening_response(Duration::from_millis(500))?; + Ok((established, server)) +} + +fn established_with_post_close_frame() -> Result> { + 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(OPENING_RESPONSE)?; + stream.write_all(&[0x88, 0x02, 0x03, 0xe8])?; + let close_reply = read_masked_control_frame(&mut stream, 0x08)?; + if close_reply != [0x03, 0xe8] { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "post-Close fixture did not receive the required Close response", + )); + } + stream.write_all(&[0x8a, 0x00])?; + Ok(()) + }); + + let endpoint = format!("ws://{local_addr}/session/{SESSION_ID}"); + let target = WebDriverBiDiWebSocketEndpoint::new(&endpoint)? + .correlate_session_id(SESSION_ID)? + .into_explicit_connect_target()?; + let connection = + WebDriverBiDiTcpConnectionPlan::new(target, Duration::from_secs(1), 1)?.connect()?; + let key = WebDriverBiDiWebSocketClientKey::new(RFC6455_SAMPLE_KEY)?; + let established = WebDriverBiDiWebSocketHandshakePlan::new(connection, key)? + .write_opening_request(Duration::from_millis(500))? + .read_opening_response(Duration::from_millis(500))?; + Ok((established, server)) +} + +fn observe( + established: originweave_network::WebDriverBiDiWebSocketEstablished, + frame_timeout: Duration, +) -> Result< + WebDriverBiDiWebSocketTransportClosureObservation, + WebDriverBiDiWebSocketTransportClosureError, +> { + WebDriverBiDiWebSocketTransportClosureObservation::observe( + established, + &[WebDriverBiDiWebSocketMaskKey::new(PONG_MASK_KEY)], + WebDriverBiDiWebSocketMaskKey::new(CLOSE_MASK_KEY), + frame_timeout, + ) +} + +#[test] +fn close_exchange_cannot_restart_the_operation_deadline() -> Result<(), Box> { + let (start_sender, start_receiver) = std::sync::mpsc::channel(); + let (established, server) = established_with_peer_script(move |stream| { + start_receiver.recv().map_err(io::Error::other)?; + thread::sleep(Duration::from_millis(750)); + stream.write_all(&[0x88, 0x02, 0x03, 0xe8])?; + assert_eq!(read_masked_control_frame(stream, 0x08)?, [0x03, 0xe8]); + thread::sleep(Duration::from_millis(750)); + Ok(()) + })?; + start_sender.send(())?; + let result = observe(established, Duration::from_secs(1)); + server + .join() + .map_err(|_| io::Error::other("deadline peer panicked"))??; + assert!(matches!( + result, + Err(WebDriverBiDiWebSocketTransportClosureError::Frame { + source: WebDriverBiDiWebSocketFrameError::FrameReadTimedOut { bytes_read: 0, .. } + }) + )); + Ok(()) +} + +#[test] +fn ping_before_close_requires_masked_pong_and_still_waits_for_tcp_closure() +-> Result<(), Box> { + let (established, server) = established_with_ping_then_held_open_peer_close()?; + + let result = observe(established, Duration::from_millis(50)); + let server_result = server + .join() + .map_err(|_| io::Error::other("ping-before-close test server panicked"))?; + let error = result.err().ok_or_else(|| { + io::Error::other("Ping/Close traffic became final transport-closure evidence") + })?; + assert!(matches!( + &error, + WebDriverBiDiWebSocketTransportClosureError::Frame { + source: WebDriverBiDiWebSocketFrameError::FrameReadTimedOut { bytes_read: 0, .. } + } + )); + server_result?; + Ok(()) +} + +#[test] +fn validated_peer_close_requires_masked_reply_and_tcp_eof() -> Result<(), Box> { + let (established, server) = + established_with_peer_close_then_eof(&[0x88, 0x04, 0x03, 0xe8, b'o', b'k'], &[0x03, 0xe8])?; + + let observation = observe(established, Duration::from_millis(500))?; + + server + .join() + .map_err(|_| io::Error::other("transport-close test server panicked"))??; + assert_eq!( + observation.kind(), + WebDriverBiDiWebSocketTransportClosureKind::PeerCloseThenEof + ); + assert_eq!(observation.peer_close_status_code(), Some(1000)); + Ok(()) +} + +#[test] +fn empty_peer_close_uses_empty_masked_reply_before_tcp_eof() -> Result<(), Box> { + let (established, server) = established_with_peer_close_then_eof(&[0x88, 0x00], &[])?; + + let observation = observe(established, Duration::from_millis(500))?; + + server + .join() + .map_err(|_| io::Error::other("empty-close test server panicked"))??; + assert_eq!( + observation.kind(), + WebDriverBiDiWebSocketTransportClosureKind::PeerCloseThenEof + ); + assert_eq!(observation.peer_close_status_code(), None); + Ok(()) +} + +#[test] +fn peer_close_without_tcp_eof_is_not_transport_closure_evidence() -> Result<(), Box> { + let (established, server) = established_with_held_open_peer_close()?; + + let result = observe(established, Duration::from_millis(50)); + + server + .join() + .map_err(|_| io::Error::other("held-open close test server panicked"))??; + let error = result.err().ok_or_else(|| { + io::Error::other("peer Close alone became final transport-closure evidence") + })?; + assert!(matches!( + &error, + WebDriverBiDiWebSocketTransportClosureError::Frame { + source: WebDriverBiDiWebSocketFrameError::FrameReadTimedOut { bytes_read: 0, .. } + } + )); + assert_eq!( + error.to_string(), + "WebDriver BiDi transport closure could not be observed safely" + ); + Ok(()) +} + +#[test] +fn one_unsolicited_pong_before_close_does_not_block_closure_observation() +-> Result<(), Box> { + let (established, server) = established_with_pong_then_peer_close()?; + + let observation = observe(established, Duration::from_millis(500))?; + + server + .join() + .map_err(|_| io::Error::other("pong-before-close test server panicked"))??; + assert_eq!( + observation.kind(), + WebDriverBiDiWebSocketTransportClosureKind::PeerCloseThenEof + ); + assert_eq!(observation.peer_close_status_code(), Some(1000)); + Ok(()) +} + +#[test] +fn repeated_pong_frames_preserve_clean_eof_observation() -> Result<(), Box> { + let (established, server) = established_with_server_frame(Some(&[0x8a, 0x00, 0x8a, 0x00]))?; + + let result = observe(established, Duration::from_millis(500)); + + server + .join() + .map_err(|_| io::Error::other("repeated-pong test server panicked"))??; + assert_eq!( + result?.kind(), + WebDriverBiDiWebSocketTransportClosureKind::PeerEof + ); + Ok(()) +} + +#[test] +fn ping_after_pre_close_pong_is_answered_before_closure() -> Result<(), Box> { + let (established, server) = established_with_peer_script(|stream| { + stream.write_all(&[0x8a, 0x00, 0x89, 0x01, b'p'])?; + assert_eq!(read_masked_control_frame(stream, 0xa)?, b"p"); + stream.write_all(&[0x88, 0])?; + assert_eq!(read_masked_control_frame(stream, 0x8)?, b""); + Ok(()) + })?; + let result = observe(established, Duration::from_millis(500)); + server + .join() + .map_err(|_| io::Error::other("Pong/Ping peer panicked"))??; + assert_eq!( + result?.kind(), + WebDriverBiDiWebSocketTransportClosureKind::PeerCloseThenEof + ); + Ok(()) +} + +#[test] +fn repeated_pings_use_distinct_keys_without_pong_consuming_one() -> Result<(), Box> { + let (established, server) = established_with_peer_script(|stream| { + stream.write_all(&[0x89, 1, b'p'])?; + let mut reply = [0; 7]; + stream.read_exact(&mut reply)?; + assert_eq!(reply, [0x8a, 0x81, 5, 6, 7, 8, 0x75]); + stream.write_all(&[0x8a, 0, 0x89, 1, b'q'])?; + stream.read_exact(&mut reply)?; + assert_eq!(reply, [0x8a, 0x81, 13, 14, 15, 16, 0x7c]); + stream.write_all(&[0x88, 0])?; + assert_eq!(read_masked_control_frame(stream, 0x8)?, b""); + Ok(()) + })?; + let result = WebDriverBiDiWebSocketTransportClosureObservation::observe( + established, + &[ + WebDriverBiDiWebSocketMaskKey::new(PONG_MASK_KEY), + WebDriverBiDiWebSocketMaskKey::new([13, 14, 15, 16]), + ], + WebDriverBiDiWebSocketMaskKey::new(CLOSE_MASK_KEY), + Duration::from_millis(500), + ); + server + .join() + .map_err(|_| io::Error::other("repeated Ping peer panicked"))??; + assert_eq!( + result?.kind(), + WebDriverBiDiWebSocketTransportClosureKind::PeerCloseThenEof + ); + Ok(()) +} + +#[test] +fn repeated_ping_rejects_missing_or_reused_keys_without_second_reply() -> Result<(), Box> +{ + let key = WebDriverBiDiWebSocketMaskKey::new(PONG_MASK_KEY); + for (keys, expected) in [ + (vec![key], "PongMaskingKeysExhausted"), + ( + vec![key, key], + "Frame { source: MalformedFrame { reason: \"client masking key was reused for consecutive frames on this established WebSocket\" } }", + ), + ] { + let (established, server) = established_with_peer_script(|stream| { + stream.write_all(&[0x89, 1, b'p'])?; + assert_eq!(read_masked_control_frame(stream, 0xa)?, b"p"); + stream.write_all(&[0x89, 1, b'q'])?; + let mut reply = [0]; + assert_eq!(stream.read(&mut reply)?, 0); + Ok(()) + })?; + let result = WebDriverBiDiWebSocketTransportClosureObservation::observe( + established, + &keys, + WebDriverBiDiWebSocketMaskKey::new(CLOSE_MASK_KEY), + Duration::from_millis(500), + ); + server + .join() + .map_err(|_| io::Error::other("missing/reused key peer panicked"))??; + assert_eq!(format!("{result:?}"), format!("Err({expected})")); + } + Ok(()) +} + +#[test] +fn peer_close_rejects_zero_frame_timeout_before_writing() -> Result<(), Box> { + let (established, server) = established_with_server_frame(Some(&[0x88, 0x02, 0x03, 0xe8]))?; + let result = observe(established, Duration::ZERO); + let _ = server.join(); + assert!(matches!( + result, + Err(WebDriverBiDiWebSocketTransportClosureError::Frame { .. }) + )); + Ok(()) +} + +#[test] +fn peer_close_rejects_excessive_frame_timeout_before_writing() -> Result<(), Box> { + let (established, server) = established_with_server_frame(Some(&[0x88, 0x02, 0x03, 0xe8]))?; + let result = observe(established, Duration::from_secs(6)); + let _ = server.join(); + assert!(matches!( + result, + Err(WebDriverBiDiWebSocketTransportClosureError::Frame { .. }) + )); + Ok(()) +} + +#[test] +fn clean_peer_eof_yields_transport_observation_without_inventing_close_status() +-> Result<(), Box> { + let (established, server) = established_with_server_frame(None)?; + + let observation = observe(established, Duration::from_millis(500))?; + + server + .join() + .map_err(|_| io::Error::other("transport-eof test server panicked"))??; + assert_eq!( + observation.kind(), + WebDriverBiDiWebSocketTransportClosureKind::PeerEof + ); + assert_eq!(observation.peer_close_status_code(), None); + Ok(()) +} + +#[test] +fn post_close_frame_is_not_transport_closure_evidence() -> Result<(), Box> { + let (established, server) = established_with_post_close_frame()?; + + let Err(error) = observe(established, Duration::from_millis(500)) else { + return Err( + io::Error::other("post-Close frame unexpectedly became closure evidence").into(), + ); + }; + + server + .join() + .map_err(|_| io::Error::other("post-Close frame test server panicked"))??; + assert!(matches!( + error, + WebDriverBiDiWebSocketTransportClosureError::UnexpectedFrame { opcode: 0xa } + )); + Ok(()) +} + +#[test] +fn application_frame_after_teardown_does_not_become_closure_evidence() -> Result<(), Box> +{ + let (established, server) = established_with_server_frame(Some(&[0x81, 0x02, b'o', b'k']))?; + + let Err(error) = observe(established, Duration::from_millis(500)) else { + return Err(io::Error::other( + "application frame unexpectedly became transport-closure evidence", + ) + .into()); + }; + + server + .join() + .map_err(|_| io::Error::other("unexpected-frame test server panicked"))??; + assert!(matches!( + &error, + WebDriverBiDiWebSocketTransportClosureError::UnexpectedFrame { opcode: 0x1 } + )); + assert_eq!( + error.to_string(), + "WebDriver BiDi peer sent non-closure traffic instead of closing" + ); + assert!(error.source().is_none()); + Ok(()) +} + +#[test] +fn unassigned_protocol_close_codes_emit_no_reply_or_evidence() -> Result<(), Box> { + for status_code in [1016_u16, 2000, 2999] { + let (established, server) = established_with_peer_script(move |stream| { + let [high, low] = status_code.to_be_bytes(); + stream.write_all(&[0x88, 0x02, high, low])?; + let mut reply = [0_u8; 1]; + assert_eq!(stream.read(&mut reply)?, 0, "invalid Close was echoed"); + Ok(()) + })?; + let result = observe(established, Duration::from_millis(500)); + server + .join() + .map_err(|_| io::Error::other("reserved-code test server panicked"))??; + assert_eq!( + format!("{result:?}"), + "Err(Frame { source: MalformedFrame { reason: \"Close frame status code is not valid on the wire\" } })" + ); + } + Ok(()) +} + +#[test] +fn malformed_close_frame_remains_a_typed_frame_failure() -> Result<(), Box> { + let (established, server) = established_with_server_frame(Some(&[0x88, 0x01, 0x00]))?; + + let Err(error) = observe(established, Duration::from_millis(500)) else { + return Err(io::Error::other( + "malformed close frame unexpectedly became transport-closure evidence", + ) + .into()); + }; + + server + .join() + .map_err(|_| io::Error::other("malformed-close test server panicked"))??; + assert!(matches!( + &error, + WebDriverBiDiWebSocketTransportClosureError::Frame { .. } + )); + assert_eq!( + error.to_string(), + "WebDriver BiDi transport closure could not be observed safely" + ); + assert!(error.source().is_some()); + Ok(()) +} + +#[test] +fn reused_pong_mask_key_preserves_typed_frame_failure() -> Result<(), Box> { + let (established, server) = established_with_peer_script(|stream| { + let mut seed = [0; 8]; + stream.read_exact(&mut seed)?; + assert_eq!(seed, [0x81, 0x82, 5, 6, 7, 8, 0x7e, 0x7b]); + stream.write_all(&[0x89, 0x01, b'p'])?; + let mut byte = [0]; + assert_eq!(stream.read(&mut byte)?, 0); + Ok(()) + })?; + let reused = WebDriverBiDiWebSocketMaskKey::new(PONG_MASK_KEY); + let established = established.write_text_frame("{}", reused, Duration::from_millis(500))?; + + let result = WebDriverBiDiWebSocketTransportClosureObservation::observe( + established, + &[reused], + WebDriverBiDiWebSocketMaskKey::new(CLOSE_MASK_KEY), + Duration::from_millis(500), + ); + assert!(matches!( + result, + Err(WebDriverBiDiWebSocketTransportClosureError::Frame { + source: WebDriverBiDiWebSocketFrameError::MalformedFrame { + reason: "client masking key was reused for consecutive frames on this established WebSocket" + } + }) + )); + server + .join() + .map_err(|_| io::Error::other("mask-reuse peer panicked"))??; + Ok(()) +} + +#[test] +fn reused_close_mask_key_preserves_typed_frame_failure() -> Result<(), Box> { + let (established, server) = established_with_peer_script(|stream| { + stream.write_all(&[0x89, 0x01, b'p'])?; + let mut pong = [0; 7]; + stream.read_exact(&mut pong)?; + assert_eq!(pong, [0x8a, 0x81, 9, 10, 11, 12, 0x79]); + stream.write_all(&[0x88, 0x02, 0x03, 0xe8])?; + let mut byte = [0]; + assert_eq!(stream.read(&mut byte)?, 0); + Ok(()) + })?; + let reused = WebDriverBiDiWebSocketMaskKey::new(CLOSE_MASK_KEY); + + let result = WebDriverBiDiWebSocketTransportClosureObservation::observe( + established, + &[reused], + reused, + Duration::from_millis(500), + ); + assert!(matches!( + result, + Err(WebDriverBiDiWebSocketTransportClosureError::Frame { + source: WebDriverBiDiWebSocketFrameError::MalformedFrame { + reason: "client masking key was reused for consecutive frames on this established WebSocket" + } + }) + )); + server + .join() + .map_err(|_| io::Error::other("mask-reuse peer panicked"))??; + Ok(()) +} + +#[test] +fn reused_close_mask_key_on_peer_close_preserves_typed_frame_failure() -> Result<(), Box> +{ + let (established, server) = established_with_peer_script(|stream| { + let mut seed = [0; 8]; + stream.read_exact(&mut seed)?; + assert_eq!(seed, [0x81, 0x82, 9, 10, 11, 12, 0x72, 0x77]); + stream.write_all(&[0x88, 0x02, 0x03, 0xe8])?; + let mut byte = [0]; + assert_eq!(stream.read(&mut byte)?, 0); + Ok(()) + })?; + let reused = WebDriverBiDiWebSocketMaskKey::new(CLOSE_MASK_KEY); + let established = established.write_text_frame("{}", reused, Duration::from_millis(500))?; + + let result = WebDriverBiDiWebSocketTransportClosureObservation::observe( + established, + &[WebDriverBiDiWebSocketMaskKey::new(PONG_MASK_KEY)], + reused, + Duration::from_millis(500), + ); + assert!(matches!( + result, + Err(WebDriverBiDiWebSocketTransportClosureError::Frame { + source: WebDriverBiDiWebSocketFrameError::MalformedFrame { + reason: "client masking key was reused for consecutive frames on this established WebSocket" + } + }) + )); + server + .join() + .map_err(|_| io::Error::other("mask-reuse peer panicked"))??; + Ok(()) +} diff --git a/crates/originweave-network/tests/webdriver_bidi_transport_close_role_validation.rs b/crates/originweave-network/tests/webdriver_bidi_transport_close_role_validation.rs new file mode 100644 index 000000000..6f6117ffb --- /dev/null +++ b/crates/originweave-network/tests/webdriver_bidi_transport_close_role_validation.rs @@ -0,0 +1,169 @@ +use std::{ + error::Error, + io::{self, Read, Write}, + net::{TcpListener, TcpStream}, + thread, + time::Duration, +}; + +use originweave_core::WebDriverBiDiWebSocketEndpoint; +use originweave_network::{ + WebDriverBiDiTcpConnectionPlan, WebDriverBiDiWebSocketClientKey, + WebDriverBiDiWebSocketHandshakePlan, WebDriverBiDiWebSocketMaskKey, + WebDriverBiDiWebSocketTransportClosureError, WebDriverBiDiWebSocketTransportClosureKind, + WebDriverBiDiWebSocketTransportClosureObservation, +}; + +const SESSION_ID: &str = "01234567-89ab-cdef-0123-456789abcdef"; +const RFC6455_SAMPLE_KEY: &str = "dGhlIHNhbXBsZSBub25jZQ=="; +const 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"; +const CLOSE_MASK_KEY: [u8; 4] = [9, 10, 11, 12]; +type EstablishedPeer = ( + originweave_network::WebDriverBiDiWebSocketEstablished, + thread::JoinHandle>>>, +); + +fn read_opening_request(stream: &mut TcpStream) -> io::Result<()> { + stream.set_read_timeout(Some(Duration::from_secs(2)))?; + let mut request = Vec::new(); + let mut buffer = [0_u8; 512]; + while !request.ends_with(b"\r\n\r\n") { + let count = stream.read(&mut buffer)?; + if count == 0 { + return Err(io::Error::new( + io::ErrorKind::UnexpectedEof, + "client opening request ended before the header terminator", + )); + } + request.extend_from_slice(&buffer[..count]); + } + Ok(()) +} + +fn read_optional_masked_close(stream: &mut TcpStream) -> io::Result>> { + stream.set_read_timeout(Some(Duration::from_millis(200)))?; + let mut header = [0_u8; 2]; + match stream.read(&mut header[..1]) { + Ok(0) => return Ok(None), + Ok(1) => {} + Ok(_) => unreachable!("one-byte header read returned more than one byte"), + Err(source) + if matches!( + source.kind(), + io::ErrorKind::TimedOut | io::ErrorKind::WouldBlock + ) => + { + return Ok(None); + } + Err(source) => return Err(source), + } + stream.read_exact(&mut header[1..])?; + if header[0] != 0x88 || header[1] & 0x80 == 0 { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "expected one final masked client Close frame", + )); + } + let payload_len = usize::from(header[1] & 0x7f); + if payload_len > 125 { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "client Close payload exceeded RFC 6455 control-frame bound", + )); + } + let mut mask = [0_u8; 4]; + stream.read_exact(&mut mask)?; + let mut payload = vec![0_u8; payload_len]; + stream.read_exact(&mut payload)?; + for (index, byte) in payload.iter_mut().enumerate() { + *byte ^= mask[index % mask.len()]; + } + Ok(Some(payload)) +} + +fn established_with_server_close(status_code: u16) -> Result> { + 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(OPENING_RESPONSE)?; + let [high, low] = status_code.to_be_bytes(); + stream.write_all(&[0x88, 0x02, high, low])?; + read_optional_masked_close(&mut stream) + }); + + let endpoint = format!("ws://{local_addr}/session/{SESSION_ID}"); + let target = WebDriverBiDiWebSocketEndpoint::new(&endpoint)? + .correlate_session_id(SESSION_ID)? + .into_explicit_connect_target()?; + let connection = + WebDriverBiDiTcpConnectionPlan::new(target, Duration::from_secs(1), 1)?.connect()?; + let key = WebDriverBiDiWebSocketClientKey::new(RFC6455_SAMPLE_KEY)?; + let established = WebDriverBiDiWebSocketHandshakePlan::new(connection, key)? + .write_opening_request(Duration::from_millis(500))? + .read_opening_response(Duration::from_millis(500))?; + Ok((established, server)) +} + +fn observe( + established: originweave_network::WebDriverBiDiWebSocketEstablished, +) -> Result< + WebDriverBiDiWebSocketTransportClosureObservation, + WebDriverBiDiWebSocketTransportClosureError, +> { + WebDriverBiDiWebSocketTransportClosureObservation::observe( + established, + &[], + WebDriverBiDiWebSocketMaskKey::new(CLOSE_MASK_KEY), + Duration::from_millis(250), + ) +} + +#[test] +fn server_close_1010_is_rejected_before_reply_or_closure_evidence() -> Result<(), Box> { + let (established, server) = established_with_server_close(1010)?; + let error = match observe(established) { + Ok(_) => { + return Err(io::Error::other("server Close 1010 produced closure evidence").into()); + } + Err(error) => error, + }; + let reply = server + .join() + .map_err(|_| io::Error::other("role-invalid Close peer panicked"))??; + + assert!( + reply.is_none(), + "server Close(1010) was mirrored by the client" + ); + assert!(matches!( + &error, + WebDriverBiDiWebSocketTransportClosureError::PeerCloseStatusNotAllowed { + status_code: 1010 + } + )); + assert_eq!( + error.to_string(), + "WebDriver BiDi server sent a Close status reserved for clients" + ); + assert!(error.source().is_none()); + Ok(()) +} + +#[test] +fn server_close_1011_remains_valid_and_is_mirrored() -> Result<(), Box> { + let (established, server) = established_with_server_close(1011)?; + let observation = observe(established)?; + let reply = server + .join() + .map_err(|_| io::Error::other("valid server Close peer panicked"))??; + + assert_eq!(reply, Some(1011_u16.to_be_bytes().to_vec())); + assert_eq!( + observation.kind(), + WebDriverBiDiWebSocketTransportClosureKind::PeerCloseThenEof + ); + assert_eq!(observation.peer_close_status_code(), Some(1011)); + Ok(()) +} diff --git a/docs/doctoring.md b/docs/doctoring.md index 42a289dd1..9d98997c7 100644 --- a/docs/doctoring.md +++ b/docs/doctoring.md @@ -18,6 +18,48 @@ UTS #39 Revision 32 is the current Unicode security-mechanisms standard and mark 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. RFC 6455 section 5.3 additionally requires every client-to-server frame to use a fresh unpredictable 32-bit masking key derived from strong entropy. OriginWeave's frame owner enforces that normative masking requirement and also rejects immediate reuse of the preceding key as a local fail-closed stuck-randomness defense; the adjacent-reuse rule is stronger local policy, not an RFC 6455 requirement. +### Transport closure deadline + +The experimental BiDi transport-closure observer applies one local operation-wide +deadline across pre-Close control traffic, the masked Close response, and final TCP +EOF. This is an OriginWeave resource policy, not an RFC 6455 timeout value. Each I/O +step receives only the remaining budget, and final closure evidence is rejected if +the monotonic deadline has expired. It prevents per-frame budget renewal without +claiming a hard real-time bound on host scheduling or proving process exit/profile +cleanup. Tests combine a delayed real peer with deterministic clock transitions. + +RFC 6455 sections 5.5.2–5.5.3 permit repeated Ping and unsolicited Pong traffic; +a Ping response preserves its exact payload and an unsolicited Pong needs no reply. +The closure observer now admits up to 64 pre-Close Ping/Pong frames under the same +deadline. This count is a local resource budget, not a protocol limit. The caller +supplies a borrowed masking-key slice consumed once per Ping and a separate Close +key; no entropy provider or callback runs inside the observer. Missing keys fail +before the corresponding response and the existing adjacent-key guard still applies. +The unpublished `observe` API now accepts that slice instead of a single Pong key; +all in-tree consumers are migrated. Boundary tests cover exactly 64 controls, the +65th control, mixed Ping/Pong, exhaustion, and literal independently masked replies. + +### WebSocket Close code admission + +RFC 6455 section 7.4.2 reserves 1000–2999 for protocol and extension definitions. +The IANA registry retrieved on 8 September 2026 lists 1016–2999 as unassigned; +1012, 1013, and 1014 are assigned. The current adapter therefore rejects the +unassigned protocol range alongside 1004–1006 and 1015, which cannot appear as +ordinary wire status codes. It retains 3000–3999 for application codes and +4000–4999 for private use without interpreting their meanings or treating them +as proof of successful browser work. This is a reviewed static admission policy, +not an online registry check. New protocol assignments require a reviewed update. +The shared frame validator runs before role-specific handling. RFC 6455 section +7.4.1 assigns status 1010 to clients and says servers do not use it, so the +client-side transport-closure state machine rejects a peer 1010 before any +Close echo or closure evidence. Server status 1011 remains admissible. + +Internet Assigned Numbers Authority. (2026). *WebSocket protocol registries*. +Retrieved September 8, 2026, from https://www.iana.org/assignments/websocket + +Fette, I., & Melnikov, A. (2011). *The WebSocket protocol* (RFC 6455, §§ 7.4.1–7.4.2). +Internet Engineering Task Force. https://www.rfc-editor.org/rfc/rfc6455 + ### 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`. @@ -132,6 +174,18 @@ The #251 integration adopts #250 `ec433b844a121f8554c062f92267991af9cacb6f` by o PR #252 ordinarily adopts current command parent `f02af6d0dd01708d495cc08dec785675f3d58898` while preserving the response implementation, public exports and four real-loopback response tests from `2015259529ada99af836989079cc85a15779a2d8`. The pre-integration native loader again collected zero correlation release-record checks; adopting the existing parent TestCase makes that contract executable without a new framework or copied owner fix. Both sides of the changelog-only conflict are retained. The response boundary still validates the entire bounded envelope before consuming exact typed correlation, retains remote errors as failures and leaves malformed or mismatched responses unable to consume another command. The acknowledgment remains unbound to received-connection provenance in this layer and does not prove process exit, profile removal or operational teardown. Later connection-bound evidence belongs to its own owner stack; local quality results, hosted checks and protected delivery remain separate. +### Teardown-assessment parent integration and acceptance limit + +PR #253 ordinarily adopts response parent `6569bf40b6595ac74c2f0a997d202137f07ba1db` and preserves its assessment implementation, public exports and existing loopback tests from `0d72082e595c0e1fcc03d609ba337896ed14e2fc`. Native release-record discovery first failed with zero collected checks; the canonical parent provides the existing executable TestCase and synchronized opening fixtures. Both changelog records remain. No new assessment API or runtime-evidence producer is introduced by this integration. + +The retained assessment still accepts three caller-supplied booleans and can label them `OperationallyComplete`; that calculation authenticates none of the observations and is not trusted operational-completion evidence. This known product gap remains open despite passing local structural tests or numerical coverage. The later #255 owner removes raw process/profile completion claims and binds received-response and closure provenance. That owner repair must remain intact when this dependency chain is integrated and independently reverified. No release or protected-main acceptance of caller claims is justified by this intermediate parent adoption. + +### Transport-closure parent integration and provenance limit + +PR #254 ordinarily adopts teardown parent `afb623e4449b7cbf926fdcef7225ceaca822cfcf` while retaining the closure implementation, exports and six loopback tests from `cbaf50dcc97753cc73135497ea8225e8b18de190`. The native correlation release-record contract first collected zero checks; the existing parent TestCase makes it execute. Both conflicting release records remain, and the synchronized opening fixtures are inherited without duplicating their repair. + +Review `5120077272` remains actionable: this intermediate closure value retains kind/status but no connection-generation identity, so it cannot prove that a particular acknowledged connection closed. The existing #255 owner carries and compares received-response and closure provenance and removes raw process/profile completion claims. Preserve that repair during subsequent integration; no local test or coverage result for this parent adoption establishes that the known provenance gap is fixed. The two-read closure envelope, rejection behavior and lack of reciprocal-handshake, process-exit or profile-removal proof are unchanged. + ### Current session-end sender adoption of status-reply provenance On 7 September 2026, the session-end sender stack reproduced the inherited status-reply gap at `edec535b8d8b3c843f9f9baf8562f31ff5c7af23`: two real loopback connections using the same session and command id allowed the replacement reply to complete the original pending status request. The ordinary merge of #250 `bbdc6ace7a5932adf24836700f806850e6b230bc` preserves the session-end sender and both original sender integration-test files byte-for-byte while inheriting sealed status receipts and connection-aware correlation. The regression now requires the exact connection-mismatch error, unchanged pending count, and subsequent completion using the original received reply. The server fixtures have already finished; this proves retained reply/correlation recovery, not liveness of an open original stream or same-endpoint replacement coverage. diff --git a/docs/doctoring/webdriver-bidi-received-response-connection-provenance.md b/docs/doctoring/webdriver-bidi-received-response-connection-provenance.md new file mode 100644 index 000000000..b94d1faff --- /dev/null +++ b/docs/doctoring/webdriver-bidi-received-response-connection-provenance.md @@ -0,0 +1,55 @@ +# WebDriver BiDi received-response connection provenance + +Status: Draft implementation evidence for PR #255. This note is not an ADR and does not grant browser, policy, process, profile, or Agent authority. + +## Problem + +A WebDriver BiDi command id is local-end correlation data, not transport identity. Before this repair, `session.end` stored the private generation of connection A when sending command id `7`, but response parsing accepted a bare assembled text message. A valid response with id `7` read from a separately verified connection B could therefore consume A's outstanding command and inherit A's stored generation. Later teardown evidence could appear internally consistent even though the protocol acknowledgment actually arrived on B. + +The test-first commit `0a94765e0c631a043febec0973fd043304877537` reproduces this with two real loopback TCP/WebSocket connections using the same WebDriver session id and command id. On the then-current implementation the foreign response was accepted and the A correlation was consumed. This is source-level RED lineage; no hosted RED is claimed because later repair commits superseded that exact generation before terminal hosted execution. + +## Constraints and rejected alternatives + +WebDriver BiDi uses the command `id` only so the local end can identify a command response; commands may finish out of order, and the id is opaque to the remote end. The 3 September 2026 W3C Working Draft's command-processing algorithms retain the WebSocket connection as an explicit input/output boundary. OriginWeave therefore needs local evidence of the receiving connection in addition to the response id. + +Passing a caller-supplied connection generation into response parsing was rejected because it would turn provenance into forgeable metadata. Passing a bare assembled text message plus a separately supplied established connection was also rejected because callers could accidentally pair a message from B with connection A. Adding a connection id to JSON was rejected because it would invent protocol data not defined by WebDriver BiDi. + +## Selected boundary + +`WebDriverBiDiWebSocketMessageReader` consumes one established WebSocket and owns its message assembler. Every text fragment admitted by that reader is therefore read from the same non-cloneable verified transport. `Pending` and interleaved control outcomes retain the reader so fragmented state cannot be moved onto another connection. Once a text message is complete, the reader emits `WebDriverBiDiReceivedTextMessage` carrying the private process-local connection generation and returns the established transport separately because no partial fragments remain. + +`session.end` response admission requires that received-message type. Correlation validates command kind, stored connection provenance, and received connection generation before removing the outstanding id. Missing provenance and a different received connection both fail closed without consuming the command. Protocol ACK remains distinct from transport closure, Chromium process exit, profile deletion, and task completion. + +This placement also preserves RFC 6455 fragmentation semantics: a message may span multiple frames, control frames may appear between fragments, recipients must support fragmented and unfragmented messages, and message fragments are delivered in order on the WebSocket connection. The reader therefore binds the assembled-message transaction to one connection rather than treating individual application payloads as free-standing evidence. + +## Quality-gate follow-up + +A complete local verification snapshot at `e7bfec4488b7cb4776df7b546cacb46c8c9eb13e` established a second RED lane: the behavioral repair passed its focused connection-provenance tests, but rustfmt and strict Clippy were not clean and production coverage still missed connection-bound event/null-id errors, connection-generation exhaustion propagation, and a post-correlation provenance fallback that could not be reached after successful connection-bound correlation. Those failures are quality-gate evidence, not hosted GREEN. + +The follow-up keeps the security invariant and removes the gate-specific causes rather than excluding them. `feba41f13fb39e83fc8b377b8a55bd7c27348f9d` exercises event and null-id envelopes through the real `session.end` response path without consuming the outstanding command. The teardown provenance comparison is expressed without the nested `if` rejected by strict Clippy. Successful connection-bound correlation now carries the already-validated received-message generation directly into `WebDriverBiDiSessionEndResult`, eliminating the impossible optional-provenance fallback instead of excluding it from coverage. Connection establishment accepts a private generation-counter seam so the exhaustion error can be exercised after exact peer verification, and the allocator uses `AtomicU64::try_update`; Rust documents `try_update` as available since 1.95.0, which remains compatible with OriginWeave's declared Rust 1.97 MSRV. The added exhaustion regression uses a real loopback TCP stream plus a verified peer result and requires `ConnectionGenerationExhausted` after one connect and one peer-inspection call. + +Formatting-only test-name repairs preserve the realistic transport scenarios while restoring conventional rustfmt layout. Exact hosted CI/coverage/security results still belong to the live PR head and must be read from GitHub; predecessor results do not transfer. + +The `ebac126d1632c94775c2454423575275eec45def` follow-up removed the unused private correlated-response generation accessor, not the stored provenance or its validation. Rust 1.97.1 rustfmt, strict Clippy, workspace tests, rustdoc with `-D warnings`, and all 141 Python contracts then passed locally. The existing production coverage checker reported 1082/1082 functions, 11025/11025 lines, 14071/14071 regions, and 1202/1202 branches. However, pinned `cargo-llvm-cov 0.8.6` emitted `warning: --branch option is unstable and it may be changed in the future`. These numerical coverage results are not warning-free acceptance evidence. Keep the warning visible and the gate outstanding; do not suppress it, remove branch measurement, or change the denominator to manufacture a clean result. + +## Current-parent integration + +The ordinary integration of closure parent `b11b6c9bccd8335b58a4fb599f8ad29ac419637f` retains the complete `ebac126d...` connection-provenance repair. The parent restores the previously uncollected command-correlation release test and keeps peer fixtures alive until opening-exchange assertions finish. Native discovery first reported zero inherited release tests and failed the expected-one assertion; the integrated tree collects and passes that test. Sender registration, connection-bound message assembly, response admission before correlation consumption, closure-generation comparison and exhaustion handling are preserved without replacing their production or child-test blobs. + +The earlier quality paragraph now identifies `ebac126d...`, the actual commit that removed the unused accessor. Its predecessor `63cbca0...` still contained that accessor and unformatted statements; the later repair's passing results cannot be assigned backward. Fresh verification of the integrated tree remains separate from those predecessor measurements, queued hosted checks, warning-free instrumentation and operational process/profile cleanup evidence. + +Fresh integration verification executes 13 focused received-message, response and teardown tests, all 142 Python contracts, compileall, and the complete Rust 1.97.1 format/check/workspace-test/strict-Clippy/rustdoc gates. Pinned-nightly numerical coverage is 1082/1082 functions, 11032/11032 lines, 14086/14086 regions and 1202/1202 branches; the branch-instrumentation warning remains visible. These exact-tree local results neither authenticate a Chromium process nor establish hosted security acceptance, counted review approval or complete operational teardown. + +## Evidence and remaining risk + +The repair includes a realistic two-connection regression for a foreign `session.end` success response and focused reader coverage for fragmented text, an interleaved control frame, malformed server framing, binary-message rejection, event/null-id correlation, and connection-generation exhaustion after verified peer admission. The original inline response-substitution and closure-substitution findings are resolved by the current source boundary, but exact-head CI and security workflows remain authoritative before any integration claim. + +The current scope deliberately does not infer browser-process ownership from the WebSocket connection and does not make protocol success an operational teardown post-condition. Process-exit and profile-removal evidence remain unavailable until their canonical runtime owners provide non-forgeable contracts. + +## References + +Fette, I., & Melnikov, A. (2011). *The WebSocket Protocol* (RFC 6455). Internet Engineering Task Force. https://www.rfc-editor.org/rfc/rfc6455 + +The Rust Project Developers. (n.d.). *AtomicU64 in std::sync::atomic*. Rust standard library documentation. Retrieved September 5, 2026, from https://doc.rust-lang.org/std/sync/atomic/type.AtomicU64.html + +W3C Browser Testing and Tools Working Group. (2026, September 3). *WebDriver BiDi* (W3C Working Draft). World Wide Web Consortium. https://www.w3.org/TR/2026/WD-webdriver-bidi-20260903/