diff --git a/src/api/config_edit.rs b/src/api/config_edit.rs index 0f46236..cb9dfeb 100644 --- a/src/api/config_edit.rs +++ b/src/api/config_edit.rs @@ -15,7 +15,7 @@ use super::model::ApiFailure; use crate::config::ProxyConfig; use crate::config::hot_reload::classify_config_changes; use crate::maestro::reload::{ReloadAccepted, ReloadRequest, ReloadSubmitError}; -use crate::maestro::runtime_build::{deferred_process_fields, resolve_reload_config}; +use crate::maestro::runtime_build::resolve_reload_config; use serde::Serialize; use std::path::{Path, PathBuf}; use std::sync::Arc; @@ -206,8 +206,7 @@ async fn prepare_patch_to_path( let revision = compute_snapshot_revision(&candidate); let new_cfg = candidate.config; let class = classify_config_changes(&old_cfg, &new_cfg); - let deferred_process_fields = - deferred_process_fields(&old_cfg, &new_cfg).map_err(ApiFailure::bad_request)?; + let resolved = resolve_reload_config(&old_cfg, &new_cfg).map_err(ApiFailure::bad_request)?; Ok(PreparedConfigPatch { owner_path, @@ -216,9 +215,9 @@ async fn prepare_patch_to_path( response: PatchConfigResponse { revision, restart_required: class.restart_required, - runtime_reload_required: class.restart_required, - process_restart_required: !deferred_process_fields.is_empty(), - deferred_process_fields, + runtime_reload_required: resolved.runtime_changed, + process_restart_required: !resolved.deferred_process_fields.is_empty(), + deferred_process_fields: resolved.deferred_process_fields, changed: class.changed, reload: None, }, diff --git a/src/api/config_edit/tests.rs b/src/api/config_edit/tests.rs index 816718b..4bbba8f 100644 --- a/src/api/config_edit/tests.rs +++ b/src/api/config_edit/tests.rs @@ -131,6 +131,24 @@ async fn patch_web_debug_is_hot_and_limits_are_process_deferred() { ); } +#[tokio::test] +async fn patch_web_decoy_fasttrack_requires_only_process_restart() { + let (path, _directory) = temp_config("[web]\nenabled = false\n"); + let patch: Json = serde_json::json!({ + "web": {"decoy_fasttrack_mode": "shadow"} + }); + + let response = apply_patch_to_path(&path, &patch, None).await.unwrap(); + + assert!(response.restart_required); + assert!(!response.runtime_reload_required); + assert!(response.process_restart_required); + assert_eq!( + response.deferred_process_fields, + vec!["web.decoy_fasttrack_mode".to_string()] + ); +} + #[tokio::test] async fn invalid_web_patch_does_not_modify_the_source() { let (path, _directory) = temp_config("[web]\nenabled = false\n"); diff --git a/src/api/web_runtime.rs b/src/api/web_runtime.rs index 4aa5e6e..f60c554 100644 --- a/src/api/web_runtime.rs +++ b/src/api/web_runtime.rs @@ -20,8 +20,8 @@ mod request; // Ingress, capacity, and decoy telemetry remain separate availability planes. mod observability; use observability::{ - WebCapacityStatus, WebCarrierNegotiationStatus, WebDecoyUpstreamStatus, WebIngressStatus, - WebLifecycleCountersStatus, + WebCapacityStatus, WebCarrierNegotiationStatus, WebDecoyFastTrackStatus, + WebDecoyUpstreamStatus, WebIngressStatus, WebLifecycleCountersStatus, }; use request::{ CloseRequest, DrainRequest, RuntimeInstanceRequest, parse_session_query, parse_session_ref, @@ -280,6 +280,7 @@ struct WebStatusData { ingress: WebIngressStatus, capacity: WebCapacityStatus, decoy_upstream: WebDecoyUpstreamStatus, + decoy_fasttrack: WebDecoyFastTrackStatus, carrier_negotiation: WebCarrierNegotiationStatus, lifecycle_counters: WebLifecycleCountersStatus, #[serde(skip_serializing_if = "Option::is_none")] @@ -315,6 +316,7 @@ impl WebStatusData { let ingress = WebIngressStatus::new(&publication, runtime.is_some()); let capacity = WebCapacityStatus::new(&publication, runtime, config); let decoy_upstream = WebDecoyUpstreamStatus::new(&publication); + let decoy_fasttrack = WebDecoyFastTrackStatus::new(&publication, config); let carrier_negotiation = WebCarrierNegotiationStatus::new(&publication); let lifecycle_counters = WebLifecycleCountersStatus::new(&publication, config); Self { @@ -332,6 +334,7 @@ impl WebStatusData { ingress, capacity, decoy_upstream, + decoy_fasttrack, carrier_negotiation, lifecycle_counters, operator_lifecycle, diff --git a/src/api/web_runtime/observability.rs b/src/api/web_runtime/observability.rs index 971fc8f..a675b07 100644 --- a/src/api/web_runtime/observability.rs +++ b/src/api/web_runtime/observability.rs @@ -1,11 +1,12 @@ use serde::Serialize; -use crate::config::{ProxyConfig, WebHttpConnectionCapacityAction}; +use crate::config::{ProxyConfig, WebDecoyFastTrackMode, WebHttpConnectionCapacityAction}; use crate::web::control::{WebRuntimeLifecycle, WebRuntimePublication}; use crate::web::manager::{WebCapacityResourceStatus, WebCapacitySnapshot, WebProcessRuntime}; use crate::web::telemetry::{ WebBridgeRecoveryCounter, WebCarrierFailureCounter, WebCarrierLearningCounter, - WebCarrierSelectionCounter, WebSessionCloseCounter, WebSessionLifecycleObservationCounter, + WebCarrierSelectionCounter, WebDecoyFastTrackCounter, WebSessionCloseCounter, + WebSessionLifecycleObservationCounter, }; use crate::web::telemetry::{WebOutcomeCounter, WebRejectionCounter}; @@ -121,6 +122,27 @@ impl WebDecoyUpstreamStatus { } } +/// Fixed-cardinality process-lifetime decoy capability-routing counters. +#[derive(Serialize)] +pub(super) struct WebDecoyFastTrackStatus { + mode: WebDecoyFastTrackMode, + requests: Vec, + shadow_mismatches_total: u64, +} + +impl WebDecoyFastTrackStatus { + /// Builds effective policy and counters without requiring the runtime manager. + pub(super) fn new(publication: &WebRuntimePublication, config: &ProxyConfig) -> Self { + Self { + mode: config.web.decoy_fasttrack_mode, + requests: publication.telemetry.decoy_fasttrack_counters(), + shadow_mismatches_total: publication + .telemetry + .decoy_fasttrack_shadow_mismatches(), + } + } +} + /// Fixed-cardinality process-lifetime carrier negotiation counters. #[derive(Serialize)] pub(super) struct WebCarrierNegotiationStatus { @@ -187,6 +209,11 @@ mod tests { serde_json::to_value(super::WebCapacityStatus::new(&publication, None, &config)) .unwrap(); let decoy = serde_json::to_value(super::WebDecoyUpstreamStatus::new(&publication)).unwrap(); + let fasttrack = serde_json::to_value(super::WebDecoyFastTrackStatus::new( + &publication, + &config, + )) + .unwrap(); let carrier = serde_json::to_value(super::WebCarrierNegotiationStatus::new(&publication)).unwrap(); let lifecycle = serde_json::to_value(super::WebLifecycleCountersStatus::new( @@ -210,6 +237,12 @@ mod tests { decoy["outcomes"].as_array().unwrap().len(), crate::web::telemetry::WebDecoyUpstreamOutcome::ALL.len() ); + assert_eq!(fasttrack["mode"], "off"); + assert_eq!( + fasttrack["requests"].as_array().unwrap().len(), + crate::web::telemetry::WebDecoyFastTrackDisposition::ALL.len() + ); + assert_eq!(fasttrack["shadow_mismatches_total"], 0); assert_eq!(capacity["partial"][0], "runtime"); assert_eq!( carrier["selections"].as_array().unwrap().len(), diff --git a/src/config/hot_reload/diff.rs b/src/config/hot_reload/diff.rs index 575c4fe..3e62595 100644 --- a/src/config/hot_reload/diff.rs +++ b/src/config/hot_reload/diff.rs @@ -85,6 +85,10 @@ pub(super) fn warn_non_hot_changes(old: &ProxyConfig, new: &ProxyConfig, non_hot warned = true; warn!("config reload: server listener settings changed; restart required"); } + if old.web.decoy_fasttrack_mode != new.web.decoy_fasttrack_mode { + warned = true; + warn!("config reload: web.decoy_fasttrack_mode changed; restart required"); + } if old.censorship.tls_domain != new.censorship.tls_domain || old.censorship.tls_domains != new.censorship.tls_domains || old.censorship.tls_fetch_scope != new.censorship.tls_fetch_scope diff --git a/src/config/hot_reload/fields.rs b/src/config/hot_reload/fields.rs index 8ca7747..a8b2ac9 100644 --- a/src/config/hot_reload/fields.rs +++ b/src/config/hot_reload/fields.rs @@ -343,8 +343,10 @@ pub(super) fn overlay_hot_fields(old: &ProxyConfig, new: &ProxyConfig) -> ProxyC cfg.access.user_max_unique_ips_mode = new.access.user_max_unique_ips_mode; cfg.access.user_max_unique_ips_window_secs = new.access.user_max_unique_ips_window_secs; let process_limits = cfg.web.limits.clone(); + let decoy_fasttrack_mode = cfg.web.decoy_fasttrack_mode; cfg.web = new.web.clone(); cfg.web.limits = process_limits; + cfg.web.decoy_fasttrack_mode = decoy_fasttrack_mode; if cfg.web.carrier_negotiation_enabled() && cfg.web.carrier_learning && cfg.web.limits.max_carrier_learning_entries < WEB_CARRIER_LEARNING_MIN_ENTRIES diff --git a/src/config/hot_reload/tests.rs b/src/config/hot_reload/tests.rs index 4bb6869..802935b 100644 --- a/src/config/hot_reload/tests.rs +++ b/src/config/hot_reload/tests.rs @@ -123,6 +123,21 @@ fn web_debug_policy_is_hot_while_debug_capacity_is_process_owned() { ); } +#[test] +fn decoy_fasttrack_mode_is_deferred_until_restart() { + let old = sample_config(); + let mut new = old.clone(); + new.web.decoy_fasttrack_mode = crate::config::WebDecoyFastTrackMode::Enforce; + + let applied = overlay_hot_fields(&old, &new); + + assert_eq!( + applied.web.decoy_fasttrack_mode, + old.web.decoy_fasttrack_mode + ); + assert_eq!(HotFields::from_config(&old), HotFields::from_config(&applied)); +} + #[test] fn hot_overlay_defers_learning_that_requires_new_process_capacity() { let mut old = sample_config(); diff --git a/src/config/load/runtime_web.rs b/src/config/load/runtime_web.rs index ab31b6a..159a5fa 100644 --- a/src/config/load/runtime_web.rs +++ b/src/config/load/runtime_web.rs @@ -36,6 +36,7 @@ pub(super) fn rebuild(config: &mut ProxyConfig) -> Result<()> { &mut static_bytes, )?; let mut profiles = Vec::with_capacity(vhost.profiles.len()); + let mut capability_table = Vec::with_capacity(vhost.profiles.len()); let mut capabilities = HashSet::with_capacity(vhost.profiles.len()); for profile in &vhost.profiles { let user_id = auth.user_id_by_name(&profile.user).ok_or_else(|| { @@ -84,6 +85,7 @@ pub(super) fn rebuild(config: &mut ProxyConfig) -> Result<()> { .max_streams_per_session .unwrap_or(config.web.limits.max_streams_per_session), }); + capability_table.push(capability); profiles.push(Arc::clone(&runtime_profile)); runtime_profiles.push(runtime_profile); } @@ -91,9 +93,11 @@ pub(super) fn rebuild(config: &mut ProxyConfig) -> Result<()> { vhost.host.clone(), Arc::new(WebRuntimeVhost { host: vhost.host.clone(), + decoy_fasttrack_mode: config.web.decoy_fasttrack_mode, decoy, decoy_header_secs: config.web.timeouts.decoy_header_secs, profiles, + capabilities: capability_table.into_boxed_slice(), }), ); } diff --git a/src/config/load/strict_keys.rs b/src/config/load/strict_keys.rs index ed3756f..69dc383 100644 --- a/src/config/load/strict_keys.rs +++ b/src/config/load/strict_keys.rs @@ -266,6 +266,7 @@ const WEB_CONFIG_KEYS: &[&str] = &[ "carriers", "carrier_learning", "carrier_negotiation_aggressiveness", + "decoy_fasttrack_mode", "http_connection_capacity_action", "debug", "limits", diff --git a/src/config/load/validate_web/memory.rs b/src/config/load/validate_web/memory.rs index 7f5ef4b..7dca7a6 100644 --- a/src/config/load/validate_web/memory.rs +++ b/src/config/load/validate_web/memory.rs @@ -6,6 +6,7 @@ const WEB_DEBUG_GROUP_SCRATCH_BYTES: usize = 4 * 1024 * 1024; const WEB_CARRIER_LEARNING_ENTRY_BYTES: usize = 512; const WEB_LANE_STATE_BYTES: usize = 512; const WEB_OVERLOAD_CONNECTION_BYTES: usize = 4 * 1024; +const WEB_CAPABILITY_INDEX_ENTRY_BYTES: usize = 32; /// Validates process-wide body, header, queue, static, and debug reservations. pub(super) fn validate(limits: &WebLimitsConfig) -> Result<()> { @@ -65,6 +66,12 @@ pub(super) fn validate(limits: &WebLimitsConfig) -> Result<()> { .ok_or_else(|| { ProxyError::Config("web.carrier learning reservation overflowed usize".to_string()) })?; + let capability_index_reservation = limits + .max_profiles + .checked_mul(WEB_CAPABILITY_INDEX_ENTRY_BYTES) + .ok_or_else(|| { + ProxyError::Config("web capability index reservation overflowed usize".to_string()) + })?; let lane_state_reservation = limits .max_streams_per_session .checked_add(limits.max_tombstones_per_session) @@ -82,6 +89,7 @@ pub(super) fn validate(limits: &WebLimitsConfig) -> Result<()> { .and_then(|value| value.checked_add(status_pages)) .and_then(|value| value.checked_add(debug_reservation)) .and_then(|value| value.checked_add(carrier_learning_reservation)) + .and_then(|value| value.checked_add(capability_index_reservation)) .and_then(|value| value.checked_add(lane_state_reservation)) .and_then(|value| value.checked_add(http_header_reservation)) .and_then(|value| value.checked_add(overload_connection_reservation)) @@ -111,4 +119,14 @@ mod tests { }; assert!(validate(&previous_envelope).is_err()); } + + #[test] + fn capability_index_reservation_rejects_size_overflow() { + let limits = WebLimitsConfig { + max_profiles: usize::MAX, + ..WebLimitsConfig::default() + }; + let error = validate(&limits).unwrap_err().to_string(); + assert!(error.contains("web capability index reservation overflowed usize")); + } } diff --git a/src/config/tests/load_basic_tests/web_tests.rs b/src/config/tests/load_basic_tests/web_tests.rs index 5c1246a..665b953 100644 --- a/src/config/tests/load_basic_tests/web_tests.rs +++ b/src/config/tests/load_basic_tests/web_tests.rs @@ -41,6 +41,12 @@ fn web_config_builds_canonical_runtime_snapshot() { .get("proxy.example.com") .expect("canonical WEB vhost"); assert_eq!(vhost.profiles.len(), 1); + assert_eq!(vhost.capabilities.len(), vhost.profiles.len()); + assert_eq!(vhost.capabilities[0], vhost.profiles[0].capability); + assert_eq!( + vhost.decoy_fasttrack_mode, + WebDecoyFastTrackMode::Off + ); assert_eq!(vhost.profiles[0].user, "alice"); assert_eq!(vhost.profiles[0].secret_mode, WebSecretMode::Dd); assert_eq!(vhost.profiles[0].carrier, WebCarrier::HttpsLanes); @@ -56,6 +62,45 @@ fn web_config_builds_canonical_runtime_snapshot() { ); } +#[test] +fn web_decoy_fasttrack_mode_is_typed_and_defaults_off() { + let defaults = ProxyConfig::default(); + assert_eq!( + defaults.web.decoy_fasttrack_mode, + WebDecoyFastTrackMode::Off + ); + + for (token, expected) in [ + ("shadow", WebDecoyFastTrackMode::Shadow), + ("enforce", WebDecoyFastTrackMode::Enforce), + ] { + let configured = WEB_CONFIG.replace( + "carrier = \"https-lanes\"", + &format!( + "carrier = \"https-lanes\"\ndecoy_fasttrack_mode = \"{token}\"" + ), + ); + let config = load_config_from_temp_toml(&configured); + assert_eq!(config.web.decoy_fasttrack_mode, expected); + assert_eq!( + config + .web + .runtime + .as_ref() + .unwrap() + .vhosts["proxy.example.com"] + .decoy_fasttrack_mode, + expected + ); + } + + let invalid = WEB_CONFIG.replace( + "carrier = \"https-lanes\"", + "carrier = \"https-lanes\"\ndecoy_fasttrack_mode = \"automatic\"", + ); + assert!(load_config_error_from_temp_toml(&invalid).contains("decoy_fasttrack_mode")); +} + #[test] fn web_http_connection_capacity_policy_is_bounded_and_configurable() { let configured = WEB_CONFIG diff --git a/src/config/types.rs b/src/config/types.rs index 12bd690..dfc0026 100644 --- a/src/config/types.rs +++ b/src/config/types.rs @@ -53,7 +53,7 @@ pub use server::{ }; #[allow(unused_imports)] pub use web::{ - WebCarrierNegotiationAggressiveness, WebConfig, WebDecoyConfig, + WebCarrierNegotiationAggressiveness, WebConfig, WebDecoyConfig, WebDecoyFastTrackMode, WebHttpConnectionCapacityAction, WebLimitsConfig, WebProfileConfig, WebSecretMode, WebTimeoutsConfig, WebVhostConfig, }; diff --git a/src/config/types/web.rs b/src/config/types/web.rs index 224e5ae..49592e1 100644 --- a/src/config/types/web.rs +++ b/src/config/types/web.rs @@ -12,6 +12,9 @@ use super::web_debug::WebDebugConfig; // Serialized WEB defaults remain separate from the runtime data model. mod defaults; use defaults::*; +// Decoy fast-track policy remains isolated from the bulky WEB data model. +mod fasttrack; +pub use fasttrack::WebDecoyFastTrackMode; // Accepted-socket overload policy remains separate from the bulky WEB data model. mod overload; pub use overload::WebHttpConnectionCapacityAction; @@ -416,6 +419,9 @@ pub struct WebConfig { /// Controls the evidence thresholds used by automatic carrier ranking. #[serde(default)] pub carrier_negotiation_aggressiveness: WebCarrierNegotiationAggressiveness, + /// Restart-only capability-scan policy for structurally impossible bridge requests. + #[serde(default)] + pub decoy_fasttrack_mode: WebDecoyFastTrackMode, /// Action applied when accepted HTTP connection capacity is exhausted. #[serde(default)] pub http_connection_capacity_action: WebHttpConnectionCapacityAction, @@ -465,6 +471,7 @@ impl Default for WebConfig { carriers: WebCarriers::default(), carrier_learning: default_web_carrier_learning(), carrier_negotiation_aggressiveness: WebCarrierNegotiationAggressiveness::default(), + decoy_fasttrack_mode: WebDecoyFastTrackMode::default(), http_connection_capacity_action: WebHttpConnectionCapacityAction::default(), limits: WebLimitsConfig::default(), debug: WebDebugConfig::default(), diff --git a/src/config/types/web/fasttrack.rs b/src/config/types/web/fasttrack.rs new file mode 100644 index 0000000..2cf7097 --- /dev/null +++ b/src/config/types/web/fasttrack.rs @@ -0,0 +1,28 @@ +use serde::{Deserialize, Serialize}; + +/// Capability-scan policy for structurally impossible WEB bridge requests. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "lowercase")] +pub enum WebDecoyFastTrackMode { + /// Preserve the legacy full scan without collecting fast-track decisions. + #[default] + Off, + /// Record eligible requests while preserving the legacy full scan. + Shadow, + /// Skip the scan only when the public request shape cannot open a bridge. + Enforce, +} + +impl WebDecoyFastTrackMode { + /// Complete fixed mode set in stable API and metric order. + pub const ALL: [Self; 3] = [Self::Off, Self::Shadow, Self::Enforce]; + + /// Returns the stable serialized mode token. + pub const fn as_str(self) -> &'static str { + match self { + Self::Off => "off", + Self::Shadow => "shadow", + Self::Enforce => "enforce", + } + } +} diff --git a/src/config/types/web/runtime.rs b/src/config/types/web/runtime.rs index c59e8b7..0485ee2 100644 --- a/src/config/types/web/runtime.rs +++ b/src/config/types/web/runtime.rs @@ -14,12 +14,16 @@ pub(crate) struct WebRuntimeConfig { pub(crate) struct WebRuntimeVhost { /// Canonical lowercase ACE hostname. pub(crate) host: String, + /// Restart-frozen decoy capability-scan policy. + pub(crate) decoy_fasttrack_mode: WebDecoyFastTrackMode, /// Immutable ordinary-site fallback snapshot. pub(crate) decoy: WebRuntimeDecoy, /// Upstream connect and response-head deadline. pub(crate) decoy_header_secs: u64, /// Exact capability profiles accepted by this host. pub(crate) profiles: Vec>, + /// Contiguous capability table aligned one-to-one with `profiles`. + pub(crate) capabilities: Box<[[u8; 32]]>, } /// Precomputed exact-user capability entry. diff --git a/src/maestro/runtime_build.rs b/src/maestro/runtime_build.rs index d08cb4d..638478f 100644 --- a/src/maestro/runtime_build.rs +++ b/src/maestro/runtime_build.rs @@ -427,6 +427,10 @@ pub(crate) fn resolve_reload_config( fields.push("web.limits".to_string()); effective.web.limits = old.web.limits.clone(); } + if old.web.decoy_fasttrack_mode != desired.web.decoy_fasttrack_mode { + fields.push("web.decoy_fasttrack_mode".to_string()); + effective.web.decoy_fasttrack_mode = old.web.decoy_fasttrack_mode; + } if effective.web.carrier_negotiation_enabled() && effective.web.carrier_learning && effective.web.limits.max_carrier_learning_entries < WEB_CARRIER_LEARNING_MIN_ENTRIES diff --git a/src/maestro/runtime_build_tests.rs b/src/maestro/runtime_build_tests.rs index a0db12e..7fe965b 100644 --- a/src/maestro/runtime_build_tests.rs +++ b/src/maestro/runtime_build_tests.rs @@ -170,6 +170,27 @@ fn web_allocation_limits_are_deferred_until_restart() { assert!(!resolved.runtime_changed); } +#[test] +fn web_decoy_fasttrack_mode_is_deferred_without_runtime_publication() { + let mut old = ProxyConfig::default(); + old.rebuild_runtime_user_auth().unwrap(); + old.rebuild_runtime_web().unwrap(); + let mut desired = old.clone(); + desired.web.decoy_fasttrack_mode = crate::config::WebDecoyFastTrackMode::Enforce; + + let resolved = resolve_reload_config(&old, &desired).unwrap(); + + assert_eq!( + resolved.deferred_process_fields, + vec!["web.decoy_fasttrack_mode".to_string()] + ); + assert_eq!( + resolved.effective.web.decoy_fasttrack_mode, + old.web.decoy_fasttrack_mode + ); + assert!(!resolved.runtime_changed); +} + #[test] fn enabling_learning_is_deferred_when_retained_capacity_is_too_small() { let mut old = ProxyConfig::default(); diff --git a/src/metrics/web.rs b/src/metrics/web.rs index 11d4cf8..7ae2ac0 100644 --- a/src/metrics/web.rs +++ b/src/metrics/web.rs @@ -10,6 +10,8 @@ use crate::web::telemetry::{ WebDecoyUpstreamOutcome, WebHttpConnectionOverloadOutcome, WebRejectionReason, }; +// Decoy fast-track metrics stay isolated from the main WEB renderer. +mod fasttrack; // Session lifecycle and aggregate families stay isolated from capacity rendering. mod lifecycle; @@ -198,6 +200,7 @@ pub(super) fn render(out: &mut String, publication: &WebRuntimePublication, conf ); } + fasttrack::render(out, publication, config); render_carrier_negotiation(out, publication, runtime.as_deref(), config); lifecycle::render(out, publication, config); } diff --git a/src/metrics/web/fasttrack.rs b/src/metrics/web/fasttrack.rs new file mode 100644 index 0000000..77525a0 --- /dev/null +++ b/src/metrics/web/fasttrack.rs @@ -0,0 +1,93 @@ +use std::fmt::Write; + +use crate::config::{ProxyConfig, WebDecoyFastTrackMode}; +use crate::web::control::WebRuntimePublication; +use crate::web::telemetry::WebDecoyFastTrackDisposition; + +/// Renders fixed-cardinality decoy capability-routing metrics. +pub(super) fn render( + out: &mut String, + publication: &WebRuntimePublication, + config: &ProxyConfig, +) { + let _ = writeln!( + out, + "# HELP telemt_web_decoy_fasttrack_mode Effective restart-frozen decoy fast-track mode" + ); + let _ = writeln!(out, "# TYPE telemt_web_decoy_fasttrack_mode gauge"); + for mode in WebDecoyFastTrackMode::ALL { + let _ = writeln!( + out, + "telemt_web_decoy_fasttrack_mode{{mode=\"{}\"}} {}", + mode.as_str(), + u8::from(config.web.decoy_fasttrack_mode == mode) + ); + } + + let _ = writeln!( + out, + "# HELP telemt_web_decoy_fasttrack_requests_total WEB root requests classified by decoy capability-routing work" + ); + let _ = writeln!( + out, + "# TYPE telemt_web_decoy_fasttrack_requests_total counter" + ); + for disposition in WebDecoyFastTrackDisposition::ALL { + let _ = writeln!( + out, + "telemt_web_decoy_fasttrack_requests_total{{disposition=\"{}\"}} {}", + disposition.as_str(), + publication.telemetry.decoy_fasttrack_total(disposition) + ); + } + + let _ = writeln!( + out, + "# HELP telemt_web_decoy_fasttrack_shadow_mismatches_total Shadow decisions that disagreed with legacy bridge eligibility" + ); + let _ = writeln!( + out, + "# TYPE telemt_web_decoy_fasttrack_shadow_mismatches_total counter" + ); + let _ = writeln!( + out, + "telemt_web_decoy_fasttrack_shadow_mismatches_total {}", + publication.telemetry.decoy_fasttrack_shadow_mismatches() + ); +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::web::control::WebRuntimeControl; + + #[test] + fn renderer_emits_one_hot_mode_and_complete_counters() { + let control = WebRuntimeControl::new(); + control + .telemetry() + .record_decoy_fasttrack(WebDecoyFastTrackDisposition::EnforceFastTrack); + control + .telemetry() + .record_decoy_fasttrack_shadow_mismatch(); + let publication = control.subscribe().borrow().clone(); + let mut config = ProxyConfig::default(); + config.web.decoy_fasttrack_mode = WebDecoyFastTrackMode::Enforce; + let mut output = String::new(); + + render(&mut output, &publication, &config); + + assert!(output.contains("telemt_web_decoy_fasttrack_mode{mode=\"off\"} 0")); + assert!(output.contains("telemt_web_decoy_fasttrack_mode{mode=\"enforce\"} 1")); + assert_eq!( + output + .matches("telemt_web_decoy_fasttrack_requests_total{") + .count(), + WebDecoyFastTrackDisposition::ALL.len() + ); + assert!(output.contains( + "telemt_web_decoy_fasttrack_requests_total{disposition=\"enforce_fasttrack\"} 1" + )); + assert!(output.contains("telemt_web_decoy_fasttrack_shadow_mismatches_total 1")); + } +} diff --git a/src/web/http.rs b/src/web/http.rs index fcd5c27..69f95a1 100644 --- a/src/web/http.rs +++ b/src/web/http.rs @@ -16,15 +16,18 @@ use ipnetwork::IpNetwork; use tokio::net::TcpStream; use tokio_util::sync::CancellationToken; -use crate::config::{WebClientIpSource, WebRuntimeVhost}; +use crate::config::{WebClientIpSource, WebDecoyFastTrackMode, WebRuntimeVhost}; use crate::maestro::generation::RuntimeGeneration; use crate::web::bridge; use crate::web::manager::{ManagerError, WebProcessRuntime}; +use crate::web::telemetry::WebDecoyFastTrackDisposition; // Response-body activity keeps connection idle accounting lifecycle-correct. mod activity; // Body collection retains allocation permits through request processing. mod body; +// Canonical capability parsing and complete scans remain isolated from HTTP routing. +mod capability; // Decoy routing and upstream proxying are isolated from carrier authentication. mod decoy; // Downlink long-poll handling remains isolated from request routing. @@ -49,11 +52,12 @@ mod trace_tests; use crate::web::trace::{HttpTraceExchange, TraceDirection, TraceLifecycleEvent, TraceRoute}; use activity::{ActivityBody, ConnectionActivity, RequestActivity, RequestDeadlineHandle}; use body::{CollectBodyError, CollectedBody, RequestBody, collect_body}; +use capability::bridge_candidate; use decoy::serve_decoy; use down::handle_down; use request::{ - bearer_token_hash, binary_content_type, bridge_candidate, canonical_request_host, - canonical_u64_header, client_ip, compatible_cookie_header, match_profile, + bearer_token_hash, binary_content_type, canonical_request_host, canonical_u64_header, + client_ip, compatible_cookie_header, match_profile, }; use response::{ bad_gateway, carrier_empty, carrier_headers, carrier_lane, full_response, generic_not_found, @@ -230,10 +234,48 @@ async fn handle_root( 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 candidate = bridge_candidate(request.uri().query()); + let canonical = candidate.is_canonical(); + let plausible_candidate = canonical && request.method() == Method::GET; + let fasttrack_mode = vhost.decoy_fasttrack_mode; + match fasttrack_mode { + WebDecoyFastTrackMode::Off => {} + WebDecoyFastTrackMode::Shadow => { + runtime.telemetry().record_decoy_fasttrack(if plausible_candidate { + WebDecoyFastTrackDisposition::ShadowCandidateFullScan + } else { + WebDecoyFastTrackDisposition::ShadowWouldFastTrack + }); + } + WebDecoyFastTrackMode::Enforce if !plausible_candidate => { + runtime + .telemetry() + .record_decoy_fasttrack(WebDecoyFastTrackDisposition::EnforceFastTrack); + let recovery_requested = + matches!(representation, recovery::RootRepresentation::Recovery(_)); + if recovery_requested { + strip_query(&mut request); + } + return serve_decoy(request, vhost, recovery_requested, &runtime).await; + } + WebDecoyFastTrackMode::Enforce => { + runtime.telemetry().record_decoy_fasttrack( + WebDecoyFastTrackDisposition::EnforceCandidateFullScan, + ); + } + } + let matched_profile = match_profile(&vhost, candidate.scan_bytes()); let recovery_requested = matches!(representation, recovery::RootRepresentation::Recovery(_)); - let Some(profile) = profile.filter(|_| canonical && request.method() == Method::GET) else { + let profile = matched_profile.filter(|_| canonical && request.method() == Method::GET); + if fasttrack_mode == WebDecoyFastTrackMode::Shadow + && !plausible_candidate + && profile.is_some() + { + runtime + .telemetry() + .record_decoy_fasttrack_shadow_mismatch(); + } + let Some(profile) = profile else { if recovery_requested { strip_query(&mut request); } diff --git a/src/web/http/capability.rs b/src/web/http/capability.rs new file mode 100644 index 0000000..f8fc353 --- /dev/null +++ b/src/web/http/capability.rs @@ -0,0 +1,94 @@ +use base64::Engine as _; +use subtle::{Choice, ConditionallySelectable, ConstantTimeEq}; + +const NON_CANONICAL_BRIDGE_CANDIDATE: [u8; 32] = [0; 32]; + +/// Parsed public bridge query without an allocated credential string. +#[derive(Clone, Copy)] +pub(super) enum BridgeCandidate { + /// The query cannot authenticate a bridge under the public request grammar. + NonCanonical, + /// Exact canonical base64url capability bytes. + Canonical([u8; 32]), +} + +impl BridgeCandidate { + /// Returns whether this query can authenticate a bridge. + pub(super) const fn is_canonical(self) -> bool { + matches!(self, Self::Canonical(_)) + } + + /// Returns the candidate used by the legacy full-scan path. + pub(super) fn scan_bytes(&self) -> &[u8; 32] { + match self { + Self::NonCanonical => &NON_CANONICAL_BRIDGE_CANDIDATE, + Self::Canonical(candidate) => candidate, + } + } +} + +/// Decodes an exact canonical bridge query without allocating credential strings. +pub(super) fn bridge_candidate(query: Option<&str>) -> BridgeCandidate { + let Some(value) = query.and_then(|query| query.strip_prefix("bridge=")) else { + return BridgeCandidate::NonCanonical; + }; + if value.len() != 43 { + return BridgeCandidate::NonCanonical; + } + let mut decoded = [0u8; 32]; + let Ok(decoded_len) = + base64::engine::general_purpose::URL_SAFE_NO_PAD.decode_slice(value, &mut decoded) + else { + return BridgeCandidate::NonCanonical; + }; + let mut canonical = [0u8; 43]; + let Ok(encoded_len) = + base64::engine::general_purpose::URL_SAFE_NO_PAD.encode_slice(decoded, &mut canonical) + else { + return BridgeCandidate::NonCanonical; + }; + if decoded_len != decoded.len() + || encoded_len != canonical.len() + || !bool::from(canonical.ct_eq(value.as_bytes())) + { + return BridgeCandidate::NonCanonical; + } + BridgeCandidate::Canonical(decoded) +} + +/// Internal result of one complete capability-table scan. +pub(super) struct CapabilityScan { + /// Whether any capability matched. + pub(super) matched: Choice, + /// Matching table position selected without a data-dependent branch. + pub(super) matched_index: u64, + #[cfg(test)] + /// Exact comparison count exposed only to deterministic unit tests. + pub(super) comparisons: usize, +} + +/// Scans every configured capability without candidate-dependent control flow. +pub(super) fn scan_capabilities( + capabilities: &[[u8; 32]], + candidate: &[u8; 32], +) -> CapabilityScan { + let mut matched = Choice::from(0); + let mut matched_index = 0u64; + #[cfg(test)] + let mut comparisons = 0usize; + for (index, capability) in capabilities.iter().enumerate() { + let equal = capability.ct_eq(candidate); + matched_index = u64::conditional_select(&matched_index, &(index as u64), equal); + matched |= equal; + #[cfg(test)] + { + comparisons += 1; + } + } + CapabilityScan { + matched, + matched_index, + #[cfg(test)] + comparisons, + } +} diff --git a/src/web/http/decoy_fasttrack_tests.rs b/src/web/http/decoy_fasttrack_tests.rs new file mode 100644 index 0000000..daefbcd --- /dev/null +++ b/src/web/http/decoy_fasttrack_tests.rs @@ -0,0 +1,293 @@ +use super::*; + +use crate::config::WebDecoyFastTrackMode; +use crate::web::telemetry::WebDecoyFastTrackDisposition; + +const RECOVERY_TYPE: &str = "application/vnd.telemt.web-recovery+json"; + +struct Observation { + response: Vec, + counters: [u64; WebDecoyFastTrackDisposition::ALL.len()], + shadow_mismatches: u64, +} + +async fn capture_origin_request(listener: &TcpListener) -> Vec { + let (mut stream, _) = listener.accept().await.unwrap(); + let mut captured = Vec::new(); + let mut buffer = [0u8; 4096]; + loop { + let read = stream.read(&mut buffer).await.unwrap(); + assert_ne!(read, 0, "decoy origin connection closed before the request completed"); + captured.extend_from_slice(&buffer[..read]); + let Some(header_end) = captured.windows(4).position(|window| window == b"\r\n\r\n") + else { + continue; + }; + let headers = std::str::from_utf8(&captured[..header_end]).unwrap(); + let content_length = headers + .lines() + .filter_map(|line| line.split_once(':')) + .find_map(|(name, value)| { + name.eq_ignore_ascii_case("content-length") + .then(|| value.trim().parse::().unwrap()) + }) + .unwrap_or(0); + if captured.len() >= header_end + 4 + content_length { + break; + } + } + stream + .write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 6\r\nConnection: close\r\n\r\norigin") + .await + .unwrap(); + captured +} + +async fn observe_http_origin( + mode: WebDecoyFastTrackMode, + capability: [u8; 32], + request_bytes: Vec, + origin: &TcpListener, +) -> (Observation, Vec) { + let mut config = runtime_config_with_fasttrack(capability, WebCarrier::Https, mode); + let runtime_config = Arc::get_mut(config.web.runtime.as_mut().unwrap()).unwrap(); + let vhost = Arc::get_mut( + runtime_config + .vhosts + .get_mut("proxy.example.com") + .unwrap(), + ) + .unwrap(); + vhost.decoy = WebRuntimeDecoy::HttpUpstream { + addr: origin.local_addr().unwrap(), + authority: "decoy.example".to_string(), + }; + let generation = test_runtime_generation(1, config); + let active_runtime = Arc::new(ArcSwap::from(Arc::clone(&generation))); + let runtime = WebProcessRuntime::start(active_runtime); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let (response, captured) = tokio::join!( + request(&listener, &runtime, request_bytes), + capture_origin_request(origin) + ); + let counters = WebDecoyFastTrackDisposition::ALL + .map(|disposition| runtime.telemetry().decoy_fasttrack_total(disposition)); + let shadow_mismatches = runtime.telemetry().decoy_fasttrack_shadow_mismatches(); + + runtime.shutdown().await; + generation.stop_sessions().await; + generation.stop_background_tasks().await; + + ( + Observation { + response, + counters, + shadow_mismatches, + }, + captured, + ) +} + +async fn observe(mode: WebDecoyFastTrackMode, capability: [u8; 32], request_bytes: Vec) -> Observation { + let generation = test_runtime_generation( + 1, + runtime_config_with_fasttrack(capability, WebCarrier::Https, mode), + ); + let active_runtime = Arc::new(ArcSwap::from(Arc::clone(&generation))); + let runtime = WebProcessRuntime::start(active_runtime); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let response = request(&listener, &runtime, request_bytes).await; + let counters = WebDecoyFastTrackDisposition::ALL + .map(|disposition| runtime.telemetry().decoy_fasttrack_total(disposition)); + let shadow_mismatches = runtime.telemetry().decoy_fasttrack_shadow_mismatches(); + + runtime.shutdown().await; + generation.stop_sessions().await; + generation.stop_background_tasks().await; + + Observation { + response, + counters, + shadow_mismatches, + } +} + +fn root_request(method: &str, query: &str) -> Vec { + format!( + "{method} /{query} HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.90\r\nConnection: close\r\n\r\n" + ) + .into_bytes() +} + +#[tokio::test] +async fn impossible_root_shapes_preserve_decoy_bytes_and_follow_the_selected_mode() { + let capability = [70u8; 32]; + let encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(capability); + for request_bytes in [ + root_request("GET", ""), + root_request("GET", "?bridge=not-canonical"), + root_request("HEAD", &format!("?bridge={encoded}")), + ] { + let off = observe(WebDecoyFastTrackMode::Off, capability, request_bytes.clone()).await; + let shadow = observe( + WebDecoyFastTrackMode::Shadow, + capability, + request_bytes.clone(), + ) + .await; + let enforce = observe(WebDecoyFastTrackMode::Enforce, capability, request_bytes).await; + + assert_eq!(shadow.response, off.response); + assert_eq!(enforce.response, off.response); + assert_eq!(off.counters, [0, 0, 0, 0]); + assert_eq!(shadow.counters, [1, 0, 0, 0]); + assert_eq!(enforce.counters, [0, 0, 1, 0]); + assert_eq!(shadow.shadow_mismatches, 0); + } +} + +#[tokio::test] +async fn canonical_hit_and_miss_always_retain_the_full_scan() { + let capability = [71u8; 32]; + for candidate in [capability, [72u8; 32]] { + let encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(candidate); + let request_bytes = root_request("GET", &format!("?bridge={encoded}")); + + for (mode, expected_counters) in [ + (WebDecoyFastTrackMode::Off, [0, 0, 0, 0]), + (WebDecoyFastTrackMode::Shadow, [0, 1, 0, 0]), + (WebDecoyFastTrackMode::Enforce, [0, 0, 0, 1]), + ] { + let observation = observe(mode, capability, request_bytes.clone()).await; + assert!(observation.response.starts_with(b"HTTP/1.1 200")); + assert_eq!(observation.counters, expected_counters); + assert_eq!(observation.shadow_mismatches, 0); + if candidate == capability { + assert!( + observation + .response + .windows(b"bootstrap=\"".len()) + .any(|window| window == b"bootstrap=\"") + ); + } else { + let (_, body) = split_response(&observation.response); + assert_eq!(body, b"decoy"); + } + } + } +} + +#[tokio::test] +async fn recovery_sanitization_precedes_fasttrack_classification() { + let capability = [73u8; 32]; + let request_bytes = format!( + "GET /?bridge=not-canonical HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.90\r\nAccept: {RECOVERY_TYPE}\r\nAuthorization: Bearer {}\r\nContent-Length: 4\r\nConnection: close\r\n\r\nbody", + "U".repeat(43) + ) + .into_bytes(); + + let off = observe(WebDecoyFastTrackMode::Off, capability, request_bytes.clone()).await; + let shadow = observe( + WebDecoyFastTrackMode::Shadow, + capability, + request_bytes.clone(), + ) + .await; + let enforce = observe(WebDecoyFastTrackMode::Enforce, capability, request_bytes).await; + + assert_eq!(shadow.response, off.response); + assert_eq!(enforce.response, off.response); + assert_eq!(response_header(split_response(&off.response).0, "cache-control"), "no-store"); + assert_eq!(off.counters, [0, 0, 0, 0]); + assert_eq!(shadow.counters, [1, 0, 0, 0]); + assert_eq!(enforce.counters, [0, 0, 1, 0]); + assert_eq!(shadow.shadow_mismatches, 0); +} + +#[tokio::test] +async fn http_origin_forwarding_is_identical_across_modes() { + let capability = [74u8; 32]; + let origin = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let mut request_bytes = b"GET /?bridge=not-canonical HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.90\r\nX-Ordinary: preserved\r\nContent-Length: 4\r\nConnection: close\r\n\r\n".to_vec(); + request_bytes.extend_from_slice(b"body"); + + let (off, off_upstream) = observe_http_origin( + WebDecoyFastTrackMode::Off, + capability, + request_bytes.clone(), + &origin, + ) + .await; + let (shadow, shadow_upstream) = observe_http_origin( + WebDecoyFastTrackMode::Shadow, + capability, + request_bytes.clone(), + &origin, + ) + .await; + let (enforce, enforce_upstream) = observe_http_origin( + WebDecoyFastTrackMode::Enforce, + capability, + request_bytes, + &origin, + ) + .await; + + assert_eq!(shadow.response, off.response); + assert_eq!(enforce.response, off.response); + assert_eq!(shadow_upstream, off_upstream); + assert_eq!(enforce_upstream, off_upstream); + assert!(off_upstream.starts_with(b"GET /?bridge=not-canonical HTTP/1.1\r\n")); + assert!(off_upstream.windows(21).any(|window| window == b"x-ordinary: preserved")); + assert!(off_upstream.ends_with(b"\r\n\r\nbody")); +} + +#[tokio::test] +async fn http_origin_recovery_sanitization_is_identical_across_modes() { + let capability = [75u8; 32]; + let origin = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let request_bytes = format!( + "GET /?bridge=not-canonical HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.90\r\nAccept: {RECOVERY_TYPE}\r\nAuthorization: Bearer {}\r\nContent-Type: application/octet-stream\r\nContent-Length: 4\r\nX-Up-Seq: 9\r\nConnection: close\r\n\r\nbody", + "U".repeat(43) + ) + .into_bytes(); + + let (off, off_upstream) = observe_http_origin( + WebDecoyFastTrackMode::Off, + capability, + request_bytes.clone(), + &origin, + ) + .await; + let (shadow, shadow_upstream) = observe_http_origin( + WebDecoyFastTrackMode::Shadow, + capability, + request_bytes.clone(), + &origin, + ) + .await; + let (enforce, enforce_upstream) = observe_http_origin( + WebDecoyFastTrackMode::Enforce, + capability, + request_bytes, + &origin, + ) + .await; + + assert_eq!(shadow.response, off.response); + assert_eq!(enforce.response, off.response); + assert_eq!(shadow_upstream, off_upstream); + assert_eq!(enforce_upstream, off_upstream); + assert!(off_upstream.starts_with(b"GET / HTTP/1.1\r\n")); + let lowercase = String::from_utf8_lossy(&off_upstream).to_ascii_lowercase(); + for forbidden in [ + "authorization:", + "accept:", + "content-type:", + "content-length:", + "x-up-seq:", + ] { + assert!(!lowercase.contains(forbidden)); + } + assert!(off_upstream.ends_with(b"\r\n\r\n")); +} diff --git a/src/web/http/request.rs b/src/web/http/request.rs index c150096..9e2a1c5 100644 --- a/src/web/http/request.rs +++ b/src/web/http/request.rs @@ -13,53 +13,27 @@ use crate::web::manager::{ const USER_AGENT_CONTEXT: &[u8] = b"telemt-web-carrier-user-agent-v1\0"; +use super::capability::scan_capabilities; + // Canonical host and forwarded-address provenance remain isolated from credentials. mod identity; pub(super) use identity::{canonical_request_host, carrier_ip_learning_eligible, client_ip}; -/// Decodes an exact canonical bridge query without allocating credential strings. -pub(super) fn bridge_candidate(query: Option<&str>) -> ([u8; 32], bool) { - let mut candidate = [0u8; 32]; - let Some(value) = query.and_then(|query| query.strip_prefix("bridge=")) else { - return (candidate, false); - }; - if value.len() != 43 { - return (candidate, false); - } - let mut decoded = [0u8; 32]; - let Ok(decoded_len) = - base64::engine::general_purpose::URL_SAFE_NO_PAD.decode_slice(value, &mut decoded) - else { - return (candidate, false); - }; - let mut canonical = [0u8; 43]; - let Ok(encoded_len) = - base64::engine::general_purpose::URL_SAFE_NO_PAD.encode_slice(decoded, &mut canonical) - else { - return (candidate, false); - }; - if decoded_len != decoded.len() - || encoded_len != canonical.len() - || !bool::from(canonical.ct_eq(value.as_bytes())) - { - return (candidate, false); - } - candidate = decoded; - (candidate, true) -} - -/// Matches a capability in constant time across every profile of one virtual host. +/// Matches one capability after a complete branchless scan of the virtual host table. pub(super) fn match_profile( vhost: &WebRuntimeVhost, candidate: &[u8; 32], ) -> Option> { - let mut matched = None; - for profile in &vhost.profiles { - if bool::from(profile.capability.ct_eq(candidate)) { - matched = Some(Arc::clone(profile)); - } + debug_assert_eq!(vhost.capabilities.len(), vhost.profiles.len()); + let scan = scan_capabilities(&vhost.capabilities, candidate); + if bool::from(scan.matched) { + usize::try_from(scan.matched_index) + .ok() + .and_then(|index| vhost.profiles.get(index)) + .map(Arc::clone) + } else { + None } - matched } /// Validates and hashes one canonical bearer credential for map lookup. diff --git a/src/web/http/request/tests.rs b/src/web/http/request/tests.rs index f1b536e..e43665d 100644 --- a/src/web/http/request/tests.rs +++ b/src/web/http/request/tests.rs @@ -1,14 +1,45 @@ use super::*; use ipnetwork::IpNetwork; +use proptest::prelude::*; use crate::config::{WebCarrier, WebClientIpSource}; +use crate::web::http::capability::{bridge_candidate, scan_capabilities}; #[test] fn canonical_bridge_query_rejects_aliases() { let token = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode([7u8; 32]); - assert!(bridge_candidate(Some(&format!("bridge={token}"))).1); - assert!(!bridge_candidate(Some(&format!("x=1&bridge={token}"))).1); - assert!(!bridge_candidate(Some(&format!("bridge={token}="))).1); + assert!(bridge_candidate(Some(&format!("bridge={token}"))).is_canonical()); + assert!(!bridge_candidate(Some(&format!("x=1&bridge={token}"))).is_canonical()); + assert!(!bridge_candidate(Some(&format!("bridge={token}="))).is_canonical()); +} + +#[test] +fn capability_scan_checks_every_entry_independent_of_match_position() { + let capabilities = [[1u8; 32], [2u8; 32], [3u8; 32], [4u8; 32]]; + for (expected_index, candidate) in capabilities.iter().enumerate() { + let scan = scan_capabilities(&capabilities, candidate); + assert!(bool::from(scan.matched)); + assert_eq!(scan.matched_index as usize, expected_index); + assert_eq!(scan.comparisons, capabilities.len()); + } + + let miss = scan_capabilities(&capabilities, &[99u8; 32]); + assert!(!bool::from(miss.matched)); + assert_eq!(miss.comparisons, capabilities.len()); + + let empty = scan_capabilities(&[], &[99u8; 32]); + assert!(!bool::from(empty.matched)); + assert_eq!(empty.comparisons, 0); +} + +proptest! { + #[test] + fn every_capability_has_one_canonical_bridge_query(capability in any::<[u8; 32]>()) { + let token = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(capability); + let candidate = bridge_candidate(Some(&format!("bridge={token}"))); + prop_assert!(candidate.is_canonical()); + prop_assert_eq!(candidate.scan_bytes(), &capability); + } } #[test] diff --git a/src/web/http/tests.rs b/src/web/http/tests.rs index 505b2e2..c9545e2 100644 --- a/src/web/http/tests.rs +++ b/src/web/http/tests.rs @@ -11,7 +11,8 @@ use tokio_util::sync::CancellationToken; use super::serve_connection; use crate::config::{ ProxyConfig, WebCarrier, WebCarriers, WebClientIpSource, WebRuntimeConfig, WebRuntimeDecoy, - WebRuntimeProfile, WebRuntimeVhost, WebSecretMode, WebStaticAsset, WebStaticSite, + WebDecoyFastTrackMode, WebRuntimeProfile, WebRuntimeVhost, WebSecretMode, WebStaticAsset, + WebStaticSite, }; use crate::maestro::generation::test_runtime_generation; use crate::web::frame::{self, FrameType}; @@ -39,6 +40,9 @@ mod operator_lifecycle_tests; // Positive-only recovery representation coverage remains isolated from ordinary root routing. #[path = "recovery_tests.rs"] mod recovery_tests; +// Decoy fast-track routing and telemetry remain isolated from carrier protocol scenarios. +#[path = "decoy_fasttrack_tests.rs"] +mod decoy_fasttrack_tests; const TEST_CARRIER_DEADLINES_SECS: [u64; 4] = [3, 5, 8, 12]; @@ -46,6 +50,21 @@ pub(super) fn runtime_config(capability: [u8; 32], carrier: WebCarrier) -> Proxy runtime_config_with_carriers(capability, carrier, false, true, Arc::from([carrier])) } +/// Builds a static-decoy runtime with one restart-frozen fast-track mode. +pub(super) fn runtime_config_with_fasttrack( + capability: [u8; 32], + carrier: WebCarrier, + mode: WebDecoyFastTrackMode, +) -> ProxyConfig { + let mut config = runtime_config(capability, carrier); + config.web.decoy_fasttrack_mode = mode; + let runtime = Arc::get_mut(config.web.runtime.as_mut().unwrap()).unwrap(); + for vhost in runtime.vhosts.values_mut() { + Arc::get_mut(vhost).unwrap().decoy_fasttrack_mode = mode; + } + config +} + pub(super) fn negotiation_runtime_config( capability: [u8; 32], carrier: WebCarrier, @@ -128,9 +147,11 @@ fn runtime_config_with_carriers_and_deadlines( }); let vhost = Arc::new(WebRuntimeVhost { host: "proxy.example.com".to_string(), + decoy_fasttrack_mode: WebDecoyFastTrackMode::Off, decoy: WebRuntimeDecoy::StaticDirectory(Arc::clone(&site)), decoy_header_secs: 1, profiles: vec![Arc::clone(&profile)], + capabilities: vec![capability].into_boxed_slice(), }); let mut vhosts = BTreeMap::new(); vhosts.insert("proxy.example.com".to_string(), vhost); @@ -138,9 +159,11 @@ fn runtime_config_with_carriers_and_deadlines( "other.example.com".to_string(), Arc::new(WebRuntimeVhost { host: "other.example.com".to_string(), + decoy_fasttrack_mode: WebDecoyFastTrackMode::Off, decoy: WebRuntimeDecoy::StaticDirectory(site), decoy_header_secs: 1, profiles: Vec::new(), + capabilities: Vec::new().into_boxed_slice(), }), ); let mut config = ProxyConfig::default(); diff --git a/src/web/telemetry.rs b/src/web/telemetry.rs index 4310630..0145382 100644 --- a/src/web/telemetry.rs +++ b/src/web/telemetry.rs @@ -10,6 +10,9 @@ pub(crate) use carrier::{ WebCarrierFailureCounter, WebCarrierFailurePhase, WebCarrierLearningCounter, WebCarrierLearningOutcome, WebCarrierSelectionCounter, WebCarrierSelectionDisposition, }; +mod fasttrack; +use fasttrack::DECOY_FASTTRACK_SLOTS; +pub(crate) use fasttrack::{WebDecoyFastTrackCounter, WebDecoyFastTrackDisposition}; mod lifecycle; use lifecycle::{SESSION_CLOSE_SLOTS, SESSION_OBSERVATION_SLOTS}; pub(crate) use lifecycle::{ @@ -316,6 +319,8 @@ pub(crate) struct WebTelemetry { carrier_selections: [AtomicU64; CARRIER_SELECTION_SLOTS], carrier_failures: [AtomicU64; CARRIER_FAILURE_SLOTS], carrier_learning_outcomes: [AtomicU64; CARRIER_LEARNING_SLOTS], + decoy_fasttrack_requests: [AtomicU64; DECOY_FASTTRACK_SLOTS], + decoy_fasttrack_shadow_mismatches: AtomicU64, session_closures: [AtomicU64; SESSION_CLOSE_SLOTS], session_observations: [AtomicU64; SESSION_OBSERVATION_SLOTS], bridge_recovery_events: [AtomicU64; WebBridgeRecoveryEvent::ALL.len()], @@ -343,6 +348,8 @@ 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)), + decoy_fasttrack_requests: std::array::from_fn(|_| AtomicU64::new(0)), + decoy_fasttrack_shadow_mismatches: 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)), diff --git a/src/web/telemetry/fasttrack.rs b/src/web/telemetry/fasttrack.rs new file mode 100644 index 0000000..3b1f776 --- /dev/null +++ b/src/web/telemetry/fasttrack.rs @@ -0,0 +1,85 @@ +use std::sync::atomic::Ordering; + +use serde::Serialize; + +use super::WebTelemetry; + +/// Terminal capability-routing work selected for one WEB root request. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +#[repr(usize)] +pub(crate) enum WebDecoyFastTrackDisposition { + /// Shadow mode identified a request that enforce mode would bypass. + ShadowWouldFastTrack, + /// Shadow mode retained a full scan for a plausible capability request. + ShadowCandidateFullScan, + /// Enforce mode bypassed a structurally impossible capability request. + EnforceFastTrack, + /// Enforce mode retained a full scan for a plausible capability request. + EnforceCandidateFullScan, +} + +impl WebDecoyFastTrackDisposition { + /// Complete fixed disposition set in stable API and metric order. + pub(crate) const ALL: [Self; 4] = [ + Self::ShadowWouldFastTrack, + Self::ShadowCandidateFullScan, + Self::EnforceFastTrack, + Self::EnforceCandidateFullScan, + ]; + + /// Returns the stable API and Prometheus label token. + pub(crate) const fn as_str(self) -> &'static str { + match self { + Self::ShadowWouldFastTrack => "shadow_would_fasttrack", + Self::ShadowCandidateFullScan => "shadow_candidate_full_scan", + Self::EnforceFastTrack => "enforce_fasttrack", + Self::EnforceCandidateFullScan => "enforce_candidate_full_scan", + } + } +} + +pub(super) const DECOY_FASTTRACK_SLOTS: usize = WebDecoyFastTrackDisposition::ALL.len(); + +/// API-safe fixed decoy fast-track counter. +#[derive(Clone, Serialize)] +pub(crate) struct WebDecoyFastTrackCounter { + /// Stable capability-routing disposition token. + pub(crate) disposition: &'static str, + /// Process-lifetime event count. + pub(crate) total: u64, +} + +impl WebTelemetry { + /// Records one shadow or enforce capability-routing disposition. + pub(crate) fn record_decoy_fasttrack(&self, disposition: WebDecoyFastTrackDisposition) { + self.decoy_fasttrack_requests[disposition as usize].fetch_add(1, Ordering::Relaxed); + } + + /// Returns one fixed decoy fast-track counter. + pub(crate) fn decoy_fasttrack_total(&self, disposition: WebDecoyFastTrackDisposition) -> u64 { + self.decoy_fasttrack_requests[disposition as usize].load(Ordering::Relaxed) + } + + /// Captures the complete decoy fast-track counter set. + pub(crate) fn decoy_fasttrack_counters(&self) -> Vec { + WebDecoyFastTrackDisposition::ALL + .into_iter() + .map(|disposition| WebDecoyFastTrackCounter { + disposition: disposition.as_str(), + total: self.decoy_fasttrack_total(disposition), + }) + .collect() + } + + /// Records a shadow decision that disagreed with legacy bridge eligibility. + pub(crate) fn record_decoy_fasttrack_shadow_mismatch(&self) { + self.decoy_fasttrack_shadow_mismatches + .fetch_add(1, Ordering::Relaxed); + } + + /// Returns shadow decisions that disagreed with legacy bridge eligibility. + pub(crate) fn decoy_fasttrack_shadow_mismatches(&self) -> u64 { + self.decoy_fasttrack_shadow_mismatches + .load(Ordering::Relaxed) + } +} diff --git a/src/web/telemetry/tests.rs b/src/web/telemetry/tests.rs index 30f7587..613f6b4 100644 --- a/src/web/telemetry/tests.rs +++ b/src/web/telemetry/tests.rs @@ -18,6 +18,7 @@ fn fixed_counter_sets_and_acceptor_guard_are_exact() { CarrierFailure::Network, ); telemetry.record_carrier_learning(WebCarrier::Https, WebCarrierLearningOutcome::Recorded); + telemetry.record_decoy_fasttrack(WebDecoyFastTrackDisposition::ShadowWouldFastTrack); telemetry.record_session_closed(WebCarrier::Https, SessionCloseReason::ApiClose); telemetry.record_session_observation( WebCarrier::Https, @@ -48,6 +49,15 @@ fn fixed_counter_sets_and_acceptor_guard_are_exact() { telemetry.carrier_learning_counters().len(), WebCarrier::ALL.len() * WebCarrierLearningOutcome::ALL.len() ); + assert_eq!( + telemetry.decoy_fasttrack_counters().len(), + WebDecoyFastTrackDisposition::ALL.len() + ); + assert_eq!( + telemetry.decoy_fasttrack_total(WebDecoyFastTrackDisposition::ShadowWouldFastTrack), + 1 + ); + assert_eq!(telemetry.decoy_fasttrack_shadow_mismatches(), 0); assert_eq!( telemetry.session_close_counters().len(), WebCarrier::ALL.len() * SessionCloseReason::ALL.len()