Restart-gated Decoy Fast-Track

Co-Authored-By: brekotis <93345790+brekotis@users.noreply.github.com>
This commit is contained in:
Alexey
2026-09-06 18:51:13 +03:00
parent 106b26a5b7
commit 20a4d50524
19 changed files with 343 additions and 171 deletions
+10 -16
View File
@@ -241,11 +241,13 @@ async fn handle_root(
match fasttrack_mode {
WebDecoyFastTrackMode::Off => {}
WebDecoyFastTrackMode::Shadow => {
runtime.telemetry().record_decoy_fasttrack(if plausible_candidate {
WebDecoyFastTrackDisposition::ShadowCandidateFullScan
} else {
WebDecoyFastTrackDisposition::ShadowWouldFastTrack
});
runtime
.telemetry()
.record_decoy_fasttrack(if plausible_candidate {
WebDecoyFastTrackDisposition::ShadowCandidateFullScan
} else {
WebDecoyFastTrackDisposition::ShadowWouldFastTrack
});
}
WebDecoyFastTrackMode::Enforce if !plausible_candidate => {
runtime
@@ -259,22 +261,14 @@ async fn handle_root(
return serve_decoy(request, vhost, recovery_requested, &runtime).await;
}
WebDecoyFastTrackMode::Enforce => {
runtime.telemetry().record_decoy_fasttrack(
WebDecoyFastTrackDisposition::EnforceCandidateFullScan,
);
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 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);
+1 -4
View File
@@ -68,10 +68,7 @@ pub(super) struct CapabilityScan {
}
/// Scans every configured capability without candidate-dependent control flow.
pub(super) fn scan_capabilities(
capabilities: &[[u8; 32]],
candidate: &[u8; 32],
) -> CapabilityScan {
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)]
+112 -34
View File
@@ -8,7 +8,20 @@ const RECOVERY_TYPE: &str = "application/vnd.telemt.web-recovery+json";
struct Observation {
response: Vec<u8>,
counters: [u64; WebDecoyFastTrackDisposition::ALL.len()],
shadow_mismatches: u64,
}
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
}
async fn capture_origin_request(listener: &TcpListener) -> Vec<u8> {
@@ -17,10 +30,12 @@ async fn capture_origin_request(listener: &TcpListener) -> Vec<u8> {
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");
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 {
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();
@@ -51,13 +66,7 @@ async fn observe_http_origin(
) -> (Observation, Vec<u8>) {
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();
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(),
@@ -72,23 +81,19 @@ async fn observe_http_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,
)
(Observation { response, counters }, captured)
}
async fn observe(mode: WebDecoyFastTrackMode, capability: [u8; 32], request_bytes: Vec<u8>) -> Observation {
async fn observe(
mode: WebDecoyFastTrackMode,
capability: [u8; 32],
request_bytes: Vec<u8>,
) -> Observation {
let generation = test_runtime_generation(
1,
runtime_config_with_fasttrack(capability, WebCarrier::Https, mode),
@@ -99,17 +104,12 @@ async fn observe(mode: WebDecoyFastTrackMode, capability: [u8; 32], request_byte
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,
}
Observation { response, counters }
}
fn root_request(method: &str, query: &str) -> Vec<u8> {
@@ -128,7 +128,12 @@ async fn impossible_root_shapes_preserve_decoy_bytes_and_follow_the_selected_mod
root_request("GET", "?bridge=not-canonical"),
root_request("HEAD", &format!("?bridge={encoded}")),
] {
let off = observe(WebDecoyFastTrackMode::Off, capability, request_bytes.clone()).await;
let off = observe(
WebDecoyFastTrackMode::Off,
capability,
request_bytes.clone(),
)
.await;
let shadow = observe(
WebDecoyFastTrackMode::Shadow,
capability,
@@ -142,7 +147,6 @@ async fn impossible_root_shapes_preserve_decoy_bytes_and_follow_the_selected_mod
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);
}
}
@@ -161,7 +165,6 @@ async fn canonical_hit_and_miss_always_retain_the_full_scan() {
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
@@ -186,7 +189,12 @@ async fn recovery_sanitization_precedes_fasttrack_classification() {
)
.into_bytes();
let off = observe(WebDecoyFastTrackMode::Off, capability, request_bytes.clone()).await;
let off = observe(
WebDecoyFastTrackMode::Off,
capability,
request_bytes.clone(),
)
.await;
let shadow = observe(
WebDecoyFastTrackMode::Shadow,
capability,
@@ -197,11 +205,42 @@ async fn recovery_sanitization_precedes_fasttrack_classification() {
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!(
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 canonical_recovery_hit_and_miss_always_retain_the_full_scan() {
let capability = [77u8; 32];
let bearer = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode([78u8; 32]);
for candidate in [capability, [79u8; 32]] {
let encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(candidate);
let request_bytes = format!(
"GET /?bridge={encoded} HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.90\r\nAccept: {RECOVERY_TYPE}\r\nAuthorization: Bearer {bearer}\r\nConnection: close\r\n\r\n"
)
.into_bytes();
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;
let (headers, body) = split_response(&observation.response);
assert_eq!(observation.counters, expected_counters);
if candidate == capability {
assert_eq!(response_header(headers, "content-type"), RECOVERY_TYPE);
assert!(serde_json::from_slice::<serde_json::Value>(body).is_ok());
} else {
assert_eq!(body, b"<!doctype html><title>decoy</title>");
}
}
}
}
#[tokio::test]
@@ -238,7 +277,11 @@ async fn http_origin_forwarding_is_identical_across_modes() {
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
.windows(21)
.any(|window| window == b"x-ordinary: preserved")
);
assert!(off_upstream.ends_with(b"\r\n\r\nbody"));
}
@@ -291,3 +334,38 @@ async fn http_origin_recovery_sanitization_is_identical_across_modes() {
}
assert!(off_upstream.ends_with(b"\r\n\r\n"));
}
#[tokio::test]
async fn fasttrack_counters_remain_process_owned_across_generation_swap() {
let capability = [76u8; 32];
let initial = test_runtime_generation(
1,
runtime_config_with_fasttrack(capability, WebCarrier::Https, WebDecoyFastTrackMode::Shadow),
);
let active_runtime = Arc::new(ArcSwap::from(Arc::clone(&initial)));
let runtime = WebProcessRuntime::start(Arc::clone(&active_runtime));
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let first = request(&listener, &runtime, root_request("GET", "")).await;
assert!(first.starts_with(b"HTTP/1.1 200"));
let replacement = test_runtime_generation(
2,
runtime_config_with_fasttrack(capability, WebCarrier::Https, WebDecoyFastTrackMode::Shadow),
);
active_runtime.store(Arc::clone(&replacement));
let second = request(&listener, &runtime, root_request("GET", "")).await;
assert!(second.starts_with(b"HTTP/1.1 200"));
assert_eq!(
runtime
.telemetry()
.decoy_fasttrack_total(WebDecoyFastTrackDisposition::ShadowWouldFastTrack),
2
);
runtime.shutdown().await;
initial.stop_sessions().await;
initial.stop_background_tasks().await;
replacement.stop_sessions().await;
replacement.stop_background_tasks().await;
}
+6 -1
View File
@@ -24,12 +24,17 @@ pub(super) fn match_profile(
vhost: &WebRuntimeVhost,
candidate: &[u8; 32],
) -> Option<Arc<WebRuntimeProfile>> {
debug_assert_eq!(vhost.capabilities.len(), vhost.profiles.len());
// Runtime construction owns alignment; any future constructor drift must fail closed.
if vhost.capabilities.len() != vhost.profiles.len() {
return None;
}
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))
// Revalidate the selected identity after the index-only complete scan.
.filter(|profile| bool::from(profile.capability.ct_eq(candidate)))
.map(Arc::clone)
} else {
None
+14
View File
@@ -32,6 +32,20 @@ fn capability_scan_checks_every_entry_independent_of_match_position() {
assert_eq!(empty.comparisons, 0);
}
#[test]
fn capability_profile_table_drift_fails_closed() {
let capability = [5u8; 32];
let mut config = crate::web::http::tests::runtime_config(capability, WebCarrier::Https);
let runtime = Arc::get_mut(config.web.runtime.as_mut().unwrap()).unwrap();
let vhost = Arc::get_mut(runtime.vhosts.get_mut("proxy.example.com").unwrap()).unwrap();
assert!(match_profile(vhost, &capability).is_some());
vhost.capabilities[0] = [6u8; 32];
assert!(match_profile(vhost, &[6u8; 32]).is_none());
vhost.capabilities = Vec::new().into_boxed_slice();
assert!(match_profile(vhost, &capability).is_none());
}
proptest! {
#[test]
fn every_capability_has_one_canonical_bridge_query(capability in any::<[u8; 32]>()) {
+18
View File
@@ -0,0 +1,18 @@
/// Splits one complete raw HTTP response into head and body slices.
pub(in crate::web::http) fn split_response(response: &[u8]) -> (&[u8], &[u8]) {
let separator = response
.windows(4)
.position(|window| window == b"\r\n\r\n")
.unwrap();
(&response[..separator], &response[separator + 4..])
}
/// Returns one case-insensitive header value from a raw HTTP response head.
pub(in crate::web::http) fn response_header<'a>(headers: &'a [u8], name: &str) -> &'a str {
std::str::from_utf8(headers)
.unwrap()
.lines()
.filter_map(|line| line.split_once(':'))
.find_map(|(header, value)| header.eq_ignore_ascii_case(name).then_some(value.trim()))
.unwrap()
}
+12 -35
View File
@@ -10,9 +10,9 @@ use tokio_util::sync::CancellationToken;
use super::serve_connection;
use crate::config::{
ProxyConfig, WebCarrier, WebCarriers, WebClientIpSource, WebRuntimeConfig, WebRuntimeDecoy,
WebDecoyFastTrackMode, WebRuntimeProfile, WebRuntimeVhost, WebSecretMode, WebStaticAsset,
WebStaticSite,
ProxyConfig, WebCarrier, WebCarriers, WebClientIpSource, WebDecoyFastTrackMode,
WebRuntimeConfig, WebRuntimeDecoy, WebRuntimeProfile, WebRuntimeVhost, WebSecretMode,
WebStaticAsset, WebStaticSite,
};
use crate::maestro::generation::test_runtime_generation;
use crate::web::frame::{self, FrameType};
@@ -43,28 +43,20 @@ mod recovery_tests;
// Decoy fast-track routing and telemetry remain isolated from carrier protocol scenarios.
#[path = "decoy_fasttrack_tests.rs"]
mod decoy_fasttrack_tests;
// Raw response parsing helpers are shared by the HTTP integration test modules.
#[path = "response_test_support.rs"]
mod response_test_support;
pub(super) use response_test_support::{response_header, split_response};
const TEST_CARRIER_DEADLINES_SECS: [u64; 4] = [3, 5, 8, 12];
/// Builds the default static-decoy runtime used by WEB HTTP tests.
pub(super) fn runtime_config(capability: [u8; 32], carrier: WebCarrier) -> ProxyConfig {
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
}
/// Builds a negotiation-enabled static-decoy runtime for carrier tests.
pub(super) fn negotiation_runtime_config(
capability: [u8; 32],
carrier: WebCarrier,
@@ -91,6 +83,7 @@ fn runtime_config_with_carriers(
)
}
/// Builds a negotiation runtime with explicit cumulative carrier deadlines.
pub(super) fn negotiation_runtime_config_with_deadlines(
capability: [u8; 32],
carrier: WebCarrier,
@@ -185,6 +178,7 @@ fn runtime_config_with_carriers_and_deadlines(
config
}
/// Serves one raw HTTP request through the private WEB listener test harness.
pub(super) async fn request(
listener: &TcpListener,
runtime: &Arc<WebProcessRuntime>,
@@ -211,23 +205,6 @@ pub(super) async fn request(
response
}
pub(super) fn split_response(response: &[u8]) -> (&[u8], &[u8]) {
let separator = response
.windows(4)
.position(|window| window == b"\r\n\r\n")
.unwrap();
(&response[..separator], &response[separator + 4..])
}
pub(super) fn response_header<'a>(headers: &'a [u8], name: &str) -> &'a str {
std::str::from_utf8(headers)
.unwrap()
.lines()
.filter_map(|line| line.split_once(':'))
.find_map(|(header, value)| header.eq_ignore_ascii_case(name).then_some(value.trim()))
.unwrap()
}
#[tokio::test]
async fn https_carrier_bootstraps_and_closes_one_session() {
let capability = [7u8; 32];
+1 -2
View File
@@ -10,6 +10,7 @@ pub(crate) use carrier::{
WebCarrierFailureCounter, WebCarrierFailurePhase, WebCarrierLearningCounter,
WebCarrierLearningOutcome, WebCarrierSelectionCounter, WebCarrierSelectionDisposition,
};
// Decoy fast-track counters retain a fixed process-owned disposition set.
mod fasttrack;
use fasttrack::DECOY_FASTTRACK_SLOTS;
pub(crate) use fasttrack::{WebDecoyFastTrackCounter, WebDecoyFastTrackDisposition};
@@ -320,7 +321,6 @@ pub(crate) struct WebTelemetry {
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()],
@@ -349,7 +349,6 @@ impl WebTelemetry {
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)),
+1 -12
View File
@@ -38,6 +38,7 @@ impl WebDecoyFastTrackDisposition {
}
}
/// Fixed storage width for process-owned decoy fast-track counters.
pub(super) const DECOY_FASTTRACK_SLOTS: usize = WebDecoyFastTrackDisposition::ALL.len();
/// API-safe fixed decoy fast-track counter.
@@ -70,16 +71,4 @@ impl WebTelemetry {
})
.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)
}
}
-1
View File
@@ -57,7 +57,6 @@ fn fixed_counter_sets_and_acceptor_guard_are_exact() {
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()