diff --git a/src/web/http.rs b/src/web/http.rs index 1472026..74e81aa 100644 --- a/src/web/http.rs +++ b/src/web/http.rs @@ -1,6 +1,6 @@ use std::convert::Infallible; use std::error::Error; -use std::net::{IpAddr, SocketAddr}; +use std::net::SocketAddr; use std::sync::Arc; use std::time::{Duration, Instant}; @@ -19,7 +19,6 @@ use tokio_util::sync::CancellationToken; use crate::config::{WebClientIpSource, WebRuntimeVhost}; use crate::web::bridge; -use crate::web::frame::{self, FrameType}; use crate::web::manager::{ManagerError, WebProcessRuntime}; // Response-body activity keeps connection idle accounting lifecycle-correct. @@ -34,6 +33,8 @@ mod down; mod request; // Carrier response construction and lane-header helpers are shared by handlers. mod response; +// Session creation and replacement negotiation remain separate from request routing. +mod session; // RFC 6455 upgrade validation and carrier drivers remain isolated from HTTP routing. #[cfg(test)] mod tests; @@ -50,19 +51,18 @@ use decoy::serve_decoy; use down::handle_down; use request::{ bearer_token_hash, binary_content_type, bridge_candidate, canonical_request_host, - canonical_u64_header, carrier_ip_learning_eligible, carrier_request, client_ip, - compatible_cookie_header, match_profile, + canonical_u64_header, client_ip, compatible_cookie_header, match_profile, }; use response::{ bad_gateway, carrier_empty, carrier_headers, carrier_lane, full_response, generic_not_found, insert_header, service_unavailable, }; +use session::handle_session; type BoxError = Box; type HttpBody = UnsyncBoxBody; type HttpResponse = Response; -const CREATE_BODY_LIMIT: usize = 64; const TRANSPORT_PATHS: [&str; 3] = ["/api/v1/session", "/api/v1/up", "/api/v1/down"]; const WEBSOCKET_PATH: &str = "/api/v1/ws"; @@ -317,199 +317,6 @@ async fn handle_api( } } -async fn handle_session( - request: Request, - runtime: Arc, - vhost: Arc, - token_hash: crate::web::manager::TokenHash, - client_ip: IpAddr, -) -> HttpResponse { - if request.headers().contains_key("x-lane-id") { - return serve_decoy(request, vhost, true, &runtime).await; - } - if request.method() == Method::DELETE { - if request.headers().contains_key(header::CONTENT_TYPE) { - return serve_decoy(request, vhost, true, &runtime).await; - } - if let Some(trace) = request_trace(&request) - && let Ok(session) = runtime.get_session(token_hash, &vhost.host) - { - trace.set_route(TraceRoute::Session); - trace.bind_identity(session.trace_identity()); - } - let CollectedBody { - request, - body, - _body_budget, - } = match collect_body(request, &runtime, 1, true).await { - Ok(result) => result, - Err(CollectBodyError::Limit) => return service_unavailable(), - Err(CollectBodyError::Invalid(request)) => { - return serve_decoy(request, vhost, true, &runtime).await; - } - }; - if !body.is_empty() || runtime.close_token(token_hash, &vhost.host).is_err() { - return serve_decoy(request, vhost, true, &runtime).await; - } - return carrier_empty(StatusCode::NO_CONTENT); - } - if request.method() != Method::POST || !binary_content_type(&request) { - return serve_decoy(request, vhost, true, &runtime).await; - } - let Some(carrier_request) = carrier_request(&request, &vhost.host) else { - return serve_decoy(request, vhost, true, &runtime).await; - }; - let ip_learning_eligible = carrier_ip_learning_eligible(&request, client_ip); - let Some((trace_session_id, profile)) = - runtime.bootstrap_trace_identity(token_hash, &vhost.host) - else { - return serve_decoy(request, vhost, true, &runtime).await; - }; - if let Some(trace) = request_trace(&request) { - trace.set_route(TraceRoute::Session); - trace.bind_profile(&profile, trace_session_id); - } - let CollectedBody { - request, - body, - _body_budget, - } = match collect_body(request, &runtime, CREATE_BODY_LIMIT, false).await { - Ok(result) => result, - Err(CollectBodyError::Limit) => return service_unavailable(), - Err(CollectBodyError::Invalid(request)) => { - return serve_decoy(request, vhost, true, &runtime).await; - } - }; - if let Some(trace) = request_trace(&request) { - trace.record_frames( - TraceDirection::Request, - &body, - &runtime.active_generation().config().web.limits, - ); - } - match runtime.create_session( - token_hash, - &vhost.host, - client_ip, - &body, - carrier_request, - ip_learning_eligible, - ) { - Ok(result) => { - let welcome = frame::encode(FrameType::Welcome, 0, &[]); - if let Some(trace) = request_trace(&request) { - trace.register_redaction(result.token.as_bytes()); - trace.record_frames( - TraceDirection::Response, - &welcome, - &runtime.active_generation().config().web.limits, - ); - } - let mut response = full_response(StatusCode::OK, welcome); - carrier_headers(&mut response); - insert_header( - &mut response, - HeaderName::from_static("x-session-token"), - &result.token, - ); - response.headers_mut().insert( - HeaderName::from_static("x-carrier-mode"), - HeaderValue::from_static(result.carrier.as_str()), - ); - response.headers_mut().insert( - HeaderName::from_static("x-down-cursor"), - HeaderValue::from_static("0"), - ); - if let Some(attempt) = result.attempt { - insert_header( - &mut response, - HeaderName::from_static("x-carrier-attempt"), - &attempt.to_string(), - ); - } - if let Some(candidate_count) = result.candidate_count { - insert_header( - &mut response, - HeaderName::from_static("x-carrier-candidate-count"), - &candidate_count.to_string(), - ); - insert_header( - &mut response, - HeaderName::from_static("x-carrier-deadline"), - &result.deadline_secs.unwrap_or_default().to_string(), - ); - if let Some(state) = result.carrier_state { - insert_header( - &mut response, - HeaderName::from_static("x-carrier-state"), - state, - ); - } - } - response - } - Err(ManagerError::Committed) => { - let mut response = carrier_empty(StatusCode::CONFLICT); - if let Some(echo) = runtime.carrier_echo( - token_hash, - &vhost.host, - client_ip, - carrier_request, - ) { - response.headers_mut().insert( - HeaderName::from_static("x-carrier-mode"), - HeaderValue::from_static(echo.carrier.as_str()), - ); - insert_header( - &mut response, - HeaderName::from_static("x-carrier-attempt"), - &echo.attempt.to_string(), - ); - insert_header( - &mut response, - HeaderName::from_static("x-carrier-candidate-count"), - &echo.candidate_count.to_string(), - ); - insert_header( - &mut response, - HeaderName::from_static("x-carrier-deadline"), - &echo.deadline_secs.to_string(), - ); - insert_header( - &mut response, - HeaderName::from_static("x-carrier-state"), - echo.state, - ); - } - response - } - Err( - error @ (ManagerError::Limit | ManagerError::Backpressure | ManagerError::Concurrent), - ) => { - runtime.trace().record_profile_lifecycle( - client_ip, - Some(trace_session_id), - &profile, - TraceLifecycleEvent::SessionRejected, - None, - Some(error.as_str()), - ); - service_unavailable() - } - Err(error) => { - runtime.trace().record_profile_lifecycle( - client_ip, - Some(trace_session_id), - &profile, - TraceLifecycleEvent::SessionRejected, - None, - Some(error.as_str()), - ); - serve_decoy(request, vhost, true, &runtime).await - } - } -} - async fn handle_up( request: Request, runtime: Arc, diff --git a/src/web/http/session.rs b/src/web/http/session.rs new file mode 100644 index 0000000..2695cc2 --- /dev/null +++ b/src/web/http/session.rs @@ -0,0 +1,213 @@ +use std::net::IpAddr; +use std::sync::Arc; + +use hyper::header::{self, HeaderName, HeaderValue}; +use hyper::{Method, Request, StatusCode}; + +use super::body::{CollectBodyError, CollectedBody, RequestBody, collect_body}; +use super::decoy::serve_decoy; +use super::request::{binary_content_type, carrier_ip_learning_eligible, carrier_request}; +use super::response::{ + carrier_empty, carrier_headers, full_response, insert_header, service_unavailable, +}; +use super::{HttpResponse, request_trace}; +use crate::config::WebRuntimeVhost; +use crate::web::frame::{self, FrameType}; +use crate::web::manager::{ManagerError, TokenHash, WebProcessRuntime}; +use crate::web::trace::{TraceDirection, TraceLifecycleEvent, TraceRoute}; + +const CREATE_BODY_LIMIT: usize = 64; + +/// Handles session creation, replacement replay, and authenticated closure. +pub(super) async fn handle_session( + request: Request, + runtime: Arc, + vhost: Arc, + token_hash: TokenHash, + client_ip: IpAddr, +) -> HttpResponse { + if request.headers().contains_key("x-lane-id") { + return serve_decoy(request, vhost, true, &runtime).await; + } + if request.method() == Method::DELETE { + if request.headers().contains_key(header::CONTENT_TYPE) { + return serve_decoy(request, vhost, true, &runtime).await; + } + if let Some(trace) = request_trace(&request) + && let Ok(session) = runtime.get_session(token_hash, &vhost.host) + { + trace.set_route(TraceRoute::Session); + trace.bind_identity(session.trace_identity()); + } + let CollectedBody { + request, + body, + _body_budget, + } = match collect_body(request, &runtime, 1, true).await { + Ok(result) => result, + Err(CollectBodyError::Limit) => return service_unavailable(), + Err(CollectBodyError::Invalid(request)) => { + return serve_decoy(request, vhost, true, &runtime).await; + } + }; + if !body.is_empty() || runtime.close_token(token_hash, &vhost.host).is_err() { + return serve_decoy(request, vhost, true, &runtime).await; + } + return carrier_empty(StatusCode::NO_CONTENT); + } + if request.method() != Method::POST || !binary_content_type(&request) { + return serve_decoy(request, vhost, true, &runtime).await; + } + let Some(carrier_request) = carrier_request(&request, &vhost.host) else { + return serve_decoy(request, vhost, true, &runtime).await; + }; + let ip_learning_eligible = carrier_ip_learning_eligible(&request, client_ip); + let Some((trace_session_id, profile)) = + runtime.bootstrap_trace_identity(token_hash, &vhost.host) + else { + return serve_decoy(request, vhost, true, &runtime).await; + }; + if let Some(trace) = request_trace(&request) { + trace.set_route(TraceRoute::Session); + trace.bind_profile(&profile, trace_session_id); + } + let CollectedBody { + request, + body, + _body_budget, + } = match collect_body(request, &runtime, CREATE_BODY_LIMIT, false).await { + Ok(result) => result, + Err(CollectBodyError::Limit) => return service_unavailable(), + Err(CollectBodyError::Invalid(request)) => { + return serve_decoy(request, vhost, true, &runtime).await; + } + }; + if let Some(trace) = request_trace(&request) { + trace.record_frames( + TraceDirection::Request, + &body, + &runtime.active_generation().config().web.limits, + ); + } + match runtime.create_session( + token_hash, + &vhost.host, + client_ip, + &body, + carrier_request, + ip_learning_eligible, + ) { + Ok(result) => { + let welcome = frame::encode(FrameType::Welcome, 0, &[]); + if let Some(trace) = request_trace(&request) { + trace.register_redaction(result.token.as_bytes()); + trace.record_frames( + TraceDirection::Response, + &welcome, + &runtime.active_generation().config().web.limits, + ); + } + let mut response = full_response(StatusCode::OK, welcome); + carrier_headers(&mut response); + insert_header( + &mut response, + HeaderName::from_static("x-session-token"), + &result.token, + ); + response.headers_mut().insert( + HeaderName::from_static("x-carrier-mode"), + HeaderValue::from_static(result.carrier.as_str()), + ); + response.headers_mut().insert( + HeaderName::from_static("x-down-cursor"), + HeaderValue::from_static("0"), + ); + if let Some(attempt) = result.attempt { + insert_header( + &mut response, + HeaderName::from_static("x-carrier-attempt"), + &attempt.to_string(), + ); + } + if let Some(candidate_count) = result.candidate_count { + insert_header( + &mut response, + HeaderName::from_static("x-carrier-candidate-count"), + &candidate_count.to_string(), + ); + insert_header( + &mut response, + HeaderName::from_static("x-carrier-deadline"), + &result.deadline_secs.unwrap_or_default().to_string(), + ); + if let Some(state) = result.carrier_state { + insert_header( + &mut response, + HeaderName::from_static("x-carrier-state"), + state, + ); + } + } + response + } + Err(ManagerError::Committed) => { + let mut response = carrier_empty(StatusCode::CONFLICT); + if let Some(echo) = runtime.carrier_echo( + token_hash, + &vhost.host, + client_ip, + carrier_request, + ) { + response.headers_mut().insert( + HeaderName::from_static("x-carrier-mode"), + HeaderValue::from_static(echo.carrier.as_str()), + ); + insert_header( + &mut response, + HeaderName::from_static("x-carrier-attempt"), + &echo.attempt.to_string(), + ); + insert_header( + &mut response, + HeaderName::from_static("x-carrier-candidate-count"), + &echo.candidate_count.to_string(), + ); + insert_header( + &mut response, + HeaderName::from_static("x-carrier-deadline"), + &echo.deadline_secs.to_string(), + ); + insert_header( + &mut response, + HeaderName::from_static("x-carrier-state"), + echo.state, + ); + } + response + } + Err( + error @ (ManagerError::Limit | ManagerError::Backpressure | ManagerError::Concurrent), + ) => { + runtime.trace().record_profile_lifecycle( + client_ip, + Some(trace_session_id), + &profile, + TraceLifecycleEvent::SessionRejected, + None, + Some(error.as_str()), + ); + service_unavailable() + } + Err(error) => { + runtime.trace().record_profile_lifecycle( + client_ip, + Some(trace_session_id), + &profile, + TraceLifecycleEvent::SessionRejected, + None, + Some(error.as_str()), + ); + serve_decoy(request, vhost, true, &runtime).await + } + } +} diff --git a/src/web/http/websocket/driver.rs b/src/web/http/websocket/driver.rs index 2d4e918..2658c00 100644 --- a/src/web/http/websocket/driver.rs +++ b/src/web/http/websocket/driver.rs @@ -21,7 +21,10 @@ const WRITE_BUFFER_BYTES: usize = 64 * 1024; // Cancellation-safe message I/O and budget retries remain separate from carrier loops. mod io; -use io::{flush, process_lane, process_multiplex, read_message, record_message, reserve_data, send}; +// Per-lane carrier state remains isolated from the multiplexed driver. +mod lane; +use io::{flush, process_multiplex, read_message, record_message, reserve_data, send}; +use lane::run_lane; pub(super) async fn run_upgraded( on_upgrade: hyper::upgrade::OnUpgrade, @@ -349,260 +352,6 @@ async fn run_multiplex( } } -#[allow(clippy::too_many_arguments)] -async fn run_lane( - socket: &mut CarrierSocket, - runtime: &Arc, - session: &Arc, - connection: &WebSocketConnection, - reservation: &mut WebSocketLaneReservation, - cancellation: CancellationToken, - trace: Option<&TraceWebSocketContext>, - acknowledge_commit: bool, -) -> Result<(), ()> { - let mut sequence = 1u64; - let mut cursor = 0u64; - // Lane reads use the same cancellation-safe fragmented-message ownership. - let mut read_budget = None; - let liveness_interval = connection.liveness_interval(); - let mut next_ping = Instant::now() + liveness_interval; - let open_deadline = - Instant::now() + Duration::from_secs(session.timeouts().websocket_open_secs); - let backpressure_timeout = - Duration::from_secs(session.timeouts().websocket_backpressure_secs); - let write_timeout = Duration::from_secs(session.timeouts().websocket_write_secs); - let maximum_message = session.limits().carrier_batch_bytes; - let mut active = false; - loop { - let down = session.poll_down_lane(reservation.lane_id(), cursor); - tokio::pin!(down); - let event = tokio::select! { - _ = cancellation.cancelled() => return Err(()), - _ = tokio::time::sleep_until(open_deadline.into()), if !active => return Err(()), - _ = tokio::time::sleep_until(next_ping.into()) => DriverEvent::Liveness, - incoming = read_message( - socket, - runtime, - session.profile_key(), - &cancellation, - &mut read_budget, - maximum_message, - backpressure_timeout, - ) => { - DriverEvent::Incoming(incoming?) - } - down = &mut down => DriverEvent::Down(down.map_err(|_| ())?), - }; - match event { - DriverEvent::Incoming((message, _budget)) => match message { - Message::Binary(body) => { - let started = Instant::now(); - let result = process_lane( - runtime, - session, - reservation, - sequence, - &body, - &cancellation, - backpressure_timeout, - ) - .await; - record_message( - runtime, - trace, - TraceDirection::Request, - "binary", - &body, - started, - ); - let progressed = result?; - if acknowledge_commit && sequence == 1 { - if !session.needs_websocket_commit_ack(connection.id()) { - return Err(()); - } - let started = Instant::now(); - if send( - socket, - Message::Binary(Bytes::new()), - &cancellation, - write_timeout, - ) - .await - .is_err() - { - session.close(); - return Err(()); - } - record_message( - runtime, - trace, - TraceDirection::Response, - "carrier-ack", - &[], - started, - ); - if !session.websocket_commit_ack_written(connection.id()) { - session.close(); - return Err(()); - } - } else if acknowledge_commit && sequence > 1 && progressed { - if !session.websocket_peer_after_commit_ack(connection.id()) { - return Err(()); - } - } - if !active && progressed { - if !connection.mark_active() { - return Err(()); - } - active = true; - } - sequence = sequence.checked_add(1).ok_or(())?; - connection.mark_peer_activity(); - next_ping = Instant::now() + liveness_interval; - } - Message::Pong(payload) => { - record_message( - runtime, - trace, - TraceDirection::Request, - "pong", - &payload, - Instant::now(), - ); - connection.mark_peer_activity(); - next_ping = Instant::now() + liveness_interval; - } - Message::Ping(payload) => { - let started = Instant::now(); - flush(socket, &cancellation, write_timeout).await?; - record_message( - runtime, - trace, - TraceDirection::Request, - "ping", - &payload, - started, - ); - record_message( - runtime, - trace, - TraceDirection::Response, - "pong", - &payload, - started, - ); - connection.mark_peer_activity(); - next_ping = Instant::now() + liveness_interval; - } - Message::Close(_) => { - record_message( - runtime, - trace, - TraceDirection::Request, - "close", - &[], - Instant::now(), - ); - return Ok(()); - } - Message::Text(text) => { - record_message( - runtime, - trace, - TraceDirection::Request, - "text", - text.as_bytes(), - Instant::now(), - ); - return Err(()); - } - Message::Frame(_) => return Err(()), - }, - DriverEvent::Down(result) => { - if result.lane_closed { - return Ok(()); - } - if result.body.is_empty() { - let started = Instant::now(); - send( - socket, - Message::Ping(Bytes::new()), - &cancellation, - write_timeout, - ) - .await?; - record_message( - runtime, - trace, - TraceDirection::Response, - "ping", - &[], - started, - ); - next_ping = Instant::now() + liveness_interval; - } else { - let _budget = reserve_data( - runtime, - session.profile_key(), - result.body.len(), - &cancellation, - backpressure_timeout, - ) - .await?; - let body = result.body; - let started = Instant::now(); - if trace.is_some() { - send( - socket, - Message::Binary(body.clone()), - &cancellation, - write_timeout, - ) - .await?; - record_message( - runtime, - trace, - TraceDirection::Response, - "binary", - &body, - started, - ); - } else { - send( - socket, - Message::Binary(body), - &cancellation, - write_timeout, - ) - .await?; - } - connection.mark_progress(); - } - cursor = result.next_cursor; - } - DriverEvent::Liveness => { - let started = Instant::now(); - send( - socket, - Message::Ping(Bytes::new()), - &cancellation, - write_timeout, - ) - .await?; - record_message( - runtime, - trace, - TraceDirection::Response, - "ping", - &[], - started, - ); - next_ping = Instant::now() + liveness_interval; - } - } - } -} - enum DriverEvent { Incoming((Message, Option)), Down(crate::web::session::PollResult), diff --git a/src/web/http/websocket/driver/lane.rs b/src/web/http/websocket/driver/lane.rs new file mode 100644 index 0000000..fc24f93 --- /dev/null +++ b/src/web/http/websocket/driver/lane.rs @@ -0,0 +1,274 @@ +use std::sync::Arc; +use std::time::{Duration, Instant}; + +use bytes::Bytes; +use tokio_tungstenite::tungstenite::protocol::Message; +use tokio_util::sync::CancellationToken; + +use super::CarrierSocket; +use super::io::{flush, process_lane, read_message, record_message, reserve_data, send}; +use crate::web::manager::{ + WebProcessRuntime, WebSocketBudgetLease, WebSocketConnection, +}; +use crate::web::session::{WebSession, WebSocketLaneReservation}; +use crate::web::trace::{TraceDirection, TraceWebSocketContext}; + +#[allow(clippy::too_many_arguments)] +pub(super) async fn run_lane( + socket: &mut CarrierSocket, + runtime: &Arc, + session: &Arc, + connection: &WebSocketConnection, + reservation: &mut WebSocketLaneReservation, + cancellation: CancellationToken, + trace: Option<&TraceWebSocketContext>, + acknowledge_commit: bool, +) -> Result<(), ()> { + let mut sequence = 1u64; + let mut cursor = 0u64; + // Lane reads use the same cancellation-safe fragmented-message ownership. + let mut read_budget = None; + let liveness_interval = connection.liveness_interval(); + let mut next_ping = Instant::now() + liveness_interval; + let open_deadline = + Instant::now() + Duration::from_secs(session.timeouts().websocket_open_secs); + let backpressure_timeout = + Duration::from_secs(session.timeouts().websocket_backpressure_secs); + let write_timeout = Duration::from_secs(session.timeouts().websocket_write_secs); + let maximum_message = session.limits().carrier_batch_bytes; + let mut active = false; + loop { + let down = session.poll_down_lane(reservation.lane_id(), cursor); + tokio::pin!(down); + let event = tokio::select! { + _ = cancellation.cancelled() => return Err(()), + _ = tokio::time::sleep_until(open_deadline.into()), if !active => return Err(()), + _ = tokio::time::sleep_until(next_ping.into()) => DriverEvent::Liveness, + incoming = read_message( + socket, + runtime, + session.profile_key(), + &cancellation, + &mut read_budget, + maximum_message, + backpressure_timeout, + ) => { + DriverEvent::Incoming(incoming?) + } + down = &mut down => DriverEvent::Down(down.map_err(|_| ())?), + }; + match event { + DriverEvent::Incoming((message, _budget)) => match message { + Message::Binary(body) => { + let started = Instant::now(); + let result = process_lane( + runtime, + session, + reservation, + sequence, + &body, + &cancellation, + backpressure_timeout, + ) + .await; + record_message( + runtime, + trace, + TraceDirection::Request, + "binary", + &body, + started, + ); + let progressed = result?; + if acknowledge_commit && sequence == 1 { + if !session.needs_websocket_commit_ack(connection.id()) { + return Err(()); + } + let started = Instant::now(); + if send( + socket, + Message::Binary(Bytes::new()), + &cancellation, + write_timeout, + ) + .await + .is_err() + { + session.close(); + return Err(()); + } + record_message( + runtime, + trace, + TraceDirection::Response, + "carrier-ack", + &[], + started, + ); + if !session.websocket_commit_ack_written(connection.id()) { + session.close(); + return Err(()); + } + } else if acknowledge_commit && sequence > 1 && progressed { + if !session.websocket_peer_after_commit_ack(connection.id()) { + return Err(()); + } + } + if !active && progressed { + if !connection.mark_active() { + return Err(()); + } + active = true; + } + sequence = sequence.checked_add(1).ok_or(())?; + connection.mark_peer_activity(); + next_ping = Instant::now() + liveness_interval; + } + Message::Pong(payload) => { + record_message( + runtime, + trace, + TraceDirection::Request, + "pong", + &payload, + Instant::now(), + ); + connection.mark_peer_activity(); + next_ping = Instant::now() + liveness_interval; + } + Message::Ping(payload) => { + let started = Instant::now(); + flush(socket, &cancellation, write_timeout).await?; + record_message( + runtime, + trace, + TraceDirection::Request, + "ping", + &payload, + started, + ); + record_message( + runtime, + trace, + TraceDirection::Response, + "pong", + &payload, + started, + ); + connection.mark_peer_activity(); + next_ping = Instant::now() + liveness_interval; + } + Message::Close(_) => { + record_message( + runtime, + trace, + TraceDirection::Request, + "close", + &[], + Instant::now(), + ); + return Ok(()); + } + Message::Text(text) => { + record_message( + runtime, + trace, + TraceDirection::Request, + "text", + text.as_bytes(), + Instant::now(), + ); + return Err(()); + } + Message::Frame(_) => return Err(()), + }, + DriverEvent::Down(result) => { + if result.lane_closed { + return Ok(()); + } + if result.body.is_empty() { + let started = Instant::now(); + send( + socket, + Message::Ping(Bytes::new()), + &cancellation, + write_timeout, + ) + .await?; + record_message( + runtime, + trace, + TraceDirection::Response, + "ping", + &[], + started, + ); + next_ping = Instant::now() + liveness_interval; + } else { + let _budget = reserve_data( + runtime, + session.profile_key(), + result.body.len(), + &cancellation, + backpressure_timeout, + ) + .await?; + let body = result.body; + let started = Instant::now(); + if trace.is_some() { + send( + socket, + Message::Binary(body.clone()), + &cancellation, + write_timeout, + ) + .await?; + record_message( + runtime, + trace, + TraceDirection::Response, + "binary", + &body, + started, + ); + } else { + send( + socket, + Message::Binary(body), + &cancellation, + write_timeout, + ) + .await?; + } + connection.mark_progress(); + } + cursor = result.next_cursor; + } + DriverEvent::Liveness => { + let started = Instant::now(); + send( + socket, + Message::Ping(Bytes::new()), + &cancellation, + write_timeout, + ) + .await?; + record_message( + runtime, + trace, + TraceDirection::Response, + "ping", + &[], + started, + ); + next_ping = Instant::now() + liveness_interval; + } + } + } +} + +enum DriverEvent { + Incoming((Message, Option)), + Down(crate::web::session::PollResult), + Liveness, +} diff --git a/src/web/manager.rs b/src/web/manager.rs index 4456304..f4f5fb2 100644 --- a/src/web/manager.rs +++ b/src/web/manager.rs @@ -191,9 +191,7 @@ impl WebProcessRuntime { trace, http_connections: Arc::new(Semaphore::new(limits.max_http_connections)), http_handlers: Arc::new(Semaphore::new(limits.max_http_handlers)), - lane_polls: Arc::new(Semaphore::new( - lane_poll_limit, - )), + lane_polls: Arc::new(Semaphore::new(lane_poll_limit)), lane_aux_polls: Arc::new(Semaphore::new(lane_aux_poll_limit)), body_readers: Arc::new(Semaphore::new(limits.max_body_readers)), body_bytes: Arc::new(Semaphore::new(limits.max_body_bytes_global)), diff --git a/src/web/manager/carrier_outcome.rs b/src/web/manager/carrier_outcome.rs index 70a6c38..c77fb64 100644 --- a/src/web/manager/carrier_outcome.rs +++ b/src/web/manager/carrier_outcome.rs @@ -80,7 +80,7 @@ impl WebProcessRuntime { class: CarrierClientClass, client_ip: IpAddr, identity: TraceIdentity, - ) { + ) -> bool { let mut state = self.state.lock(); let scores = state.bootstraps.get_mut(&bootstrap_hash).and_then(|entry| { if entry.carrier_attempt == attempt @@ -97,7 +97,7 @@ impl WebProcessRuntime { } }); drop(state); - let Some(scores) = scores else { return }; + let Some(scores) = scores else { return false }; self.trace.record_carrier_lifecycle( client_ip, identity.clone(), @@ -108,6 +108,7 @@ impl WebProcessRuntime { scores, None, ); + true } /// Promotes one exact committed attempt after transport-specific health evidence. diff --git a/src/web/manager/session_creation.rs b/src/web/manager/session_creation.rs index 61bbb12..0786385 100644 --- a/src/web/manager/session_creation.rs +++ b/src/web/manager/session_creation.rs @@ -276,13 +276,13 @@ impl WebProcessRuntime { now + Duration::from_secs(profile.carrier_negotiation_deadlines_secs[3]), ); let learning_context = learning_epoch.map(|epoch| CarrierLearningContext { - profile_key, - client_ip, - class: carrier_request.class(), - user_agent_hash: carrier_request.user_agent_hash(), - epoch, - ip_learning_eligible, - }); + profile_key, + client_ip, + class: carrier_request.class(), + user_agent_hash: carrier_request.user_agent_hash(), + epoch, + ip_learning_eligible, + }); let session = WebSession::new( Arc::downgrade(self), session_hash, @@ -540,5 +540,4 @@ impl WebProcessRuntime { ); Ok(result) } - } diff --git a/src/web/session.rs b/src/web/session.rs index aab018e..ed92b64 100644 --- a/src/web/session.rs +++ b/src/web/session.rs @@ -155,6 +155,7 @@ struct SessionState { carrier_health_activity_at: Option, carrier_health_uplink: bool, carrier_health_downlink: bool, + carrier_commit_published: bool, carrier_health_reported: bool, websocket_carrier_active: bool, websocket_commit_ack_pending: bool, @@ -281,6 +282,7 @@ impl WebSession { carrier_health_activity_at: None, carrier_health_uplink: false, carrier_health_downlink: false, + carrier_commit_published: false, carrier_health_reported: false, websocket_carrier_active: false, websocket_commit_ack_pending: false, diff --git a/src/web/session/backend.rs b/src/web/session/backend.rs index 2125374..c751b2e 100644 --- a/src/web/session/backend.rs +++ b/src/web/session/backend.rs @@ -118,11 +118,11 @@ impl WebSession { .then(|| state.streams.remove(&stream.id)) .flatten() .map(|stream_state| { - let (bytes, items) = inbound_queue_cost(&stream_state.inbound); - self.release_locked(&mut state, bytes, items, false); - self.remember_closed_locked(&mut state, stream.id); - self.queue_control_locked(&mut state, FrameType::Close, stream.id, &[]) - }) + let (bytes, items) = inbound_queue_cost(&stream_state.inbound); + self.release_locked(&mut state, bytes, items, false); + self.remember_closed_locked(&mut state, stream.id); + self.queue_control_locked(&mut state, FrameType::Close, stream.id, &[]) + }) }; if queued.is_some_and(|queued| !queued) { self.close(); @@ -137,12 +137,15 @@ impl WebSession { .streams .get(&stream.id) .is_some_and(|state| state.instance == stream.instance); - let queued = current.then(|| state.streams.remove(&stream.id)).flatten().map(|stream_state| { - let (bytes, items) = inbound_queue_cost(&stream_state.inbound); - self.release_locked(&mut state, bytes, items, false); - self.remember_closed_locked(&mut state, stream.id); - self.queue_control_locked(&mut state, FrameType::Close, stream.id, &[]) - }); + let queued = current + .then(|| state.streams.remove(&stream.id)) + .flatten() + .map(|stream_state| { + let (bytes, items) = inbound_queue_cost(&stream_state.inbound); + self.release_locked(&mut state, bytes, items, false); + self.remember_closed_locked(&mut state, stream.id); + self.queue_control_locked(&mut state, FrameType::Close, stream.id, &[]) + }); if state.closing_streams.get(&stream.id) == Some(&stream.instance) { state.closing_streams.remove(&stream.id); self.remember_closed_locked(&mut state, stream.id); diff --git a/src/web/session/backend_tests.rs b/src/web/session/backend_tests.rs index 0ebe828..077ca85 100644 --- a/src/web/session/backend_tests.rs +++ b/src/web/session/backend_tests.rs @@ -373,3 +373,17 @@ async fn cancellation_while_waiting_for_data_releases_stream_ownership() { ); runtime.shutdown().await; } + +#[tokio::test] +async fn exhausted_stream_identity_does_not_acquire_synthetic_port_ownership() { + let runtime = test_runtime(WebCarrier::Https, 1); + runtime.session.state.lock().next_stream_instance = u64::MAX; + + assert_eq!( + runtime.process_frame(1, 1, FrameType::Open, &[]), + Err(ManagerError::Closed) + ); + assert!(runtime.session.state.lock().active_peer_ports.is_empty()); + + runtime.shutdown().await; +} diff --git a/src/web/session/lane_uplink.rs b/src/web/session/lane_uplink.rs index 60d6380..70f1806 100644 --- a/src/web/session/lane_uplink.rs +++ b/src/web/session/lane_uplink.rs @@ -55,7 +55,14 @@ impl WebSession { .is_some_and(|value| value.frame_type != FrameType::Open) && only_late_frames(&frames) { - return Ok(sequence); + return if self.automatic_carrier + && state.negotiation_phase + != super::SessionNegotiationPhase::Committed + { + Err(ManagerError::Backpressure) + } else { + Ok(sequence) + }; } if lane_id == 0 || frames @@ -160,6 +167,10 @@ impl WebSession { drop(opened); return result; } + if self.automatic_carrier && !self.is_carrier_committed() { + self.lane_open_notify.notify_waiters(); + return Err(ManagerError::Backpressure); + } if committed { self.finish_carrier_commit(); } diff --git a/src/web/session/lanes.rs b/src/web/session/lanes.rs index c90a614..99e5a55 100644 --- a/src/web/session/lanes.rs +++ b/src/web/session/lanes.rs @@ -3,6 +3,7 @@ use std::time::{Duration, Instant}; use bytes::{BufMut, Bytes, BytesMut}; use tokio::sync::OwnedSemaphorePermit; + use super::lane_downlink::take_lane_down_batch; use super::{ PendingClass, PollResult, QUEUE_ITEM_COST, QueuedFrame, SessionState, WebSession, diff --git a/src/web/session/lanes/tests.rs b/src/web/session/lanes/tests.rs index 27aeb96..6de0f4c 100644 --- a/src/web/session/lanes/tests.rs +++ b/src/web/session/lanes/tests.rs @@ -20,6 +20,14 @@ fn session_with_limits(limits: WebLimitsConfig) -> Arc { fn new_session( limits: WebLimitsConfig, manager: std::sync::Weak, +) -> Arc { + new_session_with_automatic(limits, manager, false) +} + +fn new_session_with_automatic( + limits: WebLimitsConfig, + manager: std::sync::Weak, + automatic: bool, ) -> Arc { let profile = Arc::new(WebRuntimeProfile { host: "proxy.example.com".to_string(), @@ -48,9 +56,13 @@ fn new_session( 1, [3; 32], None, - crate::web::manager::CarrierClientClass::Legacy, + if automatic { + crate::web::manager::CarrierClientClass::Bridge + } else { + crate::web::manager::CarrierClientClass::Legacy + }, None, - false, + automatic, limits, WebTimeoutsConfig::default(), ) @@ -259,3 +271,20 @@ fn tombstone_eviction_releases_lane_budget_and_accepts_late_frames() { assert_eq!(session.process_up_lane(7, 7, &late), Ok(7)); assert!(!session.state.lock().closed); } + +#[test] +fn automatic_lane_does_not_ack_a_missing_lane_without_real_progress() { + let session = new_session_with_automatic( + WebLimitsConfig::default(), + std::sync::Weak::new(), + true, + ); + let late = frame::encode(FrameType::Data, 7, b"late"); + + assert_eq!( + session.process_up_lane(7, 1, &late), + Err(ManagerError::Backpressure) + ); + assert!(!session.is_carrier_committed()); + assert!(!session.state.lock().closed); +} diff --git a/src/web/session/negotiation.rs b/src/web/session/negotiation.rs index fc88320..8bd1836 100644 --- a/src/web/session/negotiation.rs +++ b/src/web/session/negotiation.rs @@ -31,7 +31,7 @@ impl WebSession { /// Publishes the already-linearized session commit to process state. pub(super) fn finish_carrier_commit(&self) { - if let Some(manager) = self.manager.upgrade() { + let published = self.manager.upgrade().is_some_and(|manager| { manager.carrier_committed( self.bootstrap_hash, self.token_hash, @@ -40,7 +40,22 @@ impl WebSession { self.carrier_class, self.client_ip, self.trace_identity(), - ); + ) + }); + if !published { + return; + } + let healthy = { + let mut state = self.state.lock(); + if state.closed || state.negotiation_phase != SessionNegotiationPhase::Committed { + false + } else { + state.carrier_commit_published = true; + self.carrier_health_ready_locked(&mut state, Instant::now()) + } + }; + if healthy { + self.finish_carrier_health(); } } @@ -97,7 +112,9 @@ impl WebSession { now: Instant, ) -> bool { if !self.automatic_carrier + || state.closed || state.negotiation_phase != SessionNegotiationPhase::Committed + || !state.carrier_commit_published || state.carrier_health_reported || state.carrier_health_due_at.is_none_or(|due| now < due) { @@ -233,6 +250,7 @@ mod tests { let now = Instant::now(); let mut state = session.state.lock(); state.negotiation_phase = SessionNegotiationPhase::Committed; + state.carrier_commit_published = true; state.carrier_health_due_at = Some(now - Duration::from_secs(1)); state.carrier_health_uplink = true; state.carrier_health_downlink = true; @@ -251,6 +269,7 @@ mod tests { let now = Instant::now(); let mut state = session.state.lock(); state.negotiation_phase = SessionNegotiationPhase::Committed; + state.carrier_commit_published = true; state.carrier_health_due_at = Some(now - Duration::from_secs(1)); state.websocket_carrier_active = true; state.websocket_commit_ack_owner = Some(7); @@ -261,6 +280,21 @@ mod tests { assert!(session.carrier_health_ready_locked(&mut state, now)); } + #[test] + fn health_waits_for_manager_commit_publication() { + let session = session(WebCarrier::Https, Instant::now() + Duration::from_secs(60)); + let now = Instant::now(); + let mut state = session.state.lock(); + state.negotiation_phase = SessionNegotiationPhase::Committed; + state.carrier_health_due_at = Some(now - Duration::from_secs(1)); + state.carrier_health_uplink = true; + state.carrier_health_downlink = true; + state.carrier_health_activity_at = Some(now); + + assert!(!session.carrier_health_ready_locked(&mut state, now)); + assert!(!state.carrier_health_reported); + } + #[test] fn commit_and_supersede_have_one_session_lock_winner() { let committed = session(WebCarrier::Https, Instant::now() + Duration::from_secs(60)); diff --git a/src/web/session/uplink.rs b/src/web/session/uplink.rs index 19af5e3..c268ef7 100644 --- a/src/web/session/uplink.rs +++ b/src/web/session/uplink.rs @@ -34,8 +34,11 @@ impl WebSession { sequence: u64, body: &[u8], ) -> Result { - self.process_up_inner(sequence, body) - .map(|(acknowledged, _)| acknowledged) + let (acknowledged, progressed) = self.process_up_inner(sequence, body)?; + if self.automatic_carrier && !progressed && !self.is_carrier_committed() { + return Err(ManagerError::Backpressure); + } + Ok(acknowledged) } /// Applies one WebSocket uplink batch and reports actual carrier progress. @@ -182,6 +185,9 @@ impl WebSession { || state.closing_streams.contains_key(&value.stream_id); match value.frame_type { FrameType::Open => { + let Some(stream) = next_stream_identity(state, value.stream_id) else { + return false; + }; let peer_port = match reserved_open.take() { Some((reserved_stream_id, peer_port)) if reserved_stream_id == value.stream_id => @@ -208,9 +214,6 @@ impl WebSession { peer_port } }; - let Some(stream) = next_stream_identity(state, value.stream_id) else { - return false; - }; state.streams.insert( value.stream_id, StreamState { @@ -420,6 +423,10 @@ mod tests { use crate::web::manager::WebProcessRuntime; fn session() -> Arc { + session_with_automatic(false) + } + + fn session_with_automatic(automatic: bool) -> Arc { let profile = Arc::new(WebRuntimeProfile { host: "proxy.example.com".to_string(), public_addr: SocketAddr::from(([203, 0, 113, 10], 443)), @@ -447,9 +454,13 @@ mod tests { 1, [3; 32], None, - crate::web::manager::CarrierClientClass::Legacy, + if automatic { + crate::web::manager::CarrierClientClass::Bridge + } else { + crate::web::manager::CarrierClientClass::Legacy + }, None, - false, + automatic, WebLimitsConfig::default(), WebTimeoutsConfig::default(), ) @@ -515,4 +526,14 @@ mod tests { assert_eq!(session.process_up(2, &body), Err(ManagerError::Protocol)); assert!(session.state.lock().closed); } + + #[test] + fn automatic_uplink_does_not_ack_a_batch_without_real_progress() { + let session = session_with_automatic(true); + let body = frame::encode(FrameType::Pong, 0, &[]); + + assert_eq!(session.process_up(1, &body), Err(ManagerError::Backpressure)); + assert!(!session.is_carrier_committed()); + assert!(!session.state.lock().closed); + } }