Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
7 changes: 3 additions & 4 deletions src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1255,7 +1255,7 @@ impl Client {
Ok(crate::transport::TransportEvent::DataReceived(data)) => {
// Update dead-socket timer (WA Web: deadSocketTimer reset)
self.last_data_received_ms.store(
wacore::time::now_millis() as u64,
wacore::time::now_millis().max(0) as u64,
Ordering::Relaxed,
);

Expand Down Expand Up @@ -1308,8 +1308,7 @@ impl Client {
}
},
Ok(crate::transport::TransportEvent::Disconnected) | Err(_) => {
self.cleanup_connection_state().await;
if !self.expected_disconnect.load(Ordering::Relaxed) {
if !self.expected_disconnect.load(Ordering::Relaxed) {
self.core.event_bus.dispatch(&Event::Disconnected(crate::types::events::Disconnected));
debug!("Transport disconnected unexpectedly.");
Comment on lines +1322 to 1323

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Clear connection state before firing Disconnected handlers

Event::Disconnected is now dispatched before cleanup_connection_state() runs, so synchronous handlers can observe stale connection state (is_connected, transport/noise handles, caches) and make incorrect decisions (for example, skipping reconnect logic because the client still appears connected during the callback). This regression comes from removing the in-loop cleanup call without preserving the prior cleanup-before-dispatch ordering for unexpected disconnects.

Useful? React with 👍 / 👎.

return Err(anyhow::anyhow!("Transport disconnected unexpectedly"));
Expand Down Expand Up @@ -3318,7 +3317,7 @@ impl Client {

// WA Web: callStanza → deadSocketTimer.onOrBefore(deadSocketTime, socketId)
self.last_data_sent_ms
.store(wacore::time::now_millis() as u64, Ordering::Relaxed);
.store(wacore::time::now_millis().max(0) as u64, Ordering::Relaxed);

Ok(())
}
Expand Down
6 changes: 3 additions & 3 deletions src/handlers/message.rs
Original file line number Diff line number Diff line change
Expand Up @@ -58,9 +58,9 @@ impl StanzaHandler for MessageHandler {
let tx = client
.message_queues
.get_with_by_ref(&chat_id, async {
// Create a channel with backpressure
// Increased capacity to handle high message rates without blocking
let (tx, rx) = async_channel::bounded::<Arc<Node>>(10000);
// Bounded capacity provides backpressure to prevent unbounded memory growth.
// 500 is enough for burst handling while limiting per-chat memory.
let (tx, rx) = async_channel::bounded::<Arc<Node>>(500);

let client_for_worker = client.clone();

Expand Down
15 changes: 9 additions & 6 deletions src/message.rs
Original file line number Diff line number Diff line change
Expand Up @@ -132,7 +132,9 @@ impl Client {
/// Increments the retry count for a message and returns the new count.
/// Returns `None` if max retries have been reached.
///
/// Uses get + insert for portability across cache backends.
/// Note: get-then-insert has a theoretical TOCTOU window, but messages are
/// processed sequentially per-chat (mailbox pattern in MessageHandler), so
/// concurrent increments for the same cache_key are practically impossible.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟠 Major

The TOCTOU window is still reachable here.

MessageHandler only serializes handle_incoming_message. spawn_retry_receipt() detaches before calling increment_retry_count, so two failure paths for the same {chat}:{msg_id}:{sender} can still race through get()+insert() and double-send retries/PDO. Move the increment into the serialized path or guard it with its own per-key lock.

Based on learnings: Use message_enqueue_locks to serialize per-chat incoming message processing.

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/message.rs` around lines 135 - 137, The TOCTOU risk arises because
spawn_retry_receipt() detaches before calling increment_retry_count, allowing
two failure paths for the same cache_key to race through get()+insert(); move
the retry-count increment into the serialized path or protect it with the
per-key lock used for incoming processing: acquire the same
message_enqueue_locks key inside handle_incoming_message (or before detaching)
and perform increment_retry_count while holding that lock (or alternatively wrap
increment_retry_count itself with the per-key lock), ensuring MessageHandler's
serialized processing (handle_incoming_message) covers the increment for the
{chat}:{msg_id}:{sender} cache_key and prevents double-send retries/PDO.

async fn increment_retry_count(&self, cache_key: &str) -> Option<u8> {
let current = self.message_retry_counts.get(&cache_key.to_string()).await;
match current {
Expand Down Expand Up @@ -537,10 +539,10 @@ impl Client {
session_enc_nodes.len()
);

// Skip session processing for group senders (@c.us, @g.us, @broadcast)
// Groups don't use 1:1 Signal Protocol sessions
let is_group_sender = sender_encryption_jid.server.contains(".us")
|| sender_encryption_jid.server.contains("broadcast");
// Skip session processing for group/broadcast JIDs — they use sender keys, not 1:1 sessions.
let is_group_sender = sender_encryption_jid.is_group()
|| sender_encryption_jid.is_broadcast_list()
|| sender_encryption_jid.is_status_broadcast();

let (
session_decrypted_successfully,
Expand Down Expand Up @@ -1034,7 +1036,7 @@ impl Client {
enc_nodes: &[&wacore_binary::node::Node],
info: &MessageInfo,
_sender_encryption_jid: &Jid,
_decrypt_fail_mode: crate::types::events::DecryptFailMode,
decrypt_fail_mode: crate::types::events::DecryptFailMode,
) -> Result<(), DecryptionError> {
if enc_nodes.is_empty() {
return Ok(());
Expand Down Expand Up @@ -1114,6 +1116,7 @@ impl Client {
"No sender key state for group message [msg:{}] from {}: {}. Sending retry receipt.",
info.id, info.source.sender, msg
);
self.dispatch_undecryptable_event(info, decrypt_fail_mode);
self.spawn_retry_receipt(info, RetryReason::NoSession);
}
Err(e) => {
Expand Down
7 changes: 4 additions & 3 deletions src/portable_cache.rs
Original file line number Diff line number Diff line change
Expand Up @@ -51,9 +51,10 @@ where
}

fn remove_key(&mut self, key: &K) -> Option<CacheEntry<V>> {
let entry = self.map.remove(key)?;
self.insertion_order.retain(|ik| ik != key);
Some(entry)
// Lazy deletion: remove from map but leave stale key in insertion_order.
// Stale keys are skipped during FIFO eviction (map.remove returns None).
// run_pending_tasks() periodically compacts insertion_order.
self.map.remove(key)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Remove stale FIFO keys when invalidating cache entries

This lazy deletion change leaves insertion_order entries behind on every remove/invalidate, but compaction only happens in run_pending_tasks() and is not part of normal cache operations. In moka-cache-disabled builds, workloads that frequently invalidate and reinsert keys can grow insertion_order indefinitely while map remains small, causing avoidable long-term memory and eviction overhead.

Useful? React with 👍 / 👎.

Comment thread
coderabbitai[bot] marked this conversation as resolved.
Outdated
}
}

Expand Down
8 changes: 4 additions & 4 deletions src/send.rs
Original file line number Diff line number Diff line change
Expand Up @@ -871,15 +871,15 @@ impl Client {
}

let force_skdm = {
use wacore::libsignal::protocol::SenderKeyStore;
use wacore::libsignal::store::sender_key_name::SenderKeyName;
let mut device_guard = device_store_arc.write().await;
let sender_address = own_sending_jid.to_protocol_address();
let sender_key_name =
SenderKeyName::new(to_str.clone(), sender_address.to_string());

let key_exists = device_guard
.load_sender_key(&sender_key_name)
let device_guard = device_store_arc.read().await;
let key_exists = self
.signal_cache
.get_sender_key(&sender_key_name, &*device_guard.backend)
.await?
.is_some();

Expand Down
18 changes: 9 additions & 9 deletions wacore/src/protocol/keepalive.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ pub fn ms_since(timestamp_ms: u64) -> Option<u64> {
if timestamp_ms == 0 {
return None;
}
let now = crate::time::now_millis() as u64;
let now = crate::time::now_millis().max(0) as u64;
Some(now.saturating_sub(timestamp_ms))
}

Expand Down Expand Up @@ -61,14 +61,14 @@ mod tests {

#[test]
fn ms_since_recent() {
let now_ms = crate::time::now_millis() as u64;
let now_ms = crate::time::now_millis().max(0) as u64;
let elapsed = ms_since(now_ms).unwrap();
assert!(elapsed < 100, "should be near-zero, got {elapsed}ms");
}

#[test]
fn ms_since_stale() {
let thirty_sec_ago = (crate::time::now_millis() as u64).saturating_sub(30_000);
let thirty_sec_ago = (crate::time::now_millis().max(0) as u64).saturating_sub(30_000);
let elapsed = ms_since(thirty_sec_ago).unwrap();
assert!(
(29_000..=31_000).contains(&elapsed),
Expand All @@ -85,33 +85,33 @@ mod tests {

#[test]
fn dead_socket_received_after_send() {
let t = crate::time::now_millis() as u64;
let t = crate::time::now_millis().max(0) as u64;
assert!(!is_dead_socket(t, t + 1));
}

#[test]
fn dead_socket_sent_recently() {
let now = crate::time::now_millis() as u64;
let now = crate::time::now_millis().max(0) as u64;
assert!(!is_dead_socket(now, 0));
}

#[test]
fn dead_socket_sent_long_ago_no_reply() {
let thirty_ago = (crate::time::now_millis() as u64).saturating_sub(30_000);
let thirty_ago = (crate::time::now_millis().max(0) as u64).saturating_sub(30_000);
assert!(is_dead_socket(thirty_ago, 0));
}

#[test]
fn dead_socket_sent_long_ago_old_reply() {
let thirty_ago = (crate::time::now_millis() as u64).saturating_sub(30_000);
let thirty_ago = (crate::time::now_millis().max(0) as u64).saturating_sub(30_000);
let thirty_one_ago = thirty_ago.saturating_sub(1_000);
assert!(is_dead_socket(thirty_ago, thirty_one_ago));
}

#[test]
fn dead_socket_sent_long_ago_recent_reply() {
let thirty_ago = (crate::time::now_millis() as u64).saturating_sub(30_000);
let one_ago = (crate::time::now_millis() as u64).saturating_sub(1_000);
let thirty_ago = (crate::time::now_millis().max(0) as u64).saturating_sub(30_000);
let one_ago = (crate::time::now_millis().max(0) as u64).saturating_sub(1_000);
assert!(!is_dead_socket(thirty_ago, one_ago));
}

Expand Down
127 changes: 74 additions & 53 deletions wacore/src/store/signal_cache.rs
Original file line number Diff line number Diff line change
Expand Up @@ -295,73 +295,94 @@ impl SignalStoreCache {
// === Flush ===

/// Flush all dirty state to the backend in a single batch.
/// Acquires all 3 mutexes to ensure consistency (matches WhatsApp Web's pattern).
///
/// Sessions are serialized here (not on every store_session call).
/// Dirty sets are only cleared after ALL writes succeed.
/// Uses a snapshot-then-release pattern: serialize dirty data under the lock,
/// release locks, then write to the backend. This avoids blocking all
/// encrypt/decrypt operations for the duration of I/O.
///
/// Dirty sets are drained before the write phase. If a write fails, the
/// data remains in the cache and will be re-dirtied on the next modification.
pub async fn flush(&self, backend: &dyn SignalStore) -> Result<()> {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🔴 Critical

Snapshot-then-drain can permanently stale the backend.

These dirty/deleted sets are cleared before any backend write and flush() itself is not serialized. A transient write error—or flush A(old snapshot) -> mutate -> flush B(new snapshot) -> B writes -> A writes—can therefore leave sessions/identities/sender keys persisted at the old value with no dirty markers left to retry. Keep a single flush mutex and only clear/reconcile dirty state after the write phase succeeds.

Also applies to: 308-384

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@wacore/src/store/signal_cache.rs` around lines 299 - 305, The flush()
implementation snapshots and drains dirty/deleted sets before performing backend
writes, which can permanently stale the backend if concurrent flushes or write
failures occur; change flush() (wacore::store::signal_cache::flush) to serialize
all flush operations with a dedicated mutex (e.g., a flush_mutex) so only one
flush runs at a time, perform the backend write using the snapshot while
retaining dirty/deleted state, and only clear or reconcile the dirty and deleted
sets after the backend write succeeds; ensure that on write failure the
dirty/deleted markers remain untouched (or are merged back) so retries can
reattempt, and update any code paths that currently drain the sets prior to
calling backend methods on the SignalStore to instead clear them post-success.

let mut sessions = self.sessions.lock().await;
let mut identities = self.identities.lock().await;
let mut sender_keys = self.sender_keys.lock().await;

// Snapshot dirty/deleted sets WITHOUT draining — preserve on failure
let session_dirty: Vec<_> = sessions.dirty.iter().cloned().collect();
let session_deleted: Vec<_> = sessions.deleted.iter().cloned().collect();
let identity_dirty: Vec<_> = identities.dirty.iter().cloned().collect();
let identity_deleted: Vec<_> = identities.deleted.iter().cloned().collect();
let sender_key_dirty: Vec<_> = sender_keys.dirty.iter().cloned().collect();

// Persist dirty sessions — serialize only here, not on every store_session
for address in &session_dirty {
if let Some(Some(record)) = sessions.cache.get(address.as_ref()) {
let bytes = record
.serialize()
.map_err(|e| anyhow::anyhow!("session serialize for {address}: {e}"))?;
backend.put_session(address, &bytes).await?;
// Phase 1: snapshot + serialize under lock, then release.
// Collect dirty keys first, then clear, then serialize from cache.
let (session_writes, session_deletes) = {
let mut state = self.sessions.lock().await;
let dirty_keys: Vec<_> = state.dirty.drain().collect();
let deleted_keys: Vec<_> = state.deleted.drain().collect();

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Keep dirty flags until backend flush completes

flush() now drains dirty/deleted before performing backend I/O, so any mid-flush error (for example a transient SQLite write failure) returns with those markers already cleared and unchanged entries will not be retried on later flushes. That can silently drop pending session/identity/sender-key persistence and lose crypto state after process restart; the same drain-before-write pattern is repeated for all three stores in this function.

Useful? React with 👍 / 👎.

let mut writes = Vec::with_capacity(dirty_keys.len());
for address in &dirty_keys {
if let Some(Some(record)) = state.cache.get(address.as_ref()) {
let bytes = record
.serialize()
.map_err(|e| anyhow::anyhow!("session serialize for {address}: {e}"))?;
writes.push((address.clone(), bytes));
}
}
(writes, deleted_keys)
};

let (identity_writes, identity_deletes) = {
let mut state = self.identities.lock().await;
let dirty_keys: Vec<_> = state.dirty.drain().collect();
let deleted_keys: Vec<_> = state.deleted.drain().collect();
let mut writes = Vec::with_capacity(dirty_keys.len());
for address in &dirty_keys {
if let Some(Some(data)) = state.cache.get(address.as_ref()) {
let key: [u8; 32] = data.as_ref().try_into().map_err(|_| {
anyhow::anyhow!(
"Corrupted identity key for {address}: expected 32 bytes, got {}",
data.len()
)
})?;
writes.push((address.clone(), key));
}
}
(writes, deleted_keys)
};

let sender_key_ops = {
let mut state = self.sender_keys.lock().await;
let dirty_keys: Vec<_> = state.dirty.drain().collect();
let mut ops: Vec<(Arc<str>, Option<Vec<u8>>)> = Vec::with_capacity(dirty_keys.len());
for name in &dirty_keys {
match state.cache.get(name.as_ref()) {
Some(Some(record)) => {
let bytes = record
.serialize()
.map_err(|e| anyhow::anyhow!("sender key serialize for {name}: {e}"))?;
ops.push((name.clone(), Some(bytes)));
}
Some(None) => {
ops.push((name.clone(), None));
}
None => {}
}
}
ops
};

// Phase 2: write to backend without holding any locks.
for (address, bytes) in &session_writes {
backend.put_session(address, bytes).await?;
}
for address in &session_deleted {
for address in &session_deletes {
backend.delete_session(address).await?;
}

for address in &identity_dirty {
if let Some(Some(data)) = identities.cache.get(address.as_ref()) {
let key: [u8; 32] = data.as_ref().try_into().map_err(|_| {
anyhow::anyhow!(
"Corrupted identity key for {address}: expected 32 bytes, got {}",
data.len()
)
})?;
backend.put_identity(address, key).await?;
}
for (address, key) in &identity_writes {
backend.put_identity(address, *key).await?;
}
for address in &identity_deleted {
for address in &identity_deletes {
backend.delete_identity(address).await?;
}

for name in &sender_key_dirty {
match sender_keys.cache.get(name.as_ref()) {
Some(Some(record)) => {
let bytes = record
.serialize()
.map_err(|e| anyhow::anyhow!("sender key serialize for {name}: {e}"))?;
backend.put_sender_key(name, &bytes).await?;
}
Some(None) => {
// Deleted via delete_sender_key — propagate to backend
backend.delete_sender_key(name).await?;
}
None => {}
for (name, bytes_opt) in &sender_key_ops {
match bytes_opt {
Some(bytes) => backend.put_sender_key(name, bytes).await?,
None => backend.delete_sender_key(name).await?,
}
}

// All writes succeeded — clear dirty sets (matches WA Web's clearDirty())
sessions.dirty.clear();
sessions.deleted.clear();
identities.dirty.clear();
identities.deleted.clear();
sender_keys.dirty.clear();

Ok(())
Comment on lines 306 to 387

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🧹 Nitpick | 🔵 Trivial

Note partial-failure semantics for future reference.

If, say, sessions flush succeeds but identities flush fails, the sessions dirty set is cleared before the error is returned. On retry, only identities (and sender_keys) will be re-flushed since sessions are already persisted and no longer dirty. This is correct behavior since the session writes did succeed.

This differs slightly from the PR description's "clearing dirty sets only after all writes succeed" (which implies a global all-or-nothing), but per-store clearing is the more practical approach given the independent store design.

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@wacore/src/store/signal_cache.rs` around lines 306 - 387, The current flush
behavior clears each store's dirty set after that store's writes succeed, which
yields per-store partial-failure semantics (e.g., sessions cleared even if
identities later fail) and differs from the PR text claiming atomic "clear dirty
sets only after all writes succeed"; update the code/docs to match intent:
either (A) change the PR description to state per-store clearing semantics, or
(B) modify flush to only clear any dirty/deleted sets after all three sections
succeed by moving the removals out of the per-store blocks and performing them
after all backend calls complete; refer to the flush method and the per-store
states sessions, identities, and sender_keys and their state.dirty/state.deleted
manipulations when making the change.

}

Expand Down
Loading