diff --git a/src/keepalive.rs b/src/keepalive.rs index 66ec2ff46..f0eb8694f 100644 --- a/src/keepalive.rs +++ b/src/keepalive.rs @@ -202,15 +202,16 @@ impl Client { } // WA Web: deadSocketTimer is an independent 20s watchdog armed on - // every send and cancelled on every receive. We approximate this by + // the FIRST send after a receive (onOrBefore keeps the earliest + // deadline) and cancelled on every receive. We approximate this by // checking is_dead_socket on EVERY keepalive tick — not just after // a failed ping. This catches scenarios where pending IQs caused // the ping to be skipped, or where the ping "succeeded" but the // connection died immediately after. - let last_sent = self.stats.last_data_sent_ms(); + let first_send = self.stats.first_send_since_recv_ms(); let last_recv = self.stats.last_data_received_ms(); - if is_dead_socket(last_sent, last_recv) { - let elapsed = ms_since(last_sent).unwrap_or(0); + if is_dead_socket(first_send, last_recv) { + let elapsed = ms_since(first_send).unwrap_or(0); warn!( target: "Client/Keepalive", "No data received for {:.1}s after send (dead socket), forcing reconnect.", diff --git a/wacore/src/protocol/keepalive.rs b/wacore/src/protocol/keepalive.rs index 695e7d662..fdc331723 100644 --- a/wacore/src/protocol/keepalive.rs +++ b/wacore/src/protocol/keepalive.rs @@ -27,23 +27,25 @@ pub fn ms_since(timestamp_ms: u64) -> Option { Some(now.saturating_sub(timestamp_ms)) } -/// Checks the dead-socket condition: data was sent but nothing received -/// within [`DEAD_SOCKET_TIME`]. +/// Checks the dead-socket condition: [`DEAD_SOCKET_TIME`] elapsed since the timer +/// was armed without a receive cancelling it. /// -/// WA Web: `deadSocketTimer` is armed on every `callStanza` (send) and -/// cancelled on every `parseAndHandleStanza` (receive). It fires when -/// `deadSocketTime` (20 s) elapses after the last send without any receive. -pub fn is_dead_socket(last_sent_ms: u64, last_received_ms: u64) -> bool { - // Never sent anything yet -- timer not armed. - if last_sent_ms == 0 { +/// `armed_ms` is the anchor WA Web's `deadSocketTimer.onOrBefore` keeps: the FIRST +/// send after the last receive (0 when unarmed / cancelled). It must NOT be the +/// most-recent send — anchoring there lets continued outgoing traffic keep pushing +/// the deadline out and hide a half-open socket forever. The caller feeds +/// `SessionStats::first_send_since_recv_ms`, reset to 0 on every receive +/// (`parseAndHandleStanza` → `cancel()`). +pub fn is_dead_socket(armed_ms: u64, last_received_ms: u64) -> bool { + // Timer not armed (never sent since the last receive). + if armed_ms == 0 { return false; } - // Received data after (or at) the last send -- timer cancelled. - if last_received_ms >= last_sent_ms { + // Received data after (or at) the armed instant -- timer cancelled. + if last_received_ms >= armed_ms { return false; } - // Sent but no reply: check if DEAD_SOCKET_TIME has elapsed since the send. - ms_since(last_sent_ms) + ms_since(armed_ms) .map(|elapsed| elapsed > DEAD_SOCKET_TIME.as_millis() as u64) .unwrap_or(false) } diff --git a/wacore/src/stats.rs b/wacore/src/stats.rs index e89ce715f..457ef5555 100644 --- a/wacore/src/stats.rs +++ b/wacore/src/stats.rs @@ -50,6 +50,11 @@ pub struct SessionStats { /// Timestamp (ms since UNIX epoch) of the last received WebSocket data. /// WA Web: `parseAndHandleStanza` → `deadSocketTimer.cancel()`. last_data_received_ms: AtomicU64, + /// Dead-socket watchdog anchor (WA Web `deadSocketTimer.onOrBefore`): the first + /// send since the last receive, so continued traffic can't push the deadline out. + /// Treated as stale (and re-armed) once `<= last_data_received_ms`, so a send that + /// raced past a receive-reset can't leave a pre-receive value stuck here. + first_send_since_recv_ms: AtomicU64, } /// Point-in-time copy of [`SessionStats`], plus client-level counters the @@ -97,8 +102,20 @@ impl SessionStats { self.bytes_sent .fetch_add(wire_bytes as u64, Ordering::Relaxed); self.frames_sent.fetch_add(1, Ordering::Relaxed); - self.last_data_sent_ms - .store(Self::now_ms(), Ordering::Relaxed); + let now = Self::now_ms(); + self.last_data_sent_ms.store(now, Ordering::Relaxed); + // Arm the dead-socket deadline on the FIRST send after a receive (WA Web + // `onOrBefore` keeps the earliest deadline; later sends must not push it out). + // Re-arm when the anchor is unset OR stale — a receive landed after it was + // armed (`anchor <= last_received`). Guarding only on `== 0` would let a send + // that captured `now` before a concurrent receive-reset write a pre-receive + // timestamp that then sticks forever (its arm raced past the reset), silently + // disabling detection; the stale check re-arms it on the next send instead. + let last_recv = self.last_data_received_ms.load(Ordering::Relaxed); + let anchor = self.first_send_since_recv_ms.load(Ordering::Relaxed); + if anchor == 0 || anchor <= last_recv { + self.first_send_since_recv_ms.store(now, Ordering::Relaxed); + } } /// One transport data event carrying `frames` decodable frames. @@ -117,6 +134,8 @@ impl SessionStats { if frames > 1 { self.last_data_received_ms .store(Self::now_ms(), Ordering::Relaxed); + // A receive cancels the dead-socket deadline; the next send re-arms it. + self.first_send_since_recv_ms.store(0, Ordering::Relaxed); } } @@ -127,6 +146,8 @@ impl SessionStats { pub fn mark_recv_activity(&self) { self.last_data_received_ms .store(Self::now_ms(), Ordering::Relaxed); + // A receive cancels the dead-socket deadline; the next send re-arms it. + self.first_send_since_recv_ms.store(0, Ordering::Relaxed); } #[inline] @@ -161,6 +182,7 @@ impl SessionStats { pub fn reset_connection_activity(&self) { self.last_data_sent_ms.store(0, Ordering::Relaxed); self.last_data_received_ms.store(0, Ordering::Relaxed); + self.first_send_since_recv_ms.store(0, Ordering::Relaxed); } #[inline] @@ -168,6 +190,14 @@ impl SessionStats { self.last_data_sent_ms.load(Ordering::Relaxed) } + /// The dead-socket watchdog anchor: the first send since the last receive + /// (0 when unarmed). Evaluate [`is_dead_socket`] against this, not the last + /// send, so continued outgoing traffic can't hide a half-open socket. + #[inline] + pub fn first_send_since_recv_ms(&self) -> u64 { + self.first_send_since_recv_ms.load(Ordering::Relaxed) + } + #[inline] pub fn last_data_received_ms(&self) -> u64 { self.last_data_received_ms.load(Ordering::Relaxed) @@ -729,6 +759,83 @@ mod tests { assert!(snap.last_data_received_ms > 0); } + #[test] + fn dead_socket_anchor_holds_across_continued_sends() { + use crate::protocol::keepalive::is_dead_socket; + + let stats = SessionStats::new(); + assert_eq!(stats.first_send_since_recv_ms(), 0, "unarmed initially"); + + stats.record_frame_sent(10); + let armed = stats.first_send_since_recv_ms(); + assert!(armed > 0, "the first send arms the dead-socket anchor"); + + // Continued outgoing traffic must NOT push the anchor out (WA Web onOrBefore + // keeps the earliest deadline) — this is what let a half-open socket hide. + // The CAS gate holds the anchor put across further sends; a sleep makes an + // "unconditional store" regression observable without the assert depending + // on the clock actually ticking (it stays == armed either way). + std::thread::sleep(std::time::Duration::from_millis(2)); + stats.record_frame_sent(10); + stats.record_frame_sent(10); + assert_eq!( + stats.first_send_since_recv_ms(), + armed, + "later sends keep the earliest anchor" + ); + + // A dead socket is detected once DEAD_SOCKET_TIME passes the anchor, even + // though sends kept happening (anchor far in the past, no receive since). + let now = crate::time::now_millis().max(0) as u64; + let stale = now.saturating_sub(21_000); + assert!( + is_dead_socket(stale, stale.saturating_sub(5_000)), + "20s past the anchor with no receive => dead" + ); + + // A receive cancels the anchor; the next send re-arms it (non-zero again). + stats.mark_recv_activity(); + assert_eq!( + stats.first_send_since_recv_ms(), + 0, + "a receive cancels the anchor" + ); + stats.record_frame_sent(10); + assert!( + stats.first_send_since_recv_ms() > 0, + "the next send after a receive re-arms the anchor" + ); + } + + /// A send whose arm raced past a concurrent receive-reset can leave the anchor + /// at a PRE-receive timestamp. The next send must re-arm it (stale: anchor <= + /// last_received) rather than treat it as live and stick there forever, which + /// would silently disable dead-socket detection for the rest of the connection. + #[test] + fn stale_pre_receive_anchor_self_heals_on_next_send() { + use crate::protocol::keepalive::is_dead_socket; + let stats = SessionStats::new(); + let base = SessionStats::now_ms(); + // Reconstruct the race outcome directly: a receive at `base`, and an anchor + // left behind at a pre-receive instant (the lost-reset send's stale `now`). + stats.last_data_received_ms.store(base, Ordering::Relaxed); + stats + .first_send_since_recv_ms + .store(base.saturating_sub(1_000), Ordering::Relaxed); + + stats.record_frame_sent(10); + + let rearmed = stats.first_send_since_recv_ms(); + assert!( + rearmed >= base, + "a stale pre-receive anchor must re-arm to a post-receive send, got {rearmed} < {base}" + ); + assert!( + !is_dead_socket(rearmed, base), + "the re-armed anchor is after the receive, so the socket is not dead" + ); + } + #[test] fn reset_connection_activity_keeps_traffic() { let stats = SessionStats::new();