diff --git a/src/web/http/down.rs b/src/web/http/down.rs index e24a125..7e45ec4 100644 --- a/src/web/http/down.rs +++ b/src/web/http/down.rs @@ -57,7 +57,8 @@ pub(super) async fn handle_down( return serve_decoy(request, vhost, true, &runtime).await; } let _lane_poll = if lane_id.is_some() { - let Some(permit) = runtime.try_lane_poll() else { + let auxiliary = lane_id.is_some_and(|lane_id| session.lane_poll_is_auxiliary(lane_id)); + let Some(permit) = runtime.try_lane_poll(auxiliary) else { return service_unavailable(); }; Some(permit) diff --git a/src/web/manager.rs b/src/web/manager.rs index 16c50ed..1efce7a 100644 --- a/src/web/manager.rs +++ b/src/web/manager.rs @@ -107,6 +107,7 @@ pub(crate) struct WebProcessRuntime { http_connections: Arc, http_handlers: Arc, lane_polls: Arc, + lane_aux_polls: Arc, body_readers: Arc, body_bytes: Arc, stream_handshakes: Arc, @@ -146,12 +147,17 @@ impl WebProcessRuntime { let websocket_connections = limits .max_http_connections .saturating_sub(limits.websocket_http_connection_reserve); + let lane_poll_limit = limits.max_http_handlers / 2; + let lane_aux_poll_limit = (lane_poll_limit / 2).max(1); let runtime = Arc::new(Self { active_runtime, trace, http_connections: Arc::new(Semaphore::new(limits.max_http_connections)), http_handlers: Arc::new(Semaphore::new(limits.max_http_handlers)), - lane_polls: Arc::new(Semaphore::new((limits.max_http_handlers / 2).max(1))), + lane_polls: Arc::new(Semaphore::new( + lane_poll_limit.saturating_sub(lane_aux_poll_limit), + )), + lane_aux_polls: Arc::new(Semaphore::new(lane_aux_poll_limit)), body_readers: Arc::new(Semaphore::new(limits.max_body_readers)), body_bytes: Arc::new(Semaphore::new(limits.max_body_bytes_global)), stream_handshakes: Arc::new(Semaphore::new(limits.max_stream_handshakes)), @@ -225,8 +231,13 @@ impl WebProcessRuntime { } /// Reserves one parked lane poll without exhausting all HTTP handlers. - pub(crate) fn try_lane_poll(&self) -> Option { - let permit = Arc::clone(&self.lane_polls).try_acquire_owned().ok(); + pub(crate) fn try_lane_poll(&self, auxiliary: bool) -> Option { + let slots = if auxiliary { + &self.lane_aux_polls + } else { + &self.lane_polls + }; + let permit = Arc::clone(slots).try_acquire_owned().ok(); if permit.is_none() { self.record_limit_hit(); } diff --git a/src/web/manager/lifecycle.rs b/src/web/manager/lifecycle.rs index a347f3c..f1bda6e 100644 --- a/src/web/manager/lifecycle.rs +++ b/src/web/manager/lifecycle.rs @@ -119,8 +119,8 @@ impl WebProcessRuntime { remove_expired_locked(&mut state, now); state.sessions.values().cloned().collect::>() }; - for session in sessions.into_iter().filter(|session| session.is_idle(now)) { - session.close(); + for session in sessions { + session.close_if_due(now); } } } diff --git a/src/web/manager/session_creation.rs b/src/web/manager/session_creation.rs index a07fa03..d181b65 100644 --- a/src/web/manager/session_creation.rs +++ b/src/web/manager/session_creation.rs @@ -354,7 +354,9 @@ impl WebProcessRuntime { let identity = session.trace_identity(); let old_identity = replacement.old_session.trace_identity(); drop(state); - replacement.old_session.finish_carrier_supersede(); + if replacement.old_session.finish_carrier_supersede() { + session.close(); + } if let Some(context) = learning_context { self.record_carrier_outcome(context, replacement.old_session.carrier(), false); } diff --git a/src/web/session.rs b/src/web/session.rs index ae93427..ba4a02a 100644 --- a/src/web/session.rs +++ b/src/web/session.rs @@ -124,6 +124,8 @@ struct SessionState { last_up_sequence: u64, last_up_digest: TokenHash, carrier_lanes: HashMap, + lane_open_claims: HashSet, + lane_open_waits: usize, next_lane_instance: u64, websocket_lane_reservations: HashMap, pending_bytes: usize, @@ -160,6 +162,7 @@ pub(crate) struct WebSession { timeouts: WebTimeoutsConfig, state: Mutex, down_notify: Arc, + lane_open_notify: Arc, cancel: CancellationToken, tasks_live: AtomicUsize, tasks_done: Arc, @@ -228,6 +231,8 @@ impl WebSession { last_up_sequence: 0, last_up_digest: [0; 32], carrier_lanes, + lane_open_claims: HashSet::new(), + lane_open_waits: 0, next_lane_instance, websocket_lane_reservations: HashMap::new(), pending_bytes: 0, @@ -240,6 +245,7 @@ impl WebSession { closed: false, }), down_notify: Arc::new(Notify::new()), + lane_open_notify: Arc::new(Notify::new()), cancel: CancellationToken::new(), tasks_live: AtomicUsize::new(0), tasks_done: Arc::new(Notify::new()), diff --git a/src/web/session/lifecycle.rs b/src/web/session/lifecycle.rs index b6e448c..0add9a2 100644 --- a/src/web/session/lifecycle.rs +++ b/src/web/session/lifecycle.rs @@ -13,7 +13,7 @@ struct ReleasedQueues { impl WebSession { /// Closes carrier state while relay tasks retain their admission until exit. pub(crate) fn close(&self) { - let Some(released) = self.begin_close(false) else { + let Some(released) = self.begin_close(false, None) else { return; }; self.finish_close(released, false); @@ -51,11 +51,13 @@ impl WebSession { } /// Completes manager-owned replacement without unregistering the old session twice. - pub(crate) fn finish_carrier_supersede(&self) { - let Some(released) = self.begin_close(true) else { - return; + pub(crate) fn finish_carrier_supersede(&self) -> bool { + let close_requested = self.state.lock().close_requested; + let Some(released) = self.begin_close(true, None) else { + return close_requested; }; self.finish_close(released, true); + close_requested } /// Waits for all logical-stream tasks after admission has closed. @@ -69,20 +71,27 @@ impl WebSession { } } - /// Returns whether reconnect grace elapsed without activity. - pub(crate) fn is_idle(&self, now: Instant) -> bool { - let state = self.state.lock(); - !state.closed - && state.negotiation_phase != SessionNegotiationPhase::Replacing - && now.saturating_duration_since(state.last_activity) - >= Duration::from_secs(self.timeouts.reconnect_grace_secs) + /// Atomically closes a session only when reconnect grace is still due. + pub(crate) fn close_if_due(&self, now: Instant) -> bool { + let Some(released) = self.begin_close(false, Some(now)) else { + return false; + }; + self.finish_close(released, false); + true } - fn begin_close(&self, superseded: bool) -> Option { + fn begin_close(&self, superseded: bool, idle_now: Option) -> Option { let mut state = self.state.lock(); if state.closed || (superseded && state.negotiation_phase != SessionNegotiationPhase::Replacing) { return None; } + if let Some(now) = idle_now + && (state.negotiation_phase == SessionNegotiationPhase::Replacing + || now.saturating_duration_since(state.last_activity) + < Duration::from_secs(self.timeouts.reconnect_grace_secs)) + { + return None; + } if !superseded && state.negotiation_phase == SessionNegotiationPhase::Replacing { state.close_requested = true; return None;