diff --git a/Cargo.lock b/Cargo.lock
index 087c840..6be0b5f 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -2900,7 +2900,7 @@ checksum = "7b2093cf4c8eb1e67749a6762251bc9cd836b6fc171623bd0a9d324d37af2417"
[[package]]
name = "telemt"
-version = "3.5.3"
+version = "3.5.4"
dependencies = [
"aes",
"anyhow",
diff --git a/Cargo.toml b/Cargo.toml
index fd15199..c6e12e0 100644
--- a/Cargo.toml
+++ b/Cargo.toml
@@ -1,6 +1,6 @@
[package]
name = "telemt"
-version = "3.5.3"
+version = "3.5.4"
edition = "2024"
[features]
diff --git a/src/api/web_status.rs b/src/api/web_status.rs
index 41ef777..0914a97 100644
--- a/src/api/web_status.rs
+++ b/src/api/web_status.rs
@@ -18,7 +18,7 @@ mod details;
// Query parsing and matching remain independent from bounded HTML rendering.
mod query;
-use details::{push_body, push_frames, push_headers};
+use details::{push_body, push_frames, push_headers, push_lifecycle};
use query::{GroupBy, StatusQuery, client_ip, parse_query, record_matches};
struct GroupSummary {
@@ -56,11 +56,11 @@ pub(super) async fn render(
);
};
let now_millis = crate::web::trace::store_epoch_millis();
- let since_millis = query
- .record
- .is_none()
- .then(|| now_millis.saturating_sub(query.window_secs.saturating_mul(1000)))
- .unwrap_or(0);
+ let since_millis = if query.record.is_none() {
+ now_millis.saturating_sub(query.window_secs.saturating_mul(1000))
+ } else {
+ 0
+ };
let records = store.snapshot_matching(|record| record_matches(record, &query, since_millis));
let status = store.status();
let mut html = String::with_capacity(MAX_PAGE_BYTES);
@@ -367,20 +367,7 @@ fn push_record(html: &mut String, record: &TraceRecord) {
push_body(html, "message body", message.body.as_ref());
push_frames(html, &message.frames);
}
- TraceRecordKind::Lifecycle(event) => {
- html.push_str("
event: ");
- html.push_str(event.event.as_str());
- html.push_str("\nstream: ");
- html.push_str(
- &event
- .stream_id
- .map(|v| v.to_string())
- .unwrap_or_else(|| "-".to_string()),
- );
- html.push_str("\nreason: ");
- html.push_str(event.reason.unwrap_or("-"));
- html.push_str("");
- }
+ TraceRecordKind::Lifecycle(event) => push_lifecycle(html, event),
}
html.push_str("");
}
diff --git a/src/api/web_status/details.rs b/src/api/web_status/details.rs
index fd21217..90ce3a2 100644
--- a/src/api/web_status/details.rs
+++ b/src/api/web_status/details.rs
@@ -84,3 +84,34 @@ pub(super) fn push_body(
}
html.push_str("");
}
+
+pub(super) fn push_lifecycle(html: &mut String, event: &crate::web::trace::TraceLifecycleRecord) {
+ html.push_str("event: ");
+ html.push_str(event.event.as_str());
+ html.push_str("\nstream: ");
+ html.push_str(
+ &event
+ .stream_id
+ .map(|value| value.to_string())
+ .unwrap_or_else(|| "-".to_string()),
+ );
+ html.push_str("\nreason: ");
+ html.push_str(event.reason.unwrap_or("-"));
+ if let Some(carrier) = &event.carrier {
+ html.push_str("\nclient class: ");
+ html.push_str(carrier.client_class);
+ html.push_str("\ncarrier: ");
+ html.push_str(carrier.carrier.as_str());
+ html.push_str("\nattempt: ");
+ html.push_str(&carrier.attempt.to_string());
+ html.push_str("\nscores: https=");
+ html.push_str(&carrier.scores[0].to_string());
+ html.push_str(" https-lanes=");
+ html.push_str(&carrier.scores[1].to_string());
+ html.push_str(" websocket=");
+ html.push_str(&carrier.scores[2].to_string());
+ html.push_str(" websocket-lanes=");
+ html.push_str(&carrier.scores[3].to_string());
+ }
+ html.push_str("");
+}
diff --git a/src/api/web_status/tests.rs b/src/api/web_status/tests.rs
index f9dc268..b2480a3 100644
--- a/src/api/web_status/tests.rs
+++ b/src/api/web_status/tests.rs
@@ -20,11 +20,15 @@ fn html_escaping_covers_active_markup_characters() {
#[tokio::test]
async fn renderer_filters_groups_and_sets_control_plane_security_headers() {
- let mut policy = WebDebugConfig::default();
- policy.enabled = true;
- let mut limits = crate::config::WebLimitsConfig::default();
- limits.debug_records_capacity = 8;
- limits.debug_bytes_global = 16 * 1024;
+ let policy = WebDebugConfig {
+ enabled: true,
+ ..Default::default()
+ };
+ let limits = crate::config::WebLimitsConfig {
+ debug_records_capacity: 8,
+ debug_bytes_global: 16 * 1024,
+ ..Default::default()
+ };
let store = WebTraceStore::new(policy.clone(), &limits);
store.record_lifecycle(
None,
@@ -79,7 +83,7 @@ async fn render_permits_remain_owned_by_inflight_response_bodies() {
#[test]
fn page_truncation_preserves_utf8_boundary_and_cap() {
- let mut html = "я".repeat(MAX_PAGE_BYTES);
+ let mut html = "\u{044f}".repeat(MAX_PAGE_BYTES);
truncate_page(&mut html);
assert!(html.len() <= MAX_PAGE_BYTES);
assert!(html.ends_with("[page output truncated]"));
diff --git a/src/config/hot_reload.rs b/src/config/hot_reload.rs
index a8f1401..10e3d9d 100644
--- a/src/config/hot_reload.rs
+++ b/src/config/hot_reload.rs
@@ -16,6 +16,7 @@
//! | `general` | `telemetry` / `me_*_policy` | Applied immediately |
//! | `network` | `dns_overrides` | Applied immediately |
//! | `access` | All user/quota fields | Effective immediately |
+//! | `web` | Carrier, timing, and debug policy | Applied to newly issued sessions |
//! Fields that require re-binding sockets (`server.listeners`, legacy
//! `server.port`, `censorship.*`, `network.*`, `use_middle_proxy`) are **not**
//! applied; a warning is emitted. SYN limiter rules are process-owned and are
diff --git a/src/config/hot_reload/tests.rs b/src/config/hot_reload/tests.rs
index 5cc8f33..8a5a99b 100644
--- a/src/config/hot_reload/tests.rs
+++ b/src/config/hot_reload/tests.rs
@@ -30,6 +30,23 @@ fn write_reload_config(path: &Path, ad_tag: Option<&str>, server_port: Option PathBuf {
let nonce = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
@@ -264,6 +281,29 @@ fn reload_keeps_hot_apply_when_non_hot_fields_change() {
let _ = std::fs::remove_file(path);
}
+#[test]
+fn reload_publishes_web_negotiation_policy_outside_hot_field_reporting() {
+ let path = temp_config_path("telemt_web_negotiation_reload");
+
+ write_web_reload_config(&path, "false", true);
+ let initial_cfg = Arc::new(ProxyConfig::load(&path).unwrap());
+ let initial_hash = ProxyConfig::load_with_metadata(&path)
+ .unwrap()
+ .rendered_hash;
+ let (config_tx, _config_rx) = watch::channel(initial_cfg.clone());
+ let (log_tx, _log_rx) = watch::channel(initial_cfg.general.log_level.clone());
+ let mut reload_state = ReloadState::new(Some(initial_hash));
+
+ write_web_reload_config(&path, "[\"websocket\", \"https\"]", false);
+ reload_config(&path, &config_tx, &log_tx, None, None, &mut reload_state).unwrap();
+
+ let applied = config_tx.borrow().clone();
+ assert!(applied.web.carrier_negotiation_enabled());
+ assert!(!applied.web.carrier_learning);
+
+ let _ = std::fs::remove_file(path);
+}
+
#[test]
fn classify_sni_change_requires_restart() {
// censorship.* is not in overlay_hot_fields -> restart.
diff --git a/src/config/hot_reload/watcher.rs b/src/config/hot_reload/watcher.rs
index e2b3e60..838988d 100644
--- a/src/config/hot_reload/watcher.rs
+++ b/src/config/hot_reload/watcher.rs
@@ -164,7 +164,7 @@ pub(super) fn reload_config(
let old_hot = HotFields::from_config(&old_cfg);
let applied_hot = HotFields::from_config(&applied_cfg);
let non_hot_changed = !config_equal(&applied_cfg, &new_cfg);
- let hot_changed = old_hot != applied_hot;
+ let hot_changed = !config_equal(&old_cfg, &applied_cfg);
if non_hot_changed {
warn_non_hot_changes(&old_cfg, &new_cfg, non_hot_changed);
diff --git a/src/config/load/runtime_web.rs b/src/config/load/runtime_web.rs
index 5d56073..ab31b6a 100644
--- a/src/config/load/runtime_web.rs
+++ b/src/config/load/runtime_web.rs
@@ -27,6 +27,7 @@ pub(super) fn rebuild(config: &mut ProxyConfig) -> Result<()> {
let mut static_files = 0usize;
let mut static_bytes = 0usize;
+ let carrier_candidates: Arc<[WebCarrier]> = config.web.carrier_candidates().into();
for vhost in &config.web.vhosts {
let decoy = build_decoy(
vhost,
@@ -63,6 +64,14 @@ pub(super) fn rebuild(config: &mut ProxyConfig) -> Result<()> {
user: profile.user.clone(),
secret_mode: profile.secret_mode,
carrier: config.web.carrier,
+ carrier_negotiation_enabled: config.web.carrier_negotiation_enabled(),
+ carrier_learning: config.web.carrier_negotiation_enabled()
+ && config.web.carrier_learning,
+ carriers: Arc::clone(&carrier_candidates),
+ carrier_negotiation_deadlines_secs: config
+ .web
+ .timeouts
+ .carrier_negotiation_deadlines_secs,
capability,
key_fingerprint,
max_sessions: profile
diff --git a/src/config/load/strict_keys.rs b/src/config/load/strict_keys.rs
index 5c153a2..6229bb9 100644
--- a/src/config/load/strict_keys.rs
+++ b/src/config/load/strict_keys.rs
@@ -260,7 +260,15 @@ const LISTENER_CONFIG_KEYS: &[&str] = &[
];
const WEB_CONFIG_KEYS: &[&str] = &[
- "enabled", "carrier", "debug", "limits", "timeouts", "vhosts",
+ "enabled",
+ "carrier",
+ "carriers",
+ "carrier_learning",
+ "carrier_negotiation_aggressiveness",
+ "debug",
+ "limits",
+ "timeouts",
+ "vhosts",
];
const WEB_LIMITS_CONFIG_KEYS: &[&str] = &[
@@ -271,10 +279,15 @@ const WEB_LIMITS_CONFIG_KEYS: &[&str] = &[
"max_frames_per_body",
"max_http_connections",
"max_http_handlers",
+ "max_lane_open_waits_per_session",
+ "pending_bytes_per_lane",
+ "pending_items_per_lane",
"websocket_bytes_global",
"websocket_admission_watermark_pct",
"websocket_eviction_watermark_pct",
"websocket_http_connection_reserve",
+ "max_websocket_evictions_in_flight",
+ "max_carrier_learning_entries",
"max_body_readers",
"max_body_bytes_global",
"max_sessions_global",
@@ -324,10 +337,17 @@ const WEB_TIMEOUTS_CONFIG_KEYS: &[&str] = &[
"header_secs",
"body_secs",
"stream_handshake_secs",
+ "stream_first_byte_secs",
"long_poll_secs",
+ "lane_open_wait_secs",
+ "carrier_health_secs",
+ "websocket_upgrade_secs",
+ "websocket_open_secs",
"websocket_write_secs",
"websocket_backpressure_secs",
"websocket_eviction_secs",
+ "carrier_negotiation_deadlines_secs",
+ "carrier_learning_secs",
"bootstrap_lifetime_secs",
"reconnect_grace_secs",
"http_idle_secs",
diff --git a/src/config/load/validate_web.rs b/src/config/load/validate_web.rs
index 20427aa..3ea8824 100644
--- a/src/config/load/validate_web.rs
+++ b/src/config/load/validate_web.rs
@@ -6,6 +6,10 @@ use super::*;
mod debug;
// Memory-envelope arithmetic remains isolated from protocol validation.
mod memory;
+// Carrier ordering, cumulative deadlines, and fallback identity are validated together.
+mod negotiation;
+// Request and lifecycle timeout relationships are validated together.
+mod timeouts;
// WebSocket transport policy is validated independently from HTTP body policy.
mod websocket;
@@ -67,11 +71,17 @@ pub(super) fn validate(config: &mut ProxyConfig) -> Result<()> {
validate_limits(&config.web.limits)?;
debug::validate(&config.web.debug, &config.web.limits)?;
- if config.web.carrier == WebCarrier::HttpsLanes && config.web.limits.max_http_handlers < 2 {
- return config_error("web.carrier=https-lanes requires web.limits.max_http_handlers >= 2");
+ let carriers = negotiation::validate(&config.web)?;
+ if carriers.contains(&WebCarrier::Https) && config.web.limits.max_http_handlers < 2 {
+ return config_error("WEB https candidates require web.limits.max_http_handlers >= 2");
}
- validate_timeouts(&config.web.timeouts)?;
- websocket::validate(config.web.carrier, &config.web.limits, &config.web.timeouts)?;
+ if carriers.contains(&WebCarrier::HttpsLanes) && config.web.limits.max_http_handlers < 4 {
+ return config_error(
+ "WEB https-lanes candidates require web.limits.max_http_handlers >= 4",
+ );
+ }
+ timeouts::validate(&config.web.timeouts)?;
+ websocket::validate(&carriers, &config.web.limits, &config.web.timeouts)?;
validate_vhosts(config)?;
Ok(())
}
@@ -136,6 +146,11 @@ fn validate_limits(limits: &WebLimitsConfig) -> Result<()> {
if !(1..=MAX_WEB_TOMBSTONES_PER_SESSION).contains(&limits.max_tombstones_per_session) {
return config_error("web.limits.max_tombstones_per_session must be within [1, 4096]");
}
+ if limits.pending_bytes_per_lane <= WEB_FRAME_HEADER_BYTES + WEB_QUEUE_ITEM_COST {
+ return config_error(
+ "web.limits.pending_bytes_per_lane must preserve one non-empty DATA frame",
+ );
+ }
if limits.carrier_batch_bytes > limits.max_body_bytes
|| limits.carrier_batch_bytes
< limits
@@ -155,6 +170,20 @@ fn validate_limits(limits: &WebLimitsConfig) -> Result<()> {
let positive = [
("max_http_connections", limits.max_http_connections),
("max_http_handlers", limits.max_http_handlers),
+ (
+ "max_lane_open_waits_per_session",
+ limits.max_lane_open_waits_per_session,
+ ),
+ ("pending_bytes_per_lane", limits.pending_bytes_per_lane),
+ ("pending_items_per_lane", limits.pending_items_per_lane),
+ (
+ "max_websocket_evictions_in_flight",
+ limits.max_websocket_evictions_in_flight,
+ ),
+ (
+ "max_carrier_learning_entries",
+ limits.max_carrier_learning_entries,
+ ),
("max_body_readers", limits.max_body_readers),
("max_body_bytes_global", limits.max_body_bytes_global),
("max_sessions_global", limits.max_sessions_global),
@@ -227,8 +256,11 @@ fn validate_limits(limits: &WebLimitsConfig) -> Result<()> {
|| limits.max_bootstraps_per_ip > limits.max_bootstraps_global
|| limits.max_http_handlers > limits.max_http_connections
|| limits.max_body_readers > limits.max_http_handlers
+ || limits.max_lane_open_waits_per_session > limits.max_streams_per_session
|| limits.pending_bytes_per_session > limits.pending_bytes_global
|| limits.pending_items_per_session > limits.pending_items_global
+ || limits.pending_bytes_per_lane > limits.pending_bytes_per_session
+ || limits.pending_items_per_lane > limits.pending_items_per_session
|| limits.control_bytes_per_session > limits.control_bytes_global
|| limits.control_bytes_per_session > limits.pending_bytes_per_session
|| limits.control_bytes_global > limits.pending_bytes_global
@@ -324,41 +356,6 @@ fn validate_limits(limits: &WebLimitsConfig) -> Result<()> {
Ok(())
}
-fn validate_timeouts(timeouts: &WebTimeoutsConfig) -> Result<()> {
- let values = [
- ("header_secs", timeouts.header_secs),
- ("body_secs", timeouts.body_secs),
- ("stream_handshake_secs", timeouts.stream_handshake_secs),
- ("long_poll_secs", timeouts.long_poll_secs),
- ("websocket_write_secs", timeouts.websocket_write_secs),
- (
- "websocket_backpressure_secs",
- timeouts.websocket_backpressure_secs,
- ),
- ("websocket_eviction_secs", timeouts.websocket_eviction_secs),
- ("bootstrap_lifetime_secs", timeouts.bootstrap_lifetime_secs),
- ("reconnect_grace_secs", timeouts.reconnect_grace_secs),
- ("http_idle_secs", timeouts.http_idle_secs),
- ("shutdown_secs", timeouts.shutdown_secs),
- ("decoy_header_secs", timeouts.decoy_header_secs),
- ];
- if let Some((field, _)) = values
- .into_iter()
- .find(|(_, value)| !(1..=3600).contains(value))
- {
- return config_error(&format!("web.timeouts.{field} must be within [1, 3600]"));
- }
- let request_deadline = timeouts
- .header_secs
- .max(timeouts.body_secs)
- .max(timeouts.long_poll_secs)
- .max(timeouts.decoy_header_secs);
- if request_deadline >= timeouts.http_idle_secs {
- return config_error("web.timeouts request deadlines must be lower than http_idle_secs");
- }
- Ok(())
-}
-
fn validate_vhosts(config: &mut ProxyConfig) -> Result<()> {
let limits = &config.web.limits;
if config.web.vhosts.len() > limits.max_vhosts {
diff --git a/src/config/load/validate_web/memory.rs b/src/config/load/validate_web/memory.rs
index 3fb9666..81262ac 100644
--- a/src/config/load/validate_web/memory.rs
+++ b/src/config/load/validate_web/memory.rs
@@ -3,9 +3,14 @@ use super::*;
const WEB_DEBUG_RENDERERS: usize = 2;
const WEB_DEBUG_STATUS_PAGE_BYTES: usize = 8 * 1024 * 1024;
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;
/// Validates process-wide body, header, queue, static, and debug reservations.
pub(super) fn validate(limits: &WebLimitsConfig) -> Result<()> {
+ if limits.max_carrier_learning_entries == 0 {
+ return config_error("web.limits.max_carrier_learning_entries must be > 0");
+ }
let body_reservation = limits
.max_body_readers
.checked_mul(limits.max_body_bytes)
@@ -47,6 +52,21 @@ pub(super) fn validate(limits: &WebLimitsConfig) -> Result<()> {
.and_then(|scratch| value.checked_add(scratch))
})
.ok_or_else(|| ProxyError::Config("web.debug reservations overflowed usize".to_string()))?;
+ let carrier_learning_reservation = limits
+ .max_carrier_learning_entries
+ .checked_mul(WEB_CARRIER_LEARNING_ENTRY_BYTES)
+ .ok_or_else(|| {
+ ProxyError::Config("web.carrier learning reservation overflowed usize".to_string())
+ })?;
+ let lane_state_reservation = limits
+ .max_streams_per_session
+ .checked_add(limits.max_tombstones_per_session)
+ .and_then(|value| value.checked_add(1))
+ .and_then(|value| value.checked_mul(limits.max_sessions_global))
+ .and_then(|value| value.checked_mul(WEB_LANE_STATE_BYTES))
+ .ok_or_else(|| {
+ ProxyError::Config("web.limits lane state reservation overflowed usize".to_string())
+ })?;
let reserved = limits
.pending_bytes_global
.checked_add(limits.max_body_bytes_global)
@@ -54,6 +74,8 @@ pub(super) fn validate(limits: &WebLimitsConfig) -> Result<()> {
.and_then(|value| value.checked_add(debug_ring_index))
.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(lane_state_reservation))
.and_then(|value| value.checked_add(http_header_reservation))
.ok_or_else(|| ProxyError::Config("web.limits byte ceilings overflow usize".to_string()))?;
if reserved > limits.memory_envelope_bytes
@@ -65,3 +87,20 @@ pub(super) fn validate(limits: &WebLimitsConfig) -> Result<()> {
}
Ok(())
}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ #[test]
+ fn default_envelope_includes_bounded_lane_and_learning_metadata() {
+ let limits = WebLimitsConfig::default();
+ assert!(validate(&limits).is_ok());
+
+ let previous_envelope = WebLimitsConfig {
+ memory_envelope_bytes: 768 * 1024 * 1024,
+ ..limits
+ };
+ assert!(validate(&previous_envelope).is_err());
+ }
+}
diff --git a/src/config/load/validate_web/negotiation.rs b/src/config/load/validate_web/negotiation.rs
new file mode 100644
index 0000000..3544ac4
--- /dev/null
+++ b/src/config/load/validate_web/negotiation.rs
@@ -0,0 +1,123 @@
+use std::collections::HashSet;
+
+use super::*;
+
+/// Validates bounded carrier selection and learning policy.
+pub(super) fn validate(config: &WebConfig) -> Result> {
+ if let Some(carriers) = config.carriers.enabled() {
+ if carriers.is_empty() {
+ return config_error("web.carriers must contain at least one carrier");
+ }
+ let mut unique = HashSet::with_capacity(carriers.len());
+ if carriers.iter().any(|carrier| !unique.insert(*carrier)) {
+ return config_error("web.carriers must not contain duplicate carriers");
+ }
+ }
+ let candidates = config.carrier_candidates();
+ if config.carrier_negotiation_enabled()
+ && config.carrier_learning
+ && config.limits.max_carrier_learning_entries < 3
+ {
+ return config_error(
+ "web.limits.max_carrier_learning_entries must be >= 3 when carrier learning is enabled",
+ );
+ }
+ if candidates.len() > WebCarrier::ALL.len() {
+ return config_error(
+ "web.carriers and the web.carrier fallback must contain at most four carriers",
+ );
+ }
+ let deadlines = config.timeouts.carrier_negotiation_deadlines_secs;
+ if deadlines[0] == 0 || deadlines.windows(2).any(|pair| pair[0] >= pair[1]) {
+ return config_error(
+ "web.timeouts.carrier_negotiation_deadlines_secs must be non-zero and strictly increasing",
+ );
+ }
+ let retained_chain_secs = deadlines[3]
+ .checked_add(config.timeouts.carrier_health_secs)
+ .and_then(|value| value.checked_add(1));
+ if retained_chain_secs.is_none_or(|value| value >= config.timeouts.bootstrap_lifetime_secs) {
+ return config_error(
+ "web.timeouts final carrier deadline plus health and cleanup must be lower than bootstrap_lifetime_secs",
+ );
+ }
+ Ok(candidates)
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ #[test]
+ fn fallback_is_appended_once() {
+ let config = WebConfig {
+ carrier: WebCarrier::Https,
+ carriers: WebCarriers::Enabled(vec![WebCarrier::Websocket, WebCarrier::Https]),
+ ..Default::default()
+ };
+ assert_eq!(
+ validate(&config).unwrap(),
+ vec![WebCarrier::Websocket, WebCarrier::Https]
+ );
+ }
+
+ #[test]
+ fn duplicate_carriers_are_rejected() {
+ let config = WebConfig {
+ carriers: WebCarriers::Enabled(vec![WebCarrier::Websocket, WebCarrier::Websocket]),
+ ..Default::default()
+ };
+ assert!(validate(&config).is_err());
+ }
+
+ #[test]
+ fn missing_or_false_carriers_disable_negotiation() {
+ #[derive(serde::Deserialize)]
+ struct Wrapper {
+ value: WebCarriers,
+ }
+
+ let config = WebConfig::default();
+ assert!(!config.carrier_negotiation_enabled());
+ assert_eq!(validate(&config).unwrap(), [WebCarrier::Https]);
+
+ let disabled: Wrapper = toml::from_str("value = false").unwrap();
+ assert_eq!(disabled.value, WebCarriers::Disabled);
+ assert!(toml::from_str::("value = true").is_err());
+ }
+
+ #[test]
+ fn fallback_cannot_expand_the_candidate_set_beyond_four() {
+ let mut config = WebConfig {
+ carrier: WebCarrier::Https,
+ carriers: WebCarriers::Enabled(vec![
+ WebCarrier::HttpsLanes,
+ WebCarrier::Websocket,
+ WebCarrier::WebsocketLanes,
+ WebCarrier::Https,
+ ]),
+ ..Default::default()
+ };
+ assert_eq!(validate(&config).unwrap().len(), 4);
+
+ config.carriers = WebCarriers::Enabled(vec![
+ WebCarrier::HttpsLanes,
+ WebCarrier::Websocket,
+ WebCarrier::WebsocketLanes,
+ ]);
+ assert_eq!(validate(&config).unwrap().len(), 4);
+ }
+
+ #[test]
+ fn deadlines_are_cumulative_and_bounded_by_bootstrap_lifetime() {
+ let mut config = WebConfig::default();
+ config.timeouts.carrier_negotiation_deadlines_secs = [3, 3, 8, 12];
+ assert!(validate(&config).is_err());
+ config.timeouts.carrier_negotiation_deadlines_secs = [3, 5, 8, 121];
+ assert!(validate(&config).is_err());
+ config.timeouts.carrier_negotiation_deadlines_secs = [3, 5, 8, 89];
+ assert!(validate(&config).is_err());
+ config.timeouts.carrier_negotiation_deadlines_secs = [3, 5, 8, 88];
+ assert!(validate(&config).is_ok());
+ }
+}
diff --git a/src/config/load/validate_web/timeouts.rs b/src/config/load/validate_web/timeouts.rs
new file mode 100644
index 0000000..dd4b989
--- /dev/null
+++ b/src/config/load/validate_web/timeouts.rs
@@ -0,0 +1,62 @@
+use super::*;
+
+/// Validates WEB request, learning, and lifecycle timeouts.
+pub(super) fn validate(timeouts: &WebTimeoutsConfig) -> Result<()> {
+ let values = [
+ ("header_secs", timeouts.header_secs),
+ ("body_secs", timeouts.body_secs),
+ ("stream_handshake_secs", timeouts.stream_handshake_secs),
+ ("stream_first_byte_secs", timeouts.stream_first_byte_secs),
+ ("long_poll_secs", timeouts.long_poll_secs),
+ ("lane_open_wait_secs", timeouts.lane_open_wait_secs),
+ ("carrier_health_secs", timeouts.carrier_health_secs),
+ ("websocket_upgrade_secs", timeouts.websocket_upgrade_secs),
+ ("websocket_open_secs", timeouts.websocket_open_secs),
+ ("websocket_write_secs", timeouts.websocket_write_secs),
+ (
+ "websocket_backpressure_secs",
+ timeouts.websocket_backpressure_secs,
+ ),
+ ("websocket_eviction_secs", timeouts.websocket_eviction_secs),
+ ("bootstrap_lifetime_secs", timeouts.bootstrap_lifetime_secs),
+ ("reconnect_grace_secs", timeouts.reconnect_grace_secs),
+ ("http_idle_secs", timeouts.http_idle_secs),
+ ("shutdown_secs", timeouts.shutdown_secs),
+ ("decoy_header_secs", timeouts.decoy_header_secs),
+ ];
+ if let Some((field, _)) = values
+ .into_iter()
+ .find(|(_, value)| !(1..=3600).contains(value))
+ {
+ return config_error(&format!("web.timeouts.{field} must be within [1, 3600]"));
+ }
+ if !(2..=86_400).contains(&timeouts.carrier_learning_secs) {
+ return config_error("web.timeouts.carrier_learning_secs must be within [2, 86400]");
+ }
+ if timeouts.stream_first_byte_secs > 300 {
+ return config_error("web.timeouts.stream_first_byte_secs must be within [1, 300]");
+ }
+ if timeouts.websocket_upgrade_secs > 60 {
+ return config_error("web.timeouts.websocket_upgrade_secs must be within [1, 60]");
+ }
+ if timeouts.websocket_open_secs > 300 {
+ return config_error("web.timeouts.websocket_open_secs must be within [1, 300]");
+ }
+ if timeouts.lane_open_wait_secs > timeouts.long_poll_secs {
+ return config_error("web.timeouts.lane_open_wait_secs must not exceed long_poll_secs");
+ }
+ if timeouts.carrier_health_secs > timeouts.reconnect_grace_secs {
+ return config_error(
+ "web.timeouts.carrier_health_secs must not exceed reconnect_grace_secs",
+ );
+ }
+ let request_deadline = timeouts
+ .header_secs
+ .max(timeouts.body_secs)
+ .max(timeouts.long_poll_secs)
+ .max(timeouts.decoy_header_secs);
+ if request_deadline >= timeouts.http_idle_secs {
+ return config_error("web.timeouts request deadlines must be lower than http_idle_secs");
+ }
+ Ok(())
+}
diff --git a/src/config/load/validate_web/websocket.rs b/src/config/load/validate_web/websocket.rs
index 84180d1..34ea38e 100644
--- a/src/config/load/validate_web/websocket.rs
+++ b/src/config/load/validate_web/websocket.rs
@@ -7,7 +7,7 @@ const WEBSOCKET_FRAME_OVERHEAD_BYTES: usize = 14;
/// Validates WebSocket admission, memory, and deadline invariants.
pub(super) fn validate(
- carrier: WebCarrier,
+ carriers: &[WebCarrier],
limits: &WebLimitsConfig,
timeouts: &WebTimeoutsConfig,
) -> Result<()> {
@@ -27,7 +27,7 @@ pub(super) fn validate(
"web.timeouts.websocket_eviction_secs must not exceed websocket_write_secs",
);
}
- if !carrier.uses_websocket() {
+ if !carriers.iter().any(|carrier| carrier.uses_websocket()) {
return Ok(());
}
if limits.carrier_batch_bytes > MAX_WEBSOCKET_BATCH_BYTES {
@@ -42,6 +42,14 @@ pub(super) fn validate(
"WebSocket carriers require websocket_http_connection_reserve within [1, max_http_connections)",
);
}
+ let websocket_capacity = limits
+ .max_http_connections
+ .saturating_sub(limits.websocket_http_connection_reserve);
+ if limits.max_websocket_evictions_in_flight > websocket_capacity {
+ return config_error(
+ "web.limits.max_websocket_evictions_in_flight must not exceed WebSocket connection capacity",
+ );
+ }
let socket_base = WEBSOCKET_IO_BUFFER_BYTES
.checked_mul(2)
.and_then(|value| value.checked_add(WEBSOCKET_DRIVER_OVERHEAD_BYTES))
diff --git a/src/config/tests/load_basic_tests/web_tests.rs b/src/config/tests/load_basic_tests/web_tests.rs
index 4390f92..c12919a 100644
--- a/src/config/tests/load_basic_tests/web_tests.rs
+++ b/src/config/tests/load_basic_tests/web_tests.rs
@@ -49,6 +49,88 @@ fn web_config_builds_canonical_runtime_snapshot() {
assert_eq!(vhost.profiles[0].max_streams_per_session, 16);
assert_eq!(vhost.profiles[0].key_fingerprint.len(), 16);
assert_ne!(vhost.profiles[0].key_fingerprint, "0001020304050607");
+ assert!(!vhost.profiles[0].carrier_negotiation_enabled);
+ assert_eq!(
+ vhost.profiles[0].carriers.as_ref(),
+ [WebCarrier::HttpsLanes]
+ );
+}
+
+#[test]
+fn web_carriers_missing_or_false_disable_negotiation() {
+ let missing = load_config_from_temp_toml(WEB_CONFIG);
+ assert!(!missing.web.carrier_negotiation_enabled());
+ assert!(!missing.web.runtime.unwrap().profiles[0].carrier_learning);
+
+ let disabled = WEB_CONFIG.replace(
+ "carrier = \"https-lanes\"",
+ "carrier = \"https-lanes\"\ncarriers = false",
+ );
+ let disabled = load_config_from_temp_toml(&disabled);
+ assert!(!disabled.web.carrier_negotiation_enabled());
+ assert!(!disabled.web.runtime.as_ref().unwrap().profiles[0].carrier_learning);
+ assert_eq!(
+ disabled.web.runtime.unwrap().profiles[0].carriers.as_ref(),
+ [WebCarrier::HttpsLanes]
+ );
+}
+
+#[test]
+fn web_carrier_array_enables_ordered_negotiation_and_appends_fallback() {
+ let configured = WEB_CONFIG.replace(
+ "carrier = \"https-lanes\"",
+ "carrier = \"https-lanes\"\ncarriers = [\"websocket\", \"https\"]\ncarrier_learning = false",
+ );
+ let config = load_config_from_temp_toml(&configured);
+ assert!(config.web.carrier_negotiation_enabled());
+ assert!(!config.web.carrier_learning);
+ let profile = &config.web.runtime.unwrap().profiles[0];
+ assert_eq!(
+ profile.carriers.as_ref(),
+ [
+ WebCarrier::Websocket,
+ WebCarrier::Https,
+ WebCarrier::HttpsLanes
+ ]
+ );
+ assert!(!profile.carrier_learning);
+}
+
+#[test]
+fn web_carriers_reject_true_empty_and_duplicates() {
+ for value in ["true", "[]", "[\"https\", \"https\"]"] {
+ let invalid = WEB_CONFIG.replace(
+ "carrier = \"https-lanes\"",
+ &format!("carrier = \"https-lanes\"\ncarriers = {value}"),
+ );
+ assert!(load_config_error_from_temp_toml(&invalid).contains("web.carriers"));
+ }
+}
+
+#[test]
+fn web_carrier_deadlines_and_learning_window_are_configurable() {
+ let configured = WEB_CONFIG.replace(
+ "[[web.vhosts]]",
+ "[web.timeouts]\ncarrier_negotiation_deadlines_secs = [1, 2, 4, 9]\ncarrier_learning_secs = 30\n\n[[web.vhosts]]",
+ );
+ let config = load_config_from_temp_toml(&configured);
+ assert_eq!(
+ config.web.timeouts.carrier_negotiation_deadlines_secs,
+ [1, 2, 4, 9]
+ );
+ assert_eq!(config.web.timeouts.carrier_learning_secs, 30);
+}
+
+#[test]
+fn web_carrier_learning_capacity_must_remain_nonzero() {
+ let invalid = WEB_CONFIG.replace(
+ "[[web.vhosts]]",
+ "[web.limits]\nmax_carrier_learning_entries = 0\n\n[[web.vhosts]]",
+ );
+ assert!(
+ load_config_error_from_temp_toml(&invalid)
+ .contains("web.limits.max_carrier_learning_entries")
+ );
}
#[test]
@@ -106,7 +188,7 @@ fn https_lanes_requires_separate_poll_and_control_handler_capacity() {
"carrier = \"https-lanes\"\n\n[web.limits]\nmax_http_handlers = 1\nmax_body_readers = 1",
);
let error = load_config_error_from_temp_toml(&invalid);
- assert!(error.contains("web.carrier=https-lanes requires"));
+ assert!(error.contains("WEB https-lanes candidates require"));
}
#[test]
diff --git a/src/config/types.rs b/src/config/types.rs
index 8c7f3f5..562f54e 100644
--- a/src/config/types.rs
+++ b/src/config/types.rs
@@ -24,6 +24,8 @@ mod network;
mod policies;
mod server;
mod web;
+// WEB carrier tokens and fixed-slot policy helpers remain independent from bulky config types.
+mod web_carrier;
// WEB debug capture policy is reusable by config reload and process storage.
mod web_debug;
@@ -51,13 +53,15 @@ pub use server::{
};
#[allow(unused_imports)]
pub use web::{
- WebCarrier, WebConfig, WebDecoyConfig, WebLimitsConfig, WebProfileConfig, WebSecretMode,
- WebTimeoutsConfig, WebVhostConfig,
+ WebCarrierNegotiationAggressiveness, WebConfig, WebDecoyConfig, WebLimitsConfig,
+ WebProfileConfig, WebSecretMode, WebTimeoutsConfig, WebVhostConfig,
};
pub(crate) use web::{
WebRuntimeConfig, WebRuntimeDecoy, WebRuntimeProfile, WebRuntimeVhost, WebStaticAsset,
WebStaticSite,
};
+#[allow(unused_imports)]
+pub use web_carrier::{WebCarrier, WebCarriers};
pub(crate) use web_debug::web_debug_fits_limits;
pub use web_debug::{WebDebugBodyCapture, WebDebugConfig};
diff --git a/src/config/types/web.rs b/src/config/types/web.rs
index db32ec0..863b533 100644
--- a/src/config/types/web.rs
+++ b/src/config/types/web.rs
@@ -6,8 +6,13 @@ use std::sync::Arc;
use bytes::Bytes;
use serde::{Deserialize, Serialize};
+use super::web_carrier::{WebCarrier, WebCarriers};
use super::web_debug::WebDebugConfig;
+// Serialized WEB defaults remain separate from the runtime data model.
+mod defaults;
+use defaults::*;
+
/// Client-facing secret representation used to derive a WEB capability.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
@@ -18,48 +23,6 @@ pub enum WebSecretMode {
Dd,
}
-/// Carrier selected for newly issued WEB bridge sessions.
-#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash, Serialize, Deserialize)]
-#[serde(rename_all = "kebab-case")]
-pub enum WebCarrier {
- /// Serialize all logical streams through one uplink and one downlink sequence.
- #[default]
- Https,
- /// Give every logical stream independent HTTPS sequencing and polling state.
- HttpsLanes,
- /// Multiplex all logical streams over one ordered WebSocket.
- Websocket,
- /// Give every logical stream an independently owned WebSocket lane.
- WebsocketLanes,
-}
-
-impl WebCarrier {
- /// Returns the exact carrier token advertised to the browser bridge.
- pub(crate) const fn as_str(self) -> &'static str {
- match self {
- Self::Https => "https",
- Self::HttpsLanes => "https-lanes",
- Self::Websocket => "websocket",
- Self::WebsocketLanes => "websocket-lanes",
- }
- }
-
- /// Returns whether one carrier owns independent state per logical stream.
- pub(crate) const fn uses_lanes(self) -> bool {
- matches!(self, Self::HttpsLanes | Self::WebsocketLanes)
- }
-
- /// Returns whether carrier messages use RFC 6455 instead of HTTP bodies.
- pub(crate) const fn uses_websocket(self) -> bool {
- matches!(self, Self::Websocket | Self::WebsocketLanes)
- }
-
- /// Returns whether all logical streams share one carrier state machine.
- pub(crate) const fn is_multiplexed(self) -> bool {
- matches!(self, Self::Https | Self::Websocket)
- }
-}
-
/// One access user explicitly exposed through a WEB virtual host.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct WebProfileConfig {
@@ -135,6 +98,15 @@ pub struct WebLimitsConfig {
/// Process-wide concurrently executing HTTP handler ceiling.
#[serde(default = "default_web_max_http_handlers")]
pub max_http_handlers: usize,
+ /// Per-session ceiling for downlink polls waiting for a lane OPEN.
+ #[serde(default = "default_web_max_lane_open_waits_per_session")]
+ pub max_lane_open_waits_per_session: usize,
+ /// Queued and resident DATA bytes allowed for one independent lane.
+ #[serde(default = "default_web_pending_bytes_per_lane")]
+ pub pending_bytes_per_lane: usize,
+ /// Queued and resident DATA items allowed for one independent lane.
+ #[serde(default = "default_web_pending_items_per_lane")]
+ pub pending_items_per_lane: usize,
/// Process-wide transient WebSocket byte sub-budget inside pending bytes.
#[serde(default = "default_web_websocket_bytes_global")]
pub websocket_bytes_global: usize,
@@ -147,6 +119,12 @@ pub struct WebLimitsConfig {
/// Accepted HTTP connections that WebSocket upgrades must leave available.
#[serde(default = "default_web_websocket_http_connection_reserve")]
pub websocket_http_connection_reserve: usize,
+ /// Concurrent pressure-eviction claims allowed process-wide.
+ #[serde(default = "default_web_max_websocket_evictions_in_flight")]
+ pub max_websocket_evictions_in_flight: usize,
+ /// Process-wide bounded carrier-learning evidence entry ceiling.
+ #[serde(default = "default_web_max_carrier_learning_entries")]
+ pub max_carrier_learning_entries: usize,
/// Process-wide concurrently collected request body ceiling.
#[serde(default = "default_web_max_body_readers")]
pub max_body_readers: usize,
@@ -216,7 +194,7 @@ pub struct WebLimitsConfig {
/// Process-wide retained and in-flight WEB debug byte ceiling.
#[serde(default = "default_web_debug_bytes_global")]
pub debug_bytes_global: usize,
- /// Declared process envelope for HTTP heads, bodies, queues, and static snapshots.
+ /// Declared process envelope for HTTP, queues, lane state, learning, and static snapshots.
#[serde(default = "default_web_memory_envelope_bytes")]
pub memory_envelope_bytes: usize,
/// Sustained process-wide bootstrap issuance rate.
@@ -249,10 +227,15 @@ impl Default for WebLimitsConfig {
max_frames_per_body: default_web_max_frames_per_body(),
max_http_connections: default_web_max_http_connections(),
max_http_handlers: default_web_max_http_handlers(),
+ max_lane_open_waits_per_session: default_web_max_lane_open_waits_per_session(),
+ pending_bytes_per_lane: default_web_pending_bytes_per_lane(),
+ pending_items_per_lane: default_web_pending_items_per_lane(),
websocket_bytes_global: default_web_websocket_bytes_global(),
websocket_admission_watermark_pct: default_web_websocket_admission_watermark_pct(),
websocket_eviction_watermark_pct: default_web_websocket_eviction_watermark_pct(),
websocket_http_connection_reserve: default_web_websocket_http_connection_reserve(),
+ max_websocket_evictions_in_flight: default_web_max_websocket_evictions_in_flight(),
+ max_carrier_learning_entries: default_web_max_carrier_learning_entries(),
max_body_readers: default_web_max_body_readers(),
max_body_bytes_global: default_web_max_body_bytes_global(),
max_sessions_global: default_web_max_sessions_global(),
@@ -299,9 +282,24 @@ pub struct WebTimeoutsConfig {
/// Deadline from the first inner byte through MTProxy authentication.
#[serde(default = "default_web_stream_handshake_timeout_secs")]
pub stream_handshake_secs: u64,
+ /// Absolute deadline for receiving the first inner MTProxy byte.
+ #[serde(default = "default_web_stream_first_byte_secs")]
+ pub stream_first_byte_secs: u64,
/// Maximum wait for one empty downlink long poll.
#[serde(default = "default_web_long_poll_timeout_secs")]
pub long_poll_secs: u64,
+ /// Grace for a canonical downlink poll that races its lane OPEN.
+ #[serde(default = "default_web_lane_open_wait_secs")]
+ pub lane_open_wait_secs: u64,
+ /// Post-commit observation interval required before learning succeeds.
+ #[serde(default = "default_web_carrier_health_secs")]
+ pub carrier_health_secs: u64,
+ /// Maximum wait for Hyper to transfer an accepted WebSocket upgrade.
+ #[serde(default = "default_web_websocket_upgrade_secs")]
+ pub websocket_upgrade_secs: u64,
+ /// Absolute deadline for the first carrier binary message after upgrade.
+ #[serde(default = "default_web_websocket_open_secs")]
+ pub websocket_open_secs: u64,
/// Maximum wait for one WebSocket write to complete.
#[serde(default = "default_web_websocket_write_secs")]
pub websocket_write_secs: u64,
@@ -311,6 +309,12 @@ pub struct WebTimeoutsConfig {
/// Maximum graceful close wait for an evicted WebSocket.
#[serde(default = "default_web_websocket_eviction_secs")]
pub websocket_eviction_secs: u64,
+ /// Cumulative carrier-attempt deadlines for up to four unique candidates.
+ #[serde(default = "default_web_carrier_negotiation_deadlines_secs")]
+ pub carrier_negotiation_deadlines_secs: [u64; 4],
+ /// Fixed process-local carrier-learning evidence lifetime.
+ #[serde(default = "default_web_carrier_learning_secs")]
+ pub carrier_learning_secs: u64,
/// Lifetime of an unused bootstrap credential and closed-token replay marker.
#[serde(default = "default_web_bootstrap_lifetime_secs")]
pub bootstrap_lifetime_secs: u64,
@@ -334,10 +338,17 @@ impl Default for WebTimeoutsConfig {
header_secs: default_web_header_timeout_secs(),
body_secs: default_web_body_timeout_secs(),
stream_handshake_secs: default_web_stream_handshake_timeout_secs(),
+ stream_first_byte_secs: default_web_stream_first_byte_secs(),
long_poll_secs: default_web_long_poll_timeout_secs(),
+ lane_open_wait_secs: default_web_lane_open_wait_secs(),
+ carrier_health_secs: default_web_carrier_health_secs(),
+ websocket_upgrade_secs: default_web_websocket_upgrade_secs(),
+ websocket_open_secs: default_web_websocket_open_secs(),
websocket_write_secs: default_web_websocket_write_secs(),
websocket_backpressure_secs: default_web_websocket_backpressure_secs(),
websocket_eviction_secs: default_web_websocket_eviction_secs(),
+ carrier_negotiation_deadlines_secs: default_web_carrier_negotiation_deadlines_secs(),
+ carrier_learning_secs: default_web_carrier_learning_secs(),
bootstrap_lifetime_secs: default_web_bootstrap_lifetime_secs(),
reconnect_grace_secs: default_web_reconnect_grace_secs(),
http_idle_secs: default_web_http_idle_secs(),
@@ -347,15 +358,37 @@ impl Default for WebTimeoutsConfig {
}
}
+/// Sensitivity of process-local carrier-learning evidence.
+#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
+#[serde(rename_all = "lowercase")]
+pub enum WebCarrierNegotiationAggressiveness {
+ /// Require broad evidence and never rank by client IP.
+ #[default]
+ Conservative,
+ /// Use moderate User-Agent, client-IP, and profile thresholds.
+ Balanced,
+ /// React to the first bounded evidence sample.
+ Aggressive,
+}
+
/// WEB ingress, carrier, fallback, and lifecycle configuration.
-#[derive(Debug, Clone, Default, Serialize, Deserialize)]
+#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WebConfig {
/// Enables issuance of new WEB bridge and session credentials.
#[serde(default)]
pub enabled: bool,
- /// Carrier selected for newly issued WEB bridge sessions.
+ /// Sole carrier when negotiation is disabled and final fallback when enabled.
#[serde(default)]
pub carrier: WebCarrier,
+ /// Ordered carriers considered by server-side negotiation before the fallback carrier.
+ #[serde(default)]
+ pub carriers: WebCarriers,
+ /// Enables bounded process-local carrier learning for automatic sessions.
+ #[serde(default = "default_web_carrier_learning")]
+ pub carrier_learning: bool,
+ /// Controls the evidence thresholds used by automatic carrier ranking.
+ #[serde(default)]
+ pub carrier_negotiation_aggressiveness: WebCarrierNegotiationAggressiveness,
/// Hard process and protocol limits.
#[serde(default)]
pub limits: WebLimitsConfig,
@@ -373,6 +406,44 @@ pub struct WebConfig {
pub(crate) runtime: Option>,
}
+impl WebConfig {
+ /// Returns the configured negotiation order with the fallback appended once.
+ pub(crate) fn carrier_candidates(&self) -> Vec {
+ let Some(configured) = self.carriers.enabled() else {
+ return vec![self.carrier];
+ };
+ let mut candidates = configured
+ .iter()
+ .copied()
+ .filter(|carrier| *carrier != self.carrier)
+ .collect::>();
+ candidates.push(self.carrier);
+ candidates
+ }
+
+ /// Returns whether the explicit candidate list enables auto-negotiation.
+ pub(crate) fn carrier_negotiation_enabled(&self) -> bool {
+ self.carriers.enabled().is_some()
+ }
+}
+
+impl Default for WebConfig {
+ fn default() -> Self {
+ Self {
+ enabled: false,
+ carrier: WebCarrier::default(),
+ carriers: WebCarriers::default(),
+ carrier_learning: default_web_carrier_learning(),
+ carrier_negotiation_aggressiveness: WebCarrierNegotiationAggressiveness::default(),
+ limits: WebLimitsConfig::default(),
+ debug: WebDebugConfig::default(),
+ timeouts: WebTimeoutsConfig::default(),
+ vhosts: Vec::new(),
+ runtime: None,
+ }
+ }
+}
+
/// Precomputed WEB configuration consumed by listener hot paths.
#[derive(Debug)]
pub(crate) struct WebRuntimeConfig {
@@ -406,8 +477,16 @@ pub(crate) struct WebRuntimeProfile {
pub(crate) user: String,
/// Client secret representation and inner protocol policy.
pub(crate) secret_mode: WebSecretMode,
- /// Carrier frozen into bridge and session state at issuance time.
+ /// Sole carrier or final fallback frozen into the issued bridge policy.
pub(crate) carrier: WebCarrier,
+ /// Whether an explicit carrier list enabled automatic negotiation.
+ pub(crate) carrier_negotiation_enabled: bool,
+ /// Whether automatic outcomes consult and update process-local evidence.
+ pub(crate) carrier_learning: bool,
+ /// Ordered negotiation candidates including the fallback carrier exactly once.
+ pub(crate) carriers: Arc<[WebCarrier]>,
+ /// Cumulative carrier-attempt deadlines frozen when the bridge is issued.
+ pub(crate) carrier_negotiation_deadlines_secs: [u64; 4],
/// HMAC-derived bridge capability.
pub(crate) capability: [u8; 32],
/// Non-secret domain-separated client-secret fingerprint for debugging.
@@ -446,93 +525,3 @@ pub(crate) struct WebStaticAsset {
/// Strong SHA-256 entity tag.
pub(crate) etag: String,
}
-
-fn default_web_static_index() -> String {
- "index.html".to_string()
-}
-
-macro_rules! usize_default {
- ($name:ident, $value:expr) => {
- fn $name() -> usize {
- $value
- }
- };
-}
-
-macro_rules! u32_default {
- ($name:ident, $value:expr) => {
- fn $name() -> u32 {
- $value
- }
- };
-}
-
-macro_rules! u8_default {
- ($name:ident, $value:expr) => {
- fn $name() -> u8 {
- $value
- }
- };
-}
-
-macro_rules! u64_default {
- ($name:ident, $value:expr) => {
- fn $name() -> u64 {
- $value
- }
- };
-}
-
-usize_default!(default_web_max_header_bytes, 16 * 1024);
-usize_default!(default_web_max_body_bytes, 2 * 1024 * 1024);
-usize_default!(default_web_max_frame_payload_bytes, 1024 * 1024);
-usize_default!(default_web_carrier_batch_bytes, 2 * 1024 * 1024);
-usize_default!(default_web_max_frames_per_body, 4096);
-usize_default!(default_web_max_http_connections, 1024);
-usize_default!(default_web_max_http_handlers, 512);
-usize_default!(default_web_websocket_bytes_global, 256 * 1024 * 1024);
-u8_default!(default_web_websocket_admission_watermark_pct, 75);
-u8_default!(default_web_websocket_eviction_watermark_pct, 90);
-usize_default!(default_web_websocket_http_connection_reserve, 64);
-usize_default!(default_web_max_body_readers, 32);
-usize_default!(default_web_max_body_bytes_global, 64 * 1024 * 1024);
-usize_default!(default_web_max_sessions_global, 128);
-usize_default!(default_web_max_sessions_per_ip, 16);
-usize_default!(default_web_max_streams_per_session, 128);
-usize_default!(default_web_max_streams_global, 4096);
-usize_default!(default_web_max_stream_handshakes, 256);
-usize_default!(default_web_max_tombstones, 4096);
-usize_default!(default_web_pending_bytes_per_session, 32 * 1024 * 1024);
-usize_default!(default_web_pending_bytes_global, 512 * 1024 * 1024);
-usize_default!(default_web_pending_items_per_session, 16 * 1024);
-usize_default!(default_web_pending_items_global, 256 * 1024);
-usize_default!(default_web_control_bytes_per_session, 256 * 1024);
-usize_default!(default_web_control_bytes_global, 16 * 1024 * 1024);
-usize_default!(default_web_max_bootstraps_global, 512);
-usize_default!(default_web_max_bootstraps_per_ip, 64);
-usize_default!(default_web_max_vhosts, 8);
-usize_default!(default_web_max_profiles, 32);
-usize_default!(default_web_max_static_files, 4096);
-usize_default!(default_web_max_static_file_bytes, 8 * 1024 * 1024);
-usize_default!(default_web_max_static_bytes, 64 * 1024 * 1024);
-usize_default!(default_web_debug_records_capacity, 65_536);
-usize_default!(default_web_debug_bytes_global, 64 * 1024 * 1024);
-usize_default!(default_web_memory_envelope_bytes, 768 * 1024 * 1024);
-u32_default!(default_web_new_bootstraps_per_minute, 1200);
-u32_default!(default_web_new_bootstraps_burst, 256);
-u32_default!(default_web_new_sessions_per_minute, 600);
-u32_default!(default_web_new_sessions_burst, 128);
-u32_default!(default_web_new_streams_per_minute, 6000);
-u32_default!(default_web_new_streams_burst, 512);
-u64_default!(default_web_header_timeout_secs, 10);
-u64_default!(default_web_body_timeout_secs, 30);
-u64_default!(default_web_stream_handshake_timeout_secs, 10);
-u64_default!(default_web_long_poll_timeout_secs, 25);
-u64_default!(default_web_websocket_write_secs, 30);
-u64_default!(default_web_websocket_backpressure_secs, 30);
-u64_default!(default_web_websocket_eviction_secs, 1);
-u64_default!(default_web_bootstrap_lifetime_secs, 120);
-u64_default!(default_web_reconnect_grace_secs, 120);
-u64_default!(default_web_http_idle_secs, 75);
-u64_default!(default_web_shutdown_secs, 15);
-u64_default!(default_web_decoy_header_timeout_secs, 30);
diff --git a/src/config/types/web/defaults.rs b/src/config/types/web/defaults.rs
new file mode 100644
index 0000000..f85b44d
--- /dev/null
+++ b/src/config/types/web/defaults.rs
@@ -0,0 +1,106 @@
+pub(super) fn default_web_static_index() -> String {
+ "index.html".to_string()
+}
+
+macro_rules! usize_default {
+ ($name:ident, $value:expr) => {
+ pub(super) fn $name() -> usize {
+ $value
+ }
+ };
+}
+
+macro_rules! u32_default {
+ ($name:ident, $value:expr) => {
+ pub(super) fn $name() -> u32 {
+ $value
+ }
+ };
+}
+
+macro_rules! u8_default {
+ ($name:ident, $value:expr) => {
+ pub(super) fn $name() -> u8 {
+ $value
+ }
+ };
+}
+
+macro_rules! u64_default {
+ ($name:ident, $value:expr) => {
+ pub(super) fn $name() -> u64 {
+ $value
+ }
+ };
+}
+
+usize_default!(default_web_max_header_bytes, 16 * 1024);
+usize_default!(default_web_max_body_bytes, 2 * 1024 * 1024);
+usize_default!(default_web_max_frame_payload_bytes, 1024 * 1024);
+usize_default!(default_web_carrier_batch_bytes, 2 * 1024 * 1024);
+usize_default!(default_web_max_frames_per_body, 4096);
+usize_default!(default_web_max_http_connections, 1024);
+usize_default!(default_web_max_http_handlers, 512);
+usize_default!(default_web_max_lane_open_waits_per_session, 16);
+usize_default!(default_web_pending_bytes_per_lane, 8 * 1024 * 1024);
+usize_default!(default_web_pending_items_per_lane, 1024);
+usize_default!(default_web_websocket_bytes_global, 256 * 1024 * 1024);
+u8_default!(default_web_websocket_admission_watermark_pct, 75);
+u8_default!(default_web_websocket_eviction_watermark_pct, 90);
+usize_default!(default_web_websocket_http_connection_reserve, 64);
+usize_default!(default_web_max_websocket_evictions_in_flight, 8);
+usize_default!(default_web_max_carrier_learning_entries, 4096);
+usize_default!(default_web_max_body_readers, 32);
+usize_default!(default_web_max_body_bytes_global, 64 * 1024 * 1024);
+usize_default!(default_web_max_sessions_global, 128);
+usize_default!(default_web_max_sessions_per_ip, 16);
+usize_default!(default_web_max_streams_per_session, 128);
+usize_default!(default_web_max_streams_global, 4096);
+usize_default!(default_web_max_stream_handshakes, 256);
+usize_default!(default_web_max_tombstones, 4096);
+usize_default!(default_web_pending_bytes_per_session, 32 * 1024 * 1024);
+usize_default!(default_web_pending_bytes_global, 512 * 1024 * 1024);
+usize_default!(default_web_pending_items_per_session, 16 * 1024);
+usize_default!(default_web_pending_items_global, 256 * 1024);
+usize_default!(default_web_control_bytes_per_session, 256 * 1024);
+usize_default!(default_web_control_bytes_global, 16 * 1024 * 1024);
+usize_default!(default_web_max_bootstraps_global, 512);
+usize_default!(default_web_max_bootstraps_per_ip, 64);
+usize_default!(default_web_max_vhosts, 8);
+usize_default!(default_web_max_profiles, 32);
+usize_default!(default_web_max_static_files, 4096);
+usize_default!(default_web_max_static_file_bytes, 8 * 1024 * 1024);
+usize_default!(default_web_max_static_bytes, 64 * 1024 * 1024);
+usize_default!(default_web_debug_records_capacity, 65_536);
+usize_default!(default_web_debug_bytes_global, 64 * 1024 * 1024);
+usize_default!(default_web_memory_envelope_bytes, 1280 * 1024 * 1024);
+u32_default!(default_web_new_bootstraps_per_minute, 1200);
+u32_default!(default_web_new_bootstraps_burst, 256);
+u32_default!(default_web_new_sessions_per_minute, 600);
+u32_default!(default_web_new_sessions_burst, 128);
+u32_default!(default_web_new_streams_per_minute, 6000);
+u32_default!(default_web_new_streams_burst, 512);
+u64_default!(default_web_header_timeout_secs, 10);
+u64_default!(default_web_body_timeout_secs, 30);
+u64_default!(default_web_stream_handshake_timeout_secs, 10);
+u64_default!(default_web_stream_first_byte_secs, 30);
+u64_default!(default_web_long_poll_timeout_secs, 25);
+u64_default!(default_web_lane_open_wait_secs, 2);
+u64_default!(default_web_carrier_health_secs, 30);
+u64_default!(default_web_websocket_upgrade_secs, 5);
+u64_default!(default_web_websocket_open_secs, 15);
+u64_default!(default_web_websocket_write_secs, 30);
+u64_default!(default_web_websocket_backpressure_secs, 30);
+u64_default!(default_web_websocket_eviction_secs, 1);
+pub(super) fn default_web_carrier_negotiation_deadlines_secs() -> [u64; 4] {
+ [3, 5, 8, 12]
+}
+u64_default!(default_web_carrier_learning_secs, 600);
+pub(super) fn default_web_carrier_learning() -> bool {
+ true
+}
+u64_default!(default_web_bootstrap_lifetime_secs, 120);
+u64_default!(default_web_reconnect_grace_secs, 120);
+u64_default!(default_web_http_idle_secs, 75);
+u64_default!(default_web_shutdown_secs, 15);
+u64_default!(default_web_decoy_header_timeout_secs, 30);
diff --git a/src/config/types/web_carrier.rs b/src/config/types/web_carrier.rs
new file mode 100644
index 0000000..f3749e2
--- /dev/null
+++ b/src/config/types/web_carrier.rs
@@ -0,0 +1,115 @@
+use serde::{Deserialize, Serialize};
+
+/// Carrier selected for one newly issued WEB relay session.
+#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash, Serialize, Deserialize)]
+#[serde(rename_all = "kebab-case")]
+pub enum WebCarrier {
+ /// Serialize all logical streams through one uplink and one downlink sequence.
+ #[default]
+ Https,
+ /// Give every logical stream independent HTTPS sequencing and polling state.
+ HttpsLanes,
+ /// Multiplex all logical streams over one ordered WebSocket.
+ Websocket,
+ /// Give every logical stream an independently owned WebSocket lane.
+ WebsocketLanes,
+}
+
+impl WebCarrier {
+ /// Every carrier supported by the WEB v1 bridge.
+ pub(crate) const ALL: [Self; 4] = [
+ Self::Https,
+ Self::HttpsLanes,
+ Self::Websocket,
+ Self::WebsocketLanes,
+ ];
+
+ /// Returns the exact carrier token advertised to the browser bridge.
+ pub(crate) const fn as_str(self) -> &'static str {
+ match self {
+ Self::Https => "https",
+ Self::HttpsLanes => "https-lanes",
+ Self::Websocket => "websocket",
+ Self::WebsocketLanes => "websocket-lanes",
+ }
+ }
+
+ /// Returns the stable fixed-slot index used by bounded learning state.
+ pub(crate) const fn index(self) -> usize {
+ match self {
+ Self::Https => 0,
+ Self::HttpsLanes => 1,
+ Self::Websocket => 2,
+ Self::WebsocketLanes => 3,
+ }
+ }
+
+ /// Returns whether one carrier owns independent state per logical stream.
+ pub(crate) const fn uses_lanes(self) -> bool {
+ matches!(self, Self::HttpsLanes | Self::WebsocketLanes)
+ }
+
+ /// Returns whether carrier messages use RFC 6455 instead of HTTP bodies.
+ pub(crate) const fn uses_websocket(self) -> bool {
+ matches!(self, Self::Websocket | Self::WebsocketLanes)
+ }
+
+ /// Returns whether all logical streams share one carrier state machine.
+ pub(crate) const fn is_multiplexed(self) -> bool {
+ matches!(self, Self::Https | Self::Websocket)
+ }
+}
+
+/// Optional ordered carrier list that enables server-side auto-negotiation.
+#[derive(Debug, Clone, Default, PartialEq, Eq)]
+pub enum WebCarriers {
+ /// Auto-negotiation is disabled and only `web.carrier` is used.
+ #[default]
+ Disabled,
+ /// Auto-negotiation uses this ordered candidate list before the fallback.
+ Enabled(Vec),
+}
+
+impl WebCarriers {
+ /// Returns the explicit candidate list when negotiation is enabled.
+ pub fn enabled(&self) -> Option<&[WebCarrier]> {
+ match self {
+ Self::Disabled => None,
+ Self::Enabled(carriers) => Some(carriers),
+ }
+ }
+}
+
+impl Serialize for WebCarriers {
+ fn serialize(&self, serializer: S) -> Result
+ where
+ S: serde::Serializer,
+ {
+ match self {
+ Self::Disabled => false.serialize(serializer),
+ Self::Enabled(carriers) => carriers.serialize(serializer),
+ }
+ }
+}
+
+impl<'de> Deserialize<'de> for WebCarriers {
+ fn deserialize(deserializer: D) -> Result
+ where
+ D: serde::Deserializer<'de>,
+ {
+ #[derive(Deserialize)]
+ #[serde(untagged)]
+ enum Repr {
+ Flag(bool),
+ List(Vec),
+ }
+
+ match Repr::deserialize(deserializer)? {
+ Repr::Flag(false) => Ok(Self::Disabled),
+ Repr::Flag(true) => Err(serde::de::Error::custom(
+ "web.carriers accepts false or a non-empty carrier array",
+ )),
+ Repr::List(carriers) => Ok(Self::Enabled(carriers)),
+ }
+ }
+}
diff --git a/src/metrics.rs b/src/metrics.rs
index fb35b84..9bd36df 100644
--- a/src/metrics.rs
+++ b/src/metrics.rs
@@ -3572,24 +3572,24 @@ async fn render_metrics(
let _ = writeln!(out, "# TYPE telemt_user_connections_current gauge");
let _ = writeln!(
out,
- "# HELP telemt_user_octets_from_client Per-user bytes received"
+ "# HELP telemt_user_octets_from_client_total Per-user total bytes received"
);
- let _ = writeln!(out, "# TYPE telemt_user_octets_from_client counter");
+ let _ = writeln!(out, "# TYPE telemt_user_octets_from_client_total counter");
let _ = writeln!(
out,
- "# HELP telemt_user_octets_to_client Per-user bytes sent"
+ "# HELP telemt_user_octets_to_client_total Per-user total bytes sent"
);
- let _ = writeln!(out, "# TYPE telemt_user_octets_to_client counter");
+ let _ = writeln!(out, "# TYPE telemt_user_octets_to_client_total counter");
let _ = writeln!(
out,
- "# HELP telemt_user_msgs_from_client Per-user messages received"
+ "# HELP telemt_user_msgs_from_client_total Per-user total messages received"
);
- let _ = writeln!(out, "# TYPE telemt_user_msgs_from_client counter");
+ let _ = writeln!(out, "# TYPE telemt_user_msgs_from_client_total counter");
let _ = writeln!(
out,
- "# HELP telemt_user_msgs_to_client Per-user messages sent"
+ "# HELP telemt_user_msgs_to_client_total Per-user total messages sent"
);
- let _ = writeln!(out, "# TYPE telemt_user_msgs_to_client counter");
+ let _ = writeln!(out, "# TYPE telemt_user_msgs_to_client_total counter");
let _ = writeln!(
out,
"# HELP telemt_ip_reservation_rollback_total IP reservation rollbacks caused by later limit checks"
@@ -3708,28 +3708,28 @@ async fn render_metrics(
);
let _ = writeln!(
out,
- "telemt_user_octets_from_client{{user=\"{}\"}} {}",
+ "telemt_user_octets_from_client_total{{user=\"{}\"}} {}",
user,
s.octets_from_client
.load(std::sync::atomic::Ordering::Relaxed)
);
let _ = writeln!(
out,
- "telemt_user_octets_to_client{{user=\"{}\"}} {}",
+ "telemt_user_octets_to_client_total{{user=\"{}\"}} {}",
user,
s.octets_to_client
.load(std::sync::atomic::Ordering::Relaxed)
);
let _ = writeln!(
out,
- "telemt_user_msgs_from_client{{user=\"{}\"}} {}",
+ "telemt_user_msgs_from_client_total{{user=\"{}\"}} {}",
user,
s.msgs_from_client
.load(std::sync::atomic::Ordering::Relaxed)
);
let _ = writeln!(
out,
- "telemt_user_msgs_to_client{{user=\"{}\"}} {}",
+ "telemt_user_msgs_to_client_total{{user=\"{}\"}} {}",
user,
s.msgs_to_client.load(std::sync::atomic::Ordering::Relaxed)
);
@@ -3990,10 +3990,10 @@ mod tests {
assert!(output.contains("telemt_me_endpoint_quarantine_draining_suppressed_total 1"));
assert!(output.contains("telemt_user_connections_total{user=\"alice\"} 1"));
assert!(output.contains("telemt_user_connections_current{user=\"alice\"} 1"));
- assert!(output.contains("telemt_user_octets_from_client{user=\"alice\"} 1024"));
- assert!(output.contains("telemt_user_octets_to_client{user=\"alice\"} 2048"));
- assert!(output.contains("telemt_user_msgs_from_client{user=\"alice\"} 1"));
- assert!(output.contains("telemt_user_msgs_to_client{user=\"alice\"} 2"));
+ assert!(output.contains("telemt_user_octets_from_client_total{user=\"alice\"} 1024"));
+ assert!(output.contains("telemt_user_octets_to_client_total{user=\"alice\"} 2048"));
+ assert!(output.contains("telemt_user_msgs_from_client_total{user=\"alice\"} 1"));
+ assert!(output.contains("telemt_user_msgs_to_client_total{user=\"alice\"} 2"));
assert!(output.contains("telemt_user_unique_ips_current{user=\"alice\"} 1"));
assert!(output.contains("telemt_user_unique_ips_recent_window{user=\"alice\"} 1"));
assert!(output.contains("telemt_user_unique_ips_limit{user=\"alice\"} 4"));
diff --git a/src/web/bridge.rs b/src/web/bridge.rs
index 67a1b57..c262351 100644
--- a/src/web/bridge.rs
+++ b/src/web/bridge.rs
@@ -1,6 +1,5 @@
use base64::Engine as _;
-use crate::config::WebCarrier;
use crate::crypto::SecureRandom;
/// Browser security policy for the transient Telegram Desktop bridge page.
@@ -14,14 +13,17 @@ pub(crate) struct BridgePage {
pub(crate) content_security_policy: String,
}
-/// Renders the selected HTTPS WEB carrier bridge with a fresh CSP nonce.
+/// Renders the bounded WEB carrier-negotiation bridge with a fresh CSP nonce.
+#[allow(clippy::too_many_arguments)]
pub(crate) fn render(
host: &str,
bootstrap: &str,
batch_limit: usize,
queue_limit: usize,
queue_items: usize,
- carrier: WebCarrier,
+ negotiation_enabled: bool,
+ candidate_count: usize,
+ carrier_deadlines: [u64; 4],
rng: &SecureRandom,
) -> BridgePage {
let mut nonce = [0u8; 18];
@@ -34,7 +36,19 @@ pub(crate) fn render(
.replace("__BATCH_LIMIT__", &batch_limit.to_string())
.replace("__QUEUE_LIMIT__", &queue_limit.to_string())
.replace("__QUEUE_ITEMS__", &queue_items.to_string())
- .replace("__CARRIER__", carrier.as_str());
+ .replace(
+ "__NEGOTIATION_ENABLED__",
+ if negotiation_enabled { "true" } else { "false" },
+ )
+ .replace("__CANDIDATE_COUNT__", &candidate_count.to_string())
+ .replace(
+ "__CARRIER_DEADLINES__",
+ &carrier_deadlines
+ .iter()
+ .map(u64::to_string)
+ .collect::>()
+ .join(","),
+ );
BridgePage {
body,
content_security_policy: format!(
@@ -54,21 +68,32 @@ const DOCUMENT: &str = r##"