diff --git a/src/api/web_status/details.rs b/src/api/web_status/details.rs index 90ce3a2..f876e0d 100644 --- a/src/api/web_status/details.rs +++ b/src/api/web_status/details.rs @@ -97,6 +97,20 @@ pub(super) fn push_lifecycle(html: &mut String, event: &crate::web::trace::Trace ); html.push_str("\nreason: "); html.push_str(event.reason.unwrap_or("-")); + html.push_str("\npeer gap ms: "); + html.push_str( + &event + .peer_gap_ms + .map(|value| value.to_string()) + .unwrap_or_else(|| "-".to_string()), + ); + html.push_str("\npredecessor session: "); + html.push_str( + &event + .predecessor_session_id + .map(|value| value.to_string()) + .unwrap_or_else(|| "-".to_string()), + ); if let Some(carrier) = &event.carrier { html.push_str("\nclient class: "); html.push_str(carrier.client_class); diff --git a/src/metrics/web.rs b/src/metrics/web.rs index 57802fe..3588941 100644 --- a/src/metrics/web.rs +++ b/src/metrics/web.rs @@ -500,5 +500,24 @@ mod tests { assert!(output.contains( "telemt_web_carrier_learning_state{state=\"unavailable\"} 1" )); + assert_eq!( + output.matches("telemt_web_session_closures_total{").count(), + crate::config::WebCarrier::ALL.len() + * crate::web::session::SessionCloseReason::ALL.len() + ); + assert_eq!( + output + .matches("telemt_web_session_lifecycle_observations_total{") + .count(), + crate::config::WebCarrier::ALL.len() + * crate::web::telemetry::WebSessionLifecycleObservation::ALL.len() + ); + assert_eq!( + output + .matches("telemt_web_bridge_recovery_events_total{") + .count(), + crate::web::telemetry::WebBridgeRecoveryEvent::ALL.len() + ); + assert!(output.contains("telemt_web_bridge_recovery_seconds 15")); } } diff --git a/src/web/bridge/request.js b/src/web/bridge/request.js index 69ff439..b613f9b 100644 --- a/src/web/bridge/request.js +++ b/src/web/bridge/request.js @@ -33,10 +33,10 @@ function create(settings){ while(attemptcontroller.abort(); + const controller=new AbortController(),abort=()=>controller.abort();let timedOut=false; if(external)external.addEventListener('abort',abort,{once:true}); const requestOptions=Object.assign({},frozenOptions,{signal:controller.signal}); - const timer=setTimeout(abort,Math.max(1,Math.min(attemptLimit,remaining))); + const timer=setTimeout(()=>{timedOut=true;controller.abort()},Math.max(1,Math.min(attemptLimit,remaining))); let response=null,wait=0; try{ const fetched=await fetch(settings.origin()+path,requestOptions); @@ -45,12 +45,18 @@ function create(settings){ }else{ const policy=responsePolicy(path,fetched.status);let body; try{body=await settings.read(fetched,policy.limit,policy.exact,controller.signal)} - catch(error){controller.abort();throw settings.failure('protocol',error&&error.message)} + catch(error){ + controller.abort(); + if(external&&external.aborted)throw error; + if(timedOut)throw settings.failure('timeout','response deadline exceeded'); + throw settings.failure('protocol',error&&error.message); + } response={status:fetched.status,headers:fetched.headers,body};return response; } }catch(error){ controller.abort(); if(settings.closed()||(external&&external.aborted))throw error; + if(timedOut)lastReason='timeout'; if(settings.reason(error,'')==='protocol')throw error; }finally{clearTimeout(timer);if(external)external.removeEventListener('abort',abort)} const after=Math.min(deadline-Date.now(),remainingBudget?remainingBudget():Infinity);if(attempt>=maxAttempts||after<=0)break; diff --git a/src/web/bridge/runtime.js b/src/web/bridge/runtime.js index c0394c3..63dfb3d 100644 --- a/src/web/bridge/runtime.js +++ b/src/web/bridge/runtime.js @@ -268,6 +268,7 @@ async function runUp(){ if(response.headers.get('X-Up-Ack')!==sequence)throw failure('protocol','uplink replay acknowledgement rejected'); }); if(!recovered||closed||lease.cancelled||sessionToken!==token)return; + break; } } if(!settleBatch(lease))return;port.postMessage({t:'traffic',up:lease.total,down:0});upSequence++;lease=null; @@ -421,6 +422,7 @@ async function runLaneUp(lane){ if(response.headers.get('X-Up-Ack')!==sequence)throw failure('protocol','lane uplink replay acknowledgement rejected'); }); if(!recovered||closed||lease.cancelled||sessionToken!==token||lanes.get(lane.id)!==lane)return; + break; } } if(!settleBatch(lease))return;port.postMessage({t:'traffic',up:lease.total,down:0});lane.sequence++;lease=null; diff --git a/src/web/bridge/tests.rs b/src/web/bridge/tests.rs index 2a07e43..4efb271 100644 --- a/src/web/bridge/tests.rs +++ b/src/web/bridge/tests.rs @@ -36,6 +36,9 @@ fn rendered_page_contains_bounded_negotiation_contract() { assert!(page.body.contains("tproxy-auto-v1.")); assert!(page.body.contains("tproxy-auto-lane-v1.")); assert!(page.body.contains("globalThis.TelemtBridgeResponse")); + assert!(page.body.contains("globalThis.TelemtBridgeRequest")); + assert!(page.body.contains("globalThis.TelemtBridgeBuffers")); + assert!(page.body.contains("globalThis.TelemtBridgeRecovery")); assert!(page.body.contains("responseBody.read")); assert!(!page.body.contains("arrayBuffer()")); assert!(page.body.contains("maxChunks=4096")); @@ -77,13 +80,13 @@ fn rendered_page_embeds_the_configured_bridge_timing_policy() { &SecureRandom::new(), ); - assert!(page.body.contains("const longPollMs=17*1000")); + assert!(page.body.contains("let longPollMs=17*1000")); assert!(page.body.contains("bridgeRequestMs=7*1000")); assert!(page.body.contains("bridgeRetryMs=41*1000")); assert!(page.body.contains("bridgeRecoveryMs=13*1000")); assert!(page.body.contains("websocketOpenMs=11*1000")); assert!(page.body.contains("reconnectGraceMs=119*1000")); - assert!(page.body.contains("const probeCoalesceMs=4")); + assert!(page.body.contains("let probeCoalesceMs=4")); assert!(page.body.contains( "helloTimer=setTimeout(()=>fail('timeout'),bridgeRequestMs)" )); @@ -137,12 +140,12 @@ fn retry_and_attempt_state_are_frozen_before_fetch() { let page = render_page("EEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEE", 4); assert!( page.body - .contains("async function request(path,frozenOptions,remainingBudget,maxAttempts)") + .contains("async function send(path,frozenOptions,remainingBudget,maxAttempts)") ); assert!(!page.body.contains("makeOptions")); assert!( page.body - .contains("if(closed||(external&&external.aborted))throw new Error('request aborted')") + .contains("if(settings.closed()||(external&&external.aborted))throw new Error('request aborted')") ); assert!(page.body.contains( "const frozen=options('POST',bootstrap,snapshot.hello,attemptHeaders(snapshot.attempt,snapshot.failure),controller.signal)" diff --git a/src/web/http.rs b/src/web/http.rs index c722c86..5c291a6 100644 --- a/src/web/http.rs +++ b/src/web/http.rs @@ -243,6 +243,12 @@ async fn handle_root( .headers() .get(header::USER_AGENT) .and_then(|value| value.to_str().ok()); + let recovery_session = match representation { + recovery::RootRepresentation::Recovery(Some(hash)) => { + runtime.bridge_recovery_session(hash, &vhost.host, &profile) + } + _ => None, + }; let bootstrap = match match representation { recovery::RootRepresentation::Bridge => runtime.issue_bootstrap_for_request( &generation, @@ -256,6 +262,9 @@ async fn handle_root( Arc::clone(&profile), client_ip, user_agent, + recovery_session + .as_ref() + .map(|session| session.trace_session_id()), ), recovery::RootRepresentation::Invalid => { strip_query(&mut request); @@ -281,10 +290,8 @@ async fn handle_root( trace.register_redaction(bootstrap.token.as_bytes()); } let config = generation.config(); - if let recovery::RootRepresentation::Recovery(previous_hash) = representation { - if let Some(session) = previous_hash.and_then(|hash| { - runtime.bridge_recovery_session(hash, &vhost.host, &profile) - }) { + if let recovery::RootRepresentation::Recovery(_) = representation { + if let Some(session) = recovery_session { let outcome = session.close(crate::web::session::SessionCloseReason::BridgeRecovery); if outcome != crate::web::session::SessionCloseOutcome::Closed && tokio::time::timeout( diff --git a/src/web/http/recovery_tests.rs b/src/web/http/recovery_tests.rs new file mode 100644 index 0000000..d8f74cb --- /dev/null +++ b/src/web/http/recovery_tests.rs @@ -0,0 +1,244 @@ +use super::*; +use sha2::{Digest, Sha256}; + +use crate::web::session::{SessionCloseOutcome, SessionCloseReason}; +use crate::web::telemetry::{WebBridgeRecoveryEvent, WebSessionLifecycleObservation}; + +const RECOVERY_TYPE: &str = "application/vnd.telemt.web-recovery+json"; + +async fn bridge_bootstrap( + listener: &TcpListener, + runtime: &Arc, + encoded_capability: &str, +) -> String { + let root = format!( + "GET /?bridge={encoded_capability} HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.40\r\nConnection: close\r\n\r\n" + ) + .into_bytes(); + let response = request(listener, runtime, root).await; + let (_, body) = split_response(&response); + std::str::from_utf8(body) + .unwrap() + .split_once("bootstrap=\"") + .and_then(|(_, suffix)| suffix.split_once('"')) + .map(|(token, _)| token.to_string()) + .unwrap() +} + +async fn create_session( + listener: &TcpListener, + runtime: &Arc, + bootstrap: &str, +) -> Vec { + let hello = frame::encode(FrameType::Hello, 0, &[1]); + let mut create = format!( + "POST /api/v1/session HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.40\r\nAuthorization: Bearer {bootstrap}\r\nContent-Type: application/octet-stream\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", + hello.len() + ) + .into_bytes(); + create.extend_from_slice(&hello); + request(listener, runtime, create).await +} + +async fn recover( + listener: &TcpListener, + runtime: &Arc, + encoded_capability: &str, + authorization: &str, +) -> Vec { + request( + listener, + runtime, + format!( + "GET /?bridge={encoded_capability} HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.40\r\nAccept: {RECOVERY_TYPE}\r\nAuthorization: {authorization}\r\nConnection: close\r\n\r\n" + ) + .into_bytes(), + ) + .await +} + +#[tokio::test] +async fn recovery_retires_current_bearer_before_single_slot_recreation() { + let capability = [40u8; 32]; + let mut config = runtime_config(capability, WebCarrier::Https); + config.web.limits.max_sessions_global = 1; + config.web.limits.max_sessions_per_ip = 1; + let generation = test_runtime_generation(1, config); + let active_runtime = Arc::new(ArcSwap::from(Arc::clone(&generation))); + let runtime = WebProcessRuntime::start(active_runtime); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(capability); + let bootstrap = bridge_bootstrap(&listener, &runtime, &encoded).await; + let create_response = create_session(&listener, &runtime, &bootstrap).await; + let (create_headers, _) = split_response(&create_response); + assert!(create_headers.starts_with(b"HTTP/1.1 200")); + let old_session = response_header(create_headers, "x-session-token").to_string(); + + let recovery_response = recover( + &listener, + &runtime, + &encoded, + &format!("Bearer {old_session}"), + ) + .await; + let (recovery_headers, recovery_body) = split_response(&recovery_response); + assert!(recovery_headers.starts_with(b"HTTP/1.1 200")); + assert_eq!(response_header(recovery_headers, "content-type"), RECOVERY_TYPE); + assert_eq!(response_header(recovery_headers, "cache-control"), "no-store"); + assert!(recovery_body.len() <= 1024); + let document: serde_json::Value = serde_json::from_slice(recovery_body).unwrap(); + assert_eq!(document["v"], 1); + assert_eq!(document["timeouts"]["bridge_recovery_secs"], 15); + let recovery_bootstrap = document["bootstrap"].as_str().unwrap(); + + let recreated = create_session(&listener, &runtime, recovery_bootstrap).await; + assert!(recreated.starts_with(b"HTTP/1.1 200")); + assert_eq!( + runtime.telemetry().session_close_total( + WebCarrier::Https, + SessionCloseReason::BridgeRecovery, + ), + 1 + ); + assert_eq!( + runtime + .telemetry() + .bridge_recovery_total(WebBridgeRecoveryEvent::BootstrapIssued), + 1 + ); + assert_eq!( + runtime + .telemetry() + .bridge_recovery_total(WebBridgeRecoveryEvent::SessionCreated), + 1 + ); + + let late = format!( + "POST /api/v1/down HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.40\r\nAuthorization: Bearer {old_session}\r\nX-Down-Cursor: 0\r\nContent-Length: 0\r\nConnection: close\r\n\r\n" + ) + .into_bytes(); + let late_response = request(&listener, &runtime, late).await; + assert!(late_response.starts_with(b"HTTP/1.1 404")); + assert_eq!( + runtime.telemetry().session_observation_total( + WebCarrier::Https, + WebSessionLifecycleObservation::RequestAfterClose, + ), + 1 + ); + + let retired_recovery = recover( + &listener, + &runtime, + &encoded, + &format!("Bearer {old_session}"), + ) + .await; + let (retired_headers, retired_body) = split_response(&retired_recovery); + assert!(retired_headers.starts_with(b"HTTP/1.1 200")); + assert_eq!(response_header(retired_headers, "content-type"), RECOVERY_TYPE); + assert!(serde_json::from_slice::(retired_body).is_ok()); + assert_eq!( + runtime + .list_sessions(SessionListRequest { + limit: 10, + cursor: None, + filter: SessionFilter::default(), + }) + .sessions + .len(), + 1 + ); + + runtime.shutdown().await; + generation.stop_sessions().await; + generation.stop_background_tasks().await; +} + +#[tokio::test] +async fn malformed_or_over_capacity_recovery_is_indistinguishable_from_decoy() { + let capability = [41u8; 32]; + let generation = test_runtime_generation(1, runtime_config(capability, WebCarrier::Https)); + let active_runtime = Arc::new(ArcSwap::from(Arc::clone(&generation))); + let runtime = WebProcessRuntime::start(active_runtime); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(capability); + + let malformed = recover(&listener, &runtime, &encoded, "Bearer malformed").await; + let (malformed_headers, malformed_body) = split_response(&malformed); + assert!(malformed_headers.starts_with(b"HTTP/1.1 200")); + assert_eq!(response_header(malformed_headers, "cache-control"), "no-store"); + assert_eq!(malformed_body, b"decoy"); + + let _held = bridge_bootstrap(&listener, &runtime, &encoded).await; + let over_capacity = recover( + &listener, + &runtime, + &encoded, + &format!("Bearer {}", "U".repeat(43)), + ) + .await; + let (capacity_headers, capacity_body) = split_response(&over_capacity); + assert!(capacity_headers.starts_with(b"HTTP/1.1 200")); + assert_eq!(response_header(capacity_headers, "cache-control"), "no-store"); + assert_eq!(capacity_body, b"decoy"); + + runtime.shutdown().await; + generation.stop_sessions().await; + generation.stop_background_tasks().await; +} + +#[tokio::test] +async fn deferred_close_preserves_the_first_reason_across_replacement_cancel() { + let capability = [42u8; 32]; + let generation = test_runtime_generation(1, runtime_config(capability, WebCarrier::Https)); + let active_runtime = Arc::new(ArcSwap::from(Arc::clone(&generation))); + let runtime = WebProcessRuntime::start(active_runtime); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(capability); + let bootstrap = bridge_bootstrap(&listener, &runtime, &encoded).await; + let created = create_session(&listener, &runtime, &bootstrap).await; + let (created_headers, _) = split_response(&created); + let bearer = response_header(created_headers, "x-session-token"); + let raw = base64::engine::general_purpose::URL_SAFE_NO_PAD + .decode(bearer) + .unwrap(); + let hash: crate::web::manager::TokenHash = Sha256::digest(raw).into(); + let session = runtime + .get_session(hash, "proxy.example.com") + .unwrap(); + let trace_session_id = session.trace_session_id(); + + assert!(session.begin_carrier_supersede()); + assert_eq!( + session.close(SessionCloseReason::ApiClose), + SessionCloseOutcome::Deferred + ); + session.cancel_carrier_supersede(); + session.wait_close_complete().await; + + assert!(matches!( + runtime.session_detail(trace_session_id), + SessionDetail::Gone { + reason: "api_close", + .. + } + )); + assert_eq!( + runtime + .telemetry() + .session_close_total(WebCarrier::Https, SessionCloseReason::ApiClose), + 1 + ); + assert_eq!( + runtime.telemetry().session_close_total( + WebCarrier::Https, + SessionCloseReason::CarrierSuperseded, + ), + 0 + ); + + runtime.shutdown().await; + generation.stop_sessions().await; + generation.stop_background_tasks().await; +} diff --git a/src/web/http/tests.rs b/src/web/http/tests.rs index 042231a..505b2e2 100644 --- a/src/web/http/tests.rs +++ b/src/web/http/tests.rs @@ -36,6 +36,9 @@ mod control_tests; // Reversible operator lifecycle coverage stays separate from terminal shutdown tests. #[path = "operator_lifecycle_tests.rs"] mod operator_lifecycle_tests; +// Positive-only recovery representation coverage remains isolated from ordinary root routing. +#[path = "recovery_tests.rs"] +mod recovery_tests; const TEST_CARRIER_DEADLINES_SECS: [u64; 4] = [3, 5, 8, 12]; @@ -276,8 +279,8 @@ async fn https_carrier_bootstraps_and_closes_one_session() { ); assert!( next_root_body - .windows(b"const negotiationEnabled=false".len()) - .any(|value| value == b"const negotiationEnabled=false") + .windows(b"let negotiationEnabled=false".len()) + .any(|value| value == b"let negotiationEnabled=false") ); let close = format!( @@ -482,7 +485,7 @@ async fn https_lanes_is_advertised_and_requires_canonical_lane_headers() { let root_response = request(&listener, &runtime, root).await; let (_, root_body) = split_response(&root_response); let root_body = std::str::from_utf8(root_body).unwrap(); - assert!(root_body.contains("const negotiationEnabled=false")); + assert!(root_body.contains("let negotiationEnabled=false")); let bootstrap = root_body .split_once("bootstrap=\"") .and_then(|(_, suffix)| suffix.split_once('"')) diff --git a/src/web/http/trace_tests.rs b/src/web/http/trace_tests.rs index ae8e2bf..166e61a 100644 --- a/src/web/http/trace_tests.rs +++ b/src/web/http/trace_tests.rs @@ -15,7 +15,7 @@ async fn enabled_debug_records_bridge_request_response_without_credentials() { let capability = [18u8; 32]; let mut config = runtime_config(capability, WebCarrier::Https); config.web.debug.enabled = true; - config.web.debug.body_capture = WebDebugBodyCapture::Prefix; + config.web.debug.body_capture = WebDebugBodyCapture::Full; config.web.debug.body_prefix_bytes = 4096; let generation = test_runtime_generation(1, config); let active_runtime = Arc::new(ArcSwap::from(Arc::clone(&generation))); diff --git a/src/web/http/websocket/driver.rs b/src/web/http/websocket/driver.rs index 701bc6e..8c9d6db 100644 --- a/src/web/http/websocket/driver.rs +++ b/src/web/http/websocket/driver.rs @@ -225,6 +225,9 @@ async fn run_multiplex( &payload, Instant::now(), ); + if !session.record_websocket_peer_activity() { + return Err(()); + } connection.mark_peer_activity(); next_ping = Instant::now() + liveness_interval; } @@ -247,6 +250,9 @@ async fn run_multiplex( &payload, started, ); + if !session.record_websocket_peer_activity() { + return Err(()); + } connection.mark_peer_activity(); next_ping = Instant::now() + liveness_interval; } diff --git a/src/web/http/websocket/driver/lane.rs b/src/web/http/websocket/driver/lane.rs index 092c19a..e635bab 100644 --- a/src/web/http/websocket/driver/lane.rs +++ b/src/web/http/websocket/driver/lane.rs @@ -133,6 +133,9 @@ pub(super) async fn run_lane( &payload, Instant::now(), ); + if !session.record_websocket_peer_activity() { + return Err(()); + } connection.mark_peer_activity(); next_ping = Instant::now() + liveness_interval; } @@ -155,6 +158,9 @@ pub(super) async fn run_lane( &payload, started, ); + if !session.record_websocket_peer_activity() { + return Err(()); + } connection.mark_peer_activity(); next_ping = Instant::now() + liveness_interval; } diff --git a/src/web/manager/credentials.rs b/src/web/manager/credentials.rs index 8c70837..e66ed62 100644 --- a/src/web/manager/credentials.rs +++ b/src/web/manager/credentials.rs @@ -24,7 +24,7 @@ impl WebProcessRuntime { client_ip: IpAddr, ) -> std::result::Result { let generation = self.active_generation(); - self.issue_bootstrap_inner(&generation, profile, client_ip, None, false) + self.issue_bootstrap_inner(&generation, profile, client_ip, None, false, None) } /// Issues one bootstrap against the generation that selected the bridge profile. @@ -35,7 +35,7 @@ impl WebProcessRuntime { profile: Arc, client_ip: IpAddr, ) -> std::result::Result { - self.issue_bootstrap_inner(generation, profile, client_ip, None, false) + self.issue_bootstrap_inner(generation, profile, client_ip, None, false, None) } /// Issues one bridge bootstrap with bounded non-secret request metadata. @@ -46,7 +46,7 @@ impl WebProcessRuntime { client_ip: IpAddr, user_agent: Option<&str>, ) -> std::result::Result { - self.issue_bootstrap_inner(generation, profile, client_ip, user_agent, false) + self.issue_bootstrap_inner(generation, profile, client_ip, user_agent, false, None) } /// Issues one recovery bootstrap with the same positive-only admission boundary. @@ -56,8 +56,16 @@ impl WebProcessRuntime { profile: Arc, client_ip: IpAddr, user_agent: Option<&str>, + predecessor_session_id: Option, ) -> std::result::Result { - self.issue_bootstrap_inner(generation, profile, client_ip, user_agent, true) + self.issue_bootstrap_inner( + generation, + profile, + client_ip, + user_agent, + true, + predecessor_session_id, + ) } fn issue_bootstrap_inner( @@ -67,6 +75,7 @@ impl WebProcessRuntime { client_ip: IpAddr, user_agent: Option<&str>, recovery: bool, + predecessor_session_id: Option, ) -> std::result::Result { let config = generation.config(); let profile = config @@ -170,6 +179,7 @@ impl WebProcessRuntime { session_ip_learning_eligible: false, used: false, recovery, + predecessor_session_id, }, ); *state.bootstraps_per_ip.entry(client_ip).or_insert(0) += 1; @@ -183,14 +193,32 @@ impl WebProcessRuntime { self.telemetry .record_bridge_recovery(WebBridgeRecoveryEvent::BootstrapIssued); } - self.trace.record_profile_lifecycle( - client_ip, - Some(trace_session_id), - &profile, - crate::web::trace::TraceLifecycleEvent::BridgeIssued, - None, - None, - ); + if recovery { + self.trace.record_lifecycle_with_context( + None, + Some(client_ip), + crate::web::trace::TraceIdentity::from_optional_profile( + Some(trace_session_id), + &profile, + ), + crate::web::trace::TraceLifecycleEvent::BridgeIssued, + None, + None, + crate::web::trace::TraceLifecycleContext { + peer_gap_ms: None, + predecessor_session_id, + }, + ); + } else { + self.trace.record_profile_lifecycle( + client_ip, + Some(trace_session_id), + &profile, + crate::web::trace::TraceLifecycleEvent::BridgeIssued, + None, + None, + ); + } Ok(BootstrapResult { token, trace_session_id, @@ -227,13 +255,28 @@ impl WebProcessRuntime { hash: TokenHash, host: &str, ) -> std::result::Result, ManagerError> { - self.state - .lock() + let state = self.state.lock(); + if let Some(session) = state .sessions .get(&hash) .cloned() .filter(|session| session.matches_host(host)) - .ok_or(ManagerError::Authentication) + { + return Ok(session); + } + let retired_carrier = state + .closed_tokens + .get(&hash) + .filter(|closed| closed.host == host) + .map(|closed| closed.carrier); + drop(state); + if let Some(carrier) = retired_carrier { + self.telemetry.record_session_observation( + carrier, + crate::web::telemetry::WebSessionLifecycleObservation::RequestAfterClose, + ); + } + Err(ManagerError::Authentication) } /// Resolves a current bearer only when it belongs to the recovering profile. diff --git a/src/web/manager/lifecycle.rs b/src/web/manager/lifecycle.rs index 2350ea0..9a4aa76 100644 --- a/src/web/manager/lifecycle.rs +++ b/src/web/manager/lifecycle.rs @@ -82,9 +82,7 @@ impl WebProcessRuntime { &mut state, hash, profile_host, - session.trace_session_id(), session.carrier(), - reason, closed_token_lifetime, self.limits.max_sessions_global.saturating_mul(16), ); diff --git a/src/web/manager/session_creation.rs b/src/web/manager/session_creation.rs index cb76114..7e473d6 100644 --- a/src/web/manager/session_creation.rs +++ b/src/web/manager/session_creation.rs @@ -339,7 +339,14 @@ impl WebProcessRuntime { state.sessions.insert(session_hash, Arc::clone(&session)); *state.sessions_per_ip.entry(client_ip).or_insert(0) += 1; *state.sessions_per_profile.entry(profile_key).or_insert(0) += 1; - let (issuance_ip, candidate_count, user_agent, user_agent_id, recovery) = { + let ( + issuance_ip, + candidate_count, + user_agent, + user_agent_id, + recovery, + predecessor_session_id, + ) = { let entry = state .bootstraps .get_mut(&bootstrap_hash) @@ -367,6 +374,7 @@ impl WebProcessRuntime { entry.user_agent.clone(), entry.user_agent_id, entry.recovery, + entry.predecessor_session_id, ) }; decrement_map(&mut state.bootstraps_per_ip, &issuance_ip); @@ -422,13 +430,17 @@ impl WebProcessRuntime { scores, None, ); - self.trace.record_lifecycle( + self.trace.record_lifecycle_with_context( None, Some(client_ip), identity, TraceLifecycleEvent::SessionCreated, None, None, + crate::web::trace::TraceLifecycleContext { + peer_gap_ms: None, + predecessor_session_id, + }, ); Ok(result) } diff --git a/src/web/manager/session_creation/replacement.rs b/src/web/manager/session_creation/replacement.rs index bd8d46c..06a6519 100644 --- a/src/web/manager/session_creation/replacement.rs +++ b/src/web/manager/session_creation/replacement.rs @@ -98,9 +98,7 @@ impl WebProcessRuntime { &mut state, old_hash, &replacement.profile.host, - replacement.old_session.trace_session_id(), replacement.old_session.carrier(), - crate::web::session::SessionCloseReason::CarrierSuperseded, Duration::from_secs(replacement.old_session.timeouts().bootstrap_lifetime_secs), self.limits.max_sessions_global.saturating_mul(16), ); @@ -126,6 +124,9 @@ impl WebProcessRuntime { self.telemetry.record_bridge_recovery( crate::web::telemetry::WebBridgeRecoveryEvent::SessionCreated, ); + self.telemetry.record_bridge_recovery( + crate::web::telemetry::WebBridgeRecoveryEvent::ClosedBeforeCommit, + ); } self.telemetry.record_session_closed( replacement.old_session.carrier(), diff --git a/src/web/manager/state.rs b/src/web/manager/state.rs index 9b93a15..67e343d 100644 --- a/src/web/manager/state.rs +++ b/src/web/manager/state.rs @@ -88,6 +88,8 @@ pub(super) struct Bootstrap { pub(super) used: bool, /// Whether this credential was issued by the post-commit recovery representation. pub(super) recovery: bool, + /// Previous logical session authenticated by the recovery request, if current. + pub(super) predecessor_session_id: Option, } /// Bounded replay marker for one explicitly or naturally closed session token. @@ -96,12 +98,8 @@ pub(super) struct ClosedToken { pub(super) expires_at: Instant, /// Canonical host that owned the session. pub(super) host: String, - /// Non-secret logical trace owner retained for exact late-request diagnostics. - pub(super) trace_session_id: u64, /// Carrier that owned the retired bearer. pub(super) carrier: WebCarrier, - /// First-writer terminal cause for the retired bearer. - pub(super) reason: crate::web::session::SessionCloseReason, } /// Current logical-session owner stored without exposing bearer credentials. @@ -363,9 +361,7 @@ pub(super) fn remember_closed_token_locked( state: &mut ManagerState, hash: TokenHash, host: &str, - trace_session_id: u64, carrier: WebCarrier, - reason: crate::web::session::SessionCloseReason, lifetime: Duration, capacity: usize, ) { @@ -374,9 +370,7 @@ pub(super) fn remember_closed_token_locked( ClosedToken { expires_at: Instant::now() + lifetime, host: host.to_string(), - trace_session_id, carrier, - reason, }, ); while state.closed_tokens.len() > capacity { diff --git a/src/web/session/activity.rs b/src/web/session/activity.rs index 2e967ed..733b447 100644 --- a/src/web/session/activity.rs +++ b/src/web/session/activity.rs @@ -1,11 +1,76 @@ use std::time::{Duration, Instant}; +use super::{SessionState, WebSession}; +use crate::web::telemetry::WebSessionLifecycleObservation; + /// Session-local activity clocks with distinct lease and diagnostic authority. pub(super) struct SessionActivity { last_peer: Instant, last_progress: Instant, } +impl WebSession { + /// Refreshes the authenticated peer lease and records only threshold-crossing gaps. + pub(super) fn touch_peer_locked( + &self, + state: &mut SessionState, + now: Instant, + observation: WebSessionLifecycleObservation, + ) { + let gap = state.activity.touch_peer(now); + if gap >= Duration::from_secs(self.timeouts.reconnect_grace_secs) + && let Some(manager) = self.manager.upgrade() + { + manager + .telemetry() + .record_session_observation(self.carrier(), observation); + } + } + + /// Records one valid WebSocket control message as authenticated peer activity. + pub(crate) fn record_websocket_peer_activity(&self) -> bool { + let mut state = self.state.lock(); + if state.closed { + return false; + } + self.touch_peer_locked( + &mut state, + Instant::now(), + WebSessionLifecycleObservation::WebSocketActivityAfterGap, + ); + true + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn progress_does_not_extend_the_authenticated_peer_lease() { + let started = Instant::now(); + let mut activity = SessionActivity::new(started); + activity.touch_progress(started + Duration::from_secs(4)); + + assert_eq!( + activity.peer_idle(started + Duration::from_secs(9)), + Duration::from_secs(9) + ); + assert_eq!( + activity.progress_idle(started + Duration::from_secs(9)), + Duration::from_secs(5) + ); + assert_eq!( + activity.touch_peer(started + Duration::from_secs(9)), + Duration::from_secs(9) + ); + assert_eq!( + activity.peer_idle(started + Duration::from_secs(10)), + Duration::from_secs(1) + ); + } +} + impl SessionActivity { /// Starts both activity clocks at the same session creation instant. pub(super) fn new(now: Instant) -> Self { diff --git a/src/web/session/downlink.rs b/src/web/session/downlink.rs index febecc2..2bd06f8 100644 --- a/src/web/session/downlink.rs +++ b/src/web/session/downlink.rs @@ -10,6 +10,7 @@ use super::{ }; use crate::web::frame::{self, FrameType}; use crate::web::manager::ManagerError; +use crate::web::telemetry::WebSessionLifecycleObservation; impl WebSession { /// Polls pending downlink frames with cursor replay and newest-poll-wins semantics. @@ -22,14 +23,19 @@ impl WebSession { if state.closed { return Err(ManagerError::Closed); } - state.activity.touch_peer(Instant::now()); if let Some(unacked) = &state.unacked { if cursor == unacked.base_cursor { - return Ok(PollResult { + let result = PollResult { body: unacked.body.clone(), next_cursor: unacked.next_cursor, lane_closed: false, - }); + }; + self.touch_peer_locked( + &mut state, + Instant::now(), + WebSessionLifecycleObservation::HttpActivityAfterGap, + ); + return Ok(result); } if cursor != unacked.next_cursor { drop(state); @@ -53,6 +59,11 @@ impl WebSession { return Err(ManagerError::Protocol); }; state.down_epoch = epoch; + self.touch_peer_locked( + &mut state, + Instant::now(), + WebSessionLifecycleObservation::HttpActivityAfterGap, + ); let healthy = self.carrier_health_ready_locked(&mut state, Instant::now()); (state.down_epoch, healthy) }; diff --git a/src/web/session/lane_uplink.rs b/src/web/session/lane_uplink.rs index a13b7c5..092f007 100644 --- a/src/web/session/lane_uplink.rs +++ b/src/web/session/lane_uplink.rs @@ -9,6 +9,7 @@ use super::{PendingClass, SessionCloseReason, WebSession, insert_carrier_lane}; use crate::config::WebCarrier; use crate::web::frame::{self, Frame, FrameType}; use crate::web::manager::{ManagerError, TokenHash}; +use crate::web::telemetry::WebSessionLifecycleObservation; impl WebSession { /// Applies one exactly-once uplink batch to an independent HTTPS lane. @@ -46,7 +47,6 @@ impl WebSession { return Err(ManagerError::Closed); } self.ensure_carrier_active_locked(&state)?; - state.activity.touch_peer(Instant::now()); let new_lane = !state.carrier_lanes.contains_key(&lane_id); if new_lane { if lane_id != 0 @@ -55,6 +55,11 @@ impl WebSession { .is_some_and(|value| value.frame_type != FrameType::Open) && only_late_frames(&frames) { + self.touch_peer_locked( + &mut state, + Instant::now(), + WebSessionLifecycleObservation::HttpActivityAfterGap, + ); return if self.automatic_carrier && state.negotiation_phase != super::SessionNegotiationPhase::Committed { @@ -89,6 +94,11 @@ impl WebSession { }); if sequence == last_sequence && sequence != 0 { return if bool::from(last_digest.ct_eq(&digest)) { + self.touch_peer_locked( + &mut state, + Instant::now(), + WebSessionLifecycleObservation::HttpActivityAfterGap, + ); Ok(sequence) } else { drop(state); @@ -109,6 +119,11 @@ impl WebSession { self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); } + self.touch_peer_locked( + &mut state, + Instant::now(), + WebSessionLifecycleObservation::HttpActivityAfterGap, + ); let (reserve_bytes, reserve_items) = inbound_reservation(&state, &frames); if !self.reserve_locked( &mut state, diff --git a/src/web/session/lanes.rs b/src/web/session/lanes.rs index 2fe0c10..ae65b56 100644 --- a/src/web/session/lanes.rs +++ b/src/web/session/lanes.rs @@ -11,6 +11,7 @@ use super::{ }; use crate::web::frame::{self, FrameType}; use crate::web::manager::ManagerError; +use crate::web::telemetry::WebSessionLifecycleObservation; impl WebSession { /// Polls one lane with independent cursor replay and newest-poll-wins semantics. @@ -65,8 +66,7 @@ impl WebSession { if state.closed { return Err(ManagerError::Closed); } - state.activity.touch_peer(Instant::now()); - let acknowledged = { + let (acknowledged, replay) = { let Some(lane) = state.carrier_lanes.get_mut(&lane_id) else { return Ok(PollResult { body: Bytes::new(), @@ -83,27 +83,40 @@ impl WebSession { } if let Some(unacked) = &lane.unacked { if cursor == unacked.base_cursor { - return Ok(PollResult { - body: unacked.body.clone(), - next_cursor: unacked.next_cursor, - lane_closed: false, - }); - } - if cursor != unacked.next_cursor { + ( + None, + Some(PollResult { + body: unacked.body.clone(), + next_cursor: unacked.next_cursor, + lane_closed: false, + }), + ) + } else if cursor != unacked.next_cursor { drop(state); self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); + } else { + (lane.unacked.take(), None) } - lane.unacked.take() } else { if cursor != lane.down_cursor { drop(state); self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); } - None + (None, None) } }; + if let Some(result) = replay { + if expected_instance.is_none() { + self.touch_peer_locked( + &mut state, + Instant::now(), + WebSessionLifecycleObservation::HttpActivityAfterGap, + ); + } + return Ok(result); + } if let Some(batch) = acknowledged { if let Some(lane) = state.carrier_lanes.get_mut(&lane_id) { lane.pending_bytes = lane.pending_bytes.saturating_sub(batch.data_bytes); @@ -139,6 +152,13 @@ impl WebSession { lane.down_epoch = epoch; let instance = lane.instance; let notify = Arc::clone(&lane.notify); + if expected_instance.is_none() { + self.touch_peer_locked( + &mut state, + Instant::now(), + WebSessionLifecycleObservation::HttpActivityAfterGap, + ); + } let healthy = self.carrier_health_ready_locked(&mut state, Instant::now()); (instance, epoch, notify, healthy) }; diff --git a/src/web/session/lifecycle.rs b/src/web/session/lifecycle.rs index 0f11b13..aaa2cd6 100644 --- a/src/web/session/lifecycle.rs +++ b/src/web/session/lifecycle.rs @@ -98,6 +98,7 @@ struct ReleasedQueues { control_items: usize, closed_before_health: bool, reason: SessionCloseReason, + peer_gap: Duration, } /// Deferred queue release after manager publication linearizes a supersede. @@ -248,6 +249,7 @@ impl WebSession { state: &mut super::SessionState, reason: SessionCloseReason, ) -> ReleasedQueues { + let peer_gap = state.activity.peer_idle(Instant::now()); let closed_before_health = self.automatic_carrier && state.negotiation_phase == SessionNegotiationPhase::Committed && self.reject_carrier_health_on_close(); @@ -304,6 +306,7 @@ impl WebSession { control_items, closed_before_health, reason, + peer_gap, } } @@ -337,11 +340,25 @@ impl WebSession { ); } if !self.finished.swap(true, Ordering::AcqRel) { - self.trace_lifecycle( - crate::web::trace::TraceLifecycleEvent::SessionClosed, - None, - Some(released.reason.as_str()), - ); + if let Some(manager) = &manager { + manager.trace().record_lifecycle_with_context( + None, + Some(self.client_ip), + self.trace_identity(), + crate::web::trace::TraceLifecycleEvent::SessionClosed, + None, + Some(released.reason.as_str()), + crate::web::trace::TraceLifecycleContext { + peer_gap_ms: Some( + released + .peer_gap + .as_millis() + .min(u128::from(u64::MAX)) as u64, + ), + predecessor_session_id: None, + }, + ); + } if released.reason != SessionCloseReason::CarrierSuperseded && let Some(manager) = manager { diff --git a/src/web/session/negotiation.rs b/src/web/session/negotiation.rs index 68e7bf7..8ca9a15 100644 --- a/src/web/session/negotiation.rs +++ b/src/web/session/negotiation.rs @@ -298,6 +298,7 @@ mod tests { WebCarrier, WebLimitsConfig, WebRuntimeProfile, WebSecretMode, WebTimeoutsConfig, }; use crate::web::manager::{CarrierClientClass, WebProcessRuntime}; + use crate::web::session::{SessionCloseOutcome, SessionCloseReason}; fn session(carrier: WebCarrier, deadline: Instant) -> Arc { let profile = Arc::new(WebRuntimeProfile { @@ -460,7 +461,7 @@ mod tests { let close = std::thread::spawn(move || { close_barrier.wait(); std::thread::yield_now(); - close_session.close(super::SessionCloseReason::ApiClose); + close_session.close(SessionCloseReason::ApiClose); }); barrier.wait(); health.join().unwrap(); @@ -473,8 +474,8 @@ mod tests { )); assert!(!session.publish_carrier_health()); assert_eq!( - session.close(super::SessionCloseReason::ApiClose), - super::SessionCloseOutcome::AlreadyClosing + session.close(SessionCloseReason::ApiClose), + SessionCloseOutcome::AlreadyClosing ); } } diff --git a/src/web/session/uplink.rs b/src/web/session/uplink.rs index e7b927b..7c68db5 100644 --- a/src/web/session/uplink.rs +++ b/src/web/session/uplink.rs @@ -14,6 +14,7 @@ use super::{ }; use crate::web::frame::{self, Frame, FrameType}; use crate::web::manager::{ManagerError, TokenHash}; +use crate::web::telemetry::WebSessionLifecycleObservation; #[derive(Clone, Copy, Default)] pub(super) struct AppliedProgress { @@ -34,7 +35,11 @@ impl WebSession { sequence: u64, body: &[u8], ) -> Result { - let (acknowledged, progressed) = self.process_up_inner(sequence, body)?; + let (acknowledged, progressed) = self.process_up_inner( + sequence, + body, + WebSessionLifecycleObservation::HttpActivityAfterGap, + )?; if self.automatic_carrier && !progressed && !self.is_carrier_committed() { return Err(ManagerError::Backpressure); } @@ -47,14 +52,19 @@ impl WebSession { sequence: u64, body: &[u8], ) -> Result { - self.process_up_inner(sequence, body) - .map(|(_, progress)| progress) + self.process_up_inner( + sequence, + body, + WebSessionLifecycleObservation::WebSocketActivityAfterGap, + ) + .map(|(_, progress)| progress) } fn process_up_inner( self: &Arc, sequence: u64, body: &[u8], + observation: WebSessionLifecycleObservation, ) -> Result<(u64, bool), ManagerError> { if !self.carrier().is_multiplexed() { return Err(ManagerError::Protocol); @@ -92,9 +102,9 @@ impl WebSession { return Err(ManagerError::Closed); } self.ensure_carrier_active_locked(&state)?; - state.activity.touch_peer(Instant::now()); if sequence == state.last_up_sequence && sequence != 0 { return if bool::from(state.last_up_digest.ct_eq(&digest)) { + self.touch_peer_locked(&mut state, Instant::now(), observation); Ok((sequence, false)) } else { drop(state); @@ -112,6 +122,7 @@ impl WebSession { self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); } + self.touch_peer_locked(&mut state, Instant::now(), observation); let (reserve_bytes, reserve_items) = inbound_reservation(&state, &frames); if !self.reserve_locked( &mut state, diff --git a/src/web/session/websocket.rs b/src/web/session/websocket.rs index dc803e7..c0d3c72 100644 --- a/src/web/session/websocket.rs +++ b/src/web/session/websocket.rs @@ -391,7 +391,11 @@ impl WebSession { lane.last_up_digest = digest; } } - state.activity.touch_peer(Instant::now()); + self.touch_peer_locked( + &mut state, + Instant::now(), + crate::web::telemetry::WebSessionLifecycleObservation::WebSocketActivityAfterGap, + ); if applied { (committed, healthy) = self.record_uplink_progress_locked(&mut state, progress); } diff --git a/src/web/trace/mod.rs b/src/web/trace/mod.rs index 205a2f5..02d631e 100644 --- a/src/web/trace/mod.rs +++ b/src/web/trace/mod.rs @@ -15,6 +15,6 @@ pub(crate) use store::{ }; pub(crate) use types::{ TraceBodySnapshot, TraceBodyState, TraceDirection, TraceFrame, TraceHeader, TraceIdentity, - TraceLifecycleEvent, TraceLifecycleRecord, TraceRecord, TraceRecordKind, TraceRoute, - TraceWebSocketContext, + TraceLifecycleContext, TraceLifecycleEvent, TraceLifecycleRecord, TraceRecord, TraceRecordKind, + TraceRoute, TraceWebSocketContext, }; diff --git a/src/web/trace/store.rs b/src/web/trace/store.rs index e50044f..f56d5ed 100644 --- a/src/web/trace/store.rs +++ b/src/web/trace/store.rs @@ -9,8 +9,8 @@ use tokio::sync::{OwnedSemaphorePermit, Semaphore}; use super::exchange::HttpTraceExchange; use super::types::{ - TraceCarrierDetail, TraceIdentity, TraceLifecycleEvent, TraceLifecycleRecord, TraceRecord, - TraceRecordKind, + TraceCarrierDetail, TraceIdentity, TraceLifecycleContext, TraceLifecycleEvent, + TraceLifecycleRecord, TraceRecord, TraceRecordKind, }; use crate::config::{WebDebugConfig, WebLimitsConfig}; @@ -227,6 +227,31 @@ impl WebTraceStore { stream_id, reason, None, + TraceLifecycleContext::default(), + ); + } + + /// Records one lifecycle event with bounded non-secret recovery context. + #[allow(clippy::too_many_arguments)] + pub(crate) fn record_lifecycle_with_context( + &self, + peer_ip: Option, + effective_ip: Option, + identity: TraceIdentity, + event: TraceLifecycleEvent, + stream_id: Option, + reason: Option<&'static str>, + context: TraceLifecycleContext, + ) { + self.record_lifecycle_detail( + peer_ip, + effective_ip, + identity, + event, + stream_id, + reason, + None, + context, ); } @@ -256,6 +281,7 @@ impl WebTraceStore { attempt, scores, }), + TraceLifecycleContext::default(), ); } @@ -269,6 +295,7 @@ impl WebTraceStore { stream_id: Option, reason: Option<&'static str>, carrier: Option, + context: TraceLifecycleContext, ) { if !self.enabled.load(Ordering::Acquire) { return; @@ -301,6 +328,8 @@ impl WebTraceStore { stream_id, reason, carrier, + peer_gap_ms: context.peer_gap_ms, + predecessor_session_id: context.predecessor_session_id, }), }; if !self.try_commit(record, reservation, epoch) { diff --git a/src/web/trace/store/tests.rs b/src/web/trace/store/tests.rs index 79f0434..f179911 100644 --- a/src/web/trace/store/tests.rs +++ b/src/web/trace/store/tests.rs @@ -109,3 +109,28 @@ fn explicit_clear_fences_inflight_commits_and_preserves_snapshot_leases() { drop(snapshot); assert_eq!(store.status().used_bytes, 0); } + +#[test] +fn lifecycle_context_retains_only_bounded_non_secret_correlations() { + let store = store(2, 4 * BASE_RECORD_RESERVATION); + store.record_lifecycle_with_context( + None, + Some("192.0.2.40".parse().unwrap()), + TraceIdentity::default(), + TraceLifecycleEvent::SessionClosed, + None, + Some("bridge_recovery"), + TraceLifecycleContext { + peer_gap_ms: Some(125_000), + predecessor_session_id: Some(41), + }, + ); + + let snapshot = store.snapshot_matching(|_| true); + let TraceRecordKind::Lifecycle(event) = &snapshot[0].record.kind else { + panic!("expected lifecycle record"); + }; + assert_eq!(event.reason, Some("bridge_recovery")); + assert_eq!(event.peer_gap_ms, Some(125_000)); + assert_eq!(event.predecessor_session_id, Some(41)); +} diff --git a/src/web/trace/types.rs b/src/web/trace/types.rs index 175982a..746ef78 100644 --- a/src/web/trace/types.rs +++ b/src/web/trace/types.rs @@ -313,6 +313,15 @@ pub(crate) struct TraceCarrierDetail { pub(crate) scores: [i16; 4], } +/// Optional non-secret context for one lifecycle transition. +#[derive(Clone, Copy, Debug, Default)] +pub(crate) struct TraceLifecycleContext { + /// Monotonic peer inactivity preceding the transition. + pub(crate) peer_gap_ms: Option, + /// Previous logical session replaced by this transition. + pub(crate) predecessor_session_id: Option, +} + /// One typed WEB lifecycle observation. #[derive(Debug)] pub(crate) struct TraceLifecycleRecord { @@ -324,6 +333,10 @@ pub(crate) struct TraceLifecycleRecord { pub(crate) reason: Option<&'static str>, /// Carrier negotiation detail when this is a carrier lifecycle event. pub(crate) carrier: Option, + /// Optional monotonic peer inactivity preceding the transition. + pub(crate) peer_gap_ms: Option, + /// Optional logical predecessor for a recovered session incarnation. + pub(crate) predecessor_session_id: Option, } /// Trace record payload variant.