From 1bd85535f4276ce81bb89ffcf89abe07e247d22a Mon Sep 17 00:00:00 2001 From: Alexey <247128645+axkurcom@users.noreply.github.com> Date: Sun, 23 Aug 2026 16:32:32 +0300 Subject: [PATCH 01/10] WEB Handshake OPEN -> DATA Misuse fixed Co-Authored-By: brekotis <93345790+brekotis@users.noreply.github.com> --- src/config/types/web.rs | 4 +- src/web/session/backend.rs | 54 ++++-- src/web/session/backend_tests.rs | 324 +++++++++++++++++++++++++++++++ 3 files changed, 367 insertions(+), 15 deletions(-) create mode 100644 src/web/session/backend_tests.rs diff --git a/src/config/types/web.rs b/src/config/types/web.rs index 829fdfb..73c9ecb 100644 --- a/src/config/types/web.rs +++ b/src/config/types/web.rs @@ -130,7 +130,7 @@ pub struct WebLimitsConfig { /// Process-wide live logical-stream ceiling. #[serde(default = "default_web_max_streams_global")] pub max_streams_global: usize, - /// Process-wide concurrent inner MTProxy handshake ceiling. + /// Process-wide ceiling for inner MTProxy handshakes that received a first byte. #[serde(default = "default_web_max_stream_handshakes")] pub max_stream_handshakes: usize, /// Closed stream identifiers retained by one session. @@ -249,7 +249,7 @@ pub struct WebTimeoutsConfig { /// Deadline for collecting one authenticated carrier request body. #[serde(default = "default_web_body_timeout_secs")] pub body_secs: u64, - /// Deadline for the inner MTProxy handshake on one logical stream. + /// Deadline from the first inner byte through MTProxy authentication. #[serde(default = "default_web_stream_handshake_timeout_secs")] pub stream_handshake_secs: u64, /// Maximum wait for one empty downlink long poll. diff --git a/src/web/session/backend.rs b/src/web/session/backend.rs index 430298c..293ccdf 100644 --- a/src/web/session/backend.rs +++ b/src/web/session/backend.rs @@ -9,6 +9,10 @@ use crate::web::stream::WebLogicalStream; use super::{WebSession, inbound_queue_cost}; +#[cfg(test)] +#[path = "backend_tests.rs"] +mod tests; + impl WebSession { /// Starts one owned inner handshake and relay task for an admitted stream. pub(super) fn spawn_stream(self: &Arc, stream_id: u32, peer_port: u16) { @@ -22,10 +26,6 @@ impl WebSession { self.stream_finished(stream_id, peer_port); return; }; - let Some(handshake_permit) = manager.try_stream_handshake() else { - self.stream_finished(stream_id, peer_port); - return; - }; let deps = generation.client_runtime_deps(); let replay_checker = Arc::clone(&generation.replay_checker); let session = Arc::clone(self); @@ -46,7 +46,6 @@ impl WebSession { stream, deps, replay_checker, - handshake_permit, peer_port, ) => {} } @@ -109,7 +108,6 @@ async fn run_stream( stream: WebLogicalStream, deps: crate::proxy::authenticated::ClientRuntimeDeps, replay_checker: Arc, - handshake_permit: tokio::sync::OwnedSemaphorePermit, peer_port: u16, ) { use tokio::io::AsyncReadExt; @@ -122,10 +120,25 @@ async fn run_stream( let mut handshake = [0u8; HANDSHAKE_LEN]; let peer = std::net::SocketAddr::new(session.client_ip, peer_port); deps.stats.increment_connects_all(); + + // A carrier may publish OPEN before the local MTProto socket writes its + // first byte. Session and stream quotas bound this idle phase without + // consuming the process-wide active-handshake budget. + if reader.read_exact(&mut handshake[..1]).await.is_err() { + deps.stats + .increment_connects_bad_with_class("web_mtproto_handshake_io"); + return; + } + let Some(manager) = session.manager.upgrade() else { + return; + }; + let Some(handshake_permit) = manager.try_stream_handshake() else { + return; + }; let handshake_result = tokio::time::timeout( Duration::from_secs(session.timeouts.stream_handshake_secs), async { - reader.read_exact(&mut handshake).await?; + reader.read_exact(&mut handshake[1..]).await?; Ok::<_, io::Error>( handle_mtproto_handshake_for_web_user( &handshake, @@ -144,12 +157,27 @@ async fn run_stream( ) .await; drop(handshake_permit); - let Ok(Ok(crate::error::HandshakeResult::Success((reader, writer, success)))) = - handshake_result - else { - deps.stats - .increment_connects_bad_with_class("web_mtproto_bad_client"); - return; + let (reader, writer, success) = match handshake_result { + Err(_) => { + deps.stats + .increment_connects_bad_with_class("web_mtproto_handshake_timeout"); + deps.stats.increment_handshake_timeouts(); + deps.stats.increment_handshake_failure_class("timeout"); + return; + } + Ok(Err(_)) => { + deps.stats + .increment_connects_bad_with_class("web_mtproto_handshake_io"); + return; + } + Ok(Ok(crate::error::HandshakeResult::Success((reader, writer, success)))) => { + (reader, writer, success) + } + Ok(Ok(_)) => { + deps.stats + .increment_connects_bad_with_class("web_mtproto_bad_client"); + return; + } }; let _ = run_authenticated( reader, diff --git a/src/web/session/backend_tests.rs b/src/web/session/backend_tests.rs new file mode 100644 index 0000000..245444c --- /dev/null +++ b/src/web/session/backend_tests.rs @@ -0,0 +1,324 @@ +use std::collections::BTreeMap; +use std::net::SocketAddr; +use std::sync::Arc; +use std::time::Duration; + +use arc_swap::ArcSwap; +use tokio::net::TcpListener; + +use super::*; +use crate::config::{ + ProxyConfig, UpstreamConfig, UpstreamType, WebCarrier, WebRuntimeConfig, WebRuntimeProfile, + WebSecretMode, +}; +use crate::crypto::{AesCtr, sha256}; +use crate::maestro::generation::{RuntimeGeneration, test_runtime_generation}; +use crate::protocol::constants::{ + DC_IDX_POS, HANDSHAKE_LEN, IV_LEN, PREKEY_LEN, PROTO_TAG_POS, ProtoTag, SKIP_LEN, +}; +use crate::web::frame; +use crate::web::manager::{ManagerError, WebProcessRuntime}; + +struct TestRuntime { + session: Arc, + manager: Arc, + generation: Arc, +} + +impl TestRuntime { + fn process_frame( + &self, + stream_id: u32, + sequence: u64, + frame_type: FrameType, + payload: &[u8], + ) -> Result { + let encoded = frame::encode(frame_type, stream_id, payload); + match self.session.carrier() { + WebCarrier::Https => self.session.process_up(sequence, &encoded), + WebCarrier::HttpsLanes => { + self.session + .process_up_lane(stream_id, sequence, &encoded) + } + } + } + + async fn shutdown(self) { + self.session.close(); + self.session.wait().await; + self.manager.shutdown().await; + self.generation.stop_sessions().await; + self.generation.stop_background_tasks().await; + } +} + +fn test_runtime(carrier: WebCarrier, max_stream_handshakes: usize) -> TestRuntime { + test_runtime_with_dc(carrier, max_stream_handshakes, None) +} + +fn test_runtime_with_dc( + carrier: WebCarrier, + max_stream_handshakes: usize, + dc_addr: Option, +) -> TestRuntime { + let profile = Arc::new(WebRuntimeProfile { + host: "proxy.example.com".to_string(), + public_addr: "203.0.113.10:443".parse().unwrap(), + user: "default".to_string(), + secret_mode: WebSecretMode::Plain, + carrier, + capability: [7; 32], + max_sessions: 4, + max_streams: 16, + max_streams_per_session: 4, + }); + let mut config = ProxyConfig::default(); + config.web.enabled = true; + config.web.carrier = carrier; + config.web.limits.max_stream_handshakes = max_stream_handshakes; + config.web.timeouts.stream_handshake_secs = 1; + config.web.timeouts.shutdown_secs = 1; + config.censorship.server_hello_delay_min_ms = 0; + config.censorship.server_hello_delay_max_ms = 0; + if let Some(dc_addr) = dc_addr { + config + .dc_overrides + .insert("2".to_string(), vec![dc_addr.to_string()]); + config.upstreams.push(UpstreamConfig { + upstream_type: UpstreamType::Direct { + interface: None, + bind_addresses: None, + bindtodevice: None, + }, + weight: 1, + enabled: true, + scopes: String::new(), + selected_scope: String::new(), + ipv4: Some(true), + ipv6: Some(false), + prefer: Some(4), + }); + } + config.web.runtime = Some(Arc::new(WebRuntimeConfig { + vhosts: BTreeMap::new(), + profiles: vec![Arc::clone(&profile)], + })); + config.rebuild_runtime_user_auth().unwrap(); + let limits = config.web.limits.clone(); + let timeouts = config.web.timeouts.clone(); + let generation = test_runtime_generation(1, config); + let manager = WebProcessRuntime::start(Arc::new(ArcSwap::from(Arc::clone(&generation)))); + let session = WebSession::new( + Arc::downgrade(&manager), + [8; 32], + "192.0.2.10".parse().unwrap(), + profile, + [7; 32], + limits, + timeouts, + ); + TestRuntime { + session, + manager, + generation, + } +} + +fn valid_plain_handshake() -> [u8; HANDSHAKE_LEN] { + let secret = [0u8; 16]; + let mut handshake = [0x5a; HANDSHAKE_LEN]; + for (index, byte) in handshake[SKIP_LEN..SKIP_LEN + PREKEY_LEN + IV_LEN] + .iter_mut() + .enumerate() + { + *byte = (index as u8).wrapping_add(1); + } + let dec_prekey = &handshake[SKIP_LEN..SKIP_LEN + PREKEY_LEN]; + let dec_iv_bytes = &handshake[SKIP_LEN + PREKEY_LEN..SKIP_LEN + PREKEY_LEN + IV_LEN]; + let mut dec_key_input = Vec::with_capacity(PREKEY_LEN + secret.len()); + dec_key_input.extend_from_slice(dec_prekey); + dec_key_input.extend_from_slice(&secret); + let dec_key = sha256(&dec_key_input); + let mut dec_iv = [0u8; IV_LEN]; + dec_iv.copy_from_slice(dec_iv_bytes); + let mut cipher = AesCtr::new(&dec_key, u128::from_be_bytes(dec_iv)); + let keystream = cipher.encrypt(&[0u8; HANDSHAKE_LEN]); + let mut plaintext = [0u8; HANDSHAKE_LEN]; + plaintext[PROTO_TAG_POS..PROTO_TAG_POS + 4] + .copy_from_slice(&ProtoTag::Intermediate.to_bytes()); + plaintext[DC_IDX_POS..DC_IDX_POS + 2].copy_from_slice(&2i16.to_le_bytes()); + for index in PROTO_TAG_POS..HANDSHAKE_LEN { + handshake[index] = plaintext[index] ^ keystream[index]; + } + handshake +} + +async fn settle_tasks() { + for _ in 0..4 { + tokio::task::yield_now().await; + } +} + +fn bad_class(runtime: &TestRuntime, class: &str) -> u64 { + runtime + .generation + .stats + .get_connects_bad_class_counts() + .into_iter() + .find_map(|(name, total)| (name == class).then_some(total)) + .unwrap_or(0) +} + +#[tokio::test(start_paused = true)] +async fn open_without_data_does_not_start_the_inner_handshake_timeout() { + for carrier in [WebCarrier::Https, WebCarrier::HttpsLanes] { + let runtime = test_runtime(carrier, 1); + + assert_eq!(runtime.process_frame(1, 1, FrameType::Open, &[]), Ok(1)); + settle_tasks().await; + assert!(runtime.session.state.lock().streams.contains_key(&1)); + + tokio::time::advance(Duration::from_secs(2)).await; + settle_tasks().await; + + assert!( + runtime.session.state.lock().streams.contains_key(&1), + "OPEN without DATA must remain live until session cancellation or idle expiry" + ); + runtime.shutdown().await; + } +} + +#[tokio::test(start_paused = true)] +async fn the_first_inner_byte_starts_the_handshake_timeout() { + let runtime = test_runtime(WebCarrier::Https, 1); + + assert_eq!(runtime.process_frame(1, 1, FrameType::Open, &[]), Ok(1)); + settle_tasks().await; + tokio::time::advance(Duration::from_secs(2)).await; + assert_eq!(runtime.process_frame(1, 2, FrameType::Data, &[0x5a]), Ok(2)); + settle_tasks().await; + assert!(runtime.session.state.lock().streams.contains_key(&1)); + + tokio::time::advance(Duration::from_secs(2)).await; + settle_tasks().await; + + assert!(!runtime.session.state.lock().streams.contains_key(&1)); + assert_eq!(bad_class(&runtime, "web_mtproto_handshake_timeout"), 1); + assert_eq!(runtime.generation.stats.get_handshake_timeouts(), 1); + assert_eq!( + runtime + .generation + .stats + .get_handshake_failure_class_counts(), + vec![("timeout".to_string(), 1)] + ); + runtime.shutdown().await; +} + +#[tokio::test(start_paused = true)] +async fn delayed_complete_handshake_is_classified_after_data_arrives() { + let runtime = test_runtime(WebCarrier::Https, 1); + + assert_eq!(runtime.process_frame(1, 1, FrameType::Open, &[]), Ok(1)); + settle_tasks().await; + tokio::time::advance(Duration::from_secs(2)).await; + assert_eq!( + runtime.process_frame(1, 2, FrameType::Data, &[0; 64]), + Ok(2) + ); + settle_tasks().await; + + assert!(!runtime.session.state.lock().streams.contains_key(&1)); + assert_eq!(bad_class(&runtime, "web_mtproto_bad_client"), 1); + assert_eq!(bad_class(&runtime, "web_mtproto_handshake_timeout"), 0); + assert_eq!(runtime.generation.stats.get_handshake_timeouts(), 0); + runtime.shutdown().await; +} + +#[tokio::test(start_paused = true)] +async fn delayed_valid_handshake_reaches_the_authenticated_relay() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let runtime = test_runtime_with_dc(WebCarrier::Https, 1, Some(listener.local_addr().unwrap())); + + assert_eq!(runtime.process_frame(1, 1, FrameType::Open, &[]), Ok(1)); + settle_tasks().await; + tokio::time::advance(Duration::from_secs(2)).await; + assert_eq!( + runtime.process_frame(1, 2, FrameType::Data, &valid_plain_handshake()), + Ok(2) + ); + + let accepted = tokio::time::timeout(Duration::from_secs(1), listener.accept()).await; + assert!( + accepted.is_ok(), + "valid handshake did not reach the configured upstream; bad classes: {:?}", + runtime.generation.stats.get_connects_bad_class_counts() + ); + let (upstream, _) = accepted.unwrap().unwrap(); + assert!(runtime.session.state.lock().streams.contains_key(&1)); + assert_eq!(runtime.generation.stats.get_connects_bad(), 0); + assert_eq!(runtime.generation.stats.get_handshake_timeouts(), 0); + + runtime.shutdown().await; + drop(upstream); +} + +#[tokio::test(start_paused = true)] +async fn silent_streams_do_not_consume_active_handshake_capacity() { + let runtime = test_runtime(WebCarrier::Https, 1); + + assert_eq!(runtime.process_frame(1, 1, FrameType::Open, &[]), Ok(1)); + assert_eq!(runtime.process_frame(2, 2, FrameType::Open, &[]), Ok(2)); + settle_tasks().await; + tokio::time::advance(Duration::from_secs(2)).await; + settle_tasks().await; + { + let state = runtime.session.state.lock(); + assert!(state.streams.contains_key(&1)); + assert!(state.streams.contains_key(&2)); + } + + assert_eq!(runtime.process_frame(1, 3, FrameType::Data, &[1]), Ok(3)); + settle_tasks().await; + assert_eq!(runtime.process_frame(2, 4, FrameType::Data, &[2]), Ok(4)); + settle_tasks().await; + + let state = runtime.session.state.lock(); + assert!(state.streams.contains_key(&1)); + assert!(!state.streams.contains_key(&2)); + drop(state); + runtime.shutdown().await; +} + +#[tokio::test(start_paused = true)] +async fn cancellation_while_waiting_for_data_releases_stream_ownership() { + let runtime = test_runtime(WebCarrier::Https, 1); + + assert_eq!(runtime.process_frame(1, 1, FrameType::Open, &[]), Ok(1)); + settle_tasks().await; + assert_eq!(runtime.session.tasks_live.load(Ordering::Acquire), 1); + + runtime.session.close(); + runtime.session.wait().await; + + assert_eq!(runtime.session.tasks_live.load(Ordering::Acquire), 0); + assert!(runtime.session.state.lock().active_peer_ports.is_empty()); + let peer_port = runtime + .manager + .try_acquire_stream( + runtime.session.profile_key, + runtime.session.profile.max_streams, + runtime.session.client_ip, + runtime.session.profile.public_addr, + ) + .unwrap(); + assert_eq!(peer_port, 1); + runtime.manager.release_stream( + runtime.session.profile_key, + runtime.session.client_ip, + runtime.session.profile.public_addr, + peer_port, + ); + runtime.shutdown().await; +} From 3188c8e4175e9a52200e3e959c7fb744d92c4891 Mon Sep 17 00:00:00 2001 From: Alexey <247128645+axkurcom@users.noreply.github.com> Date: Sun, 23 Aug 2026 18:05:12 +0300 Subject: [PATCH 02/10] WEB + Apple iOS Fallback Cache fixed Co-Authored-By: brekotis <93345790+brekotis@users.noreply.github.com> --- src/web/http/decoy.rs | 11 ++++++++++- src/web/http/tests.rs | 30 ++++++++++++++++++++++++++++++ 2 files changed, 40 insertions(+), 1 deletion(-) diff --git a/src/web/http/decoy.rs b/src/web/http/decoy.rs index 3d194ea..4fbe5e2 100644 --- a/src/web/http/decoy.rs +++ b/src/web/http/decoy.rs @@ -17,6 +17,7 @@ use crate::config::{WebRuntimeDecoy, WebRuntimeVhost}; use crate::web::manager::WebProcessRuntime; /// Serves the configured ordinary site after optionally removing carrier material. +/// Transport-sanitized static fallbacks remain uncacheable after query removal. pub(super) async fn serve_decoy( mut request: Request, vhost: Arc, @@ -41,7 +42,15 @@ where }; let request = Request::from_parts(parts, body); match &vhost.decoy { - WebRuntimeDecoy::StaticDirectory(site) => serve_static(request, site), + WebRuntimeDecoy::StaticDirectory(site) => { + let mut response = serve_static(request, site); + if sanitize_transport { + response + .headers_mut() + .insert(header::CACHE_CONTROL, HeaderValue::from_static("no-store")); + } + response + } WebRuntimeDecoy::HttpUpstream { addr, authority } => { proxy_to_upstream( request, diff --git a/src/web/http/tests.rs b/src/web/http/tests.rs index 9b19fd8..7ed5f02 100644 --- a/src/web/http/tests.rs +++ b/src/web/http/tests.rs @@ -209,6 +209,36 @@ async fn https_carrier_bootstraps_and_closes_one_session() { replacement.stop_background_tasks().await; } +#[tokio::test] +async fn rejected_bridge_bootstrap_falls_back_to_uncacheable_static_index() { + let capability = [15u8; 32]; + let generation = test_runtime_generation(1, runtime_config(capability, WebCarrier::Https)); + let active_runtime = Arc::new(ArcSwap::from(Arc::clone(&generation))); + let runtime = WebProcessRuntime::start(active_runtime); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(capability); + let bridge_request = || { + format!( + "GET /?bridge={encoded} HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nConnection: close\r\n\r\n" + ) + .into_bytes() + }; + + let first_response = request(&listener, &runtime, bridge_request()).await; + let (_, first_body) = split_response(&first_response); + assert!(first_body.windows(11).any(|value| value == b"bootstrap='")); + + let fallback_response = request(&listener, &runtime, bridge_request()).await; + let (fallback_headers, fallback_body) = split_response(&fallback_response); + assert!(fallback_headers.starts_with(b"HTTP/1.1 200")); + assert_eq!(response_header(fallback_headers, "cache-control"), "no-store"); + assert_eq!(fallback_body, b"decoy"); + + runtime.shutdown().await; + generation.stop_sessions().await; + generation.stop_background_tasks().await; +} + #[tokio::test] async fn bootstrap_survives_client_address_family_change() { let capability = [8u8; 32]; From 1a18f2452fc1863ce52b9d013c1a0fc1e33720cb Mon Sep 17 00:00:00 2001 From: Alexey <247128645+axkurcom@users.noreply.github.com> Date: Sun, 23 Aug 2026 20:43:22 +0300 Subject: [PATCH 03/10] Fixed routing across endpoint refresh Co-Authored-By: brekotis <93345790+brekotis@users.noreply.github.com> --- src/maestro/generation.rs | 12 +++- src/transport/middle_proxy/send/selection.rs | 38 ++++++++--- .../tests/send_adversarial_tests.rs | 63 ++++++++++++++++++- src/web/session/backend.rs | 4 ++ src/web/session/backend_tests.rs | 43 ++++++++++++- 5 files changed, 146 insertions(+), 14 deletions(-) diff --git a/src/maestro/generation.rs b/src/maestro/generation.rs index de29b7e..8a8c5de 100644 --- a/src/maestro/generation.rs +++ b/src/maestro/generation.rs @@ -305,8 +305,18 @@ impl RuntimeGeneration { #[cfg(test)] /// Builds a lightweight runtime generation without network startup tasks. pub(crate) fn test_runtime_generation(id: u64, config: ProxyConfig) -> Arc { - let (config_tx, config_rx) = watch::channel(Arc::new(config.clone())); let (_admission_tx, admission_rx) = watch::channel(true); + test_runtime_generation_with_admission(id, config, admission_rx) +} + +#[cfg(test)] +/// Builds a lightweight runtime generation with a controllable admission gate. +pub(crate) fn test_runtime_generation_with_admission( + id: u64, + config: ProxyConfig, + admission_rx: watch::Receiver, +) -> Arc { + let (config_tx, config_rx) = watch::channel(Arc::new(config.clone())); let stats = Arc::new(Stats::new()); let upstream_manager = Arc::new(UpstreamManager::new( config.upstreams, diff --git a/src/transport/middle_proxy/send/selection.rs b/src/transport/middle_proxy/send/selection.rs index 834e0c0..837062c 100644 --- a/src/transport/middle_proxy/send/selection.rs +++ b/src/transport/middle_proxy/send/selection.rs @@ -16,19 +16,37 @@ impl MePool { include_warm: bool, ) -> Vec { let preferred_snapshot = self.preferred_endpoints_by_dc.load(); - let Some(preferred) = preferred_snapshot.get(&routed_dc) else { - return Vec::new(); - }; - if preferred.is_empty() { - return Vec::new(); + let mut out = Vec::new(); + if let Some(preferred) = preferred_snapshot + .get(&routed_dc) + .filter(|preferred| !preferred.is_empty()) + { + for (idx, w) in writers.iter().enumerate() { + if !self.writer_eligible_for_selection(w, include_warm) { + continue; + } + if w.writer_dc == routed_dc && preferred.binary_search(&w.addr).is_ok() { + out.push(idx); + } + } + } + if !out.is_empty() || !include_warm { + return out; } - let mut out = Vec::new(); + // A map update publishes desired endpoints before replacement coverage is + // guaranteed. Preserve the existing same-DC writer as the final tier so + // the data plane remains available while the pool-owned reinit converges. for (idx, w) in writers.iter().enumerate() { - if !self.writer_eligible_for_selection(w, include_warm) { - continue; - } - if w.writer_dc == routed_dc && preferred.binary_search(&w.addr).is_ok() { + let family_enabled = if w.addr.is_ipv4() { + self.decision.ipv4_me + } else { + self.decision.ipv6_me + }; + if family_enabled + && w.writer_dc == routed_dc + && self.writer_eligible_for_selection(w, true) + { out.push(idx); } } diff --git a/src/transport/middle_proxy/tests/send_adversarial_tests.rs b/src/transport/middle_proxy/tests/send_adversarial_tests.rs index 3b6e718..eac095f 100644 --- a/src/transport/middle_proxy/tests/send_adversarial_tests.rs +++ b/src/transport/middle_proxy/tests/send_adversarial_tests.rs @@ -15,6 +15,12 @@ use crate::network::probe::NetworkDecision; use crate::stats::Stats; async fn make_pool() -> (Arc, Arc) { + make_pool_with_decision(NetworkDecision::default()).await +} + +async fn make_pool_with_decision( + decision: NetworkDecision, +) -> (Arc, Arc) { let general = GeneralConfig { me_route_no_writer_mode: MeRouteNoWriterMode::AsyncRecoveryFailfast, me_route_no_writer_wait_ms: 50, @@ -40,7 +46,7 @@ async fn make_pool() -> (Arc, Arc) { HashMap::new(), HashMap::new(), None, - NetworkDecision::default(), + decision, None, rng.clone(), Arc::new(Stats::default()), @@ -213,6 +219,61 @@ fn proxy_req_our_addr_from_payload(payload: &[u8]) -> SocketAddr { ) } +#[tokio::test] +async fn send_proxy_req_uses_live_same_dc_writer_while_preferred_endpoint_refills() { + let decision = NetworkDecision { + ipv4_dc: true, + ipv4_me: true, + effective_prefer: 4, + ..NetworkDecision::default() + }; + let (pool, _rng) = make_pool_with_decision(decision).await; + let old_positive_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 3, 2)), 443); + let old_negative_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 3, 3)), 443); + let mut old_positive_rx = insert_writer(&pool, 41, 2, old_positive_addr, true).await; + let _old_negative_rx = insert_writer(&pool, 42, -2, old_negative_addr, true).await; + + let new_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(192, 0, 2, 42)), 443); + pool.update_proxy_maps( + HashMap::from([(2, vec![(new_addr.ip(), new_addr.port())])]), + None, + ) + .await; + + assert!(pool.admission_ready_conditional_cast().await); + assert_eq!( + pool.preferred_endpoints_by_dc + .load() + .get(&2) + .cloned() + .unwrap_or_default(), + vec![new_addr] + ); + + let (conn_id, _rx) = pool.registry.register().await; + let result = pool + .send_proxy_req( + conn_id, + 2, + SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 30004), + SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 443), + b"cold-route", + 0, + None, + None, + ) + .await; + + assert!( + result.is_ok(), + "a live same-DC writer must bridge preferred-endpoint refill: {result:?}" + ); + assert_eq!( + recv_data_count(&mut old_positive_rx, Duration::from_millis(50)).await, + 1 + ); +} + #[tokio::test] async fn send_proxy_req_does_not_replay_when_first_bind_commit_fails() { let (pool, _rng) = make_pool().await; diff --git a/src/web/session/backend.rs b/src/web/session/backend.rs index 293ccdf..69acef6 100644 --- a/src/web/session/backend.rs +++ b/src/web/session/backend.rs @@ -21,6 +21,10 @@ impl WebSession { return; }; let generation = manager.active_generation(); + if !*generation.admission_rx.borrow() { + self.stream_finished(stream_id, peer_port); + return; + } let Ok(connection_permit) = generation.max_connections.clone().try_acquire_owned() else { manager.record_stream_rejected(); self.stream_finished(stream_id, peer_port); diff --git a/src/web/session/backend_tests.rs b/src/web/session/backend_tests.rs index 245444c..b09989d 100644 --- a/src/web/session/backend_tests.rs +++ b/src/web/session/backend_tests.rs @@ -5,6 +5,7 @@ use std::time::Duration; use arc_swap::ArcSwap; use tokio::net::TcpListener; +use tokio::sync::watch; use super::*; use crate::config::{ @@ -12,7 +13,9 @@ use crate::config::{ WebSecretMode, }; use crate::crypto::{AesCtr, sha256}; -use crate::maestro::generation::{RuntimeGeneration, test_runtime_generation}; +use crate::maestro::generation::{ + RuntimeGeneration, test_runtime_generation_with_admission, +}; use crate::protocol::constants::{ DC_IDX_POS, HANDSHAKE_LEN, IV_LEN, PREKEY_LEN, PROTO_TAG_POS, ProtoTag, SKIP_LEN, }; @@ -23,6 +26,7 @@ struct TestRuntime { session: Arc, manager: Arc, generation: Arc, + admission_tx: watch::Sender, } impl TestRuntime { @@ -106,7 +110,8 @@ fn test_runtime_with_dc( config.rebuild_runtime_user_auth().unwrap(); let limits = config.web.limits.clone(); let timeouts = config.web.timeouts.clone(); - let generation = test_runtime_generation(1, config); + let (admission_tx, admission_rx) = watch::channel(true); + let generation = test_runtime_generation_with_admission(1, config, admission_rx); let manager = WebProcessRuntime::start(Arc::new(ArcSwap::from(Arc::clone(&generation)))); let session = WebSession::new( Arc::downgrade(&manager), @@ -121,6 +126,7 @@ fn test_runtime_with_dc( session, manager, generation, + admission_tx, } } @@ -264,6 +270,39 @@ async fn delayed_valid_handshake_reaches_the_authenticated_relay() { drop(upstream); } +#[tokio::test] +async fn closed_generation_admission_rejects_web_stream_before_backend() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let runtime = test_runtime_with_dc(WebCarrier::Https, 1, Some(listener.local_addr().unwrap())); + runtime.admission_tx.send_replace(false); + + assert_eq!(runtime.process_frame(1, 1, FrameType::Open, &[]), Ok(1)); + settle_tasks().await; + assert!( + tokio::time::timeout(Duration::from_millis(50), listener.accept()) + .await + .is_err(), + "WEB stream bypassed the closed generation admission gate" + ); + assert!(!runtime.session.state.lock().streams.contains_key(&1)); + assert_eq!(runtime.generation.max_connections.available_permits(), 64); + + runtime.admission_tx.send_replace(true); + assert_eq!(runtime.process_frame(2, 2, FrameType::Open, &[]), Ok(2)); + assert_eq!( + runtime.process_frame(2, 3, FrameType::Data, &valid_plain_handshake()), + Ok(3) + ); + let (upstream, _) = tokio::time::timeout(Duration::from_secs(1), listener.accept()) + .await + .expect("a new WEB stream did not start after admission reopened") + .unwrap(); + assert_eq!(runtime.generation.max_connections.available_permits(), 63); + + runtime.shutdown().await; + drop(upstream); +} + #[tokio::test(start_paused = true)] async fn silent_streams_do_not_consume_active_handshake_capacity() { let runtime = test_runtime(WebCarrier::Https, 1); From 0779cd090136228eb3b3055231dd4e43c43e33d7 Mon Sep 17 00:00:00 2001 From: Alexey <247128645+axkurcom@users.noreply.github.com> Date: Sun, 23 Aug 2026 22:07:20 +0300 Subject: [PATCH 04/10] Compat Bootstrap Literal for iOS Co-Authored-By: brekotis <93345790+brekotis@users.noreply.github.com> --- src/web/bridge.rs | 28 +++++++++++++++++++++++++++- src/web/http/tests.rs | 28 ++++++++++++++-------------- 2 files changed, 41 insertions(+), 15 deletions(-) diff --git a/src/web/bridge.rs b/src/web/bridge.rs index 4418325..fce1113 100644 --- a/src/web/bridge.rs +++ b/src/web/bridge.rs @@ -54,7 +54,7 @@ const DOCUMENT: &str = r##"