diff --git a/src/api/web_runtime.rs b/src/api/web_runtime.rs index b825f49..4aa5e6e 100644 --- a/src/api/web_runtime.rs +++ b/src/api/web_runtime.rs @@ -21,6 +21,7 @@ mod request; mod observability; use observability::{ WebCapacityStatus, WebCarrierNegotiationStatus, WebDecoyUpstreamStatus, WebIngressStatus, + WebLifecycleCountersStatus, }; use request::{ CloseRequest, DrainRequest, RuntimeInstanceRequest, parse_session_query, parse_session_ref, @@ -92,12 +93,20 @@ pub(super) async fn handle( let trace_session_id = parse_session_ref(&runtime, session_ref)?; match runtime.session_detail(trace_session_id) { SessionDetail::Active(row) => Ok(success_response(StatusCode::OK, row, revision)), - SessionDetail::Gone { attempt } => Ok(success_response( + SessionDetail::Gone { + attempt, + carrier, + reason, + closed_age_ms, + } => Ok(success_response( StatusCode::GONE, GoneSessionData { session_ref: session_ref.to_string(), state: "closed", attempt, + carrier, + reason, + closed_age_ms, }, revision, )), @@ -272,6 +281,7 @@ struct WebStatusData { capacity: WebCapacityStatus, decoy_upstream: WebDecoyUpstreamStatus, carrier_negotiation: WebCarrierNegotiationStatus, + lifecycle_counters: WebLifecycleCountersStatus, #[serde(skip_serializing_if = "Option::is_none")] operator_lifecycle: Option, #[serde(skip_serializing_if = "Option::is_none")] @@ -306,6 +316,7 @@ impl WebStatusData { let capacity = WebCapacityStatus::new(&publication, runtime, config); let decoy_upstream = WebDecoyUpstreamStatus::new(&publication); let carrier_negotiation = WebCarrierNegotiationStatus::new(&publication); + let lifecycle_counters = WebLifecycleCountersStatus::new(&publication, config); Self { lifecycle: publication.lifecycle.as_str(), lifecycle_epoch: publication.epoch, @@ -322,6 +333,7 @@ impl WebStatusData { capacity, decoy_upstream, carrier_negotiation, + lifecycle_counters, operator_lifecycle, runtime: runtime.map(WebProcessRuntime::try_status), } @@ -333,6 +345,9 @@ struct GoneSessionData { session_ref: String, state: &'static str, attempt: u8, + carrier: crate::config::WebCarrier, + reason: &'static str, + closed_age_ms: u64, } #[derive(Serialize)] diff --git a/src/api/web_runtime/observability.rs b/src/api/web_runtime/observability.rs index a545ca7..a1d4fc7 100644 --- a/src/api/web_runtime/observability.rs +++ b/src/api/web_runtime/observability.rs @@ -5,7 +5,8 @@ use crate::web::control::{WebRuntimeLifecycle, WebRuntimePublication}; use crate::web::manager::{WebCapacityResourceStatus, WebCapacitySnapshot, WebProcessRuntime}; use crate::web::telemetry::{WebOutcomeCounter, WebRejectionCounter}; use crate::web::telemetry::{ - WebCarrierFailureCounter, WebCarrierLearningCounter, WebCarrierSelectionCounter, + WebBridgeRecoveryCounter, WebCarrierFailureCounter, WebCarrierLearningCounter, + WebCarrierSelectionCounter, WebSessionCloseCounter, WebSessionLifecycleObservationCounter, }; /// Private WEB ingress state owned by this Telemt process. @@ -139,6 +140,27 @@ impl WebCarrierNegotiationStatus { } } +/// Fixed-cardinality process-lifetime WEB lifecycle counters. +#[derive(Serialize)] +pub(super) struct WebLifecycleCountersStatus { + bridge_recovery_secs: u64, + session_closures: Vec, + session_observations: Vec, + bridge_recovery_events: Vec, +} + +impl WebLifecycleCountersStatus { + /// Builds a complete counter set from process-owned telemetry. + pub(super) fn new(publication: &WebRuntimePublication, config: &ProxyConfig) -> Self { + Self { + bridge_recovery_secs: config.web.timeouts.bridge_recovery_secs, + session_closures: publication.telemetry.session_close_counters(), + session_observations: publication.telemetry.session_observation_counters(), + bridge_recovery_events: publication.telemetry.bridge_recovery_counters(), + } + } +} + #[cfg(test)] mod tests { use crate::config::ProxyConfig; @@ -167,6 +189,11 @@ mod tests { let decoy = serde_json::to_value(super::WebDecoyUpstreamStatus::new(&publication)).unwrap(); let carrier = serde_json::to_value(super::WebCarrierNegotiationStatus::new(&publication)).unwrap(); + let lifecycle = serde_json::to_value(super::WebLifecycleCountersStatus::new( + &publication, + &config, + )) + .unwrap(); assert_eq!( capacity["rejections"].as_array().unwrap().len(), @@ -200,5 +227,22 @@ mod tests { crate::config::WebCarrier::ALL.len() * crate::web::telemetry::WebCarrierLearningOutcome::ALL.len() ); + assert_eq!( + lifecycle["session_closures"].as_array().unwrap().len(), + crate::config::WebCarrier::ALL.len() + * crate::web::session::SessionCloseReason::ALL.len() + ); + assert_eq!( + lifecycle["session_observations"].as_array().unwrap().len(), + crate::config::WebCarrier::ALL.len() + * crate::web::telemetry::WebSessionLifecycleObservation::ALL.len() + ); + assert_eq!( + lifecycle["bridge_recovery_events"] + .as_array() + .unwrap() + .len(), + crate::web::telemetry::WebBridgeRecoveryEvent::ALL.len() + ); } } diff --git a/src/config/load/strict_keys.rs b/src/config/load/strict_keys.rs index ffb3d4b..ed3756f 100644 --- a/src/config/load/strict_keys.rs +++ b/src/config/load/strict_keys.rs @@ -344,6 +344,7 @@ const WEB_TIMEOUTS_CONFIG_KEYS: &[&str] = &[ "long_poll_secs", "bridge_request_secs", "bridge_retry_secs", + "bridge_recovery_secs", "carrier_probe_coalesce_ms", "lane_open_wait_secs", "carrier_health_secs", diff --git a/src/config/load/validate_web/timeouts.rs b/src/config/load/validate_web/timeouts.rs index 66fa6e9..a301a32 100644 --- a/src/config/load/validate_web/timeouts.rs +++ b/src/config/load/validate_web/timeouts.rs @@ -42,6 +42,9 @@ pub(super) fn validate(timeouts: &WebTimeoutsConfig) -> Result<()> { if !(1..=300).contains(&timeouts.bridge_retry_secs) { return config_error("web.timeouts.bridge_retry_secs must be within [1, 300]"); } + if !(1..=60).contains(&timeouts.bridge_recovery_secs) { + return config_error("web.timeouts.bridge_recovery_secs must be within [1, 60]"); + } if timeouts.bridge_request_secs > timeouts.bridge_retry_secs { return config_error("web.timeouts.bridge_request_secs must not exceed bridge_retry_secs"); } diff --git a/src/config/tests/load_basic_tests/web_tests.rs b/src/config/tests/load_basic_tests/web_tests.rs index e8ed6b4..5c1246a 100644 --- a/src/config/tests/load_basic_tests/web_tests.rs +++ b/src/config/tests/load_basic_tests/web_tests.rs @@ -194,7 +194,7 @@ fn web_carriers_reject_true_empty_and_duplicates() { fn web_carrier_and_bridge_deadlines_are_configurable() { let configured = WEB_CONFIG.replace( "[[web.vhosts]]", - "[web.timeouts]\ncarrier_negotiation_deadlines_secs = [1, 2, 4, 9]\ncarrier_learning_secs = 30\nbridge_request_secs = 7\nbridge_retry_secs = 41\ncarrier_probe_coalesce_ms = 4\n\n[[web.vhosts]]", + "[web.timeouts]\ncarrier_negotiation_deadlines_secs = [1, 2, 4, 9]\ncarrier_learning_secs = 30\nbridge_request_secs = 7\nbridge_retry_secs = 41\nbridge_recovery_secs = 13\ncarrier_probe_coalesce_ms = 4\n\n[[web.vhosts]]", ); let config = load_config_from_temp_toml(&configured); assert_eq!( @@ -204,6 +204,7 @@ fn web_carrier_and_bridge_deadlines_are_configurable() { assert_eq!(config.web.timeouts.carrier_learning_secs, 30); assert_eq!(config.web.timeouts.bridge_request_secs, 7); assert_eq!(config.web.timeouts.bridge_retry_secs, 41); + assert_eq!(config.web.timeouts.bridge_recovery_secs, 13); assert_eq!(config.web.timeouts.carrier_probe_coalesce_ms, 4); } @@ -211,13 +212,14 @@ fn web_carrier_and_bridge_deadlines_are_configurable() { fn web_bridge_deadlines_are_known_in_strict_mode() { let configured = WEB_CONFIG.replace( "[[web.vhosts]]", - "[web.timeouts]\nbridge_request_secs = 7\nbridge_retry_secs = 41\ncarrier_probe_coalesce_ms = 4\n\n[[web.vhosts]]", + "[web.timeouts]\nbridge_request_secs = 7\nbridge_retry_secs = 41\nbridge_recovery_secs = 13\ncarrier_probe_coalesce_ms = 4\n\n[[web.vhosts]]", ); let configured = format!("[general]\nconfig_strict = true\n{configured}"); let config = load_config_from_temp_toml(&configured); assert_eq!(config.web.timeouts.bridge_request_secs, 7); assert_eq!(config.web.timeouts.bridge_retry_secs, 41); + assert_eq!(config.web.timeouts.bridge_recovery_secs, 13); assert_eq!(config.web.timeouts.carrier_probe_coalesce_ms, 4); } @@ -228,6 +230,8 @@ fn web_bridge_deadlines_are_bounded_and_ordered() { ("bridge_request_secs", "61"), ("bridge_retry_secs", "0"), ("bridge_retry_secs", "301"), + ("bridge_recovery_secs", "0"), + ("bridge_recovery_secs", "61"), ("carrier_probe_coalesce_ms", "11"), ] { let invalid = WEB_CONFIG.replace( diff --git a/src/config/types/web.rs b/src/config/types/web.rs index ade2523..9f917fc 100644 --- a/src/config/types/web.rs +++ b/src/config/types/web.rs @@ -301,6 +301,9 @@ pub struct WebTimeoutsConfig { /// Absolute generated-bridge budget for one retryable HTTP operation. #[serde(default = "default_web_bridge_retry_secs")] pub bridge_retry_secs: u64, + /// Absolute post-commit budget for one surviving bridge recovery epoch. + #[serde(default = "default_web_bridge_recovery_secs")] + pub bridge_recovery_secs: u64, /// Optional delay for coalescing the first OPEN with immediate DATA. #[serde(default = "default_web_carrier_probe_coalesce_ms")] pub carrier_probe_coalesce_ms: u64, @@ -361,6 +364,7 @@ impl Default for WebTimeoutsConfig { long_poll_secs: default_web_long_poll_timeout_secs(), bridge_request_secs: default_web_bridge_request_secs(), bridge_retry_secs: default_web_bridge_retry_secs(), + bridge_recovery_secs: default_web_bridge_recovery_secs(), carrier_probe_coalesce_ms: default_web_carrier_probe_coalesce_ms(), lane_open_wait_secs: default_web_lane_open_wait_secs(), carrier_health_secs: default_web_carrier_health_secs(), diff --git a/src/config/types/web/defaults.rs b/src/config/types/web/defaults.rs index 1616e98..10e91e6 100644 --- a/src/config/types/web/defaults.rs +++ b/src/config/types/web/defaults.rs @@ -88,6 +88,7 @@ u64_default!(default_web_stream_first_byte_secs, 30); u64_default!(default_web_long_poll_timeout_secs, 25); u64_default!(default_web_bridge_request_secs, 10); u64_default!(default_web_bridge_retry_secs, 90); +u64_default!(default_web_bridge_recovery_secs, 15); u64_default!(default_web_carrier_probe_coalesce_ms, 0); u64_default!(default_web_lane_open_wait_secs, 2); u64_default!(default_web_carrier_health_secs, 30); diff --git a/src/metrics/web.rs b/src/metrics/web.rs index 6cdbaa7..57802fe 100644 --- a/src/metrics/web.rs +++ b/src/metrics/web.rs @@ -11,6 +11,9 @@ use crate::web::telemetry::{ WebDecoyUpstreamOutcome, WebHttpConnectionOverloadOutcome, WebRejectionReason, }; +// Session lifecycle and aggregate families stay isolated from capacity rendering. +mod lifecycle; + /// Renders fixed-cardinality process-owned WEB observability families. pub(super) fn render(out: &mut String, publication: &WebRuntimePublication, config: &ProxyConfig) { let runtime = publication.runtime.upgrade(); @@ -197,7 +200,7 @@ pub(super) fn render(out: &mut String, publication: &WebRuntimePublication, conf } render_carrier_negotiation(out, publication, runtime.as_deref(), config); - render_aggregate_totals(out, publication); + lifecycle::render(out, publication, config); } fn render_carrier_negotiation( @@ -413,55 +416,6 @@ fn render_capacity(out: &mut String, snapshot: &crate::web::manager::WebCapacity } } -fn render_aggregate_totals(out: &mut String, publication: &WebRuntimePublication) { - let totals = publication.telemetry.aggregates(); - let _ = writeln!( - out, - "# HELP telemt_web_session_incarnations_total Process-owned WEB session lifecycle totals" - ); - let _ = writeln!(out, "# TYPE telemt_web_session_incarnations_total counter"); - let _ = writeln!( - out, - "telemt_web_session_incarnations_total{{event=\"created\"}} {}", - totals.sessions_created - ); - let _ = writeln!( - out, - "telemt_web_session_incarnations_total{{event=\"closed\"}} {}", - totals.sessions_closed - ); - let _ = writeln!( - out, - "# HELP telemt_web_streams_total Process-owned WEB logical stream totals" - ); - let _ = writeln!(out, "# TYPE telemt_web_streams_total counter"); - let _ = writeln!( - out, - "telemt_web_streams_total{{event=\"opened\"}} {}", - totals.streams_opened - ); - let _ = writeln!( - out, - "telemt_web_streams_total{{event=\"rejected\"}} {}", - totals.streams_rejected - ); - let _ = writeln!( - out, - "# HELP telemt_web_carrier_bytes_total Process-owned WEB carrier payload bytes" - ); - let _ = writeln!(out, "# TYPE telemt_web_carrier_bytes_total counter"); - let _ = writeln!( - out, - "telemt_web_carrier_bytes_total{{direction=\"up\"}} {}", - totals.bytes_up - ); - let _ = writeln!( - out, - "telemt_web_carrier_bytes_total{{direction=\"down\"}} {}", - totals.bytes_down - ); -} - fn operator_state_token(state: OperatorLifecycleState) -> &'static str { match state { OperatorLifecycleState::Running => "running", diff --git a/src/metrics/web/lifecycle.rs b/src/metrics/web/lifecycle.rs new file mode 100644 index 0000000..152c458 --- /dev/null +++ b/src/metrics/web/lifecycle.rs @@ -0,0 +1,130 @@ +use std::fmt::Write; + +use crate::config::{ProxyConfig, WebCarrier}; +use crate::web::control::WebRuntimePublication; +use crate::web::session::SessionCloseReason; +use crate::web::telemetry::{WebBridgeRecoveryEvent, WebSessionLifecycleObservation}; + +pub(super) fn render( + out: &mut String, + publication: &WebRuntimePublication, + config: &ProxyConfig, +) { + let _ = writeln!( + out, + "# HELP telemt_web_session_closures_total Closed WEB session incarnations by carrier and terminal reason" + ); + let _ = writeln!(out, "# TYPE telemt_web_session_closures_total counter"); + for carrier in WebCarrier::ALL { + for reason in SessionCloseReason::ALL { + let _ = writeln!( + out, + "telemt_web_session_closures_total{{carrier=\"{}\",reason=\"{}\"}} {}", + carrier.as_str(), + reason.as_str(), + publication.telemetry.session_close_total(carrier, reason) + ); + } + } + + let _ = writeln!( + out, + "# HELP telemt_web_session_lifecycle_observations_total Authenticated activity after bounded lifecycle gaps" + ); + let _ = writeln!( + out, + "# TYPE telemt_web_session_lifecycle_observations_total counter" + ); + for carrier in WebCarrier::ALL { + for observation in WebSessionLifecycleObservation::ALL { + let _ = writeln!( + out, + "telemt_web_session_lifecycle_observations_total{{carrier=\"{}\",observation=\"{}\"}} {}", + carrier.as_str(), + observation.as_str(), + publication + .telemetry + .session_observation_total(carrier, observation) + ); + } + } + + let _ = writeln!( + out, + "# HELP telemt_web_bridge_recovery_events_total Server-observed bridge recovery milestones" + ); + let _ = writeln!( + out, + "# TYPE telemt_web_bridge_recovery_events_total counter" + ); + for event in WebBridgeRecoveryEvent::ALL { + let _ = writeln!( + out, + "telemt_web_bridge_recovery_events_total{{event=\"{}\"}} {}", + event.as_str(), + publication.telemetry.bridge_recovery_total(event) + ); + } + + let _ = writeln!( + out, + "# HELP telemt_web_bridge_recovery_seconds Effective bridge recovery deadline" + ); + let _ = writeln!(out, "# TYPE telemt_web_bridge_recovery_seconds gauge"); + let _ = writeln!( + out, + "telemt_web_bridge_recovery_seconds {}", + config.web.timeouts.bridge_recovery_secs + ); + + render_aggregates(out, publication); +} + +fn render_aggregates(out: &mut String, publication: &WebRuntimePublication) { + let totals = publication.telemetry.aggregates(); + let _ = writeln!( + out, + "# HELP telemt_web_session_incarnations_total Process-owned WEB session lifecycle totals" + ); + let _ = writeln!(out, "# TYPE telemt_web_session_incarnations_total counter"); + let _ = writeln!( + out, + "telemt_web_session_incarnations_total{{event=\"created\"}} {}", + totals.sessions_created + ); + let _ = writeln!( + out, + "telemt_web_session_incarnations_total{{event=\"closed\"}} {}", + totals.sessions_closed + ); + let _ = writeln!( + out, + "# HELP telemt_web_streams_total Process-owned WEB logical stream totals" + ); + let _ = writeln!(out, "# TYPE telemt_web_streams_total counter"); + let _ = writeln!( + out, + "telemt_web_streams_total{{event=\"opened\"}} {}", + totals.streams_opened + ); + let _ = writeln!( + out, + "telemt_web_streams_total{{event=\"rejected\"}} {}", + totals.streams_rejected + ); + let _ = writeln!( + out, + "# HELP telemt_web_carrier_bytes_total Process-owned WEB carrier payload bytes" + ); + let _ = writeln!(out, "# TYPE telemt_web_carrier_bytes_total counter"); + let _ = writeln!( + out, + "telemt_web_carrier_bytes_total{{direction=\"up\"}} {}", + totals.bytes_up + ); + let _ = writeln!( + out, + "telemt_web_carrier_bytes_total{{direction=\"down\"}} {}", + totals.bytes_down + ); +} diff --git a/src/web/bridge.rs b/src/web/bridge.rs index 385ab17..beed502 100644 --- a/src/web/bridge.rs +++ b/src/web/bridge.rs @@ -21,12 +21,16 @@ pub(crate) fn render( batch_limit: usize, queue_limit: usize, queue_items: usize, + max_streams: usize, negotiation_enabled: bool, candidate_count: usize, carrier_deadlines: [u64; 4], long_poll_secs: u64, bridge_request_secs: u64, bridge_retry_secs: u64, + bridge_recovery_secs: u64, + websocket_open_secs: u64, + reconnect_grace_secs: u64, carrier_probe_coalesce_ms: u64, rng: &SecureRandom, ) -> BridgePage { @@ -35,6 +39,9 @@ pub(crate) fn render( let nonce = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(nonce); let body = DOCUMENT .replace("__RESPONSE_RUNTIME__", RESPONSE_RUNTIME) + .replace("__REQUEST_RUNTIME__", REQUEST_RUNTIME) + .replace("__BUFFER_RUNTIME__", BUFFER_RUNTIME) + .replace("__RECOVERY_RUNTIME__", RECOVERY_RUNTIME) .replace("__RUNTIME__", RUNTIME) .replace("__NONCE__", &nonce) .replace("__HOST__", host) @@ -42,6 +49,7 @@ pub(crate) fn render( .replace("__BATCH_LIMIT__", &batch_limit.to_string()) .replace("__QUEUE_LIMIT__", &queue_limit.to_string()) .replace("__QUEUE_ITEMS__", &queue_items.to_string()) + .replace("__MAX_STREAMS__", &max_streams.to_string()) .replace( "__NEGOTIATION_ENABLED__", if negotiation_enabled { "true" } else { "false" }, @@ -50,6 +58,12 @@ pub(crate) fn render( .replace("__LONG_POLL_SECS__", &long_poll_secs.to_string()) .replace("__BRIDGE_REQUEST_SECS__", &bridge_request_secs.to_string()) .replace("__BRIDGE_RETRY_SECS__", &bridge_retry_secs.to_string()) + .replace( + "__BRIDGE_RECOVERY_SECS__", + &bridge_recovery_secs.to_string(), + ) + .replace("__WEBSOCKET_OPEN_SECS__", &websocket_open_secs.to_string()) + .replace("__RECONNECT_GRACE_SECS__", &reconnect_grace_secs.to_string()) .replace( "__CARRIER_PROBE_COALESCE_MS__", &carrier_probe_coalesce_ms.to_string(), @@ -72,6 +86,9 @@ pub(crate) fn render( const DOCUMENT: &str = include_str!("bridge/document.html"); const RESPONSE_RUNTIME: &str = include_str!("bridge/response.js"); +const REQUEST_RUNTIME: &str = include_str!("bridge/request.js"); +const BUFFER_RUNTIME: &str = include_str!("bridge/buffers.js"); +const RECOVERY_RUNTIME: &str = include_str!("bridge/recovery.js"); const RUNTIME: &str = include_str!("bridge/runtime.js"); // Rendered wire-contract tests remain separate from the embedded document. diff --git a/src/web/bridge/buffers.js b/src/web/bridge/buffers.js new file mode 100644 index 0000000..3989581 --- /dev/null +++ b/src/web/bridge/buffers.js @@ -0,0 +1,143 @@ +(()=>{'use strict'; +function create(settings){ + let queuedBytes=0,queuedItems=0; + const retiredLimit=4096,activeStreams=new Set(),retiredStreams=new Set(),retiredStreamOrder=[]; + function reserve(data,lane){ + const limits=settings.limits(),buffered=settings.buffered(); + if(!data.byteLength||data.byteLength>limits.queueBytes-queuedBytes-buffered||queuedItems>=limits.queueItems)return false; + if(lane&&(data.byteLength>limits.laneBytes-lane.bytes-(lane.socket?lane.socket.bufferedAmount:0)||lane.items>=limits.laneItems))return false; + queuedBytes+=data.byteLength;queuedItems++;if(lane){lane.bytes+=data.byteLength;lane.items++}return true; + } + function release(bytes,items,lane){ + if(bytes>queuedBytes||items>queuedItems||(lane&&(bytes>lane.bytes||items>lane.items)))throw new Error('queue accounting invariant'); + queuedBytes-=bytes;queuedItems-=items;if(lane){lane.bytes-=bytes;lane.items-=items} + } + function releasePending(values,lane){ + if(!values.length)return;let bytes=0;for(const value of values)bytes+=value.byteLength; + const items=values.length;values.length=0;release(bytes,items,lane); + } + function frameBound(value,maxFrames,maxBytes){ + const view=new DataView(value);let offset=0,frames=0; + while(offset1048576||end>value.byteLength)throw new Error('invalid frame'); + if(frames>0&&(frames>=maxFrames||end>maxBytes))break; + frames++;offset=end; + } + if(!frames)throw new Error('empty frame batch'); + return {frames,bytes:offset}; + } + function splitFrames(value){ + const view=new DataView(value),result=[];let offset=0; + while(offset=4096)throw new Error('invalid frame batch'); + const type=view.getUint8(offset),id=(view.getUint8(offset+1)<<16)|(view.getUint8(offset+2)<<8)|view.getUint8(offset+3); + const size=view.getUint32(offset+4),end=offset+8+size; + if((type===2&&!size)||size>1048576||end>value.byteLength)throw new Error('invalid frame'); + result.push({type,id,data:offset===0&&end===value.byteLength?value:value.slice(offset,end)});offset=end; + } + if(!result.length)throw new Error('empty frame batch');return result; + } + function rememberStreamRetired(id){ + if(!id||retiredStreams.has(id))return; + if(retiredStreamOrder.length===retiredLimit)retiredStreams.delete(retiredStreamOrder.shift()); + retiredStreams.add(id);retiredStreamOrder.push(id); + } + function acceptNativeFrames(data){ + const values=splitFrames(data),accepted=[]; + for(const value of values){ + if(retiredStreams.has(value.id)&&!activeStreams.has(value.id))continue; + if(value.type===1){ + if(!activeStreams.has(value.id)&&activeStreams.size>=settings.maxStreams())throw settings.failure('capacity','stream capacity exhausted'); + activeStreams.add(value.id);accepted.push(value.data);continue; + } + if(value.type===3){activeStreams.delete(value.id);rememberStreamRetired(value.id)} + accepted.push(value.data); + } + if(accepted.length===values.length)return data;if(!accepted.length)return null; + let total=0;for(const value of accepted)total+=value.byteLength; + const joined=new Uint8Array(total);let offset=0; + for(const value of accepted){joined.set(new Uint8Array(value),offset);offset+=value.byteLength} + return joined.buffer; + } + function observeServerFrames(data){ + for(const value of splitFrames(data))if(value.type===3){activeStreams.delete(value.id);rememberStreamRetired(value.id)} + } + function probeFrames(){ + const result=[];let scanned=0,pending=settings.pending(),batchLimit=settings.limits().batchBytes; + for(let index=0;index=4096)throw new Error('invalid frame batch'); + const type=view.getUint8(start),id=(view.getUint8(start+1)<<16)|(view.getUint8(start+2)<<8)|view.getUint8(start+3); + const size=view.getUint32(start+4),end=start+8+size,bytes=end-start; + if((type===2&&!size)||size>1048576||end>source.byteLength)throw new Error('invalid frame'); + if(scanned+bytes>batchLimit)return result; + result.push({source,index,start,end,type,id});scanned+=bytes;start=end; + } + } + return result; + } + function findProbe(includeData){ + const frames=probeFrames(),first=frames.findIndex(frame=>frame.type===1||frame.type===2);if(first<0)return null; + const laneMode=settings.laneMode(),selected=[];let hasData=frames[first].type===2; + if(laneMode){ + selected.push(frames[first]); + if(includeData&&!hasData)for(let index=first+1;indexright-left); + for(const index of indexes){ + const source=pending[index],spans=groups.get(index).sort((left,right)=>left.start-right.start);let removed=0,offset=0; + for(const span of spans){if(span.startbatchLimit||frames+bound.frames>4096))break; + total+=values[count].byteLength;frames+=bound.frames;count++; + } + const joined=new Uint8Array(total);let offset=0; + for(const data of values.splice(0,count)){joined.set(new Uint8Array(data),offset);offset+=data.byteLength} + return {body:joined.buffer,total,count}; + } + function takeBatch(values,lane){return Object.assign(joinPending(values,lane),{lane,controller:null,cancelled:false,settled:false})} + function settleBatch(lease){if(!lease||lease.settled)return false;lease.settled=true;release(lease.total,lease.count,lease.lane);return true} + function cancelBatch(lease){if(!lease||lease.settled)return;lease.cancelled=true;if(lease.controller)lease.controller.abort();settleBatch(lease)} + function closeFrame(id){const value=new Uint8Array(8);value[0]=3;value[1]=(id>>>16)&255;value[2]=(id>>>8)&255;value[3]=id&255;return value.buffer} + function retireStream(id){activeStreams.delete(id);rememberStreamRetired(id)} + function retireAllStreams(){const ids=Array.from(activeStreams);activeStreams.clear();for(const id of ids)rememberStreamRetired(id);return ids} + function clearStreams(){activeStreams.clear();retiredStreams.clear();retiredStreamOrder.length=0} + function assertEmpty(){if(queuedBytes!==0||queuedItems!==0)throw new Error('queue accounting leak')} + return Object.freeze({reserve,release,releasePending,frameBound,splitFrames,acceptNativeFrames,observeServerFrames,findProbe,consumeProbe,takeBatch,settleBatch,cancelBatch,closeFrame,retireStream,retireAllStreams,clearStreams,assertEmpty}); +} +globalThis.TelemtBridgeBuffers=Object.freeze({create}); +})(); diff --git a/src/web/bridge/document.html b/src/web/bridge/document.html index 56dab71..ea68ec7 100644 --- a/src/web/bridge/document.html +++ b/src/web/bridge/document.html @@ -10,6 +10,15 @@ __RESPONSE_RUNTIME__ + + + diff --git a/src/web/bridge/recovery.js b/src/web/bridge/recovery.js new file mode 100644 index 0000000..b433bf2 --- /dev/null +++ b/src/web/bridge/recovery.js @@ -0,0 +1,75 @@ +(()=>{ +'use strict'; +const mediaType='application/vnd.telemt.web-recovery+json',maxBytes=1024,heartbeatMs=2500; +const exactKeys=(value,keys)=>value&&typeof value==='object'&&!Array.isArray(value)&&Object.keys(value).sort().join(',')===keys.slice().sort().join(','); +const integer=(value,min,max)=>Number.isSafeInteger(value)&&value>=min&&value<=max; +function parsePolicy(bytes){ + let value;try{value=JSON.parse(new TextDecoder('utf-8',{fatal:true}).decode(bytes))}catch(error){throw new Error('invalid recovery document')} + if(!exactKeys(value,['v','bootstrap','limits','timeouts','negotiation'])||value.v!==1||!/^[A-Za-z0-9_-]{43}$/.test(value.bootstrap))throw new Error('invalid recovery document'); + const limits=value.limits,timeouts=value.timeouts,negotiation=value.negotiation; + if(!exactKeys(limits,['carrier_batch_bytes','pending_bytes_per_session','pending_items_per_session','max_streams_per_session']) + ||!integer(limits.carrier_batch_bytes,8,16777216)||!integer(limits.pending_bytes_per_session,limits.carrier_batch_bytes,4294967296) + ||!integer(limits.pending_items_per_session,1,1048576)||!integer(limits.max_streams_per_session,1,16777215))throw new Error('invalid recovery limits'); + if(!exactKeys(timeouts,['long_poll_secs','bridge_request_secs','bridge_retry_secs','bridge_recovery_secs','websocket_open_secs','reconnect_grace_secs']) + ||!integer(timeouts.long_poll_secs,1,3600)||!integer(timeouts.bridge_request_secs,1,60)||!integer(timeouts.bridge_retry_secs,1,300) + ||!integer(timeouts.bridge_recovery_secs,1,60)||!integer(timeouts.websocket_open_secs,1,300)||!integer(timeouts.reconnect_grace_secs,1,3600))throw new Error('invalid recovery timeouts'); + if(!exactKeys(negotiation,['enabled','candidate_count','deadlines_secs','carrier_probe_coalesce_ms'])||typeof negotiation.enabled!=='boolean' + ||!integer(negotiation.candidate_count,1,4)||!Array.isArray(negotiation.deadlines_secs)||negotiation.deadlines_secs.length!==4 + ||!negotiation.deadlines_secs.every((entry,index)=>integer(entry,1,3600)&&(index===0||entry>negotiation.deadlines_secs[index-1])) + ||!integer(negotiation.carrier_probe_coalesce_ms,0,10))throw new Error('invalid recovery negotiation'); + return value; +} +function create(settings){ + let current=null,nextEpoch=1; + const remaining=owner=>Math.max(0,Math.min(owner.wall-Date.now(),owner.monotonic-performance.now())); + function stop(owner){ + if(owner.heartbeat)clearTimeout(owner.heartbeat);owner.heartbeat=null; + if(owner.deadlineTimer)clearTimeout(owner.deadlineTimer);owner.deadlineTimer=null; + } + function heartbeat(owner){ + if(current!==owner||owner.controller.signal.aborted)return; + const left=remaining(owner);settings.status(left); + if(left>0)owner.heartbeat=setTimeout(()=>heartbeat(owner),Math.min(heartbeatMs,left)); + } + async function load(owner){ + const left=remaining(owner);if(left<=0)throw new Error('recovery deadline'); + const requestController=new AbortController(),abort=()=>requestController.abort(); + owner.controller.signal.addEventListener('abort',abort,{once:true}); + const timer=setTimeout(abort,Math.max(1,Math.min(settings.requestMs(),left))); + try{ + const token=settings.token(); + const response=await fetch(settings.url(),{ + method:'GET',signal:requestController.signal,mode:'same-origin',credentials:'omit',cache:'no-store',redirect:'error',referrerPolicy:'no-referrer', + headers:Object.assign({Accept:mediaType},token?{Authorization:'Bearer '+token}:{}) + }); + if(response.status!==200||response.headers.get('Content-Type')!==mediaType){settings.cancel(response);throw new Error('recovery representation rejected')} + const bytes=await settings.read(response,maxBytes,false,requestController.signal); + return parsePolicy(bytes); + }finally{ + clearTimeout(timer);owner.controller.signal.removeEventListener('abort',abort); + } + } + async function run(owner,replay){ + heartbeat(owner); + if(replay){ + try{await replay(owner.controller.signal,()=>remaining(owner));settings.restored();return true} + catch(error){if(!settings.replaceable(error))throw error} + } + const policy=await load(owner); + if(remaining(owner)<=0)throw new Error('recovery deadline'); + await settings.replace(policy,owner.controller.signal,()=>remaining(owner),owner.epoch); + return true; + } + function recover(reason,replay){ + if(current)return current.promise; + const budget=settings.budgetMs(),owner={epoch:nextEpoch++,controller:new AbortController(),heartbeat:null,deadlineTimer:null,promise:null}; + owner.wall=Date.now()+budget;owner.monotonic=performance.now()+budget;current=owner; + owner.deadlineTimer=setTimeout(()=>owner.controller.abort(),Math.max(1,remaining(owner))); + owner.promise=run(owner,replay).catch(error=>{settings.terminal(settings.reason(error,reason));return false}).finally(()=>{stop(owner);if(current===owner)current=null}); + return owner.promise; + } + function cancel(){if(current){const owner=current;current=null;owner.controller.abort();stop(owner)}} + return Object.freeze({recover,cancel,active:()=>current!==null,remaining:()=>current?remaining(current):0}); +} +globalThis.TelemtBridgeRecovery=Object.freeze({create}); +})(); diff --git a/src/web/bridge/request.js b/src/web/bridge/request.js new file mode 100644 index 0000000..69ff439 --- /dev/null +++ b/src/web/bridge/request.js @@ -0,0 +1,66 @@ +(()=>{'use strict'; +function create(settings){ + const pause=(milliseconds,signal)=>new Promise((resolve,reject)=>{ + if(signal&&signal.aborted){reject(new Error('request aborted'));return} + const timer=setTimeout(done,milliseconds);function done(){if(signal)signal.removeEventListener('abort',abort);resolve()} + function abort(){clearTimeout(timer);signal.removeEventListener('abort',abort);reject(new Error('request aborted'))} + if(signal)signal.addEventListener('abort',abort,{once:true}); + }); + const options=(method,token,body,headers,signal,keepalive)=>({ + method,body,signal,keepalive:!!keepalive,mode:'same-origin',credentials:'omit',cache:'no-store',redirect:'error',referrerPolicy:'no-referrer', + headers:Object.assign(token?{Authorization:'Bearer '+token}:{},body?{'Content-Type':'application/octet-stream'}:{},headers||{}) + }); + function retryAfterMs(response){ + const header=response.headers.get('Retry-After'); + if(!header)return 0; + const seconds=Number(header); + if(Number.isFinite(seconds)&&seconds>=0)return Math.min(seconds*1000,30000); + const when=Date.parse(header); + if(Number.isFinite(when)){const delta=when-Date.now();return delta>0?Math.min(delta,30000):0} + return 0; + } + function retryableStatus(status){return status===408||status===429||status===502||status===503||status===504} + function responsePolicy(path,status){ + if(path==='/api/v1/session'&&status===200)return {limit:8,exact:true}; + if(path==='/api/v1/down'&&status===200)return {limit:settings.batchLimit(),exact:false}; + return {limit:0,exact:true}; + } + async function send(path,frozenOptions,remainingBudget,maxAttempts){ + let delay=250,attempt=0,lastReason='network';maxAttempts=maxAttempts||9; + const initialBudget=remainingBudget?Math.min(settings.retryMs(),remainingBudget()):settings.retryMs(); + const deadline=Date.now()+Math.max(0,initialBudget),external=frozenOptions.signal; + const attemptLimit=path==='/api/v1/down'?settings.longPollMs()+settings.requestMs():settings.requestMs(); + while(attemptcontroller.abort(); + 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))); + let response=null,wait=0; + try{ + const fetched=await fetch(settings.origin()+path,requestOptions); + if(retryableStatus(fetched.status)){ + lastReason='http';wait=retryAfterMs(fetched);settings.cancel(fetched); + }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)} + response={status:fetched.status,headers:fetched.headers,body};return response; + } + }catch(error){ + controller.abort(); + if(settings.closed()||(external&&external.aborted))throw error; + 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; + settings.retrying(); + const backoff=wait||delay+Math.floor(Math.random()*Math.max(1,delay/4)); + await pause(Math.min(backoff,after),external);delay=Math.min(delay*2,2000); + } + throw settings.failure(lastReason,'carrier retry limit reached'); + } + return Object.freeze({options,pause,send}); +} +globalThis.TelemtBridgeRequest=Object.freeze({create}); +})(); diff --git a/src/web/bridge/runtime.js b/src/web/bridge/runtime.js index 903fe7c..c0394c3 100644 --- a/src/web/bridge/runtime.js +++ b/src/web/bridge/runtime.js @@ -1,20 +1,25 @@ (()=>{'use strict'; -const bootstrap="__BOOTSTRAP__"; +let bootstrap="__BOOTSTRAP__"; const relayOrigin='https://__HOST__',carrierCapabilities='https,https-lanes,websocket,websocket-lanes'; const responseBody=globalThis.TelemtBridgeResponse;if(!responseBody)throw new Error('missing response runtime'); -const negotiationEnabled=__NEGOTIATION_ENABLED__,candidateCount=__CANDIDATE_COUNT__,candidateDeadlines=[__CARRIER_DEADLINES__]; -const longPollMs=__LONG_POLL_SECS__*1000,bridgeRequestMs=__BRIDGE_REQUEST_SECS__*1000,bridgeRetryMs=__BRIDGE_RETRY_SECS__*1000; -const probeCoalesceMs=__CARRIER_PROBE_COALESCE_MS__; +const requestSupport=globalThis.TelemtBridgeRequest;if(!requestSupport)throw new Error('missing request runtime'); +const bufferSupport=globalThis.TelemtBridgeBuffers;if(!bufferSupport)throw new Error('missing buffer runtime'); +const recoverySupport=globalThis.TelemtBridgeRecovery;if(!recoverySupport)throw new Error('missing recovery runtime'); +let negotiationEnabled=__NEGOTIATION_ENABLED__,candidateCount=__CANDIDATE_COUNT__,candidateDeadlines=[__CARRIER_DEADLINES__]; +let longPollMs=__LONG_POLL_SECS__*1000,bridgeRequestMs=__BRIDGE_REQUEST_SECS__*1000,bridgeRetryMs=__BRIDGE_RETRY_SECS__*1000; +let bridgeRecoveryMs=__BRIDGE_RECOVERY_SECS__*1000,websocketOpenMs=__WEBSOCKET_OPEN_SECS__*1000,reconnectGraceMs=__RECONNECT_GRACE_SECS__*1000; +let probeCoalesceMs=__CARRIER_PROBE_COALESCE_MS__; let negotiatedCandidateCount=candidateCount,negotiatedFinalDeadline=candidateDeadlines[3],negotiatedFrozen=false; -const batchLimit=__BATCH_LIMIT__,queueLimit=__QUEUE_LIMIT__,queueItemLimit=__QUEUE_ITEMS__; -const laneQueueLimit=Math.min(queueLimit,8388608),laneItemLimit=Math.min(queueItemLimit,1024),closedLaneLimit=4096; -const fragment=location.hash,androidNonce=/^#android=([A-Za-z0-9_-]{43})$/.exec(fragment)?.[1]||''; +let batchLimit=__BATCH_LIMIT__,queueLimit=__QUEUE_LIMIT__,queueItemLimit=__QUEUE_ITEMS__,maxStreams=__MAX_STREAMS__; +let laneQueueLimit=Math.min(queueLimit,8388608),laneItemLimit=Math.min(queueItemLimit,1024);const closedLaneLimit=4096; +const fragment=location.hash,androidNonce=/^#android=([A-Za-z0-9_-]{43})$/.exec(fragment)?.[1]||'',recoveryPath=location.pathname+location.search; history.replaceState(null,'',location.pathname); let initialized=false,closed=false,port=null,sessionToken='',cleanupToken='',createStarted=false,socket=null,socketReady=false,carrier=''; -let queuedBytes=0,queuedItems=0,upSequence=1,downCursor='0',upRunning=false,upLease=null,pollController=null; +let upSequence=1,downCursor='0',upRunning=false,upLease=null,pollController=null; let helloFrame=null,helloTimer=null,welcomeSent=false,carrierAttempt=1,carrierFailure='',carrierCommitted=false,terminalFailure=''; let negotiationStartedAt=0,carrierTimer=null,probeTimer=null,attemptController=null,attemptEpoch=1,candidateRunning=false,switching=false,currentAttempt=null; -const pending=[],upPending=[],lanes=new Map(),closedLanes=new Set(),closedLaneOrder=[]; +let recoveryController=null,recoveryCommit=null,recoveryReplaced=false,lastSchedulerWall=Date.now(),lastSchedulerMonotonic=performance.now(); +const pending=[],upPending=[],recoveryPending=[],lanes=new Map(),closedLanes=new Set(),closedLaneOrder=[]; const canonicalFailures=['timeout','network','upgrade','http','protocol']; const failure=(reason,message)=>Object.assign(new Error(message||reason),{telemtReason:reason}); const failureReason=(error,fallback)=>error&&canonicalFailures.includes(error.telemtReason)?error.telemtReason:fallback; @@ -25,182 +30,84 @@ const status=(state,phase,reason,deadlineMs)=>{ if(currentDeadline===undefined)currentDeadline=negotiationStartedAt?Math.max(0,negotiationStartedAt+negotiatedFinalDeadline*1000-Date.now()):0; port.postMessage({t:'status',state,phase:currentPhase,reason:reason||'',deadline_ms:Math.max(0,Math.ceil(currentDeadline))}); }; -const pause=(milliseconds,signal)=>new Promise((resolve,reject)=>{ - if(signal&&signal.aborted){reject(new Error('request aborted'));return} - const timer=setTimeout(done,milliseconds);function done(){if(signal)signal.removeEventListener('abort',abort);resolve()} - function abort(){clearTimeout(timer);signal.removeEventListener('abort',abort);reject(new Error('request aborted'))} - if(signal)signal.addEventListener('abort',abort,{once:true}); -}); const socketURL=()=>relayOrigin.replace(/^https:/,'wss:')+'/api/v1/ws'; -const options=(method,token,body,headers,signal,keepalive)=>({ - method,body,signal,keepalive:!!keepalive,mode:'same-origin',credentials:'omit',cache:'no-store',redirect:'error',referrerPolicy:'no-referrer', - headers:Object.assign(token?{Authorization:'Bearer '+token}:{},body?{'Content-Type':'application/octet-stream'}:{},headers||{}) +const requestClient=requestSupport.create({ + origin:()=>relayOrigin,closed:()=>closed,retryMs:()=>bridgeRetryMs,longPollMs:()=>longPollMs,requestMs:()=>bridgeRequestMs, + batchLimit:()=>batchLimit,read:(response,limit,exact,signal)=>responseBody.read(response,limit,exact,signal),cancel:responseBody.cancel, + failure,reason:failureReason,retrying:()=>status('reconnecting') }); +const options=requestClient.options,pause=requestClient.pause,request=requestClient.send; +const buffers=bufferSupport.create({ + limits:()=>({batchBytes:batchLimit,queueBytes:queueLimit,queueItems:queueItemLimit,laneBytes:laneQueueLimit,laneItems:laneItemLimit}), + buffered:()=>{let total=socket?socket.bufferedAmount:0;for(const value of lanes.values())if(value.socket)total+=value.socket.bufferedAmount;return total}, + pending:()=>pending,laneMode:()=>carrier==='https-lanes'||carrier==='websocket-lanes',maxStreams:()=>maxStreams,failure +}); +const {reserve,release,releasePending,frameBound,splitFrames,acceptNativeFrames,observeServerFrames,findProbe,consumeProbe,takeBatch,closeFrame,retireStream,retireAllStreams,clearStreams}=buffers; +function detachLease(lease){if(lease.lane){if(lease.lane.upLease===lease)lease.lane.upLease=null}else if(upLease===lease)upLease=null} +function settleBatch(lease){if(!buffers.settleBatch(lease))return false;detachLease(lease);return true} +function cancelBatch(lease){if(!lease||lease.settled)return;buffers.cancelBatch(lease);detachLease(lease)} const attemptHeaders=(attempt,failure)=>negotiationEnabled?Object.assign({'X-Carrier-Capabilities':carrierCapabilities,'X-Carrier-Attempt':String(attempt)},failure?{'X-Carrier-Failure':failure}:{}):{}; -function reserve(data,lane){ - let buffered=socket?socket.bufferedAmount:0;for(const value of lanes.values())if(value.socket)buffered+=value.socket.bufferedAmount; - if(!data.byteLength||data.byteLength>queueLimit-queuedBytes-buffered||queuedItems>=queueItemLimit)return false; - if(lane&&(data.byteLength>laneQueueLimit-lane.bytes-(lane.socket?lane.socket.bufferedAmount:0)||lane.items>=laneItemLimit))return false; - queuedBytes+=data.byteLength;queuedItems++;if(lane){lane.bytes+=data.byteLength;lane.items++}return true; +function finishOldRecovery(){ + recoveryReplaced=false;status('connected','committed','',0); + for(const data of recoveryPending.splice(0)){release(data.byteLength,1,null);queueCarrier(data)} } -function release(bytes,items,lane){ - if(bytes>queuedBytes||items>queuedItems||(lane&&(bytes>lane.bytes||items>lane.items)))throw new Error('queue accounting invariant'); - queuedBytes-=bytes;queuedItems-=items;if(lane){lane.bytes-=bytes;lane.items-=items} +function rejectRecoveryCommit(error){ + const commit=recoveryCommit;if(!commit)return;recoveryCommit=null; + commit.signal.removeEventListener('abort',commit.abort);commit.reject(error); } -function releasePending(values,lane){ - if(!values.length)return;let bytes=0;for(const value of values)bytes+=value.byteLength; - const items=values.length;values.length=0;release(bytes,items,lane); +function resolveRecoveryCommit(){ + const commit=recoveryCommit;if(!commit)return;recoveryCommit=null;recoveryReplaced=false; + commit.signal.removeEventListener('abort',commit.abort);commit.resolve(); } -function frameBound(value,maxFrames,maxBytes){ - const view=new DataView(value);let offset=0,frames=0; - while(offset1048576||end>value.byteLength)throw new Error('invalid frame'); - if(frames>0&&(frames>=maxFrames||end>maxBytes))break; - frames++;offset=end; +function retireCarrier(policy){ + recoveryReplaced=true;attemptEpoch++;if(carrierTimer)clearTimeout(carrierTimer);carrierTimer=null;clearProbeTimer(); + if(attemptController)attemptController.abort();attemptController=null;if(pollController)pollController.abort();pollController=null; + if(socket){const previous=socket;socket=null;previous.close()}socketReady=false;cancelBatch(upLease);releasePending(upPending,null); + for(const lane of lanes.values()){ + if(lane.controller)lane.controller.abort();cancelBatch(lane.upLease);releasePending(lane.pending,lane);if(lane.socket)lane.socket.close(); } - if(!frames)throw new Error('empty frame batch'); - return {frames,bytes:offset}; + lanes.clear();closedLanes.clear();closedLaneOrder.length=0;releasePending(pending,null);releasePending(recoveryPending,null); + for(const id of retireAllStreams())if(port){const frame=closeFrame(id);port.postMessage(frame,[frame])} + bootstrap=policy.bootstrap;batchLimit=policy.limits.carrier_batch_bytes;queueLimit=policy.limits.pending_bytes_per_session; + queueItemLimit=policy.limits.pending_items_per_session;maxStreams=policy.limits.max_streams_per_session; + laneQueueLimit=Math.min(queueLimit,8388608);laneItemLimit=Math.min(queueItemLimit,1024); + longPollMs=policy.timeouts.long_poll_secs*1000;bridgeRequestMs=policy.timeouts.bridge_request_secs*1000; + bridgeRetryMs=policy.timeouts.bridge_retry_secs*1000;bridgeRecoveryMs=policy.timeouts.bridge_recovery_secs*1000; + websocketOpenMs=policy.timeouts.websocket_open_secs*1000;reconnectGraceMs=policy.timeouts.reconnect_grace_secs*1000; + negotiationEnabled=policy.negotiation.enabled;candidateCount=policy.negotiation.candidate_count; + candidateDeadlines=policy.negotiation.deadlines_secs;probeCoalesceMs=policy.negotiation.carrier_probe_coalesce_ms; + negotiatedCandidateCount=candidateCount;negotiatedFinalDeadline=candidateDeadlines[3];negotiatedFrozen=false; + sessionToken='';cleanupToken='';carrier='';carrierAttempt=1;carrierFailure='';carrierCommitted=false; + candidateRunning=false;switching=false;currentAttempt=null;upSequence=1;downCursor='0';upRunning=false; } -function splitFrames(value){ - const view=new DataView(value),result=[];let offset=0; - while(offset=4096)throw new Error('invalid frame batch'); - const type=view.getUint8(offset),id=(view.getUint8(offset+1)<<16)|(view.getUint8(offset+2)<<8)|view.getUint8(offset+3); - const size=view.getUint32(offset+4),end=offset+8+size; - if((type===2&&!size)||size>1048576||end>value.byteLength)throw new Error('invalid frame'); - result.push({type,id,data:offset===0&&end===value.byteLength?value:value.slice(offset,end)});offset=end; - } - if(!result.length)throw new Error('empty frame batch');return result; +function replaceCarrier(policy,signal,remaining){ + if(closed||!helloFrame||!port)throw failure('protocol','missing recovery owner'); + retireCarrier(policy);negotiationStartedAt=Date.now();armCarrierDeadline(attemptEpoch); + return new Promise((resolve,reject)=>{ + const abort=()=>rejectRecoveryCommit(failure('timeout','recovery deadline')); + recoveryCommit={resolve,reject,signal,abort};signal.addEventListener('abort',abort,{once:true}); + if(remaining()<=0){abort();return}createSession(attemptEpoch); + }); } -function probeFrames(){ - const result=[];let scanned=0; - for(let index=0;index=4096)throw new Error('invalid frame batch'); - const type=view.getUint8(start),id=(view.getUint8(start+1)<<16)|(view.getUint8(start+2)<<8)|view.getUint8(start+3); - const size=view.getUint32(start+4),end=start+8+size,bytes=end-start; - if((type===2&&!size)||size>1048576||end>source.byteLength)throw new Error('invalid frame'); - if(scanned+bytes>batchLimit)return result; - result.push({source,index,start,end,type,id});scanned+=bytes;start=end; - } - } - return result; +function recoverTransport(error,replay){ + if(closed)return Promise.resolve(false); + const reason=failureReason(error,'network'); + if(reason==='protocol'){fail(reason);return Promise.resolve(false)} + return recoveryController.recover(reason,replay); } -function findProbe(includeData){ - const frames=probeFrames(),first=frames.findIndex(frame=>frame.type===1||frame.type===2);if(first<0)return null; - const laneMode=carrier==='https-lanes'||carrier==='websocket-lanes',selected=[];let hasData=frames[first].type===2; - if(laneMode){ - selected.push(frames[first]); - if(includeData&&!hasData)for(let index=first+1;indexright-left); - for(const index of indexes){ - const source=pending[index],spans=groups.get(index).sort((left,right)=>left.start-right.start);let removed=0,offset=0; - for(const span of spans){if(span.startbatchLimit||frames+bound.frames>4096))break; - total+=values[count].byteLength;frames+=bound.frames;count++; - } - const joined=new Uint8Array(total);let offset=0; - for(const data of values.splice(0,count)){joined.set(new Uint8Array(data),offset);offset+=data.byteLength} - return {body:joined.buffer,total,count}; -} -function takeBatch(values,lane){ - const batch=joinPending(values,lane); - return Object.assign(batch,{lane,controller:null,cancelled:false,settled:false}); -} -function settleBatch(lease){ - if(!lease||lease.settled)return false;lease.settled=true; - if(lease.lane){if(lease.lane.upLease===lease)lease.lane.upLease=null}else if(upLease===lease)upLease=null; - release(lease.total,lease.count,lease.lane);return true; -} -function cancelBatch(lease){ - if(!lease||lease.settled)return;lease.cancelled=true;if(lease.controller)lease.controller.abort();settleBatch(lease); -} -function retryAfterMs(response){ - const header=response.headers.get('Retry-After'); - if(!header)return 0; - const seconds=Number(header); - if(Number.isFinite(seconds)&&seconds>=0)return Math.min(seconds*1000,30000); - const when=Date.parse(header); - if(Number.isFinite(when)){const delta=when-Date.now();return delta>0?Math.min(delta,30000):0} - return 0; -} -function retryableStatus(status){return status===408||status===429||status===502||status===503||status===504} -function responsePolicy(path,status){ - if(path==='/api/v1/session'&&status===200)return {limit:8,exact:true}; - if(path==='/api/v1/down'&&status===200)return {limit:batchLimit,exact:false}; - return {limit:0,exact:true}; -} -async function request(path,frozenOptions){ - let delay=250,attempt=0,lastReason='network';const deadline=Date.now()+bridgeRetryMs,external=frozenOptions.signal; - const attemptLimit=path==='/api/v1/down'?longPollMs+bridgeRequestMs:bridgeRequestMs; - while(attempt<9){ - if(closed||(external&&external.aborted))throw new Error('request aborted'); - const remaining=deadline-Date.now();if(remaining<=0)break;attempt++; - const controller=new AbortController(),abort=()=>controller.abort(); - 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))); - let response=null,wait=0; - try{ - const fetched=await fetch(relayOrigin+path,requestOptions); - if(retryableStatus(fetched.status)){ - lastReason='http';wait=retryAfterMs(fetched);responseBody.cancel(fetched); - }else{ - const policy=responsePolicy(path,fetched.status);let body; - try{body=await responseBody.read(fetched,policy.limit,policy.exact,controller.signal)} - catch(error){controller.abort();throw failure('protocol',error&&error.message)} - response={status:fetched.status,headers:fetched.headers,body};return response; - } - }catch(error){ - controller.abort(); - if(closed||(external&&external.aborted))throw error; - if(failureReason(error,'')==='protocol')throw error; - }finally{clearTimeout(timer);if(external)external.removeEventListener('abort',abort)} - const after=deadline-Date.now();if(attempt>=9||after<=0)break; - status('reconnecting'); - const backoff=wait||delay+Math.floor(Math.random()*Math.max(1,delay/4)); - await pause(Math.min(backoff,after),external);delay=Math.min(delay*2,5000); - } - throw failure(lastReason,'carrier retry limit reached'); +function observeResumeTrigger(){ + const gap=schedulerGap();if(!carrierCommitted||closed)return; + if(gap>=2*longPollMs)status('reconnecting','retrying','',bridgeRecoveryMs); + if(gap>=reconnectGraceMs&&!recoveryController.active())recoveryController.recover('timeout',null); } function fail(reason){ if(closed)return;reason=reason||'protocol';if(canonicalFailures.includes(reason))terminalFailure=reason; + rejectRecoveryCommit(failure(reason)); status('failed','terminal',reason,0);if(port)port.postMessage({t:'close'});close(true); } function knownCarrier(value){return value==='https'||value==='https-lanes'||value==='websocket'||value==='websocket-lanes'} @@ -333,6 +240,7 @@ function commitCarrier(probe,epoch){ if(carrier==='https')poll(); else if(carrier==='https-lanes'){const lane=lanes.get(probe.id);if(lane&&!lane.polling)pollLane(lane)} for(const data of pending.splice(0)){release(data.byteLength,1,null);queueCarrier(data)} + resolveRecoveryCommit(); } function queueCarrier(data){ try{ @@ -346,10 +254,22 @@ async function runUp(){ if(upRunning)return;upRunning=true;let lease=null; try{ while(!closed&&sessionToken&&upPending.length){ - lease=takeBatch(upPending,null);upLease=lease;lease.controller=new AbortController();const sequence=String(upSequence); - const response=await request('/api/v1/up',options('POST',sessionToken,lease.body,{'X-Up-Seq':sequence},lease.controller.signal)); - if(response.status!==204)throw failure('http','uplink rejected'); - if(response.headers.get('X-Up-Ack')!==sequence)throw failure('protocol','uplink acknowledgement rejected'); + lease=takeBatch(upPending,null);upLease=lease;lease.controller=new AbortController();const sequence=String(upSequence),token=sessionToken; + for(;;){ + try{ + const response=await request('/api/v1/up',options('POST',token,lease.body,{'X-Up-Seq':sequence},lease.controller.signal),null,1); + if(response.status!==204)throw failure('http','uplink rejected'); + if(response.headers.get('X-Up-Ack')!==sequence)throw failure('protocol','uplink acknowledgement rejected'); + break; + }catch(error){ + const recovered=await recoverTransport(error,async(signal,remaining)=>{ + const response=await request('/api/v1/up',options('POST',token,lease.body,{'X-Up-Seq':sequence},signal),remaining,2); + if(response.status!==204)throw failure('http','uplink replay rejected'); + if(response.headers.get('X-Up-Ack')!==sequence)throw failure('protocol','uplink replay acknowledgement rejected'); + }); + if(!recovered||closed||lease.cancelled||sessionToken!==token)return; + } + } if(!settleBatch(lease))return;port.postMessage({t:'traffic',up:lease.total,down:0});upSequence++;lease=null; } }catch(error){if(!closed&&!(lease&&lease.cancelled))fail(failureReason(error,'network'))} @@ -359,13 +279,17 @@ function sendCandidateSocket(next){ const state=next.telemt;if(!state||state.sent||next.readyState!==WebSocket.OPEN||!state.probe)return; let probe=state.probe; try{const fresh=findProbe(true);if(fresh&&fresh.id===probe.id)probe=fresh;next.send(probe.data)}catch(error){advanceCarrier('upgrade',state.epoch);return} - state.probe=probe;state.sent=true;if(!negotiationEnabled)commitCarrier(probe,state.epoch); + state.probe=probe;state.sent=true;if(!negotiationEnabled){if(state.openTimer)clearTimeout(state.openTimer);state.openTimer=null;commitCarrier(probe,state.epoch)} } function openCandidateSocket(probe,laneID,epoch){ let lane=laneID===null?null:ensureLane(laneID),next=lane?lane.socket:socket; if(next){if(!next.telemt||next.telemt.epoch!==epoch){advanceCarrier('protocol',epoch);return}if(probe)next.telemt.probe=probe;sendCandidateSocket(next);return} const token=sessionToken,protocol=laneID===null?(negotiationEnabled?'tproxy-auto-v1.':'tproxy-v1.')+token:(negotiationEnabled?'tproxy-auto-lane-v1.':'tproxy-lane-v1.')+token+'.'+String(laneID); - next=new WebSocket(socketURL(),protocol);next.binaryType='arraybuffer';next.telemt={epoch,lane,probe,opened:false,sent:false}; + next=new WebSocket(socketURL(),protocol);next.binaryType='arraybuffer';next.telemt={epoch,lane,probe,opened:false,sent:false,openTimer:null}; + next.telemt.openTimer=setTimeout(()=>{ + const state=next.telemt;if(closed||state.epoch!==attemptEpoch)return; + next.close();advanceCarrier('timeout',state.epoch); + },websocketOpenMs); if(lane)lane.socket=next;else socket=next; next.onopen=()=>{ const state=next.telemt;if(closed||state.epoch!==attemptEpoch){next.close();return}state.opened=true; @@ -373,18 +297,19 @@ function openCandidateSocket(probe,laneID,epoch){ }; next.onmessage=event=>{ const state=next.telemt;if(closed||state.epoch!==attemptEpoch||!(event.data instanceof ArrayBuffer))return; + if(state.openTimer)clearTimeout(state.openTimer);state.openTimer=null; if(!carrierCommitted){if(!state.sent||event.data.byteLength!==0){advanceCarrier('protocol',state.epoch);return}commitCarrier(state.probe,state.epoch);return} try{ if(state.lane){const values=splitFrames(event.data);for(const value of values)if(value.id!==state.lane.id)throw new Error('cross-lane frame');if(values.some(value=>value.type===3))state.lane.remoteClosed=true} else{const bound=frameBound(event.data,4096,batchLimit);if(bound.bytes!==event.data.byteLength)throw new Error('invalid frame batch')} }catch(error){if(state.lane)finishLane(state.lane,true);else fail('protocol');return} - port.postMessage({t:'traffic',up:0,down:event.data.byteLength});port.postMessage(event.data,[event.data]);status('connected'); + observeServerFrames(event.data);port.postMessage({t:'traffic',up:0,down:event.data.byteLength});port.postMessage(event.data,[event.data]);status('connected'); }; next.onerror=()=>{}; next.onclose=()=>{ - const state=next.telemt;if(state.epoch!==attemptEpoch||closed)return; + const state=next.telemt;if(state.openTimer)clearTimeout(state.openTimer);state.openTimer=null;if(state.epoch!==attemptEpoch||closed)return; if(!carrierCommitted){advanceCarrier(state.opened?'network':'upgrade',state.epoch);return} - if(state.lane){state.lane.ready=false;state.lane.socket=null;finishLane(state.lane,true)}else{socketReady=false;fail('network')} + if(state.lane){state.lane.ready=false;state.lane.socket=null;finishLane(state.lane,true)}else{socketReady=false;recoveryController.recover('network',null)} }; } function queueSocket(data){if(!reserve(data,null)){fail('capacity');return}upPending.push(data);runSocketUp()} @@ -400,21 +325,30 @@ async function runSocketUp(){ await waitSocket(socket,lease.total,queueLimit,lease.controller.signal);socket.send(lease.body); if(!settleBatch(lease))return;port.postMessage({t:'traffic',up:lease.total,down:0});lease=null; } - }catch(error){if(!closed&&!(lease&&lease.cancelled))fail(failureReason(error,'network'))} + }catch(error){if(!closed&&!(lease&&lease.cancelled))recoverTransport(error,null)} finally{upRunning=false;if(!closed&&socketReady&&upPending.length)runSocketUp()} } async function poll(){ while(!closed&&sessionToken){ + const token=sessionToken,cursor=downCursor; try{ pollController=new AbortController(); - const response=await request('/api/v1/down',options('POST',sessionToken,null,{'X-Down-Cursor':downCursor},pollController.signal)); + const response=await request('/api/v1/down',options('POST',token,null,{'X-Down-Cursor':cursor},pollController.signal),null,1); if(response.status===204){status('connected');continue} if(response.status!==200)throw failure('http','downlink rejected'); const next=response.headers.get('X-Down-Cursor')||'',data=response.body; if(!next||!data.byteLength)throw failure('protocol','invalid downlink response'); if(closed)return; - port.postMessage({t:'traffic',up:0,down:data.byteLength});port.postMessage(data,[data]);downCursor=next;status('connected'); - }catch(error){if(!closed)fail(failureReason(error,'network'));return} + observeServerFrames(data);port.postMessage({t:'traffic',up:0,down:data.byteLength});port.postMessage(data,[data]);downCursor=next;status('connected'); + }catch(error){ + if(closed)return; + const recovered=await recoverTransport(error,async(signal,remaining)=>{ + const response=await request('/api/v1/down',options('POST',token,null,{'X-Down-Cursor':cursor},signal),remaining,2); + if(response.status===204)return; + if(response.status!==200||!response.body.byteLength||!response.headers.get('X-Down-Cursor'))throw failure('http','downlink replay rejected'); + }); + if(!recovered||closed||sessionToken!==token)return; + } } } function ensureLane(id){ @@ -427,13 +361,12 @@ function rememberLaneClosed(id){ if(closedLaneOrder.length===closedLaneLimit)closedLanes.delete(closedLaneOrder.shift()); closedLanes.add(id);closedLaneOrder.push(id); } -function closeFrame(id){const value=new Uint8Array(8);value[0]=3;value[1]=(id>>>16)&255;value[2]=(id>>>8)&255;value[3]=id&255;return value.buffer} function finishLane(lane,notifyClient){ if(lanes.get(lane.id)!==lane)return; if(lane.controller)lane.controller.abort();lane.controller=null;cancelBatch(lane.upLease); if(lane.socket&&lane.socket.readyState{if(!closed&&lanes.get(lane.id)===lane&&lane.socket===opened)finishLane(lane,true)},websocketOpenMs); lane.socket.onopen=()=>{if(closed||lanes.get(lane.id)!==lane)return;lane.ready=true;status('connected');runLaneSocketUp(lane)}; lane.socket.onmessage=event=>{ + clearTimeout(openTimer); if(closed||lanes.get(lane.id)!==lane||!(event.data instanceof ArrayBuffer)){finishLane(lane,true);return} let values;try{values=splitFrames(event.data);for(const value of values)if(value.id!==lane.id)throw new Error('cross-lane frame')}catch(error){finishLane(lane,true);return} if(values.some(value=>value.type===3))lane.remoteClosed=true; - port.postMessage({t:'traffic',up:0,down:event.data.byteLength});port.postMessage(event.data,[event.data]);status('connected'); + observeServerFrames(event.data);port.postMessage({t:'traffic',up:0,down:event.data.byteLength});port.postMessage(event.data,[event.data]);status('connected'); }; - lane.socket.onerror=()=>{};lane.socket.onclose=()=>{lane.ready=false;lane.socket=null;if(!closed)finishLane(lane,true)}; + lane.socket.onerror=()=>{};lane.socket.onclose=()=>{clearTimeout(openTimer);lane.ready=false;lane.socket=null;if(!closed)finishLane(lane,true)}; } async function runLaneSocketUp(lane){ if(lane.running||!lane.ready)return;lane.running=true;let lease=null; @@ -472,10 +407,22 @@ async function runLaneUp(lane){ try{ while(!closed&&sessionToken&&lane.pending.length){ lease=takeBatch(lane.pending,lane);lane.upLease=lease;lease.controller=new AbortController(); - const sequence=String(lane.sequence),laneID=String(lane.id); - const response=await request('/api/v1/up',options('POST',sessionToken,lease.body,{'X-Up-Seq':sequence,'X-Lane-ID':laneID},lease.controller.signal)); - if(response.status!==204)throw failure('http','lane uplink rejected'); - if(response.headers.get('X-Up-Ack')!==sequence)throw failure('protocol','lane uplink acknowledgement rejected'); + const sequence=String(lane.sequence),laneID=String(lane.id),token=sessionToken; + for(;;){ + try{ + const response=await request('/api/v1/up',options('POST',token,lease.body,{'X-Up-Seq':sequence,'X-Lane-ID':laneID},lease.controller.signal),null,1); + if(response.status!==204)throw failure('http','lane uplink rejected'); + if(response.headers.get('X-Up-Ack')!==sequence)throw failure('protocol','lane uplink acknowledgement rejected'); + break; + }catch(error){ + const recovered=await recoverTransport(error,async(signal,remaining)=>{ + const response=await request('/api/v1/up',options('POST',token,lease.body,{'X-Up-Seq':sequence,'X-Lane-ID':laneID},signal),remaining,2); + if(response.status!==204)throw failure('http','lane uplink replay rejected'); + 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; + } + } if(!settleBatch(lease))return;port.postMessage({t:'traffic',up:lease.total,down:0});lane.sequence++;lease=null; if(!lane.polling)pollLane(lane); } @@ -483,11 +430,12 @@ async function runLaneUp(lane){ finally{lane.running=false;if(!closed&&lanes.get(lane.id)===lane&&sessionToken&&lane.pending.length)runLaneUp(lane)} } async function pollLane(lane){ - if(!lane||lane.polling)return;lane.polling=true; + if(!lane||lane.polling)return;lane.polling=true;let restart=false,failedToken='',failedCursor='',failedLaneID=''; try{ while(!closed&&sessionToken&&lanes.get(lane.id)===lane){ - const controller=new AbortController(),laneID=String(lane.id);lane.controller=controller; - const response=await request('/api/v1/down',options('POST',sessionToken,null,{'X-Down-Cursor':lane.cursor,'X-Lane-ID':laneID},controller.signal)); + const controller=new AbortController(),laneID=String(lane.id),token=sessionToken,cursor=lane.cursor;lane.controller=controller; + failedToken=token;failedCursor=cursor;failedLaneID=laneID; + const response=await request('/api/v1/down',options('POST',token,null,{'X-Down-Cursor':cursor,'X-Lane-ID':laneID},controller.signal),null,1); if(response.status===204){ if(response.headers.get('X-Lane-Closed')==='1'){finishLane(lane,false);return} status('connected');continue; @@ -497,35 +445,55 @@ async function pollLane(lane){ if(!next||!data.byteLength)throw failure('protocol','invalid lane downlink response'); for(const value of splitFrames(data))if(value.id!==lane.id)throw new Error('cross-lane frame'); if(closed)return; - port.postMessage({t:'traffic',up:0,down:data.byteLength});port.postMessage(data,[data]);lane.cursor=next;status('connected'); + observeServerFrames(data);port.postMessage({t:'traffic',up:0,down:data.byteLength});port.postMessage(data,[data]);lane.cursor=next;status('connected'); } - }catch(error){if(!closed)fail(failureReason(error,'network'))} - finally{lane.polling=false;lane.controller=null} + }catch(error){ + if(!closed&&lanes.get(lane.id)===lane){ + const recovered=await recoverTransport(error,async(signal,remaining)=>{ + const response=await request('/api/v1/down',options('POST',failedToken,null,{'X-Down-Cursor':failedCursor,'X-Lane-ID':failedLaneID},signal),remaining,2); + if(response.status===204)return; + if(response.status!==200||!response.body.byteLength||!response.headers.get('X-Down-Cursor'))throw failure('http','lane downlink replay rejected'); + }); + restart=recovered&&!closed&&sessionToken===failedToken&&lanes.get(lane.id)===lane; + } + } + finally{lane.polling=false;lane.controller=null;if(restart)pollLane(lane)} } function deleteSession(){ const token=cleanupToken||sessionToken,headers=canonicalFailures.includes(terminalFailure)?{'X-Carrier-Failure':terminalFailure}:null; if(token)fetch(relayOrigin+'/api/v1/session',options('DELETE',token,null,headers,undefined,true)).catch(()=>{}); } function close(notifyServer){ - if(closed)return;closed=true;if(helloTimer)clearTimeout(helloTimer);helloTimer=null;if(carrierTimer)clearTimeout(carrierTimer);clearProbeTimer();if(attemptController)attemptController.abort();if(pollController)pollController.abort(); + if(closed)return;closed=true;if(recoveryController)recoveryController.cancel();rejectRecoveryCommit(failure('network','bridge closed'));if(helloTimer)clearTimeout(helloTimer);helloTimer=null;if(carrierTimer)clearTimeout(carrierTimer);clearProbeTimer();if(attemptController)attemptController.abort();if(pollController)pollController.abort(); if(socket)socket.close();cancelBatch(upLease);releasePending(upPending,null); for(const lane of lanes.values()){ if(lane.controller)lane.controller.abort();cancelBatch(lane.upLease);releasePending(lane.pending,lane);if(lane.socket)lane.socket.close(); } - if(notifyServer)deleteSession();releasePending(pending,null);lanes.clear();if(port)port.close(); - if(queuedBytes!==0||queuedItems!==0)throw new Error('queue accounting leak'); + if(notifyServer)deleteSession();releasePending(pending,null);releasePending(recoveryPending,null);lanes.clear();clearStreams();if(port)port.close(); + buffers.assertEmpty(); } function activatePort(nextPort){ initialized=true;port=nextPort; port.onmessage=message=>{ + observeResumeTrigger(); if(message.data instanceof ArrayBuffer){ if(!createStarted){createStarted=true;if(helloTimer)clearTimeout(helloTimer);helloTimer=null;helloFrame=message.data;if(negotiationEnabled){negotiationStartedAt=Date.now();armCarrierDeadline(attemptEpoch)}createSession(attemptEpoch)} - else if(!carrierCommitted){if(!reserve(message.data,null)){fail('capacity');return}pending.push(message.data);maybeStartCandidate()} - else queueCarrier(message.data); + else{ + let data;try{data=acceptNativeFrames(message.data)}catch(error){fail(error&&error.telemtReason==='capacity'?'capacity':'protocol');return}if(!data)return; + if(recoveryController.active()&&!recoveryReplaced){if(!reserve(data,null)){fail('capacity');return}recoveryPending.push(data)} + else if(!carrierCommitted){if(!reserve(data,null)){fail('capacity');return}pending.push(data);maybeStartCandidate()} + else queueCarrier(data); + } }else if(message.data&&message.data.t==='close'){status('failed','terminal','closed',0);close(true)} }; port.start();status('connecting','starting','',bridgeRequestMs);helloTimer=setTimeout(()=>fail('timeout'),bridgeRequestMs); } +recoveryController=recoverySupport.create({ + budgetMs:()=>bridgeRecoveryMs,requestMs:()=>bridgeRequestMs,url:()=>relayOrigin+recoveryPath,token:()=>cleanupToken||sessionToken, + read:(response,limit,exact,signal)=>responseBody.read(response,limit,exact,signal),cancel:responseBody.cancel,status:remaining=>status('reconnecting','retrying','',remaining), + restored:finishOldRecovery,replace:replaceCarrier,replaceable:error=>failureReason(error,'network')!=='protocol', + reason:(error,fallback)=>failureReason(error,fallback),terminal:reason=>fail(recoveryController.remaining()<=0?'timeout':reason) +}); addEventListener('message',event=>{ if(event.source!==parent)return;if(initialized){if(event.ports&&event.ports.length===1)event.ports[0].close();return} if(event.data===null||typeof event.data!=='object')return; @@ -535,8 +503,7 @@ addEventListener('message',event=>{ if(source.protocol!=='http:'||source.hostname!=='127.0.0.1'||!source.port||source.origin!==event.origin)return; activatePort(event.ports[0]); },{once:false}); -const androidBridge=globalThis.TelegramWebProxy; -if(!initialized&&androidNonce&&androidBridge&&typeof androidBridge.postMessage==='function'){ +function activateAndroid(androidBridge){ const androidPort={onmessage:null,start(){},close(){androidBridge.onmessage=null},postMessage(value){ if(value instanceof ArrayBuffer){ let frames;try{frames=splitFrames(value)}catch(error){fail('protocol');return} @@ -546,5 +513,16 @@ if(!initialized&&androidNonce&&androidBridge&&typeof androidBridge.postMessage== androidBridge.onmessage=event=>{let data=event.data;if(typeof data==='string'){try{data=JSON.parse(data)}catch(error){return}}if(androidPort.onmessage)androidPort.onmessage({data})}; activatePort(androidPort);androidBridge.postMessage(JSON.stringify({t:'tproxy-android-init',v:1,nonce:androidNonce})); } +function discoverAndroid(){ + if(!androidNonce)return;const wall=Date.now()+bridgeRequestMs,monotonic=performance.now()+bridgeRequestMs; + const probe=()=>{ + if(initialized||closed)return;const androidBridge=globalThis.TelegramWebProxy; + if(androidBridge&&typeof androidBridge.postMessage==='function'){activateAndroid(androidBridge);return} + const remaining=Math.min(wall-Date.now(),monotonic-performance.now());if(remaining>0)setTimeout(probe,Math.min(100,remaining)); + };probe(); +} +discoverAndroid(); +addEventListener('online',observeResumeTrigger); +if(globalThis.document&&typeof globalThis.document.addEventListener==='function')globalThis.document.addEventListener('visibilitychange',()=>{if(globalThis.document.visibilityState==='visible')observeResumeTrigger()}); addEventListener('pagehide',()=>fail('navigation'),{once:true}); })(); diff --git a/src/web/bridge/tests.rs b/src/web/bridge/tests.rs index 6f157c9..2a07e43 100644 --- a/src/web/bridge/tests.rs +++ b/src/web/bridge/tests.rs @@ -7,12 +7,16 @@ fn render_page(bootstrap: &str, candidate_count: usize) -> BridgePage { 2 * 1024 * 1024, 32 * 1024 * 1024, 16 * 1024, + 1024, true, candidate_count, [3, 5, 8, 12], 25, 10, 90, + 15, + 15, + 120, 0, &SecureRandom::new(), ) @@ -47,7 +51,7 @@ fn rendered_page_preserves_the_ios_bootstrap_literal() { let page = render_page(bootstrap, 2); assert!( page.body - .contains(&format!("const bootstrap=\"{bootstrap}\"")) + .contains(&format!("let bootstrap=\"{bootstrap}\"")) ); } @@ -59,12 +63,16 @@ fn rendered_page_embeds_the_configured_bridge_timing_policy() { 2 * 1024 * 1024, 32 * 1024 * 1024, 16 * 1024, + 1024, true, 4, [3, 5, 8, 12], 17, 7, 41, + 13, + 11, + 119, 4, &SecureRandom::new(), ); @@ -72,6 +80,9 @@ fn rendered_page_embeds_the_configured_bridge_timing_policy() { assert!(page.body.contains("const 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( "helloTimer=setTimeout(()=>fail('timeout'),bridgeRequestMs)" @@ -99,12 +110,16 @@ fn disabled_negotiation_does_not_arm_a_carrier_deadline() { 2 * 1024 * 1024, 32 * 1024 * 1024, 16 * 1024, + 1024, false, 1, [3, 5, 8, 12], 25, 10, 90, + 15, + 15, + 120, 0, &SecureRandom::new(), ); @@ -122,7 +137,7 @@ fn retry_and_attempt_state_are_frozen_before_fetch() { let page = render_page("EEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEE", 4); assert!( page.body - .contains("async function request(path,frozenOptions)") + .contains("async function request(path,frozenOptions,remainingBudget,maxAttempts)") ); assert!(!page.body.contains("makeOptions")); assert!( diff --git a/src/web/http.rs b/src/web/http.rs index 20f86d4..c722c86 100644 --- a/src/web/http.rs +++ b/src/web/http.rs @@ -31,6 +31,8 @@ mod decoy; mod down; // Canonical request parsing rejects ambiguous credentials before routing. mod request; +// Positive-only recovery representation stays separate from ordinary bridge rendering. +mod recovery; // Carrier response construction and lane-header helpers are shared by handlers. mod response; // Session creation and replacement negotiation remain separate from request routing. @@ -219,6 +221,11 @@ async fn handle_root( generation: Arc, vhost: Arc, ) -> HttpResponse { + let representation = recovery::classify(&request); + if matches!(representation, recovery::RootRepresentation::Invalid) { + strip_query(&mut request); + return serve_decoy(request, vhost, true, &runtime).await; + } let (candidate, canonical) = bridge_candidate(request.uri().query()); let profile = match_profile(&vhost, &candidate); let Some(profile) = profile.filter(|_| canonical && request.method() == Method::GET) else { @@ -236,12 +243,25 @@ async fn handle_root( .headers() .get(header::USER_AGENT) .and_then(|value| value.to_str().ok()); - let bootstrap = match runtime.issue_bootstrap_for_request( - &generation, - Arc::clone(&profile), - client_ip, - user_agent, - ) { + let bootstrap = match match representation { + recovery::RootRepresentation::Bridge => runtime.issue_bootstrap_for_request( + &generation, + Arc::clone(&profile), + client_ip, + user_agent, + ), + recovery::RootRepresentation::Recovery(_) => runtime + .issue_recovery_bootstrap_for_request( + &generation, + Arc::clone(&profile), + client_ip, + user_agent, + ), + recovery::RootRepresentation::Invalid => { + strip_query(&mut request); + return serve_decoy(request, vhost, true, &runtime).await; + } + } { Ok(bootstrap) => bootstrap, Err(error) => { runtime.trace().record_profile_lifecycle( @@ -261,18 +281,51 @@ 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) + }) { + let outcome = session.close(crate::web::session::SessionCloseReason::BridgeRecovery); + if outcome != crate::web::session::SessionCloseOutcome::Closed + && tokio::time::timeout( + Duration::from_secs(config.web.timeouts.bridge_request_secs), + session.wait_close_complete(), + ) + .await + .is_err() + { + strip_query(&mut request); + return serve_decoy(request, vhost, true, &runtime).await; + } + } + let Some(response) = recovery::response( + &bootstrap, + &vhost, + &profile, + &config.web.limits, + &config.web.timeouts, + ) else { + strip_query(&mut request); + return serve_decoy(request, vhost, true, &runtime).await; + }; + return response; + } let page = bridge::render( &vhost.host, &bootstrap.token, config.web.limits.carrier_batch_bytes, config.web.limits.pending_bytes_per_session, config.web.limits.pending_items_per_session, + profile.max_streams_per_session, profile.carrier_negotiation_enabled, profile.carriers.len(), profile.carrier_negotiation_deadlines_secs, config.web.timeouts.long_poll_secs, config.web.timeouts.bridge_request_secs, config.web.timeouts.bridge_retry_secs, + config.web.timeouts.bridge_recovery_secs, + config.web.timeouts.websocket_open_secs, + config.web.timeouts.reconnect_grace_secs, config.web.timeouts.carrier_probe_coalesce_ms, &generation.rng, ); diff --git a/src/web/http/decoy.rs b/src/web/http/decoy.rs index 32873de..66fb139 100644 --- a/src/web/http/decoy.rs +++ b/src/web/http/decoy.rs @@ -297,6 +297,7 @@ fn lease_deadline( fn sanitize_transport_request(request: &mut Request) { for name in [ header::AUTHORIZATION, + header::ACCEPT, header::CONTENT_LENGTH, header::CONTENT_TYPE, header::UPGRADE, diff --git a/src/web/http/recovery.rs b/src/web/http/recovery.rs new file mode 100644 index 0000000..f1594f1 --- /dev/null +++ b/src/web/http/recovery.rs @@ -0,0 +1,141 @@ +use bytes::Bytes; +use hyper::header::{self, HeaderValue}; +use hyper::{Request, StatusCode}; +use serde::Serialize; + +use super::request::{bearer_token_hash, compatible_cookie_header}; +use super::response::{full_response, insert_header}; +use super::{HttpResponse, RequestBody}; +use crate::config::{WebRuntimeProfile, WebRuntimeVhost}; +use crate::web::manager::{BootstrapResult, TokenHash}; + +pub(super) const MEDIA_TYPE: &str = "application/vnd.telemt.web-recovery+json"; +const MAX_RESPONSE_BYTES: usize = 1024; + +/// Positive representation selected by exact recovery headers. +#[derive(Clone, Copy)] +pub(super) enum RootRepresentation { + /// Render the ordinary transient bridge document. + Bridge, + /// Return fresh recovery policy and optionally retire one current bearer. + Recovery(Option), + /// Hide malformed recovery material behind the configured decoy. + Invalid, +} + +/// Classifies recovery headers without changing capability authentication. +pub(super) fn classify(request: &Request) -> RootRepresentation { + let accepts = request.headers().get_all(header::ACCEPT); + let mut values = accepts.iter(); + let first = values.next(); + let exact = first.is_some_and(|value| value.as_bytes() == MEDIA_TYPE.as_bytes()) + && values.next().is_none(); + let recovery_present = accepts + .iter() + .any(|value| value.as_bytes() == MEDIA_TYPE.as_bytes()); + let authorization_present = request.headers().contains_key(header::AUTHORIZATION); + if !exact { + return if authorization_present || recovery_present { + RootRepresentation::Invalid + } else { + RootRepresentation::Bridge + }; + } + if !compatible_cookie_header(request) { + return RootRepresentation::Invalid; + } + if !authorization_present { + return RootRepresentation::Recovery(None); + } + bearer_token_hash(request) + .map(|hash| RootRepresentation::Recovery(Some(hash))) + .unwrap_or(RootRepresentation::Invalid) +} + +/// Builds the bounded no-store recovery representation. +pub(super) fn response( + bootstrap: &BootstrapResult, + vhost: &WebRuntimeVhost, + profile: &WebRuntimeProfile, + limits: &crate::config::WebLimitsConfig, + timeouts: &crate::config::WebTimeoutsConfig, +) -> Option { + let document = RecoveryDocument { + version: 1, + bootstrap: &bootstrap.token, + limits: RecoveryLimits { + carrier_batch_bytes: limits.carrier_batch_bytes, + pending_bytes_per_session: limits.pending_bytes_per_session, + pending_items_per_session: limits.pending_items_per_session, + max_streams_per_session: profile.max_streams_per_session, + }, + timeouts: RecoveryTimeouts { + long_poll_secs: timeouts.long_poll_secs, + bridge_request_secs: timeouts.bridge_request_secs, + bridge_retry_secs: timeouts.bridge_retry_secs, + bridge_recovery_secs: timeouts.bridge_recovery_secs, + websocket_open_secs: timeouts.websocket_open_secs, + reconnect_grace_secs: timeouts.reconnect_grace_secs, + }, + negotiation: RecoveryNegotiation { + enabled: profile.carrier_negotiation_enabled, + candidate_count: profile.carriers.len(), + deadlines_secs: profile.carrier_negotiation_deadlines_secs, + carrier_probe_coalesce_ms: timeouts.carrier_probe_coalesce_ms, + }, + }; + let body = serde_json::to_vec(&document).ok()?; + if body.len() > MAX_RESPONSE_BYTES || profile.host != vhost.host { + return None; + } + let mut response = full_response(StatusCode::OK, Bytes::from(body)); + insert_header(&mut response, header::CONTENT_TYPE, MEDIA_TYPE); + response + .headers_mut() + .insert(header::CACHE_CONTROL, HeaderValue::from_static("no-store")); + response.headers_mut().insert( + header::REFERRER_POLICY, + HeaderValue::from_static("no-referrer"), + ); + response.headers_mut().insert( + header::X_CONTENT_TYPE_OPTIONS, + HeaderValue::from_static("nosniff"), + ); + Some(response) +} + +#[derive(Serialize)] +struct RecoveryDocument<'a> { + #[serde(rename = "v")] + version: u8, + bootstrap: &'a str, + limits: RecoveryLimits, + timeouts: RecoveryTimeouts, + negotiation: RecoveryNegotiation, +} + +#[derive(Serialize)] +struct RecoveryLimits { + carrier_batch_bytes: usize, + pending_bytes_per_session: usize, + pending_items_per_session: usize, + max_streams_per_session: usize, +} + +#[derive(Serialize)] +struct RecoveryTimeouts { + long_poll_secs: u64, + bridge_request_secs: u64, + bridge_retry_secs: u64, + bridge_recovery_secs: u64, + websocket_open_secs: u64, + reconnect_grace_secs: u64, +} + +#[derive(Serialize)] +struct RecoveryNegotiation { + enabled: bool, + candidate_count: usize, + deadlines_secs: [u64; 4], + carrier_probe_coalesce_ms: u64, +} diff --git a/src/web/http/websocket/driver.rs b/src/web/http/websocket/driver.rs index 72c3d14..701bc6e 100644 --- a/src/web/http/websocket/driver.rs +++ b/src/web/http/websocket/driver.rs @@ -10,7 +10,9 @@ use tokio_util::sync::CancellationToken; use super::ConnectionIo; use crate::web::http::activity::UpgradeDeadlineLease; use crate::web::manager::{WebProcessRuntime, WebSocketBudgetLease, WebSocketConnection}; -use crate::web::session::{WebSession, WebSocketLaneReservation, WebSocketProbeReservation}; +use crate::web::session::{ + SessionCloseReason, WebSession, WebSocketLaneReservation, WebSocketProbeReservation, +}; use crate::web::trace::{TraceDirection, TraceWebSocketContext}; const READ_BUFFER_BYTES: usize = 64 * 1024; @@ -103,7 +105,7 @@ pub(super) async fn run_upgraded( if let Some(reservation) = lane_reservation { session.close_websocket_lane(reservation); } else if !acknowledge_commit || session.is_carrier_committed() { - session.close(); + session.close(SessionCloseReason::WebSocketEnded); } } @@ -194,7 +196,7 @@ async fn run_multiplex( started, ); if !session.websocket_commit_ack_written(connection.id()) { - session.close(); + session.close(SessionCloseReason::Protocol); return Err(()); } } else if acknowledge_commit diff --git a/src/web/http/websocket/driver/lane.rs b/src/web/http/websocket/driver/lane.rs index d6542a2..092c19a 100644 --- a/src/web/http/websocket/driver/lane.rs +++ b/src/web/http/websocket/driver/lane.rs @@ -8,7 +8,7 @@ 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::session::{SessionCloseReason, WebSession, WebSocketLaneReservation}; use crate::web::trace::{TraceDirection, TraceWebSocketContext}; #[allow(clippy::too_many_arguments)] @@ -92,7 +92,7 @@ pub(super) async fn run_lane( .await .is_err() { - session.close(); + session.close(SessionCloseReason::Protocol); return Err(()); } record_message( @@ -104,7 +104,7 @@ pub(super) async fn run_lane( started, ); if !session.websocket_commit_ack_written(connection.id()) { - session.close(); + session.close(SessionCloseReason::Protocol); return Err(()); } } else if acknowledge_commit diff --git a/src/web/manager/carrier_outcome.rs b/src/web/manager/carrier_outcome.rs index 86a5c23..d6f237d 100644 --- a/src/web/manager/carrier_outcome.rs +++ b/src/web/manager/carrier_outcome.rs @@ -99,7 +99,7 @@ impl WebProcessRuntime { identity: TraceIdentity, ) -> bool { let mut state = self.state.lock(); - let scores = state.bootstraps.get_mut(&bootstrap_hash).and_then(|entry| { + let outcome = state.bootstraps.get_mut(&bootstrap_hash).and_then(|entry| { if entry.carrier_attempt == attempt && entry .session @@ -108,13 +108,20 @@ impl WebProcessRuntime { && entry.carrier_phase == CarrierChainPhase::Provisional { entry.carrier_phase = CarrierChainPhase::CommittedPendingHealth; - Some(entry.carrier_scores) + Some((entry.carrier_scores, entry.recovery)) } else { None } }); drop(state); - let Some(scores) = scores else { return false }; + let Some((scores, recovery)) = outcome else { + return false; + }; + if recovery { + self.telemetry.record_bridge_recovery( + crate::web::telemetry::WebBridgeRecoveryEvent::Committed, + ); + } self.trace.record_carrier_lifecycle( client_ip, identity.clone(), diff --git a/src/web/manager/control.rs b/src/web/manager/control.rs index d28a54f..4100450 100644 --- a/src/web/manager/control.rs +++ b/src/web/manager/control.rs @@ -7,6 +7,7 @@ use serde::Serialize; use super::status::immutable_matches; use super::{SessionFilter, WebProcessRuntime}; +use crate::web::session::SessionCloseReason; const OPERATION_REF_VERSION: &str = "wo1"; const OPERATION_RETENTION: usize = 32; @@ -297,7 +298,7 @@ impl WebProcessRuntime { status.matched = status.matched.saturating_add(direct.len()); }); for candidate in direct { - candidate.session.close(); + candidate.session.close(SessionCloseReason::ApiClose); self.update_operation(sequence, |status| { status.close_signalled = status.close_signalled.saturating_add(1) }); @@ -325,7 +326,7 @@ impl WebProcessRuntime { current }; if let Some(session) = session { - session.close(); + session.close(SessionCloseReason::ApiClose); self.update_operation(sequence, |status| { status.matched = status.matched.saturating_add(1); status.close_signalled = status.close_signalled.saturating_add(1); diff --git a/src/web/manager/credentials.rs b/src/web/manager/credentials.rs index 2b20a7c..8c70837 100644 --- a/src/web/manager/credentials.rs +++ b/src/web/manager/credentials.rs @@ -7,13 +7,13 @@ use zeroize::Zeroizing; use super::state::{ Bootstrap, CarrierChainPhase, allow_rate, evict_oldest_unused_bootstrap, matching_profile, - new_unique_token, remove_expired_locked, + new_unique_token, profile_key, remove_expired_locked, }; use super::{BootstrapResult, ManagerError, TOKEN_BYTES, TokenHash, WebProcessRuntime}; use crate::config::WebRuntimeProfile; use crate::maestro::generation::RuntimeGeneration; -use crate::web::session::WebSession; -use crate::web::telemetry::WebRejectionReason; +use crate::web::session::{SessionCloseReason, WebSession}; +use crate::web::telemetry::{WebBridgeRecoveryEvent, WebRejectionReason}; impl WebProcessRuntime { /// Issues a one-use bootstrap credential for an active compatible profile. @@ -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) + self.issue_bootstrap_inner(&generation, profile, client_ip, None, false) } /// 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) + self.issue_bootstrap_inner(generation, profile, client_ip, None, false) } /// Issues one bridge bootstrap with bounded non-secret request metadata. @@ -46,7 +46,18 @@ impl WebProcessRuntime { client_ip: IpAddr, user_agent: Option<&str>, ) -> std::result::Result { - self.issue_bootstrap_inner(generation, profile, client_ip, user_agent) + self.issue_bootstrap_inner(generation, profile, client_ip, user_agent, false) + } + + /// Issues one recovery bootstrap with the same positive-only admission boundary. + pub(crate) fn issue_recovery_bootstrap_for_request( + &self, + generation: &Arc, + profile: Arc, + client_ip: IpAddr, + user_agent: Option<&str>, + ) -> std::result::Result { + self.issue_bootstrap_inner(generation, profile, client_ip, user_agent, true) } fn issue_bootstrap_inner( @@ -55,6 +66,7 @@ impl WebProcessRuntime { profile: Arc, client_ip: IpAddr, user_agent: Option<&str>, + recovery: bool, ) -> std::result::Result { let config = generation.config(); let profile = config @@ -76,7 +88,11 @@ impl WebProcessRuntime { let _operator_admission = self.try_operator_admission()?; let now = Instant::now(); let mut state = self.state.lock(); - remove_expired_locked(&mut state, now); + let expired_recoveries = remove_expired_locked(&mut state, now); + self.telemetry.record_bridge_recovery_count( + WebBridgeRecoveryEvent::ExpiredUnused, + expired_recoveries, + ); state.apply_issuance_policy(generation.id, config.web.enabled); if state.closed || !state.issuance_enabled { self.record_limit_hit(); @@ -153,6 +169,7 @@ impl WebProcessRuntime { session_client_ip: None, session_ip_learning_eligible: false, used: false, + recovery, }, ); *state.bootstraps_per_ip.entry(client_ip).or_insert(0) += 1; @@ -162,6 +179,10 @@ impl WebProcessRuntime { .map(|entry| Arc::clone(&entry.profile)) .ok_or(ManagerError::Closed)?; drop(state); + if recovery { + self.telemetry + .record_bridge_recovery(WebBridgeRecoveryEvent::BootstrapIssued); + } self.trace.record_profile_lifecycle( client_ip, Some(trace_session_id), @@ -215,6 +236,24 @@ impl WebProcessRuntime { .ok_or(ManagerError::Authentication) } + /// Resolves a current bearer only when it belongs to the recovering profile. + pub(crate) fn bridge_recovery_session( + &self, + hash: TokenHash, + host: &str, + profile: &WebRuntimeProfile, + ) -> Option> { + let expected_profile = profile_key(profile); + self.state + .lock() + .sessions + .get(&hash) + .filter(|session| { + session.matches_host(host) && session.profile_key() == expected_profile + }) + .cloned() + } + /// Closes a live token and accepts bounded tombstone retries. pub(crate) fn close_token( &self, @@ -256,7 +295,7 @@ impl WebProcessRuntime { .is_some_and(|closed| closed.host == host); drop(state); if let Some(session) = session { - if session.close() + if session.close(SessionCloseReason::ClientDelete).accepted() && let (Some(failure), Some(phase)) = (failure, failure_phase) { self.telemetry diff --git a/src/web/manager/lifecycle.rs b/src/web/manager/lifecycle.rs index e924901..2350ea0 100644 --- a/src/web/manager/lifecycle.rs +++ b/src/web/manager/lifecycle.rs @@ -11,6 +11,7 @@ use super::state::{ }; use super::{ProfileKey, TokenHash, WebProcessRuntime}; use crate::maestro::generation::RuntimeGeneration; +use crate::web::session::SessionCloseReason; /// Result of draining all process-owned WEB work under one absolute deadline. #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -61,17 +62,29 @@ impl WebProcessRuntime { profile_key: ProfileKey, profile_host: &str, closed_token_lifetime: Duration, + reason: SessionCloseReason, ) { let mut state = self.state.lock(); let Some(session) = state.sessions.remove(&hash) else { return; }; + let recovery_closed_before_commit = state.bootstraps.values().any(|bootstrap| { + bootstrap.recovery + && bootstrap + .session + .as_ref() + .is_some_and(|current| current.token_hash() == hash) + && !session.is_carrier_committed() + }); decrement_map(&mut state.sessions_per_ip, &client_ip); decrement_map(&mut state.sessions_per_profile, &profile_key); remember_closed_token_locked( &mut state, hash, profile_host, + session.trace_session_id(), + session.carrier(), + reason, closed_token_lifetime, self.limits.max_sessions_global.saturating_mul(16), ); @@ -86,6 +99,8 @@ impl WebProcessRuntime { &mut state, trace_session_id, session.carrier_attempt(), + session.carrier(), + reason, closed_token_lifetime, self.limits.max_sessions_global, ); @@ -104,7 +119,13 @@ impl WebProcessRuntime { for bootstrap_hash in bootstrap_hashes { remove_bootstrap_locked(&mut state, bootstrap_hash); } - self.telemetry.record_session_closed(); + self.telemetry + .record_session_closed(session.carrier(), reason); + if recovery_closed_before_commit { + self.telemetry.record_bridge_recovery( + crate::web::telemetry::WebBridgeRecoveryEvent::ClosedBeforeCommit, + ); + } drop(state); self.notify_operator_work_changed(); } @@ -135,7 +156,7 @@ impl WebProcessRuntime { }; self.stream_admission.lock().closed = true; for session in &sessions { - session.close(); + session.close(SessionCloseReason::RuntimeShutdown); } self.tasks.close(); WebShutdownDrain { @@ -215,14 +236,18 @@ impl WebProcessRuntime { for (hash, _) in expired { remove_bootstrap_locked(&mut state, hash); } - remove_expired_locked(&mut state, now); + let expired_recoveries = remove_expired_locked(&mut state, now); + self.telemetry.record_bridge_recovery_count( + crate::web::telemetry::WebBridgeRecoveryEvent::ExpiredUnused, + expired_recoveries, + ); ( state.sessions.values().cloned().collect::>(), expired_chains, ) }; for session in expired_chains { - session.close(); + session.close(SessionCloseReason::NegotiationTimeout); } for session in sessions { session.close_if_due(now); diff --git a/src/web/manager/operator_lifecycle.rs b/src/web/manager/operator_lifecycle.rs index fed6171..499d5f6 100644 --- a/src/web/manager/operator_lifecycle.rs +++ b/src/web/manager/operator_lifecycle.rs @@ -9,6 +9,7 @@ use tokio::time::Instant as TokioInstant; use tokio_util::sync::CancellationToken; use super::WebProcessRuntime; +use crate::web::session::SessionCloseReason; // Serialized state-machine values stay separate from synchronization mechanics. mod status; @@ -447,7 +448,7 @@ impl WebProcessRuntime { } forced = true; for session in sessions { - session.close(); + session.close(SessionCloseReason::OperatorForce); } } } diff --git a/src/web/manager/session_creation.rs b/src/web/manager/session_creation.rs index 67c3674..cb76114 100644 --- a/src/web/manager/session_creation.rs +++ b/src/web/manager/session_creation.rs @@ -18,7 +18,7 @@ use super::{ }; use crate::config::{WebCarrier, WebRuntimeProfile}; use crate::web::frame; -use crate::web::session::WebSession; +use crate::web::session::{SessionCloseReason, WebSession}; use crate::web::telemetry::WebCarrierSelectionDisposition; use crate::web::trace::TraceLifecycleEvent; @@ -56,7 +56,11 @@ impl WebProcessRuntime { let config = generation.config(); let now = Instant::now(); let mut state = self.state.lock(); - remove_expired_locked(&mut state, now); + let expired_recoveries = remove_expired_locked(&mut state, now); + self.telemetry.record_bridge_recovery_count( + crate::web::telemetry::WebBridgeRecoveryEvent::ExpiredUnused, + expired_recoveries, + ); state.apply_issuance_policy(generation.id, config.web.enabled); if state.closed || !state.issuance_enabled { return Err(ManagerError::Closed); @@ -79,7 +83,7 @@ impl WebProcessRuntime { let session = entry.session.clone(); drop(state); if let Some(session) = session { - session.close(); + session.close(SessionCloseReason::NegotiationTimeout); } return Err(ManagerError::Closed); } @@ -335,7 +339,7 @@ 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) = { + let (issuance_ip, candidate_count, user_agent, user_agent_id, recovery) = { let entry = state .bootstraps .get_mut(&bootstrap_hash) @@ -362,10 +366,16 @@ impl WebProcessRuntime { u8::try_from(entry.carrier_candidates.len()).unwrap_or(4), entry.user_agent.clone(), entry.user_agent_id, + entry.recovery, ) }; decrement_map(&mut state.bootstraps_per_ip, &issuance_ip); self.telemetry.record_session_created(); + if recovery { + self.telemetry.record_bridge_recovery( + crate::web::telemetry::WebBridgeRecoveryEvent::SessionCreated, + ); + } let identity = session.trace_identity(); let result = CreateResult { token: session_token, diff --git a/src/web/manager/session_creation/replacement.rs b/src/web/manager/session_creation/replacement.rs index 985a790..bd8d46c 100644 --- a/src/web/manager/session_creation/replacement.rs +++ b/src/web/manager/session_creation/replacement.rs @@ -20,7 +20,11 @@ impl WebProcessRuntime { let config = generation.config(); let now = Instant::now(); let mut state = self.state.lock(); - remove_expired_locked(&mut state, now); + let expired_recoveries = remove_expired_locked(&mut state, now); + self.telemetry.record_bridge_recovery_count( + crate::web::telemetry::WebBridgeRecoveryEvent::ExpiredUnused, + expired_recoveries, + ); state.apply_issuance_policy(generation.id, config.web.enabled); let valid = state.bootstraps.get(&bootstrap_hash).is_some_and(|entry| { entry.carrier_transitioning @@ -85,7 +89,7 @@ impl WebProcessRuntime { let Some(supersede) = replacement.old_session.prepare_carrier_supersede() else { drop(state); self.cancel_replacement(bootstrap_hash, &replacement.old_session); - session.close(); + session.close(crate::web::session::SessionCloseReason::Protocol); return Err(ManagerError::Closed); }; let old_hash = replacement.old_session.token_hash(); @@ -94,6 +98,9 @@ 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), ); @@ -115,7 +122,15 @@ impl WebProcessRuntime { *slot = Some(replacement.old_session.carrier()); } self.telemetry.record_session_created(); - self.telemetry.record_session_closed(); + if entry.recovery { + self.telemetry.record_bridge_recovery( + crate::web::telemetry::WebBridgeRecoveryEvent::SessionCreated, + ); + } + self.telemetry.record_session_closed( + replacement.old_session.carrier(), + crate::web::session::SessionCloseReason::CarrierSuperseded, + ); let result = CreateResult { token: session_token, carrier: replacement.carrier, diff --git a/src/web/manager/state.rs b/src/web/manager/state.rs index 71c84d5..9b93a15 100644 --- a/src/web/manager/state.rs +++ b/src/web/manager/state.rs @@ -86,6 +86,8 @@ pub(super) struct Bootstrap { pub(super) session_ip_learning_eligible: bool, /// Distinguishes unused issuance quota from completed creation replay state. pub(super) used: bool, + /// Whether this credential was issued by the post-commit recovery representation. + pub(super) recovery: bool, } /// Bounded replay marker for one explicitly or naturally closed session token. @@ -94,6 +96,12 @@ 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. @@ -116,6 +124,12 @@ pub(super) struct ClosedSession { pub(super) expires_at: Instant, /// Last carrier incarnation closed for this logical session. pub(super) attempt: u8, + /// Last carrier owning this logical session. + pub(super) carrier: WebCarrier, + /// First-writer terminal close cause. + pub(super) reason: crate::web::session::SessionCloseReason, + /// Monotonic instant used only to report bounded close age. + pub(super) closed_at: Instant, } /// Token-bucket state for one process-wide creation class. @@ -303,13 +317,19 @@ pub(super) fn evict_oldest_unused_bootstrap(state: &mut ManagerState) -> bool { } /// Removes expired bootstrap and closed-token entries while the manager lock is held. -pub(super) fn remove_expired_locked(state: &mut ManagerState, now: Instant) { +pub(super) fn remove_expired_locked(state: &mut ManagerState, now: Instant) -> usize { let expired = state .bootstraps .iter() - .filter_map(|(hash, bootstrap)| (now > bootstrap.expires_at).then_some(*hash)) + .filter_map(|(hash, bootstrap)| { + (now > bootstrap.expires_at).then_some((*hash, bootstrap.recovery && !bootstrap.used)) + }) .collect::>(); - for hash in expired { + let recovery_unused = expired + .iter() + .filter(|(_, recovery_unused)| *recovery_unused) + .count(); + for (hash, _) in expired { remove_bootstrap_locked(state, hash); } state @@ -325,6 +345,7 @@ pub(super) fn remove_expired_locked(state: &mut ManagerState, now: Instant) { state.closed_sessions.remove(&trace_session_id); } } + recovery_unused } /// Removes one bootstrap and releases its per-address issuance quota when unused. @@ -342,6 +363,9 @@ 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, ) { @@ -350,6 +374,9 @@ 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 { @@ -370,6 +397,8 @@ pub(super) fn remember_closed_session_locked( state: &mut ManagerState, trace_session_id: u64, attempt: u8, + carrier: WebCarrier, + reason: crate::web::session::SessionCloseReason, lifetime: Duration, capacity: usize, ) { @@ -380,6 +409,9 @@ pub(super) fn remember_closed_session_locked( ClosedSession { expires_at: Instant::now() + lifetime, attempt, + carrier, + reason, + closed_at: Instant::now(), }, ) .is_none() diff --git a/src/web/manager/status.rs b/src/web/manager/status.rs index 297e915..49622b6 100644 --- a/src/web/manager/status.rs +++ b/src/web/manager/status.rs @@ -198,7 +198,16 @@ pub(crate) enum SessionDetail { /// One exact live-session snapshot. Active(Box), /// One bounded retained closed-session tombstone. - Gone { attempt: u8 }, + Gone { + /// Last carrier incarnation number. + attempt: u8, + /// Carrier that owned the final incarnation. + carrier: WebCarrier, + /// First-writer terminal close cause. + reason: &'static str, + /// Monotonic age since final closure. + closed_age_ms: u64, + }, /// A required short lock was contended. Busy, /// Neither a live session nor a retained tombstone exists. @@ -500,13 +509,16 @@ impl WebProcessRuntime { .map(|status| SessionDetail::Active(Box::new(self.row(candidate, status)))) .unwrap_or(SessionDetail::Busy); } - let closed = state - .closed_sessions - .get(&trace_session_id) - .map(|closed| closed.attempt); - closed.map_or(SessionDetail::NotFound, |attempt| SessionDetail::Gone { - attempt, - }) + let now = Instant::now(); + state.closed_sessions.get(&trace_session_id).map_or( + SessionDetail::NotFound, + |closed| SessionDetail::Gone { + attempt: closed.attempt, + carrier: closed.carrier, + reason: closed.reason.as_str(), + closed_age_ms: millis(now.saturating_duration_since(closed.closed_at)), + }, + ) } fn row(&self, candidate: Candidate, status: WebSessionStatus) -> SessionRow { @@ -524,3 +536,7 @@ mod reference; pub(crate) use reference::SessionRefError; pub(super) use reference::immutable_matches; use reference::permits; + +fn millis(duration: std::time::Duration) -> u64 { + duration.as_millis().min(u128::from(u64::MAX)) as u64 +} diff --git a/src/web/session.rs b/src/web/session.rs index f4b70e5..27ae1a5 100644 --- a/src/web/session.rs +++ b/src/web/session.rs @@ -20,6 +20,9 @@ use crate::web::manager::{ // Backend tasks own generation admission and authenticated MTProxy relay lifetimes. mod backend; +// Activity clocks separate authenticated peer leases from diagnostic progress. +mod activity; +use activity::SessionActivity; // Downlink queues own cursor replay, flow control, and memory reservations. mod downlink; // Response ownership keeps detached batches charged until the last body clone drops. @@ -42,6 +45,8 @@ mod negotiation; use negotiation::CarrierHealthPublicationState; // Session closure and carrier-attempt transitions share one cancellation boundary. mod lifecycle; +pub(crate) use lifecycle::{SessionCloseOutcome, SessionCloseReason}; +use lifecycle::SessionNegotiationPhase; // Uplink batches own exactly-once sequencing and client-frame validation. mod uplink; @@ -169,7 +174,7 @@ struct SessionState { pending_items: usize, pending_control_bytes: usize, pending_control_items: usize, - last_activity: Instant, + activity: SessionActivity, negotiation_phase: SessionNegotiationPhase, carrier_health_due_at: Option, carrier_health_activity_at: Option, @@ -181,18 +186,10 @@ struct SessionState { websocket_commit_ack_owner: Option, websocket_commit_ack_written: bool, websocket_probe_claimed: bool, - close_requested: bool, + close_requested: Option, closed: bool, } -#[derive(Clone, Copy, PartialEq, Eq)] -enum SessionNegotiationPhase { - Uncommitted, - Replacing, - Committed, - Superseded, -} - /// One bounded WEB carrier session containing logical MTProxy streams. pub(crate) struct WebSession { manager: std::sync::Weak, @@ -213,6 +210,8 @@ pub(crate) struct WebSession { timeouts: WebTimeoutsConfig, state: Mutex, carrier_health_publication: AtomicU8, + close_complete: AtomicBool, + close_notify: Notify, down_notify: Arc, lane_open_notify: Arc, cancel: CancellationToken, @@ -253,6 +252,7 @@ impl WebSession { limits: WebLimitsConfig, timeouts: WebTimeoutsConfig, ) -> Arc { + let created_at = Instant::now(); let mut carrier_lanes = HashMap::new(); let mut next_lane_instance = 1; if selected_carrier == WebCarrier::HttpsLanes { @@ -273,7 +273,7 @@ impl WebSession { carrier_class, learning_context, automatic_carrier, - created_at: Instant::now(), + created_at, limits, timeouts, state: Mutex::new(SessionState { @@ -298,7 +298,7 @@ impl WebSession { pending_items: 0, pending_control_bytes: 0, pending_control_items: 0, - last_activity: Instant::now(), + activity: SessionActivity::new(created_at), negotiation_phase: SessionNegotiationPhase::Uncommitted, carrier_health_due_at: None, carrier_health_activity_at: None, @@ -310,12 +310,14 @@ impl WebSession { websocket_commit_ack_owner: None, websocket_commit_ack_written: false, websocket_probe_claimed: false, - close_requested: false, + close_requested: None, closed: false, }), carrier_health_publication: AtomicU8::new( CarrierHealthPublicationState::Awaiting as u8, ), + close_complete: AtomicBool::new(false), + close_notify: Notify::new(), down_notify: Arc::new(Notify::new()), lane_open_notify: Arc::new(Notify::new()), cancel: CancellationToken::new(), @@ -431,7 +433,7 @@ impl WebSession { self.release_locked(&mut state, count + overhead, usize::from(finished), false); if !self.queue_window_locked(&mut state, stream.id, count as u32) { drop(state); - self.close(); + self.close(SessionCloseReason::Backpressure); return Poll::Ready(Err(io::Error::other( "WEB session control budget exhausted", ))); @@ -497,7 +499,7 @@ impl WebSession { ))); }; stream_state.send_credit -= count as u64; - state.last_activity = Instant::now(); + state.activity.touch_progress(Instant::now()); drop(state); if self.carrier().is_multiplexed() { self.down_notify.notify_waiters(); diff --git a/src/web/session/activity.rs b/src/web/session/activity.rs new file mode 100644 index 0000000..2e967ed --- /dev/null +++ b/src/web/session/activity.rs @@ -0,0 +1,40 @@ +use std::time::{Duration, Instant}; + +/// Session-local activity clocks with distinct lease and diagnostic authority. +pub(super) struct SessionActivity { + last_peer: Instant, + last_progress: Instant, +} + +impl SessionActivity { + /// Starts both activity clocks at the same session creation instant. + pub(super) fn new(now: Instant) -> Self { + Self { + last_peer: now, + last_progress: now, + } + } + + /// Records one validated peer operation and returns the preceding peer gap. + pub(super) fn touch_peer(&mut self, now: Instant) -> Duration { + let gap = now.saturating_duration_since(self.last_peer); + self.last_peer = now; + self.last_progress = now; + gap + } + + /// Records server-side carrier progress without extending the peer lease. + pub(super) fn touch_progress(&mut self, now: Instant) { + self.last_progress = now; + } + + /// Returns elapsed time since the latest validated peer operation. + pub(super) fn peer_idle(&self, now: Instant) -> Duration { + now.saturating_duration_since(self.last_peer) + } + + /// Returns elapsed time since the latest carrier-side progress. + pub(super) fn progress_idle(&self, now: Instant) -> Duration { + now.saturating_duration_since(self.last_progress) + } +} diff --git a/src/web/session/backend.rs b/src/web/session/backend.rs index c9ea00f..173f75c 100644 --- a/src/web/session/backend.rs +++ b/src/web/session/backend.rs @@ -7,7 +7,7 @@ use crate::proxy::shared_state::ConntrackClosePolicy; use crate::web::frame::FrameType; use crate::web::stream::WebLogicalStream; -use super::{StreamIdentity, WebSession, inbound_queue_cost}; +use super::{SessionCloseReason, StreamIdentity, WebSession, inbound_queue_cost}; #[cfg(test)] #[path = "backend_tests.rs"] @@ -119,7 +119,7 @@ impl WebSession { }) }; if queued.is_some_and(|queued| !queued) { - self.close(); + self.close(SessionCloseReason::Backpressure); } } @@ -156,7 +156,7 @@ impl WebSession { } if let Some(queued) = queued { if !queued { - self.close(); + self.close(SessionCloseReason::Backpressure); } if self.carrier().is_multiplexed() { self.down_notify.notify_waiters(); diff --git a/src/web/session/backend_tests.rs b/src/web/session/backend_tests.rs index 6c9d5a3..9a07c9c 100644 --- a/src/web/session/backend_tests.rs +++ b/src/web/session/backend_tests.rs @@ -47,7 +47,7 @@ impl TestRuntime { } async fn shutdown(self) { - self.session.close(); + self.session.close(super::SessionCloseReason::ApiClose); self.session.wait().await; self.manager.shutdown().await; self.generation.stop_sessions().await; @@ -402,7 +402,7 @@ async fn cancellation_while_waiting_for_data_releases_stream_ownership() { settle_tasks().await; assert_eq!(runtime.session.tasks_live.load(Ordering::Acquire), 1); - runtime.session.close(); + runtime.session.close(super::SessionCloseReason::ApiClose); runtime.session.wait().await; assert_eq!(runtime.session.tasks_live.load(Ordering::Acquire), 0); diff --git a/src/web/session/downlink.rs b/src/web/session/downlink.rs index b7b5f3a..febecc2 100644 --- a/src/web/session/downlink.rs +++ b/src/web/session/downlink.rs @@ -5,7 +5,8 @@ use bytes::{BufMut, Bytes, BytesMut}; use super::resident::{OwnedBatchBody, PendingCounts, PendingResponseLease}; use super::{ - DownBatch, PendingClass, PollResult, QUEUE_ITEM_COST, QueuedFrame, SessionState, WebSession, + DownBatch, PendingClass, PollResult, QUEUE_ITEM_COST, QueuedFrame, SessionCloseReason, + SessionState, WebSession, }; use crate::web::frame::{self, FrameType}; use crate::web::manager::ManagerError; @@ -21,7 +22,7 @@ impl WebSession { if state.closed { return Err(ManagerError::Closed); } - state.last_activity = Instant::now(); + state.activity.touch_peer(Instant::now()); if let Some(unacked) = &state.unacked { if cursor == unacked.base_cursor { return Ok(PollResult { @@ -32,7 +33,7 @@ impl WebSession { } if cursor != unacked.next_cursor { drop(state); - self.close(); + self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); } let carrier_health_eligible = unacked.carrier_health_eligible; @@ -43,12 +44,12 @@ impl WebSession { } } else if cursor != state.down_cursor { drop(state); - self.close(); + self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); } let Some(epoch) = state.down_epoch.checked_add(1) else { drop(state); - self.close(); + self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); }; state.down_epoch = epoch; @@ -83,7 +84,7 @@ impl WebSession { } Err(error) => { drop(state); - self.close(); + self.close(SessionCloseReason::Protocol); return Err(error); } }; @@ -110,7 +111,7 @@ impl WebSession { Err(_) => { let mut state = self.state.lock(); if state.down_epoch == epoch { - state.last_activity = Instant::now(); + state.activity.touch_progress(Instant::now()); } Ok(PollResult { body: Bytes::new(), diff --git a/src/web/session/downlink_tests.rs b/src/web/session/downlink_tests.rs index d85e9f7..69a22b5 100644 --- a/src/web/session/downlink_tests.rs +++ b/src/web/session/downlink_tests.rs @@ -69,7 +69,7 @@ async fn downlink_replays_unacknowledged_batch_byte_for_byte() { assert_eq!(first.body, replay.body); drop(first); drop(replay); - session.close(); + session.close(super::SessionCloseReason::ApiClose); manager.shutdown().await; } @@ -89,7 +89,7 @@ async fn acknowledged_response_stays_resident_until_the_last_body_clone_drops() assert!(session.resident.snapshot().bytes() > 0); drop(retained); assert_eq!(session.resident.snapshot().bytes(), 0); - session.close(); + session.close(super::SessionCloseReason::ApiClose); manager.shutdown().await; } @@ -139,6 +139,6 @@ async fn newer_poll_supersedes_older_poll_without_closing_session() { assert_eq!(superseded.next_cursor, 0); assert!(!session.state.lock().closed); second.abort(); - session.close(); + session.close(super::SessionCloseReason::ApiClose); manager.shutdown().await; } diff --git a/src/web/session/lane_uplink.rs b/src/web/session/lane_uplink.rs index 27fd7bc..a13b7c5 100644 --- a/src/web/session/lane_uplink.rs +++ b/src/web/session/lane_uplink.rs @@ -5,7 +5,7 @@ use sha2::{Digest, Sha256}; use subtle::ConstantTimeEq; use super::uplink::{AppliedProgress, inbound_reservation, validate_batch}; -use super::{PendingClass, WebSession, insert_carrier_lane}; +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}; @@ -24,7 +24,7 @@ impl WebSession { let frames = match frame::parse_all(body, &self.limits) { Ok(frames) => frames, Err(_) => { - self.close(); + self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); } }; @@ -33,7 +33,7 @@ impl WebSession { .copied() .any(|value| value.stream_id != lane_id || frame::validate_client_shape(value).is_err()) { - self.close(); + self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); } let digest: TokenHash = Sha256::digest(body).into(); @@ -46,7 +46,7 @@ impl WebSession { return Err(ManagerError::Closed); } self.ensure_carrier_active_locked(&state)?; - state.last_activity = Instant::now(); + state.activity.touch_peer(Instant::now()); let new_lane = !state.carrier_lanes.contains_key(&lane_id); if new_lane { if lane_id != 0 @@ -69,7 +69,7 @@ impl WebSession { .is_none_or(|value| value.frame_type != FrameType::Open) { drop(state); - self.close(); + self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); } let lane_limit = self @@ -92,13 +92,13 @@ impl WebSession { Ok(sequence) } else { drop(state); - self.close(); + self.close(SessionCloseReason::Protocol); Err(ManagerError::Protocol) }; } if sequence == 0 || sequence != last_sequence.saturating_add(1) { drop(state); - self.close(); + self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); } if up_active { @@ -106,7 +106,7 @@ impl WebSession { } if !validate_batch(&state, &frames) { drop(state); - self.close(); + self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); } let (reserve_bytes, reserve_items) = inbound_reservation(&state, &frames); @@ -124,7 +124,7 @@ impl WebSession { if new_lane && insert_carrier_lane(&mut state, lane_id).is_none() { self.release_locked(&mut state, reserve_bytes, reserve_items, false); drop(state); - self.close(); + self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); } let Some(lane) = state.carrier_lanes.get_mut(&lane_id) else { @@ -161,7 +161,7 @@ impl WebSession { return result; } if result.is_err() { - self.close(); + self.close(SessionCloseReason::Protocol); drop(opened); return result; } diff --git a/src/web/session/lanes.rs b/src/web/session/lanes.rs index 52f47ac..2fe0c10 100644 --- a/src/web/session/lanes.rs +++ b/src/web/session/lanes.rs @@ -6,8 +6,8 @@ use tokio::sync::OwnedSemaphorePermit; use super::lane_downlink::take_lane_down_batch; use super::{ - CarrierLaneIdentity, PendingClass, PollResult, QUEUE_ITEM_COST, QueuedFrame, SessionState, - WebSession, remember_closed, + CarrierLaneIdentity, PendingClass, PollResult, QUEUE_ITEM_COST, QueuedFrame, SessionCloseReason, + SessionState, WebSession, remember_closed, }; use crate::web::frame::{self, FrameType}; use crate::web::manager::ManagerError; @@ -65,7 +65,7 @@ impl WebSession { if state.closed { return Err(ManagerError::Closed); } - state.last_activity = Instant::now(); + state.activity.touch_peer(Instant::now()); let acknowledged = { let Some(lane) = state.carrier_lanes.get_mut(&lane_id) else { return Ok(PollResult { @@ -91,14 +91,14 @@ impl WebSession { } if cursor != unacked.next_cursor { drop(state); - self.close(); + self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); } lane.unacked.take() } else { if cursor != lane.down_cursor { drop(state); - self.close(); + self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); } None @@ -133,7 +133,7 @@ impl WebSession { .ok_or(ManagerError::Protocol)?; let Some(epoch) = lane.down_epoch.checked_add(1) else { drop(state); - self.close(); + self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); }; lane.down_epoch = epoch; @@ -195,7 +195,7 @@ impl WebSession { } Err(error) => { drop(state); - self.close(); + self.close(SessionCloseReason::Protocol); return Err(error); } }; @@ -258,7 +258,7 @@ impl WebSession { }); } if lane.down_epoch == epoch { - state.last_activity = Instant::now(); + state.activity.touch_progress(Instant::now()); } } Ok(PollResult { @@ -281,7 +281,7 @@ impl WebSession { } if cursor != 0 || lane_id == 0 { drop(state); - self.close(); + self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); } if state.closed_streams.contains(&lane_id) diff --git a/src/web/session/lanes/tests.rs b/src/web/session/lanes/tests.rs index 4adaf7d..c1a41a0 100644 --- a/src/web/session/lanes/tests.rs +++ b/src/web/session/lanes/tests.rs @@ -108,7 +108,7 @@ async fn early_down_waits_without_creating_a_provisional_lane() { assert!(!result.body.is_empty()); assert_eq!(session.state.lock().lane_open_waits, 0); drop(result); - session.close(); + session.close(super::super::SessionCloseReason::ApiClose); manager.shutdown().await; } @@ -126,7 +126,7 @@ async fn early_down_timeout_is_empty_and_releases_its_session_slot() { assert_eq!(result.next_cursor, 0); assert!(!result.lane_closed); assert_eq!(session.state.lock().lane_open_waits, 0); - session.close(); + session.close(super::super::SessionCloseReason::ApiClose); manager.shutdown().await; } @@ -156,7 +156,7 @@ async fn early_down_admission_is_bounded_and_cancellation_safe() { let _ = wait.await; } assert_eq!(session.state.lock().lane_open_waits, 0); - session.close(); + session.close(super::super::SessionCloseReason::ApiClose); manager.shutdown().await; } @@ -168,7 +168,7 @@ async fn session_close_wakes_early_down_with_closed_state() { while session.state.lock().lane_open_waits == 0 { tokio::task::yield_now().await; } - session.close(); + session.close(super::super::SessionCloseReason::ApiClose); assert!(matches!( tokio::time::timeout(Duration::from_secs(1), poll) .await @@ -229,7 +229,7 @@ async fn drained_closed_lane_replays_then_signals_completion() { assert!(finished.lane_closed); drop(first); drop(replay); - session.close(); + session.close(super::super::SessionCloseReason::ApiClose); manager.shutdown().await; } diff --git a/src/web/session/lifecycle.rs b/src/web/session/lifecycle.rs index b0c221b..0f11b13 100644 --- a/src/web/session/lifecycle.rs +++ b/src/web/session/lifecycle.rs @@ -1,7 +1,95 @@ use std::sync::atomic::Ordering; use std::time::{Duration, Instant}; -use super::{SessionNegotiationPhase, WebSession}; +use super::WebSession; + +/// Stable terminal cause assigned by the first session-close winner. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +#[repr(usize)] +pub(crate) enum SessionCloseReason { + /// An authenticated client explicitly deleted its current session. + ClientDelete, + /// A surviving bridge replaced an unreachable carrier incarnation. + BridgeRecovery, + /// No validated peer operation arrived within the frozen reconnect grace. + PeerIdle, + /// Automatic carrier negotiation exhausted its absolute deadline. + NegotiationTimeout, + /// A successful negotiation replacement retired this incarnation. + CarrierSuperseded, + /// Authenticated carrier framing or sequencing violated the protocol. + Protocol, + /// Mandatory bounded control state could not be retained. + Backpressure, + /// A committed WebSocket carrier ended. + WebSocketEnded, + /// An authenticated control-plane request selected this session. + ApiClose, + /// A graceful operator drain reached its force-close deadline. + OperatorForce, + /// Terminal process shutdown closed all remaining sessions. + RuntimeShutdown, +} + +impl SessionCloseReason { + /// Complete fixed reason set in stable API and metric order. + pub(crate) const ALL: [Self; 11] = [ + Self::ClientDelete, + Self::BridgeRecovery, + Self::PeerIdle, + Self::NegotiationTimeout, + Self::CarrierSuperseded, + Self::Protocol, + Self::Backpressure, + Self::WebSocketEnded, + Self::ApiClose, + Self::OperatorForce, + Self::RuntimeShutdown, + ]; + + /// Returns the stable API, trace, and Prometheus token. + pub(crate) const fn as_str(self) -> &'static str { + match self { + Self::ClientDelete => "client_delete", + Self::BridgeRecovery => "bridge_recovery", + Self::PeerIdle => "peer_idle", + Self::NegotiationTimeout => "negotiation_timeout", + Self::CarrierSuperseded => "carrier_superseded", + Self::Protocol => "protocol", + Self::Backpressure => "backpressure", + Self::WebSocketEnded => "websocket_ended", + Self::ApiClose => "api_close", + Self::OperatorForce => "operator_force", + Self::RuntimeShutdown => "runtime_shutdown", + } + } +} + +/// Result of one first-writer-wins close request. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) enum SessionCloseOutcome { + /// This caller closed the session synchronously. + Closed, + /// This caller owns a close deferred behind carrier replacement. + Deferred, + /// An earlier caller already owns or completed session closure. + AlreadyClosing, +} + +impl SessionCloseOutcome { + /// Returns whether this caller won the terminal cause. + pub(crate) const fn accepted(self) -> bool { + !matches!(self, Self::AlreadyClosing) + } +} + +#[derive(Clone, Copy, PartialEq, Eq)] +pub(super) enum SessionNegotiationPhase { + Uncommitted, + Replacing, + Committed, + Superseded, +} struct ReleasedQueues { data_bytes: usize, @@ -9,6 +97,7 @@ struct ReleasedQueues { control_bytes: usize, control_items: usize, closed_before_health: bool, + reason: SessionCloseReason, } /// Deferred queue release after manager publication linearizes a supersede. @@ -21,24 +110,31 @@ pub(crate) struct CarrierSupersedeCompletion<'a> { impl CarrierSupersedeCompletion<'_> { /// Releases process budgets and signals cancellation after manager locks are dropped. pub(crate) fn finish(self) { - self.session.finish_close(self.released, true); + self.session.finish_close(self.released); } } impl WebSession { /// Closes carrier state while relay tasks retain their admission until exit. - pub(crate) fn close(&self) -> bool { - let Some(released) = self.begin_close(false, None) else { - return false; - }; - self.finish_close(released, false); - true + pub(crate) fn close(&self, reason: SessionCloseReason) -> SessionCloseOutcome { + let mut state = self.state.lock(); + if state.closed || state.close_requested.is_some() { + return SessionCloseOutcome::AlreadyClosing; + } + if state.negotiation_phase == SessionNegotiationPhase::Replacing { + state.close_requested = Some(reason); + return SessionCloseOutcome::Deferred; + } + let released = self.release_on_close_locked(&mut state, reason); + drop(state); + self.finish_close(released); + SessionCloseOutcome::Closed } /// Atomically prevents first-frame commit while one successor is prepared. pub(crate) fn begin_carrier_supersede(&self) -> bool { let mut state = self.state.lock(); - if state.closed || state.close_requested { + if state.closed || state.close_requested.is_some() { return false; } match state.negotiation_phase { @@ -54,21 +150,32 @@ impl WebSession { /// Restores an uncommitted attempt after successor admission failed. pub(crate) fn cancel_carrier_supersede(&self) { - let close_requested = { + let released = { let mut state = self.state.lock(); if !state.closed && state.negotiation_phase == SessionNegotiationPhase::Replacing { state.negotiation_phase = SessionNegotiationPhase::Uncommitted; } - state.close_requested + state + .close_requested + .filter(|_| !state.closed) + .map(|reason| self.release_on_close_locked(&mut state, reason)) }; - if close_requested { - self.close(); + if let Some(released) = released { + self.finish_close(released); } } /// Linearizes manager publication against close requests on the old token. pub(crate) fn prepare_carrier_supersede(&self) -> Option> { - let released = self.begin_close(true, None)?; + let mut state = self.state.lock(); + if state.closed + || state.negotiation_phase != SessionNegotiationPhase::Replacing + || state.close_requested.is_some() + { + return None; + } + let released = + self.release_on_close_locked(&mut state, SessionCloseReason::CarrierSuperseded); Some(CarrierSupersedeCompletion { session: self, released, @@ -88,6 +195,19 @@ impl WebSession { } } + /// Waits until registry removal and close telemetry have completed. + pub(crate) async fn wait_close_complete(&self) { + loop { + let notified = self.close_notify.notified(); + tokio::pin!(notified); + notified.as_mut().enable(); + if self.close_complete.load(Ordering::Acquire) { + return; + } + notified.await; + } + } + /// Returns the current number of registered logical-stream tasks. pub(crate) fn tasks_live(&self) -> usize { self.tasks_live.load(Ordering::Acquire) @@ -102,38 +222,38 @@ impl WebSession { if let Some(claim) = healthy { self.finish_carrier_health(claim); } - let Some(released) = self.begin_close(false, Some(now)) else { + let Some(released) = self.begin_idle_close(now) else { return false; }; - self.finish_close(released, false); + self.finish_close(released); true } - fn begin_close(&self, superseded: bool, idle_now: Option) -> Option { + fn begin_idle_close(&self, now: Instant) -> Option { let mut state = self.state.lock(); - if state.closed - || (superseded - && (state.negotiation_phase != SessionNegotiationPhase::Replacing - || state.close_requested)) + if state.closed || state.close_requested.is_some() { + return None; + } + if state.negotiation_phase == SessionNegotiationPhase::Replacing + || state.activity.peer_idle(now) + < Duration::from_secs(self.timeouts.reconnect_grace_secs) { return None; } - if let Some(now) = idle_now - && (state.negotiation_phase == SessionNegotiationPhase::Replacing - || now.saturating_duration_since(state.last_activity) - < Duration::from_secs(self.timeouts.reconnect_grace_secs)) - { - return None; - } - if !superseded && state.negotiation_phase == SessionNegotiationPhase::Replacing { - state.close_requested = true; - return None; - } + Some(self.release_on_close_locked(&mut state, SessionCloseReason::PeerIdle)) + } + + fn release_on_close_locked( + &self, + state: &mut super::SessionState, + reason: SessionCloseReason, + ) -> ReleasedQueues { let closed_before_health = self.automatic_carrier && state.negotiation_phase == SessionNegotiationPhase::Committed && self.reject_carrier_health_on_close(); + state.close_requested = Some(reason); state.closed = true; - if superseded { + if reason == SessionCloseReason::CarrierSuperseded { state.negotiation_phase = SessionNegotiationPhase::Superseded; } for stream in state.streams.values_mut() { @@ -149,8 +269,8 @@ impl WebSession { state.pending_windows.clear(); if let Some(batch) = state.unacked.take() { batch.lease.detach(); - self.release_local_locked(&mut state, batch.data_bytes, batch.data_items, false); - self.release_local_locked(&mut state, batch.control_bytes, batch.control_items, true); + self.release_local_locked(state, batch.data_bytes, batch.data_items, false); + self.release_local_locked(state, batch.control_bytes, batch.control_items, true); } let mut lane_data_bytes = 0usize; let mut lane_data_items = 0usize; @@ -166,8 +286,8 @@ impl WebSession { lane_control_items = lane_control_items.saturating_add(batch.control_items); } } - self.release_local_locked(&mut state, lane_data_bytes, lane_data_items, false); - self.release_local_locked(&mut state, lane_control_bytes, lane_control_items, true); + self.release_local_locked(state, lane_data_bytes, lane_data_items, false); + self.release_local_locked(state, lane_control_bytes, lane_control_items, true); state.carrier_lanes.clear(); let control_bytes = state.pending_control_bytes; let control_items = state.pending_control_items; @@ -177,16 +297,17 @@ impl WebSession { state.pending_items = 0; state.pending_control_bytes = 0; state.pending_control_items = 0; - Some(ReleasedQueues { + ReleasedQueues { data_bytes, data_items, control_bytes, control_items, closed_before_health, - }) + reason, + } } - fn finish_close(&self, released: ReleasedQueues, superseded: bool) { + fn finish_close(&self, released: ReleasedQueues) { self.cancel.cancel(); if self.carrier().is_multiplexed() { self.down_notify.notify_waiters(); @@ -194,7 +315,8 @@ impl WebSession { if self.carrier().uses_lanes() { self.lane_open_notify.notify_waiters(); } - if let Some(manager) = self.manager.upgrade() { + let manager = self.manager.upgrade(); + if let Some(manager) = &manager { if released.closed_before_health { manager.telemetry().record_carrier_learning( self.selected_carrier, @@ -213,22 +335,27 @@ impl WebSession { released.control_items, true, ); - if !self.finished.swap(true, Ordering::AcqRel) { - self.trace_lifecycle( - crate::web::trace::TraceLifecycleEvent::SessionClosed, - None, - Some(if superseded { "superseded" } else { "closed" }), + } + if !self.finished.swap(true, Ordering::AcqRel) { + self.trace_lifecycle( + crate::web::trace::TraceLifecycleEvent::SessionClosed, + None, + Some(released.reason.as_str()), + ); + if released.reason != SessionCloseReason::CarrierSuperseded + && let Some(manager) = manager + { + manager.session_finished( + self.token_hash, + self.client_ip, + self.profile_key, + &self.profile.host, + Duration::from_secs(self.timeouts.bootstrap_lifetime_secs), + released.reason, ); - if !superseded { - manager.session_finished( - self.token_hash, - self.client_ip, - self.profile_key, - &self.profile.host, - Duration::from_secs(self.timeouts.bootstrap_lifetime_secs), - ); - } } + self.close_complete.store(true, Ordering::Release); + self.close_notify.notify_waiters(); } } } diff --git a/src/web/session/negotiation.rs b/src/web/session/negotiation.rs index a213b30..68e7bf7 100644 --- a/src/web/session/negotiation.rs +++ b/src/web/session/negotiation.rs @@ -460,7 +460,7 @@ mod tests { let close = std::thread::spawn(move || { close_barrier.wait(); std::thread::yield_now(); - close_session.close(); + close_session.close(super::SessionCloseReason::ApiClose); }); barrier.wait(); health.join().unwrap(); @@ -472,7 +472,10 @@ mod tests { | CarrierHealthPublicationState::Rejected )); assert!(!session.publish_carrier_health()); - assert!(!session.close()); + assert_eq!( + session.close(super::SessionCloseReason::ApiClose), + super::SessionCloseOutcome::AlreadyClosing + ); } } diff --git a/src/web/session/status.rs b/src/web/session/status.rs index cd55cda..7e1a8f7 100644 --- a/src/web/session/status.rs +++ b/src/web/session/status.rs @@ -54,6 +54,12 @@ pub(crate) struct WebSessionStatus { pub(crate) age_ms: u64, /// Monotonic age since the latest carrier activity. pub(crate) idle_ms: u64, + /// Monotonic age since the latest validated peer operation. + pub(crate) peer_idle_ms: u64, + /// Session-frozen authenticated peer inactivity allowance. + pub(crate) reconnect_grace_ms: u64, + /// Remaining time before peer inactivity makes the session eligible for cleanup. + pub(crate) peer_deadline_remaining_ms: u64, /// Remaining automatic negotiation deadline. #[serde(skip_serializing_if = "Option::is_none")] pub(crate) negotiation_remaining_ms: Option, @@ -66,7 +72,7 @@ impl WebSession { let resident = self.resident.snapshot(); let state_name = if state.closed { "closed" - } else if state.close_requested { + } else if state.close_requested.is_some() { "closing" } else if self.carrier_health_publication_state() == CarrierHealthPublicationState::Published @@ -112,7 +118,14 @@ impl WebSession { .pending_control_items .saturating_add(resident.control_items), age_ms: millis(now.saturating_duration_since(self.created_at)), - idle_ms: millis(now.saturating_duration_since(state.last_activity)), + idle_ms: millis(state.activity.progress_idle(now)), + peer_idle_ms: millis(state.activity.peer_idle(now)), + reconnect_grace_ms: self.timeouts.reconnect_grace_secs.saturating_mul(1_000), + peer_deadline_remaining_ms: self + .timeouts + .reconnect_grace_secs + .saturating_mul(1_000) + .saturating_sub(millis(state.activity.peer_idle(now))), negotiation_remaining_ms: self .carrier_deadline_at .map(|deadline| millis(deadline.saturating_duration_since(now))), diff --git a/src/web/session/uplink.rs b/src/web/session/uplink.rs index 91f0a0d..e7b927b 100644 --- a/src/web/session/uplink.rs +++ b/src/web/session/uplink.rs @@ -9,8 +9,8 @@ use subtle::ConstantTimeEq; use super::backend::StreamCompletion; use super::{ - InboundChunk, PendingClass, QUEUE_ITEM_COST, SessionState, StreamIdentity, StreamState, - WebSession, inbound_queue_cost, + InboundChunk, PendingClass, QUEUE_ITEM_COST, SessionCloseReason, SessionState, StreamIdentity, + StreamState, WebSession, inbound_queue_cost, }; use crate::web::frame::{self, Frame, FrameType}; use crate::web::manager::{ManagerError, TokenHash}; @@ -70,7 +70,7 @@ impl WebSession { let frames = match frame::parse_all(body, &self.limits) { Ok(frames) => frames, Err(_) => { - self.close(); + self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); } }; @@ -79,7 +79,7 @@ impl WebSession { .copied() .any(|value| frame::validate_client_shape(value).is_err()) { - self.close(); + self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); } let digest: TokenHash = Sha256::digest(body).into(); @@ -92,24 +92,24 @@ impl WebSession { return Err(ManagerError::Closed); } self.ensure_carrier_active_locked(&state)?; - state.last_activity = Instant::now(); + 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)) { Ok((sequence, false)) } else { drop(state); - self.close(); + self.close(SessionCloseReason::Protocol); Err(ManagerError::Protocol) }; } if sequence == 0 || sequence != state.last_up_sequence.saturating_add(1) { drop(state); - self.close(); + self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); } if !validate_batch(&state, &frames) { drop(state); - self.close(); + self.close(SessionCloseReason::Protocol); return Err(ManagerError::Protocol); } let (reserve_bytes, reserve_items) = inbound_reservation(&state, &frames); @@ -147,7 +147,7 @@ impl WebSession { return result; } if result.is_err() { - self.close(); + self.close(SessionCloseReason::Protocol); drop(opened); return result; } diff --git a/src/web/session/websocket.rs b/src/web/session/websocket.rs index b1b6b26..dc803e7 100644 --- a/src/web/session/websocket.rs +++ b/src/web/session/websocket.rs @@ -391,7 +391,7 @@ impl WebSession { lane.last_up_digest = digest; } } - state.last_activity = Instant::now(); + state.activity.touch_peer(Instant::now()); if applied { (committed, healthy) = self.record_uplink_progress_locked(&mut state, progress); } diff --git a/src/web/session/websocket/tests.rs b/src/web/session/websocket/tests.rs index 45f964b..af05cfe 100644 --- a/src/web/session/websocket/tests.rs +++ b/src/web/session/websocket/tests.rs @@ -18,7 +18,7 @@ struct TestRuntime { impl TestRuntime { async fn shutdown(self) { - self.session.close(); + self.session.close(super::super::SessionCloseReason::ApiClose); self.session.wait().await; self.manager.shutdown().await; self.generation.stop_sessions().await; @@ -181,7 +181,9 @@ async fn closed_session_releases_bound_lane_quota_on_reservation_drop() { let mut reservation = runtime.session.reserve_websocket_lane(7).unwrap(); reservation.bind(1).unwrap(); - runtime.session.close(); + runtime + .session + .close(super::super::SessionCloseReason::ApiClose); drop(reservation); assert!(runtime.session.state.lock().active_peer_ports.is_empty()); @@ -220,7 +222,9 @@ async fn closed_session_releases_transferred_rejected_lane_quota() { WebSocketLaneReservationPhase::Transferred ); - runtime.session.close(); + runtime + .session + .close(super::super::SessionCloseReason::ApiClose); drop(reservation); assert!(runtime.session.state.lock().active_peer_ports.is_empty()); @@ -259,7 +263,9 @@ async fn closed_session_keeps_stream_owned_quota_until_task_completion() { WebSocketLaneReservationPhase::StreamOwned ); - runtime.session.close(); + runtime + .session + .close(super::super::SessionCloseReason::ApiClose); drop(reservation); assert!( diff --git a/src/web/telemetry.rs b/src/web/telemetry.rs index 1b3e4d3..fbce47d 100644 --- a/src/web/telemetry.rs +++ b/src/web/telemetry.rs @@ -10,6 +10,12 @@ pub(crate) use carrier::{ WebCarrierLearningOutcome, WebCarrierSelectionCounter, WebCarrierSelectionDisposition, }; use carrier::{CARRIER_FAILURE_SLOTS, CARRIER_LEARNING_SLOTS, CARRIER_SELECTION_SLOTS}; +mod lifecycle; +pub(crate) use lifecycle::{ + WebBridgeRecoveryCounter, WebBridgeRecoveryEvent, WebSessionCloseCounter, + WebSessionLifecycleObservation, WebSessionLifecycleObservationCounter, +}; +use lifecycle::{SESSION_CLOSE_SLOTS, SESSION_OBSERVATION_SLOTS}; const LAST_DECOY_OUTCOME_BITS: u32 = 4; const LAST_DECOY_OUTCOME_MASK: u64 = (1 << LAST_DECOY_OUTCOME_BITS) - 1; @@ -310,6 +316,9 @@ pub(crate) struct WebTelemetry { carrier_selections: [AtomicU64; CARRIER_SELECTION_SLOTS], carrier_failures: [AtomicU64; CARRIER_FAILURE_SLOTS], carrier_learning_outcomes: [AtomicU64; CARRIER_LEARNING_SLOTS], + session_closures: [AtomicU64; SESSION_CLOSE_SLOTS], + session_observations: [AtomicU64; SESSION_OBSERVATION_SLOTS], + bridge_recovery_events: [AtomicU64; WebBridgeRecoveryEvent::ALL.len()], last_decoy: AtomicU64, sessions_created: AtomicU64, sessions_closed: AtomicU64, @@ -334,6 +343,9 @@ impl WebTelemetry { carrier_selections: std::array::from_fn(|_| AtomicU64::new(0)), carrier_failures: std::array::from_fn(|_| AtomicU64::new(0)), carrier_learning_outcomes: std::array::from_fn(|_| AtomicU64::new(0)), + session_closures: std::array::from_fn(|_| AtomicU64::new(0)), + session_observations: std::array::from_fn(|_| AtomicU64::new(0)), + bridge_recovery_events: std::array::from_fn(|_| AtomicU64::new(0)), last_decoy: AtomicU64::new(0), sessions_created: AtomicU64::new(0), sessions_closed: AtomicU64::new(0), @@ -467,11 +479,6 @@ impl WebTelemetry { self.sessions_created.fetch_add(1, Ordering::Relaxed); } - /// Records one closed session incarnation. - pub(crate) fn record_session_closed(&self) { - self.sessions_closed.fetch_add(1, Ordering::Relaxed); - } - /// Records one admitted logical stream. pub(crate) fn record_stream_opened(&self) { self.streams_opened.fetch_add(1, Ordering::Relaxed); diff --git a/src/web/telemetry/lifecycle.rs b/src/web/telemetry/lifecycle.rs new file mode 100644 index 0000000..826af91 --- /dev/null +++ b/src/web/telemetry/lifecycle.rs @@ -0,0 +1,236 @@ +use std::sync::atomic::{AtomicU64, Ordering}; + +use serde::Serialize; + +use crate::config::WebCarrier; +use crate::web::session::SessionCloseReason; + +use super::WebTelemetry; + +pub(super) const SESSION_CLOSE_SLOTS: usize = + WebCarrier::ALL.len() * SessionCloseReason::ALL.len(); +pub(super) const SESSION_OBSERVATION_SLOTS: usize = + WebCarrier::ALL.len() * WebSessionLifecycleObservation::ALL.len(); + +/// Stable observation emitted after an authenticated session lifecycle gap. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +#[repr(usize)] +pub(crate) enum WebSessionLifecycleObservation { + /// A valid HTTP carrier request resumed after a suspicious peer gap. + HttpActivityAfterGap, + /// A valid WebSocket message resumed after a suspicious peer gap. + WebSocketActivityAfterGap, + /// A valid retained bearer was used after its session closed. + RequestAfterClose, +} + +impl WebSessionLifecycleObservation { + /// Complete fixed observation set in stable metric order. + pub(crate) const ALL: [Self; 3] = [ + Self::HttpActivityAfterGap, + Self::WebSocketActivityAfterGap, + Self::RequestAfterClose, + ]; + + /// Returns the stable API and Prometheus token. + pub(crate) const fn as_str(self) -> &'static str { + match self { + Self::HttpActivityAfterGap => "http_activity_after_gap", + Self::WebSocketActivityAfterGap => "websocket_activity_after_gap", + Self::RequestAfterClose => "request_after_close", + } + } +} + +/// Stable server-side milestone for one bridge recovery incarnation. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +#[repr(usize)] +pub(crate) enum WebBridgeRecoveryEvent { + /// A recovery request received a fresh bootstrap. + BootstrapIssued, + /// A fresh recovery session was registered. + SessionCreated, + /// The recovery session committed a real carrier. + Committed, + /// An issued recovery bootstrap expired without session creation. + ExpiredUnused, + /// A recovery session closed before carrier commit. + ClosedBeforeCommit, +} + +impl WebBridgeRecoveryEvent { + /// Complete fixed recovery event set in stable metric order. + pub(crate) const ALL: [Self; 5] = [ + Self::BootstrapIssued, + Self::SessionCreated, + Self::Committed, + Self::ExpiredUnused, + Self::ClosedBeforeCommit, + ]; + + /// Returns the stable API and Prometheus token. + pub(crate) const fn as_str(self) -> &'static str { + match self { + Self::BootstrapIssued => "bootstrap_issued", + Self::SessionCreated => "session_created", + Self::Committed => "committed", + Self::ExpiredUnused => "expired_unused", + Self::ClosedBeforeCommit => "closed_before_commit", + } + } +} + +/// API-safe typed session-close counter. +#[derive(Clone, Serialize)] +pub(crate) struct WebSessionCloseCounter { + /// Fixed carrier owning the closed incarnation. + pub(crate) carrier: &'static str, + /// Fixed terminal close cause. + pub(crate) reason: &'static str, + /// Monotonic process-lifetime count. + pub(crate) total: u64, +} + +/// API-safe typed lifecycle-observation counter. +#[derive(Clone, Serialize)] +pub(crate) struct WebSessionLifecycleObservationCounter { + /// Fixed carrier associated with the observation. + pub(crate) carrier: &'static str, + /// Fixed lifecycle observation. + pub(crate) observation: &'static str, + /// Monotonic process-lifetime count. + pub(crate) total: u64, +} + +/// API-safe typed bridge-recovery counter. +#[derive(Clone, Serialize)] +pub(crate) struct WebBridgeRecoveryCounter { + /// Fixed recovery milestone. + pub(crate) event: &'static str, + /// Monotonic process-lifetime count. + pub(crate) total: u64, +} + +pub(super) const fn session_close_slot( + carrier: WebCarrier, + reason: SessionCloseReason, +) -> usize { + carrier as usize * SessionCloseReason::ALL.len() + reason as usize +} + +pub(super) const fn session_observation_slot( + carrier: WebCarrier, + observation: WebSessionLifecycleObservation, +) -> usize { + carrier as usize * WebSessionLifecycleObservation::ALL.len() + observation as usize +} + +pub(super) fn load(counter: &AtomicU64) -> u64 { + counter.load(Ordering::Relaxed) +} + +impl WebTelemetry { + /// Records one closed session incarnation and its exact terminal cause. + pub(crate) fn record_session_closed( + &self, + carrier: WebCarrier, + reason: SessionCloseReason, + ) { + self.sessions_closed.fetch_add(1, Ordering::Relaxed); + self.session_closures[session_close_slot(carrier, reason)] + .fetch_add(1, Ordering::Relaxed); + } + + /// Returns one fixed session-close counter. + pub(crate) fn session_close_total( + &self, + carrier: WebCarrier, + reason: SessionCloseReason, + ) -> u64 { + load(&self.session_closures[session_close_slot(carrier, reason)]) + } + + /// Captures the complete fixed session-close matrix for API serialization. + pub(crate) fn session_close_counters(&self) -> Vec { + WebCarrier::ALL + .into_iter() + .flat_map(|carrier| { + SessionCloseReason::ALL.into_iter().map(move |reason| { + WebSessionCloseCounter { + carrier: carrier.as_str(), + reason: reason.as_str(), + total: self.session_close_total(carrier, reason), + } + }) + }) + .collect() + } + + /// Records one fixed authenticated lifecycle observation. + pub(crate) fn record_session_observation( + &self, + carrier: WebCarrier, + observation: WebSessionLifecycleObservation, + ) { + self.session_observations[session_observation_slot(carrier, observation)] + .fetch_add(1, Ordering::Relaxed); + } + + /// Returns one fixed authenticated lifecycle observation counter. + pub(crate) fn session_observation_total( + &self, + carrier: WebCarrier, + observation: WebSessionLifecycleObservation, + ) -> u64 { + load(&self.session_observations[session_observation_slot(carrier, observation)]) + } + + /// Captures the complete fixed lifecycle-observation matrix. + pub(crate) fn session_observation_counters( + &self, + ) -> Vec { + WebCarrier::ALL + .into_iter() + .flat_map(|carrier| { + WebSessionLifecycleObservation::ALL + .into_iter() + .map(move |observation| WebSessionLifecycleObservationCounter { + carrier: carrier.as_str(), + observation: observation.as_str(), + total: self.session_observation_total(carrier, observation), + }) + }) + .collect() + } + + /// Records one fixed bridge-recovery milestone. + pub(crate) fn record_bridge_recovery(&self, event: WebBridgeRecoveryEvent) { + self.bridge_recovery_events[event as usize].fetch_add(1, Ordering::Relaxed); + } + + /// Adds a bounded batch of identical recovery milestones. + pub(crate) fn record_bridge_recovery_count( + &self, + event: WebBridgeRecoveryEvent, + count: usize, + ) { + self.bridge_recovery_events[event as usize] + .fetch_add(u64::try_from(count).unwrap_or(u64::MAX), Ordering::Relaxed); + } + + /// Returns one fixed bridge-recovery counter. + pub(crate) fn bridge_recovery_total(&self, event: WebBridgeRecoveryEvent) -> u64 { + load(&self.bridge_recovery_events[event as usize]) + } + + /// Captures the complete fixed bridge-recovery event set. + pub(crate) fn bridge_recovery_counters(&self) -> Vec { + WebBridgeRecoveryEvent::ALL + .into_iter() + .map(|event| WebBridgeRecoveryCounter { + event: event.as_str(), + total: self.bridge_recovery_total(event), + }) + .collect() + } +}