diff --git a/src/web/http/websocket/driver.rs b/src/web/http/websocket/driver.rs index 8c9d6db..f34df62 100644 --- a/src/web/http/websocket/driver.rs +++ b/src/web/http/websocket/driver.rs @@ -134,7 +134,7 @@ async fn run_multiplex( let maximum_message = session.limits().carrier_batch_bytes; let mut active = false; loop { - let down = session.poll_down(cursor); + let down = session.poll_down_websocket(cursor); tokio::pin!(down); let event = tokio::select! { _ = cancellation.cancelled() => return Err(()), diff --git a/src/web/session/downlink.rs b/src/web/session/downlink.rs index 2bd06f8..6447610 100644 --- a/src/web/session/downlink.rs +++ b/src/web/session/downlink.rs @@ -15,6 +15,14 @@ use crate::web::telemetry::WebSessionLifecycleObservation; impl WebSession { /// Polls pending downlink frames with cursor replay and newest-poll-wins semantics. pub(crate) async fn poll_down(&self, cursor: u64) -> Result { + self.poll_down_inner(cursor, true).await + } + + async fn poll_down_inner( + &self, + cursor: u64, + peer_activity: bool, + ) -> Result { if !self.carrier().is_multiplexed() { return Err(ManagerError::Protocol); } @@ -30,11 +38,13 @@ impl WebSession { next_cursor: unacked.next_cursor, lane_closed: false, }; - self.touch_peer_locked( - &mut state, - Instant::now(), - WebSessionLifecycleObservation::HttpActivityAfterGap, - ); + if peer_activity { + self.touch_peer_locked( + &mut state, + Instant::now(), + WebSessionLifecycleObservation::HttpActivityAfterGap, + ); + } return Ok(result); } if cursor != unacked.next_cursor { @@ -59,11 +69,13 @@ impl WebSession { return Err(ManagerError::Protocol); }; state.down_epoch = epoch; - self.touch_peer_locked( - &mut state, - Instant::now(), - WebSessionLifecycleObservation::HttpActivityAfterGap, - ); + if peer_activity { + self.touch_peer_locked( + &mut state, + Instant::now(), + WebSessionLifecycleObservation::HttpActivityAfterGap, + ); + } let healthy = self.carrier_health_ready_locked(&mut state, Instant::now()); (state.down_epoch, healthy) }; @@ -133,6 +145,14 @@ impl WebSession { } } + /// Polls multiplexed WebSocket downlink without renewing the peer lease. + pub(crate) async fn poll_down_websocket( + &self, + cursor: u64, + ) -> Result { + self.poll_down_inner(cursor, false).await + } + /// Reserves session and process queue capacity while the session lock is held. pub(super) fn reserve_locked( &self, diff --git a/src/web/session/downlink_tests.rs b/src/web/session/downlink_tests.rs index 69a22b5..35a0e8f 100644 --- a/src/web/session/downlink_tests.rs +++ b/src/web/session/downlink_tests.rs @@ -142,3 +142,19 @@ async fn newer_poll_supersedes_older_poll_without_closing_session() { session.close(super::SessionCloseReason::ApiClose); manager.shutdown().await; } + +#[tokio::test] +async fn websocket_downlink_poll_does_not_extend_the_peer_lease() { + let (session, manager) = session(); + session + .state + .lock() + .activity + .touch_peer(Instant::now() - Duration::from_secs(121)); + queue_close(&session); + + session.poll_down_websocket(0).await.unwrap(); + + assert!(session.close_if_due(Instant::now())); + manager.shutdown().await; +}