Skip to content
Merged
Show file tree
Hide file tree
Changes from 2 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 5 additions & 4 deletions src/keepalive.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.",
Expand Down
26 changes: 14 additions & 12 deletions wacore/src/protocol/keepalive.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,23 +27,25 @@ pub fn ms_since(timestamp_ms: u64) -> Option<u64> {
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)
}
Expand Down
79 changes: 77 additions & 2 deletions wacore/src/stats.rs
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,9 @@ 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.
first_send_since_recv_ms: AtomicU64,
}

/// Point-in-time copy of [`SessionStats`], plus client-level counters the
Expand Down Expand Up @@ -97,8 +100,19 @@ 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();
Comment thread
cubic-dev-ai[bot] marked this conversation as resolved.
self.last_data_sent_ms.store(now, Ordering::Relaxed);
// Arm the dead-socket deadline on the FIRST send after a receive only;
// later sends must not push it out (WA Web `onOrBefore` keeps the earliest
// deadline). The load fast-paths the common already-armed case.
if self.first_send_since_recv_ms.load(Ordering::Relaxed) == 0 {
let _ = self.first_send_since_recv_ms.compare_exchange(
0,
now,
Ordering::Relaxed,
Ordering::Relaxed,
);
}
}

/// One transport data event carrying `frames` decodable frames.
Expand All @@ -117,6 +131,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);
}
}

Expand All @@ -127,6 +143,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]
Expand Down Expand Up @@ -161,13 +179,22 @@ 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]
pub fn last_data_sent_ms(&self) -> u64 {
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)
Expand Down Expand Up @@ -729,6 +756,54 @@ 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"
);
}

#[test]
fn reset_connection_activity_keeps_traffic() {
let stats = SessionStats::new();
Expand Down
Loading