From 14e8d10ad37f821a4ec1843828184f88f0dbef9c Mon Sep 17 00:00:00 2001 From: Alexey <247128645+axkurcom@users.noreply.github.com> Date: Thu, 27 Aug 2026 00:06:37 +0300 Subject: [PATCH] Rustfmt --- src/api/web_status/details.rs | 5 +- src/config/load/validate_web/timeouts.rs | 4 +- .../tests/load_basic_tests/web_tests.rs | 11 ++- src/config/types.rs | 4 +- src/config/types/web.rs | 9 +- src/web/bridge/tests.rs | 48 ++++++----- src/web/http.rs | 6 +- src/web/http/down.rs | 6 +- src/web/http/negotiation_tests.rs | 69 ++++++++-------- src/web/http/request.rs | 25 ++---- src/web/http/session.rs | 9 +- src/web/http/tests.rs | 11 +-- src/web/http/websocket.rs | 34 ++++---- src/web/http/websocket/driver.rs | 19 +---- src/web/http/websocket/driver/io.rs | 12 +-- src/web/http/websocket/driver/lane.rs | 15 +--- src/web/http/websocket/tests.rs | 3 +- src/web/manager.rs | 7 +- src/web/manager/carrier_learning.rs | 30 ++++--- src/web/manager/negotiation.rs | 9 +- src/web/manager/session_creation.rs | 82 ++++++++++--------- src/web/manager/websocket.rs | 11 +-- src/web/session.rs | 1 - src/web/session/downlink.rs | 6 +- src/web/session/downlink_tests.rs | 3 +- src/web/session/lane_uplink.rs | 6 +- src/web/session/lanes.rs | 17 ++-- src/web/session/lanes/tests.rs | 14 ++-- src/web/session/lifecycle.rs | 14 +--- src/web/session/negotiation.rs | 5 +- src/web/session/resident.rs | 12 ++- src/web/session/uplink.rs | 8 +- src/web/session/websocket.rs | 3 +- src/web/trace/sanitize.rs | 2 +- 34 files changed, 221 insertions(+), 299 deletions(-) diff --git a/src/api/web_status/details.rs b/src/api/web_status/details.rs index e48775b..90ce3a2 100644 --- a/src/api/web_status/details.rs +++ b/src/api/web_status/details.rs @@ -85,10 +85,7 @@ pub(super) fn push_body( html.push_str(""); } -pub(super) fn push_lifecycle( - html: &mut String, - event: &crate::web::trace::TraceLifecycleRecord, -) { +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: ");
diff --git a/src/config/load/validate_web/timeouts.rs b/src/config/load/validate_web/timeouts.rs
index 73360a0..dd4b989 100644
--- a/src/config/load/validate_web/timeouts.rs
+++ b/src/config/load/validate_web/timeouts.rs
@@ -43,9 +43,7 @@ pub(super) fn validate(timeouts: &WebTimeoutsConfig) -> Result<()> {
         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",
-        );
+        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(
diff --git a/src/config/tests/load_basic_tests/web_tests.rs b/src/config/tests/load_basic_tests/web_tests.rs
index 28200bb..c12919a 100644
--- a/src/config/tests/load_basic_tests/web_tests.rs
+++ b/src/config/tests/load_basic_tests/web_tests.rs
@@ -50,7 +50,10 @@ fn web_config_builds_canonical_runtime_snapshot() {
     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]);
+    assert_eq!(
+        vhost.profiles[0].carriers.as_ref(),
+        [WebCarrier::HttpsLanes]
+    );
 }
 
 #[test]
@@ -95,11 +98,7 @@ fn web_carrier_array_enables_ordered_negotiation_and_appends_fallback() {
 
 #[test]
 fn web_carriers_reject_true_empty_and_duplicates() {
-    for value in [
-        "true",
-        "[]",
-        "[\"https\", \"https\"]",
-    ] {
+    for value in ["true", "[]", "[\"https\", \"https\"]"] {
         let invalid = WEB_CONFIG.replace(
             "carrier = \"https-lanes\"",
             &format!("carrier = \"https-lanes\"\ncarriers = {value}"),
diff --git a/src/config/types.rs b/src/config/types.rs
index 0b9f907..562f54e 100644
--- a/src/config/types.rs
+++ b/src/config/types.rs
@@ -56,12 +56,12 @@ pub use web::{
     WebCarrierNegotiationAggressiveness, WebConfig, WebDecoyConfig, WebLimitsConfig,
     WebProfileConfig, WebSecretMode, WebTimeoutsConfig, WebVhostConfig,
 };
-#[allow(unused_imports)]
-pub use web_carrier::{WebCarrier, WebCarriers};
 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 2a8fc58..863b533 100644
--- a/src/config/types/web.rs
+++ b/src/config/types/web.rs
@@ -234,8 +234,7 @@ impl Default for WebLimitsConfig {
             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_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(),
@@ -348,8 +347,7 @@ impl Default for WebTimeoutsConfig {
             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_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(),
@@ -436,8 +434,7 @@ impl Default for WebConfig {
             carrier: WebCarrier::default(),
             carriers: WebCarriers::default(),
             carrier_learning: default_web_carrier_learning(),
-            carrier_negotiation_aggressiveness:
-                WebCarrierNegotiationAggressiveness::default(),
+            carrier_negotiation_aggressiveness: WebCarrierNegotiationAggressiveness::default(),
             limits: WebLimitsConfig::default(),
             debug: WebDebugConfig::default(),
             timeouts: WebTimeoutsConfig::default(),
diff --git a/src/web/bridge/tests.rs b/src/web/bridge/tests.rs
index aa674d2..58fa517 100644
--- a/src/web/bridge/tests.rs
+++ b/src/web/bridge/tests.rs
@@ -16,10 +16,7 @@ fn render_page(bootstrap: &str, candidate_count: usize) -> BridgePage {
 
 #[test]
 fn rendered_page_contains_bounded_negotiation_contract() {
-    let page = render_page(
-        "AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA",
-        4,
-    );
+    let page = render_page("AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA", 4);
     assert!(!page.body.contains("__"));
     assert!(!page.body.contains("bridge="));
     assert!(page.body.contains("X-Carrier-Capabilities"));
@@ -40,17 +37,23 @@ fn rendered_page_contains_bounded_negotiation_contract() {
 fn rendered_page_preserves_the_ios_bootstrap_literal() {
     let bootstrap = "BBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBB";
     let page = render_page(bootstrap, 2);
-    assert!(page.body.contains(&format!("const bootstrap=\"{bootstrap}\"")));
+    assert!(
+        page.body
+            .contains(&format!("const bootstrap=\"{bootstrap}\""))
+    );
 }
 
 #[test]
 fn effective_deadline_formula_uses_the_final_checkpoint() {
-    let page = render_page(
-        "CCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCC",
-        3,
+    let page = render_page("CCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCCC", 3);
+    assert!(
+        page.body
+            .contains("negotiatedFinalDeadline=candidateDeadlines[3]")
+    );
+    assert!(
+        page.body
+            .contains("carrierAttempt>=negotiatedCandidateCount?negotiatedFinalDeadline")
     );
-    assert!(page.body.contains("negotiatedFinalDeadline=candidateDeadlines[3]"));
-    assert!(page.body.contains("carrierAttempt>=negotiatedCandidateCount?negotiatedFinalDeadline"));
 }
 
 #[test]
@@ -69,18 +72,19 @@ fn disabled_negotiation_does_not_arm_a_carrier_deadline() {
     assert!(page.body.contains(
         "if(negotiationEnabled){negotiationStartedAt=Date.now();armCarrierDeadline(attemptEpoch)}"
     ));
-    assert!(page.body.contains(
-        "negotiationEnabled?'tproxy-auto-v1.':'tproxy-v1.'"
-    ));
+    assert!(
+        page.body
+            .contains("negotiationEnabled?'tproxy-auto-v1.':'tproxy-v1.'")
+    );
 }
 
 #[test]
 fn retry_and_attempt_state_are_frozen_before_fetch() {
-    let page = render_page(
-        "EEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEE",
-        4,
+    let page = render_page("EEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEEE", 4);
+    assert!(
+        page.body
+            .contains("async function request(path,frozenOptions)")
     );
-    assert!(page.body.contains("async function request(path,frozenOptions)"));
     assert!(!page.body.contains("makeOptions"));
     assert!(
         page.body
@@ -93,14 +97,14 @@ fn retry_and_attempt_state_are_frozen_before_fetch() {
 
 #[test]
 fn ambiguous_commit_is_resolved_before_carrier_advance() {
-    let page = render_page(
-        "FFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFF",
-        4,
-    );
+    let page = render_page("FFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFF", 4);
     assert!(page.body.contains("resolveAttempt(reason,epoch,snapshot)"));
     assert!(page.body.contains(
         "sessionEcho(response,snapshot.attempt,['provisional','committed','healthy'],true)"
     ));
-    assert!(page.body.contains("if(echo.state!=='provisional'){switching=false;fail();return}"));
+    assert!(
+        page.body
+            .contains("if(echo.state!=='provisional'){switching=false;fail();return}")
+    );
     assert!(page.body.contains("const token=cleanupToken||sessionToken"));
 }
diff --git a/src/web/http.rs b/src/web/http.rs
index 74e81aa..c7f1946 100644
--- a/src/web/http.rs
+++ b/src/web/http.rs
@@ -356,11 +356,7 @@ async fn handle_up(
         }
     };
     if let Some(trace) = request_trace(&request) {
-        trace.record_frames(
-            TraceDirection::Request,
-            &body,
-            session.limits(),
-        );
+        trace.record_frames(TraceDirection::Request, &body, session.limits());
     }
     let result = match lane_id {
         Some(lane_id) => session.process_up_lane(lane_id, sequence, &body),
diff --git a/src/web/http/down.rs b/src/web/http/down.rs
index a60fb44..4796aaa 100644
--- a/src/web/http/down.rs
+++ b/src/web/http/down.rs
@@ -89,11 +89,7 @@ pub(super) async fn handle_down(
         }
         Ok(result) => {
             if let Some(trace) = request_trace(&request) {
-                trace.record_frames(
-                    TraceDirection::Response,
-                    &result.body,
-                    session.limits(),
-                );
+                trace.record_frames(TraceDirection::Response, &result.body, session.limits());
             }
             let mut response = full_response(StatusCode::OK, result.body);
             carrier_headers(&mut response);
diff --git a/src/web/http/negotiation_tests.rs b/src/web/http/negotiation_tests.rs
index 47fa1b3..2ce9244 100644
--- a/src/web/http/negotiation_tests.rs
+++ b/src/web/http/negotiation_tests.rs
@@ -130,13 +130,7 @@ async fn metadata_free_native_client_can_use_each_fixed_carrier() {
         let response = request(
             &listener,
             &runtime,
-            create_request_with_headers(
-                &bootstrap,
-                &hello,
-                None,
-                None,
-                NATIVE_USER_AGENT_HEADER,
-            ),
+            create_request_with_headers(&bootstrap, &hello, None, None, NATIVE_USER_AGENT_HEADER),
         )
         .await;
         let (headers, _) = split_response(&response);
@@ -170,13 +164,7 @@ async fn metadata_free_native_client_uses_fallback_when_candidates_are_enabled()
     let response = request(
         &listener,
         &runtime,
-        create_request_with_headers(
-            &bootstrap,
-            &hello,
-            None,
-            None,
-            NATIVE_USER_AGENT_HEADER,
-        ),
+        create_request_with_headers(&bootstrap, &hello, None, None, NATIVE_USER_AGENT_HEADER),
     )
     .await;
     let (headers, _) = split_response(&response);
@@ -209,13 +197,7 @@ async fn explicit_native_capabilities_participate_in_automatic_selection() {
     let response = request(
         &listener,
         &runtime,
-        create_request_with_headers(
-            &bootstrap,
-            &hello,
-            Some(1),
-            None,
-            NATIVE_USER_AGENT_HEADER,
-        ),
+        create_request_with_headers(&bootstrap, &hello, Some(1), None, NATIVE_USER_AGENT_HEADER),
     )
     .await;
     let (headers, _) = split_response(&response);
@@ -261,9 +243,18 @@ async fn negotiation_replays_replaces_and_freezes_after_carrier_commit() {
 
     let replay = request(&listener, &runtime, first_request).await;
     let (replay_headers, _) = split_response(&replay);
-    assert_eq!(response_header(replay_headers, "x-session-token"), first_token);
-    assert_eq!(response_header(replay_headers, "x-carrier-candidate-count"), "3");
-    assert_eq!(response_header(replay_headers, "x-carrier-state"), "provisional");
+    assert_eq!(
+        response_header(replay_headers, "x-session-token"),
+        first_token
+    );
+    assert_eq!(
+        response_header(replay_headers, "x-carrier-candidate-count"),
+        "3"
+    );
+    assert_eq!(
+        response_header(replay_headers, "x-carrier-state"),
+        "provisional"
+    );
 
     let second_request = create_request(&bootstrap, &hello, Some(2), Some("timeout"));
     let second = request(&listener, &runtime, second_request.clone()).await;
@@ -346,11 +337,20 @@ async fn negotiation_replays_replaces_and_freezes_after_carrier_commit() {
     let (third_headers, _) = split_response(&third);
     assert!(third_headers.starts_with(b"HTTP/1.1 409"));
     assert!(optional_response_header(third_headers, "x-session-token").is_none());
-    assert_eq!(response_header(third_headers, "x-carrier-mode"), "https-lanes");
+    assert_eq!(
+        response_header(third_headers, "x-carrier-mode"),
+        "https-lanes"
+    );
     assert_eq!(response_header(third_headers, "x-carrier-attempt"), "2");
-    assert_eq!(response_header(third_headers, "x-carrier-candidate-count"), "3");
+    assert_eq!(
+        response_header(third_headers, "x-carrier-candidate-count"),
+        "3"
+    );
     assert_eq!(response_header(third_headers, "x-carrier-deadline"), "12");
-    assert_eq!(response_header(third_headers, "x-carrier-state"), "committed");
+    assert_eq!(
+        response_header(third_headers, "x-carrier-state"),
+        "committed"
+    );
     assert!(
         runtime
             .get_session(token_hash(&second_token), "proxy.example.com")
@@ -397,7 +397,10 @@ async fn timed_out_attempt_replays_before_successor_own_deadline() {
         response_header(replay_headers, "x-session-token"),
         first_token
     );
-    assert_eq!(response_header(replay_headers, "x-carrier-state"), "provisional");
+    assert_eq!(
+        response_header(replay_headers, "x-carrier-state"),
+        "provisional"
+    );
 
     let second = request(
         &listener,
@@ -453,9 +456,8 @@ async fn https_lane_downlink_can_arrive_before_its_uplink_open() {
     .into_bytes();
     let down_listener = Arc::clone(&listener);
     let down_runtime = Arc::clone(&runtime);
-    let down = tokio::spawn(async move {
-        request(&down_listener, &down_runtime, down_request).await
-    });
+    let down =
+        tokio::spawn(async move { request(&down_listener, &down_runtime, down_request).await });
     tokio::task::yield_now().await;
 
     let open = frame::encode(FrameType::Open, 7, &[]);
@@ -477,10 +479,7 @@ async fn https_lane_downlink_can_arrive_before_its_uplink_open() {
         .unwrap()
         .unwrap();
     let (down_headers, _) = split_response(&down);
-    assert!(
-        down_headers.starts_with(b"HTTP/1.1 200")
-            || down_headers.starts_with(b"HTTP/1.1 204")
-    );
+    assert!(down_headers.starts_with(b"HTTP/1.1 200") || down_headers.starts_with(b"HTTP/1.1 204"));
     assert!(optional_response_header(down_headers, "x-down-cursor").is_some());
 
     runtime.shutdown().await;
diff --git a/src/web/http/request.rs b/src/web/http/request.rs
index 239dea5..4c677da 100644
--- a/src/web/http/request.rs
+++ b/src/web/http/request.rs
@@ -15,9 +15,7 @@ const USER_AGENT_CONTEXT: &[u8] = b"telemt-web-carrier-user-agent-v1\0";
 
 // Canonical host and forwarded-address provenance remain isolated from credentials.
 mod identity;
-pub(super) use identity::{
-    canonical_request_host, carrier_ip_learning_eligible, client_ip,
-};
+pub(super) use identity::{canonical_request_host, carrier_ip_learning_eligible, client_ip};
 
 /// Decodes an exact canonical bridge query without allocating credential strings.
 pub(super) fn bridge_candidate(query: Option<&str>) -> ([u8; 32], bool) {
@@ -205,17 +203,13 @@ fn parse_capabilities(value: &str) -> Option {
 }
 
 fn strict_browser_hint(request: &Request, host: &str) -> bool {
-    single_header(request, header::ORIGIN)
-        .is_some_and(|value| value == format!("https://{host}"))
+    single_header(request, header::ORIGIN).is_some_and(|value| value == format!("https://{host}"))
         && single_header(request, "sec-fetch-site") == Some("same-origin")
         && single_header(request, "sec-fetch-mode") == Some("cors")
         && single_header(request, "sec-fetch-dest") == Some("empty")
 }
 
-fn optional_canonical_u8_header(
-    request: &Request,
-    name: &'static str,
-) -> Option> {
+fn optional_canonical_u8_header(request: &Request, name: &'static str) -> Option> {
     if !request.headers().contains_key(name) {
         return Some(None);
     }
@@ -452,9 +446,11 @@ mod tests {
             .header(header::USER_AGENT, "Native")
             .body(())
             .unwrap();
-        assert!(!carrier_request(&legacy, "proxy.example.com")
-            .unwrap()
-            .is_automatic());
+        assert!(
+            !carrier_request(&legacy, "proxy.example.com")
+                .unwrap()
+                .is_automatic()
+        );
 
         let reordered = Request::builder()
             .header("x-carrier-capabilities", "websocket,https")
@@ -479,10 +475,7 @@ mod tests {
         assert!(!parsed.uses_capabilities());
 
         let automatic = Request::builder()
-            .header(
-                "x-carrier-capabilities",
-                "https,https-lanes",
-            )
+            .header("x-carrier-capabilities", "https,https-lanes")
             .header("x-carrier-attempt", "1")
             .header(
                 header::USER_AGENT,
diff --git a/src/web/http/session.rs b/src/web/http/session.rs
index 2695cc2..cd3aa94 100644
--- a/src/web/http/session.rs
+++ b/src/web/http/session.rs
@@ -152,12 +152,9 @@ pub(super) async fn handle_session(
         }
         Err(ManagerError::Committed) => {
             let mut response = carrier_empty(StatusCode::CONFLICT);
-            if let Some(echo) = runtime.carrier_echo(
-                token_hash,
-                &vhost.host,
-                client_ip,
-                carrier_request,
-            ) {
+            if let Some(echo) =
+                runtime.carrier_echo(token_hash, &vhost.host, client_ip, carrier_request)
+            {
                 response.headers_mut().insert(
                     HeaderName::from_static("x-carrier-mode"),
                     HeaderValue::from_static(echo.carrier.as_str()),
diff --git a/src/web/http/tests.rs b/src/web/http/tests.rs
index d1ee7e6..f1493e8 100644
--- a/src/web/http/tests.rs
+++ b/src/web/http/tests.rs
@@ -134,8 +134,7 @@ fn runtime_config_with_carriers_and_deadlines(
         WebCarriers::Disabled
     };
     config.web.carrier_learning = carrier_learning;
-    config.web.timeouts.carrier_negotiation_deadlines_secs =
-        carrier_negotiation_deadlines_secs;
+    config.web.timeouts.carrier_negotiation_deadlines_secs = carrier_negotiation_deadlines_secs;
     config.web.limits.max_bootstraps_per_ip = 1;
     config.web.timeouts.shutdown_secs = 1;
     config.web.runtime = Some(Arc::new(WebRuntimeConfig {
@@ -260,9 +259,11 @@ async fn https_carrier_bootstraps_and_closes_one_session() {
             .windows(11)
             .any(|value| value == b"bootstrap=\"")
     );
-    assert!(next_root_body
-        .windows(b"const negotiationEnabled=false".len())
-        .any(|value| value == b"const negotiationEnabled=false"));
+    assert!(
+        next_root_body
+            .windows(b"const negotiationEnabled=false".len())
+            .any(|value| value == b"const negotiationEnabled=false")
+    );
 
     let close = format!(
         "DELETE /api/v1/session HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nAuthorization: Bearer {session}\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"
diff --git a/src/web/http/websocket.rs b/src/web/http/websocket.rs
index ff0fea6..762b229 100644
--- a/src/web/http/websocket.rs
+++ b/src/web/http/websocket.rs
@@ -389,21 +389,20 @@ fn parse_upgrade(request: &Request) -> Option {
     {
         return None;
     }
-    let (token, carrier, acknowledge_commit) = if let Some(token) =
-        protocol.strip_prefix("tproxy-auto-v1.")
-    {
-        (token, ParsedCarrier::Multiplex, true)
-    } else if let Some(lane) = protocol.strip_prefix("tproxy-auto-lane-v1.") {
-        let (token, lane_id) = parse_lane_protocol(lane)?;
-        (token, ParsedCarrier::Lane(lane_id), true)
-    } else if let Some(token) = protocol.strip_prefix("tproxy-v1.") {
-        (token, ParsedCarrier::Multiplex, false)
-    } else if let Some(lane) = protocol.strip_prefix("tproxy-lane-v1.") {
-        let (token, lane_id) = parse_lane_protocol(lane)?;
-        (token, ParsedCarrier::Lane(lane_id), false)
-    } else {
-        return None;
-    };
+    let (token, carrier, acknowledge_commit) =
+        if let Some(token) = protocol.strip_prefix("tproxy-auto-v1.") {
+            (token, ParsedCarrier::Multiplex, true)
+        } else if let Some(lane) = protocol.strip_prefix("tproxy-auto-lane-v1.") {
+            let (token, lane_id) = parse_lane_protocol(lane)?;
+            (token, ParsedCarrier::Lane(lane_id), true)
+        } else if let Some(token) = protocol.strip_prefix("tproxy-v1.") {
+            (token, ParsedCarrier::Multiplex, false)
+        } else if let Some(lane) = protocol.strip_prefix("tproxy-lane-v1.") {
+            let (token, lane_id) = parse_lane_protocol(lane)?;
+            (token, ParsedCarrier::Lane(lane_id), false)
+        } else {
+            return None;
+        };
     let raw_token = base64::engine::general_purpose::URL_SAFE_NO_PAD
         .decode(token)
         .ok()?;
@@ -441,10 +440,7 @@ fn parse_lane_protocol(value: &str) -> Option<(&str, u32)> {
     Some((token, lane_id))
 }
 
-fn single_header(
-    request: &Request,
-    name: impl hyper::header::AsHeaderName,
-) -> Option<&str> {
+fn single_header(request: &Request, name: impl hyper::header::AsHeaderName) -> Option<&str> {
     let mut values = request.headers().get_all(name).iter();
     let value = values.next()?.to_str().ok()?;
     values.next().is_none().then_some(value)
diff --git a/src/web/http/websocket/driver.rs b/src/web/http/websocket/driver.rs
index 2658c00..7ff4932 100644
--- a/src/web/http/websocket/driver.rs
+++ b/src/web/http/websocket/driver.rs
@@ -8,12 +8,8 @@ use tokio_tungstenite::tungstenite::protocol::{Message, Role, WebSocketConfig};
 use tokio_util::sync::CancellationToken;
 
 use super::ConnectionIo;
-use crate::web::manager::{
-    WebProcessRuntime, WebSocketBudgetLease, WebSocketConnection,
-};
-use crate::web::session::{
-    WebSession, WebSocketLaneReservation, WebSocketProbeReservation,
-};
+use crate::web::manager::{WebProcessRuntime, WebSocketBudgetLease, WebSocketConnection};
+use crate::web::session::{WebSession, WebSocketLaneReservation, WebSocketProbeReservation};
 use crate::web::trace::{TraceDirection, TraceWebSocketContext};
 
 const READ_BUFFER_BYTES: usize = 64 * 1024;
@@ -126,8 +122,7 @@ async fn run_multiplex(
     let mut next_ping = Instant::now() + liveness_interval;
     let open_deadline =
         Instant::now() + Duration::from_secs(session.timeouts().websocket_open_secs);
-    let backpressure_timeout =
-        Duration::from_secs(session.timeouts().websocket_backpressure_secs);
+    let backpressure_timeout = Duration::from_secs(session.timeouts().websocket_backpressure_secs);
     let write_timeout = Duration::from_secs(session.timeouts().websocket_write_secs);
     let maximum_message = session.limits().carrier_batch_bytes;
     let mut active = false;
@@ -317,13 +312,7 @@ async fn run_multiplex(
                             started,
                         );
                     } else {
-                        send(
-                            socket,
-                            Message::Binary(body),
-                            &cancellation,
-                            write_timeout,
-                        )
-                        .await?;
+                        send(socket, Message::Binary(body), &cancellation, write_timeout).await?;
                     }
                     connection.mark_progress();
                 }
diff --git a/src/web/http/websocket/driver/io.rs b/src/web/http/websocket/driver/io.rs
index 132c040..5d3a350 100644
--- a/src/web/http/websocket/driver/io.rs
+++ b/src/web/http/websocket/driver/io.rs
@@ -24,16 +24,8 @@ pub(super) async fn read_message(
         ready = socket.get_ref().readable() => ready.map_err(|_| ())?,
     }
     if retained_budget.is_none() {
-        *retained_budget = Some(
-            reserve_data(
-                runtime,
-                owner,
-                maximum,
-                cancellation,
-                backpressure_timeout,
-            )
-            .await?,
-        );
+        *retained_budget =
+            Some(reserve_data(runtime, owner, maximum, cancellation, backpressure_timeout).await?);
     }
     let message = tokio::select! {
         _ = cancellation.cancelled() => return Err(()),
diff --git a/src/web/http/websocket/driver/lane.rs b/src/web/http/websocket/driver/lane.rs
index fc24f93..18fc899 100644
--- a/src/web/http/websocket/driver/lane.rs
+++ b/src/web/http/websocket/driver/lane.rs
@@ -7,9 +7,7 @@ use tokio_util::sync::CancellationToken;
 
 use super::CarrierSocket;
 use super::io::{flush, process_lane, read_message, record_message, reserve_data, send};
-use crate::web::manager::{
-    WebProcessRuntime, WebSocketBudgetLease, WebSocketConnection,
-};
+use crate::web::manager::{WebProcessRuntime, WebSocketBudgetLease, WebSocketConnection};
 use crate::web::session::{WebSession, WebSocketLaneReservation};
 use crate::web::trace::{TraceDirection, TraceWebSocketContext};
 
@@ -32,8 +30,7 @@ pub(super) async fn run_lane(
     let mut next_ping = Instant::now() + liveness_interval;
     let open_deadline =
         Instant::now() + Duration::from_secs(session.timeouts().websocket_open_secs);
-    let backpressure_timeout =
-        Duration::from_secs(session.timeouts().websocket_backpressure_secs);
+    let backpressure_timeout = Duration::from_secs(session.timeouts().websocket_backpressure_secs);
     let write_timeout = Duration::from_secs(session.timeouts().websocket_write_secs);
     let maximum_message = session.limits().carrier_batch_bytes;
     let mut active = false;
@@ -232,13 +229,7 @@ pub(super) async fn run_lane(
                             started,
                         );
                     } else {
-                        send(
-                            socket,
-                            Message::Binary(body),
-                            &cancellation,
-                            write_timeout,
-                        )
-                        .await?;
+                        send(socket, Message::Binary(body), &cancellation, write_timeout).await?;
                     }
                     connection.mark_progress();
                 }
diff --git a/src/web/http/websocket/tests.rs b/src/web/http/websocket/tests.rs
index fac2c98..029d14e 100644
--- a/src/web/http/websocket/tests.rs
+++ b/src/web/http/websocket/tests.rs
@@ -452,8 +452,7 @@ async fn failed_automatic_multiplex_socket_remains_supersedable() {
         Arc::from([WebCarrier::Websocket, WebCarrier::Https]),
     );
     let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
-    let (bootstrap_hash, hello, session, session_hash) =
-        create_automatic_session(&live.runtime);
+    let (bootstrap_hash, hello, session, session_hash) = create_automatic_session(&live.runtime);
     let protocol = format!("tproxy-auto-v1.{session}");
     let mut socket = upgrade(&listener, &live.runtime, &protocol).await;
     socket.close(None).await.unwrap();
diff --git a/src/web/manager.rs b/src/web/manager.rs
index f4f5fb2..3730b6f 100644
--- a/src/web/manager.rs
+++ b/src/web/manager.rs
@@ -39,10 +39,10 @@ mod budget;
 mod websocket;
 pub(crate) use budget::WebSocketBudgetLease;
 use budget::{WebDataBudget, WebSocketBudgetClass};
-use state::{ManagerState, StreamAdmissionState};
 pub(crate) use negotiation::{
     CarrierCapabilities, CarrierClientClass, CarrierFailure, CarrierLearningContext, CarrierRequest,
 };
+use state::{ManagerState, StreamAdmissionState};
 pub(crate) use websocket::{WebSocketConnection, WebSocketKind};
 
 const TOKEN_BYTES: usize = 32;
@@ -327,10 +327,7 @@ impl WebProcessRuntime {
     }
 
     /// Reserves transient bytes while one downlink batch replaces queued frames.
-    pub(crate) fn try_downlink_staging_budget(
-        &self,
-        bytes: usize,
-    ) -> Option {
+    pub(crate) fn try_downlink_staging_budget(&self, bytes: usize) -> Option {
         let bytes = u32::try_from(bytes).ok()?;
         let permit = Arc::clone(&self.body_bytes)
             .try_acquire_many_owned(bytes)
diff --git a/src/web/manager/carrier_learning.rs b/src/web/manager/carrier_learning.rs
index d36b93b..991c92d 100644
--- a/src/web/manager/carrier_learning.rs
+++ b/src/web/manager/carrier_learning.rs
@@ -4,8 +4,8 @@ use std::time::{Duration, Instant};
 
 use sha2::{Digest, Sha256};
 
-use super::negotiation::{CarrierClientClass, CarrierLearningContext};
 use super::ProfileKey;
+use super::negotiation::{CarrierClientClass, CarrierLearningContext};
 use crate::config::{WebCarrier, WebCarrierNegotiationAggressiveness};
 
 const PROFILE_WEIGHT: i16 = 32;
@@ -85,9 +85,7 @@ impl Evidence {
             }
             aggregate.outcomes = aggregate.outcomes.saturating_add(bucket.outcomes);
             for (score, value) in aggregate.scores.iter_mut().zip(bucket.scores) {
-                *score = score
-                    .saturating_add(value)
-                    .clamp(SCORE_MIN, SCORE_MAX);
+                *score = score.saturating_add(value).clamp(SCORE_MIN, SCORE_MAX);
             }
             for cohort in bucket.cohorts.iter().flatten() {
                 if !aggregate.cohorts.contains(&Some(*cohort))
@@ -262,10 +260,9 @@ impl CarrierLearning {
         let user_agent_ready = user_agent
             .as_ref()
             .is_some_and(|entry| entry.outcomes >= thresholds.user_agent);
-        let ip_ready = thresholds.ip.is_some_and(|minimum| {
-            ip.as_ref()
-                .is_some_and(|entry| entry.outcomes >= minimum)
-        });
+        let ip_ready = thresholds
+            .ip
+            .is_some_and(|minimum| ip.as_ref().is_some_and(|entry| entry.outcomes >= minimum));
         let mut scores = [0i16; 4];
         for carrier in WebCarrier::ALL {
             let index = carrier.index();
@@ -274,11 +271,9 @@ impl CarrierLearning {
                     * PROFILE_WEIGHT;
             }
             if user_agent_ready {
-                scores[index] += i16::from(
-                    user_agent
-                        .as_ref()
-                        .map_or(0, |value| value.scores[index]),
-                ) * USER_AGENT_WEIGHT;
+                scores[index] +=
+                    i16::from(user_agent.as_ref().map_or(0, |value| value.scores[index]))
+                        * USER_AGENT_WEIGHT;
             }
             if ip_ready {
                 scores[index] +=
@@ -362,7 +357,11 @@ impl CarrierLearning {
             if !current {
                 continue;
             }
-            if self.entries.get(&key).is_some_and(|entry| entry.is_live(slot)) {
+            if self
+                .entries
+                .get(&key)
+                .is_some_and(|entry| entry.is_live(slot))
+            {
                 self.insertion_order.push_back((key, sequence));
             } else {
                 self.entries.remove(&key);
@@ -411,8 +410,7 @@ impl CarrierLearning {
         let Some(insertion_sequence) = self.next_insertion_sequence() else {
             return;
         };
-        self.entries
-            .insert(key, Evidence::new(insertion_sequence));
+        self.entries.insert(key, Evidence::new(insertion_sequence));
         self.insertion_order.push_back((key, insertion_sequence));
         if let Some(entry) = self.entries.get_mut(&key) {
             entry.update(slot, deltas, cohort);
diff --git a/src/web/manager/negotiation.rs b/src/web/manager/negotiation.rs
index 0267eca..13b3321 100644
--- a/src/web/manager/negotiation.rs
+++ b/src/web/manager/negotiation.rs
@@ -189,9 +189,7 @@ impl CarrierRequest {
 
     /// Checks the complete idempotent identity of one exact attempt request.
     pub(crate) fn matches_attempt(self, other: Self) -> bool {
-        self.matches_client(other)
-            && self.attempt == other.attempt
-            && self.failure == other.failure
+        self.matches_client(other) && self.attempt == other.attempt && self.failure == other.failure
     }
 
     fn capabilities_bits(self) -> Option {
@@ -252,7 +250,10 @@ mod tests {
     #[test]
     fn invalid_candidate_or_attempt_counts_have_no_deadline_slot() {
         for (candidate_count, attempt) in [(0, 1), (5, 1), (1, 0), (1, 2), (3, 4)] {
-            assert_eq!(carrier_attempt_deadline_index(candidate_count, attempt), None);
+            assert_eq!(
+                carrier_attempt_deadline_index(candidate_count, attempt),
+                None
+            );
         }
     }
 }
diff --git a/src/web/manager/session_creation.rs b/src/web/manager/session_creation.rs
index 24cdb03..6207a2e 100644
--- a/src/web/manager/session_creation.rs
+++ b/src/web/manager/session_creation.rs
@@ -7,14 +7,15 @@ use sha2::{Digest, Sha256};
 use subtle::ConstantTimeEq;
 use zeroize::Zeroizing;
 
+use super::negotiation::carrier_attempt_deadline_index;
+use super::session_admission::admit_initial;
 use super::state::{
     CarrierChainPhase, decrement_map, matching_profile, new_unique_token, profile_key,
     remember_closed_token_locked, remove_expired_locked,
 };
-use super::negotiation::carrier_attempt_deadline_index;
-use super::session_admission::admit_initial;
 use super::{
-    CarrierLearningContext, CarrierRequest, CreateResult, ManagerError, TokenHash, WebProcessRuntime,
+    CarrierLearningContext, CarrierRequest, CreateResult, ManagerError, TokenHash,
+    WebProcessRuntime,
 };
 use crate::config::{WebCarrier, WebRuntimeProfile};
 use crate::web::frame;
@@ -62,7 +63,9 @@ impl WebProcessRuntime {
             return Err(ManagerError::Authentication);
         }
         if entry.used {
-            if entry.carrier_deadline_at.is_some_and(|deadline| now >= deadline)
+            if entry
+                .carrier_deadline_at
+                .is_some_and(|deadline| now >= deadline)
                 && entry
                     .session
                     .as_ref()
@@ -112,9 +115,8 @@ impl WebProcessRuntime {
                     token: entry.session_token.as_str().to_owned(),
                     carrier: session.carrier(),
                     attempt: carrier_request.attempt(),
-                    candidate_count: automatic.then(|| {
-                        u8::try_from(entry.carrier_candidates.len()).unwrap_or(4)
-                    }),
+                    candidate_count: automatic
+                        .then(|| u8::try_from(entry.carrier_candidates.len()).unwrap_or(4)),
                     deadline_secs: automatic
                         .then_some(entry.profile.carrier_negotiation_deadlines_secs[3]),
                     carrier_state: automatic.then_some(carrier_state),
@@ -139,14 +141,16 @@ impl WebProcessRuntime {
                     CarrierChainPhase::CommittedPendingHealth | CarrierChainPhase::Healthy
                 )
             {
-                return Err(if matches!(
-                    entry.carrier_phase,
-                    CarrierChainPhase::CommittedPendingHealth | CarrierChainPhase::Healthy
-                ) {
-                    ManagerError::Committed
-                } else {
-                    ManagerError::Protocol
-                });
+                return Err(
+                    if matches!(
+                        entry.carrier_phase,
+                        CarrierChainPhase::CommittedPendingHealth | CarrierChainPhase::Healthy
+                    ) {
+                        ManagerError::Committed
+                    } else {
+                        ManagerError::Protocol
+                    },
+                );
             }
             let Some(carrier) = entry
                 .carrier_candidates
@@ -155,8 +159,8 @@ impl WebProcessRuntime {
             else {
                 return Err(ManagerError::Protocol);
             };
-            let candidate_count = u8::try_from(entry.carrier_candidates.len())
-                .map_err(|_| ManagerError::Protocol)?;
+            let candidate_count =
+                u8::try_from(entry.carrier_candidates.len()).map_err(|_| ManagerError::Protocol)?;
             let deadline_index = carrier_attempt_deadline_index(candidate_count, next_attempt)
                 .ok_or(ManagerError::Protocol)?;
             if entry.carrier_started_at.is_some_and(|started| {
@@ -211,8 +215,8 @@ impl WebProcessRuntime {
         if carrier_request.is_automatic() && !profile.carrier_negotiation_enabled {
             return Err(ManagerError::Protocol);
         }
-        let capability_selection = carrier_request.uses_capabilities()
-            && profile.carrier_negotiation_enabled;
+        let capability_selection =
+            carrier_request.uses_capabilities() && profile.carrier_negotiation_enabled;
         let learning_policy = (
             config.web.carrier_negotiation_enabled() && config.web.carrier_learning,
             config.web.carrier_negotiation_aggressiveness,
@@ -222,11 +226,9 @@ impl WebProcessRuntime {
             && profile.carrier_learning
         {
             let learning = self.learning.lock();
-            if let Some(epoch) = learning.epoch_for_policy(
-                learning_policy.0,
-                learning_policy.1,
-                learning_policy.2,
-            ) {
+            if let Some(epoch) =
+                learning.epoch_for_policy(learning_policy.0, learning_policy.1, learning_policy.2)
+            {
                 let (candidates, scores) = learning.rank(
                     now,
                     &profile.carriers,
@@ -259,8 +261,7 @@ impl WebProcessRuntime {
                 [0; 4],
                 None,
             )
-        } else if carrier_request.uses_capabilities()
-            && !carrier_request.supports(profile.carrier)
+        } else if carrier_request.uses_capabilities() && !carrier_request.supports(profile.carrier)
         {
             return Err(ManagerError::Protocol);
         } else {
@@ -276,9 +277,9 @@ impl WebProcessRuntime {
             self.limit_hits.fetch_add(1, Ordering::Relaxed);
             return Err(ManagerError::Limit);
         };
-        let carrier_deadline_at = carrier_request.is_automatic().then_some(
-            now + Duration::from_secs(profile.carrier_negotiation_deadlines_secs[3]),
-        );
+        let carrier_deadline_at = carrier_request
+            .is_automatic()
+            .then_some(now + Duration::from_secs(profile.carrier_negotiation_deadlines_secs[3]));
         let learning_context = learning_epoch.map(|epoch| CarrierLearningContext {
             profile_key,
             client_ip,
@@ -336,9 +337,7 @@ impl WebProcessRuntime {
             token: session_token,
             carrier,
             attempt: carrier_request.attempt(),
-            candidate_count: carrier_request
-                .is_automatic()
-                .then_some(candidate_count),
+            candidate_count: carrier_request.is_automatic().then_some(candidate_count),
             deadline_secs: carrier_request
                 .is_automatic()
                 .then_some(profile.carrier_negotiation_deadlines_secs[3]),
@@ -494,9 +493,7 @@ impl WebProcessRuntime {
             token: session_token,
             carrier: replacement.carrier,
             attempt: Some(replacement.attempt),
-            candidate_count: Some(
-                u8::try_from(entry.carrier_candidates.len()).unwrap_or(4),
-            ),
+            candidate_count: Some(u8::try_from(entry.carrier_candidates.len()).unwrap_or(4)),
             deadline_secs: Some(entry.profile.carrier_negotiation_deadlines_secs[3]),
             carrier_state: Some(CarrierChainPhase::Provisional.as_str()),
         };
@@ -512,7 +509,10 @@ impl WebProcessRuntime {
             replacement.old_session.carrier(),
             replacement.attempt - 1,
             replacement.scores,
-            replacement.request.failure().map(|failure| failure.as_str()),
+            replacement
+                .request
+                .failure()
+                .map(|failure| failure.as_str()),
         );
         self.trace.record_carrier_lifecycle(
             client_ip,
@@ -522,7 +522,10 @@ impl WebProcessRuntime {
             replacement.old_session.carrier(),
             replacement.attempt - 1,
             replacement.scores,
-            replacement.request.failure().map(|failure| failure.as_str()),
+            replacement
+                .request
+                .failure()
+                .map(|failure| failure.as_str()),
         );
         self.trace.record_carrier_lifecycle(
             client_ip,
@@ -540,7 +543,10 @@ impl WebProcessRuntime {
             identity,
             TraceLifecycleEvent::SessionCreated,
             None,
-            replacement.request.failure().map(|failure| failure.as_str()),
+            replacement
+                .request
+                .failure()
+                .map(|failure| failure.as_str()),
         );
         Ok(result)
     }
diff --git a/src/web/manager/websocket.rs b/src/web/manager/websocket.rs
index 78042cd..cee8df2 100644
--- a/src/web/manager/websocket.rs
+++ b/src/web/manager/websocket.rs
@@ -357,9 +357,7 @@ fn select_victim(
     let requester_usage = runtime.data_budget.owner_usage(owner);
     let now = runtime.websocket_tick();
     let mut registry = runtime.websockets.lock();
-    if claim
-        && registry.evictions_in_flight >= runtime.limits.max_websocket_evictions_in_flight
-    {
+    if claim && registry.evictions_in_flight >= runtime.limits.max_websocket_evictions_in_flight {
         return None;
     }
     let selected = registry
@@ -408,9 +406,7 @@ fn select_pressure_victim(
     claim: bool,
 ) -> Option> {
     let mut registry = runtime.websockets.lock();
-    if claim
-        && registry.evictions_in_flight >= runtime.limits.max_websocket_evictions_in_flight
-    {
+    if claim && registry.evictions_in_flight >= runtime.limits.max_websocket_evictions_in_flight {
         return None;
     }
     let selected = registry
@@ -447,8 +443,7 @@ fn claim_stale_victims(runtime: &WebProcessRuntime, now: u64) -> Vec= dead_after(entry)
+            now.saturating_sub(entry.last_peer_tick.load(Ordering::Acquire)) >= dead_after(entry)
         })
         .take(available)
         .cloned()
diff --git a/src/web/session.rs b/src/web/session.rs
index ed92b64..1db78b8 100644
--- a/src/web/session.rs
+++ b/src/web/session.rs
@@ -482,7 +482,6 @@ impl WebSession {
             .upgrade()
             .map(|manager| manager.budget_notify())
     }
-
 }
 
 fn inbound_queue_cost(queue: &VecDeque) -> (usize, usize) {
diff --git a/src/web/session/downlink.rs b/src/web/session/downlink.rs
index 94a3492..13a749d 100644
--- a/src/web/session/downlink.rs
+++ b/src/web/session/downlink.rs
@@ -156,10 +156,8 @@ impl WebSession {
         let fits = if control {
             bytes <= self.limits.control_bytes_per_session
                 && items <= item_reserve
-                && pending_bytes
-                    <= self.limits.pending_bytes_per_session.saturating_sub(bytes)
-                && pending_items
-                    <= self.limits.pending_items_per_session.saturating_sub(items)
+                && pending_bytes <= self.limits.pending_bytes_per_session.saturating_sub(bytes)
+                && pending_items <= self.limits.pending_items_per_session.saturating_sub(items)
                 && pending_control_bytes
                     <= self.limits.control_bytes_per_session.saturating_sub(bytes)
                 && pending_control_items <= item_reserve.saturating_sub(items)
diff --git a/src/web/session/downlink_tests.rs b/src/web/session/downlink_tests.rs
index 6f0b3a8..b1b229c 100644
--- a/src/web/session/downlink_tests.rs
+++ b/src/web/session/downlink_tests.rs
@@ -5,8 +5,7 @@ use std::sync::Arc;
 use arc_swap::ArcSwap;
 
 use crate::config::{
-    ProxyConfig, WebCarrier, WebLimitsConfig, WebRuntimeProfile, WebSecretMode,
-    WebTimeoutsConfig,
+    ProxyConfig, WebCarrier, WebLimitsConfig, WebRuntimeProfile, WebSecretMode, WebTimeoutsConfig,
 };
 use crate::maestro::generation::test_runtime_generation;
 use crate::web::manager::WebProcessRuntime;
diff --git a/src/web/session/lane_uplink.rs b/src/web/session/lane_uplink.rs
index 70f1806..157c601 100644
--- a/src/web/session/lane_uplink.rs
+++ b/src/web/session/lane_uplink.rs
@@ -56,8 +56,7 @@ impl WebSession {
                     && only_late_frames(&frames)
                 {
                     return if self.automatic_carrier
-                        && state.negotiation_phase
-                            != super::SessionNegotiationPhase::Committed
+                        && state.negotiation_phase != super::SessionNegotiationPhase::Committed
                     {
                         Err(ManagerError::Backpressure)
                     } else {
@@ -154,8 +153,7 @@ impl WebSession {
                 }
             }
             if applied {
-                (committed, healthy) =
-                    self.record_uplink_progress_locked(&mut state, progress);
+                (committed, healthy) = self.record_uplink_progress_locked(&mut state, progress);
             }
             applied.then_some(sequence).ok_or(ManagerError::Closed)
         };
diff --git a/src/web/session/lanes.rs b/src/web/session/lanes.rs
index 99e5a55..7755d66 100644
--- a/src/web/session/lanes.rs
+++ b/src/web/session/lanes.rs
@@ -119,8 +119,7 @@ impl WebSession {
                         return Err(ManagerError::Closed);
                     }
                     let carrier_health_eligible = lane_id != 0
-                        && state.negotiation_phase
-                            == super::SessionNegotiationPhase::Committed;
+                        && state.negotiation_phase == super::SessionNegotiationPhase::Committed;
                     let Some(lane) = state.carrier_lanes.get_mut(&lane_id) else {
                         return Ok(PollResult {
                             body: Bytes::new(),
@@ -231,11 +230,7 @@ impl WebSession {
         }
     }
 
-    async fn wait_for_lane_open(
-        &self,
-        lane_id: u32,
-        cursor: u64,
-    ) -> Result {
+    async fn wait_for_lane_open(&self, lane_id: u32, cursor: u64) -> Result {
         let wait = {
             let mut state = self.state.lock();
             if state.closed {
@@ -355,10 +350,10 @@ impl WebSession {
                 let resident = lane.resident.snapshot();
                 payload.len() > self.limits.pending_bytes_per_lane
                     || lane.pending_bytes.saturating_add(resident.data_bytes)
-                    > self
-                        .limits
-                        .pending_bytes_per_lane
-                        .saturating_sub(payload.len())
+                        > self
+                            .limits
+                            .pending_bytes_per_lane
+                            .saturating_sub(payload.len())
             }) {
                 return false;
             }
diff --git a/src/web/session/lanes/tests.rs b/src/web/session/lanes/tests.rs
index 6de0f4c..4adaf7d 100644
--- a/src/web/session/lanes/tests.rs
+++ b/src/web/session/lanes/tests.rs
@@ -6,8 +6,7 @@ use bytes::BytesMut;
 
 use super::*;
 use crate::config::{
-    ProxyConfig, WebCarrier, WebLimitsConfig, WebRuntimeProfile, WebSecretMode,
-    WebTimeoutsConfig,
+    ProxyConfig, WebCarrier, WebLimitsConfig, WebRuntimeProfile, WebSecretMode, WebTimeoutsConfig,
 };
 use crate::maestro::generation::test_runtime_generation;
 use crate::web::manager::WebProcessRuntime;
@@ -187,7 +186,9 @@ fn lane_uplink_sequences_are_independent_and_exactly_once() {
     {
         let mut state = session.state.lock();
         for lane_id in [51, 52] {
-            state.carrier_lanes.insert(lane_id, CarrierLane::new(u64::from(lane_id)));
+            state
+                .carrier_lanes
+                .insert(lane_id, CarrierLane::new(u64::from(lane_id)));
             state.closed_streams.insert(lane_id);
         }
     }
@@ -274,11 +275,8 @@ fn tombstone_eviction_releases_lane_budget_and_accepts_late_frames() {
 
 #[test]
 fn automatic_lane_does_not_ack_a_missing_lane_without_real_progress() {
-    let session = new_session_with_automatic(
-        WebLimitsConfig::default(),
-        std::sync::Weak::new(),
-        true,
-    );
+    let session =
+        new_session_with_automatic(WebLimitsConfig::default(), std::sync::Weak::new(), true);
     let late = frame::encode(FrameType::Data, 7, b"late");
 
     assert_eq!(
diff --git a/src/web/session/lifecycle.rs b/src/web/session/lifecycle.rs
index 630320d..bde6f1a 100644
--- a/src/web/session/lifecycle.rs
+++ b/src/web/session/lifecycle.rs
@@ -138,12 +138,7 @@ impl WebSession {
         if let Some(batch) = state.unacked.take() {
             batch.lease.detach();
             self.release_local_locked(&mut state, batch.data_bytes, batch.data_items, false);
-            self.release_local_locked(
-                &mut state,
-                batch.control_bytes,
-                batch.control_items,
-                true,
-            );
+            self.release_local_locked(&mut state, batch.control_bytes, batch.control_items, true);
         }
         let mut lane_data_bytes = 0usize;
         let mut lane_data_items = 0usize;
@@ -160,12 +155,7 @@ impl WebSession {
             }
         }
         self.release_local_locked(&mut state, lane_data_bytes, lane_data_items, false);
-        self.release_local_locked(
-            &mut state,
-            lane_control_bytes,
-            lane_control_items,
-            true,
-        );
+        self.release_local_locked(&mut state, lane_control_bytes, lane_control_items, true);
         state.carrier_lanes.clear();
         let control_bytes = state.pending_control_bytes;
         let control_items = state.pending_control_items;
diff --git a/src/web/session/negotiation.rs b/src/web/session/negotiation.rs
index 8bd1836..26e10d6 100644
--- a/src/web/session/negotiation.rs
+++ b/src/web/session/negotiation.rs
@@ -93,9 +93,8 @@ impl WebSession {
         let now = Instant::now();
         let committed = if state.negotiation_phase == SessionNegotiationPhase::Uncommitted {
             state.negotiation_phase = SessionNegotiationPhase::Committed;
-            state.carrier_health_due_at = Some(
-                now + Duration::from_secs(self.timeouts.carrier_health_secs),
-            );
+            state.carrier_health_due_at =
+                Some(now + Duration::from_secs(self.timeouts.carrier_health_secs));
             true
         } else {
             false
diff --git a/src/web/session/resident.rs b/src/web/session/resident.rs
index e6e59f9..155e959 100644
--- a/src/web/session/resident.rs
+++ b/src/web/session/resident.rs
@@ -43,8 +43,10 @@ impl ResidentCounters {
     }
 
     fn add(&self, counts: PendingCounts) {
-        self.data_bytes.fetch_add(counts.data_bytes, Ordering::AcqRel);
-        self.data_items.fetch_add(counts.data_items, Ordering::AcqRel);
+        self.data_bytes
+            .fetch_add(counts.data_bytes, Ordering::AcqRel);
+        self.data_items
+            .fetch_add(counts.data_items, Ordering::AcqRel);
         self.control_bytes
             .fetch_add(counts.control_bytes, Ordering::AcqRel);
         self.control_items
@@ -52,8 +54,10 @@ impl ResidentCounters {
     }
 
     fn remove(&self, counts: PendingCounts) {
-        self.data_bytes.fetch_sub(counts.data_bytes, Ordering::AcqRel);
-        self.data_items.fetch_sub(counts.data_items, Ordering::AcqRel);
+        self.data_bytes
+            .fetch_sub(counts.data_bytes, Ordering::AcqRel);
+        self.data_items
+            .fetch_sub(counts.data_items, Ordering::AcqRel);
         self.control_bytes
             .fetch_sub(counts.control_bytes, Ordering::AcqRel);
         self.control_items
diff --git a/src/web/session/uplink.rs b/src/web/session/uplink.rs
index c268ef7..ca4b824 100644
--- a/src/web/session/uplink.rs
+++ b/src/web/session/uplink.rs
@@ -139,8 +139,7 @@ impl WebSession {
             } else {
                 state.last_up_sequence = sequence;
                 state.last_up_digest = digest;
-                (committed, healthy) =
-                    self.record_uplink_progress_locked(&mut state, progress);
+                (committed, healthy) = self.record_uplink_progress_locked(&mut state, progress);
                 Ok((sequence, progress.any()))
             }
         };
@@ -532,7 +531,10 @@ mod tests {
         let session = session_with_automatic(true);
         let body = frame::encode(FrameType::Pong, 0, &[]);
 
-        assert_eq!(session.process_up(1, &body), Err(ManagerError::Backpressure));
+        assert_eq!(
+            session.process_up(1, &body),
+            Err(ManagerError::Backpressure)
+        );
         assert!(!session.is_carrier_committed());
         assert!(!session.state.lock().closed);
     }
diff --git a/src/web/session/websocket.rs b/src/web/session/websocket.rs
index cd6299c..3ef0edd 100644
--- a/src/web/session/websocket.rs
+++ b/src/web/session/websocket.rs
@@ -273,8 +273,7 @@ impl WebSession {
             }
             state.last_activity = Instant::now();
             if applied {
-                (committed, healthy) =
-                    self.record_uplink_progress_locked(&mut state, progress);
+                (committed, healthy) = self.record_uplink_progress_locked(&mut state, progress);
             }
             applied
                 .then_some(progress.any())
diff --git a/src/web/trace/sanitize.rs b/src/web/trace/sanitize.rs
index de7999b..1396b63 100644
--- a/src/web/trace/sanitize.rs
+++ b/src/web/trace/sanitize.rs
@@ -49,7 +49,7 @@ pub(super) fn response_dynamic_bytes(
     } else {
         0
     }
-        .saturating_add(sensitive_value_bytes(response.headers(), None))
+    .saturating_add(sensitive_value_bytes(response.headers(), None))
 }
 
 /// Copies header names and only allowlisted bounded values.