From b1acfcba800f70695603d5b179f6b68d886b8b74 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jo=C3=A3o=20Lucas?= Date: Sat, 28 Mar 2026 20:00:06 -0300 Subject: [PATCH 1/4] fix: address 9 audit findings across correctness, safety, and performance MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit HIGH: 1. Sender key check now reads through signal_cache with a read lock instead of bypassing cache and taking a write lock (src/send.rs) 2. Dispatch UndecryptableMessage event for group NoSenderKeyState decrypt failures — matches the session-based path (src/message.rs) 3. Remove double cleanup_connection_state call on transport disconnect — the call in run() already covers it (src/client.rs) MEDIUM: 4. Document TOCTOU window in increment_retry_count — mitigated by per-chat mailbox serialization (src/message.rs) 5. Signal cache flush uses snapshot-then-release pattern — locks are held only during serialization, not during I/O (signal_cache.rs) 6. Reduce per-chat message queue capacity from 10000 to 500 to limit memory amplification (src/handlers/message.rs) LOW: 7. PortableCache::remove_key uses lazy deletion (O(1)) instead of O(n) VecDeque scan (src/portable_cache.rs) 8. Guard now_millis() casts with .max(0) before as u64 to prevent silent wrap on negative clock values (keepalive.rs, client.rs) 9. Replace fragile .contains(".us") heuristic with proper JID type methods is_group()/is_broadcast_list()/is_status_broadcast() --- src/client.rs | 7 +- src/handlers/message.rs | 6 +- src/message.rs | 15 ++-- src/portable_cache.rs | 7 +- src/send.rs | 8 +- wacore/src/protocol/keepalive.rs | 18 ++--- wacore/src/store/signal_cache.rs | 127 ++++++++++++++++++------------- 7 files changed, 106 insertions(+), 82 deletions(-) diff --git a/src/client.rs b/src/client.rs index 91a014f43..feff1a534 100644 --- a/src/client.rs +++ b/src/client.rs @@ -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, ); @@ -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."); return Err(anyhow::anyhow!("Transport disconnected unexpectedly")); @@ -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(()) } diff --git a/src/handlers/message.rs b/src/handlers/message.rs index 63d1b1f9f..3698dc488 100644 --- a/src/handlers/message.rs +++ b/src/handlers/message.rs @@ -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::>(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::>(500); let client_for_worker = client.clone(); diff --git a/src/message.rs b/src/message.rs index d1f7517a4..ed4a84814 100644 --- a/src/message.rs +++ b/src/message.rs @@ -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. async fn increment_retry_count(&self, cache_key: &str) -> Option { let current = self.message_retry_counts.get(&cache_key.to_string()).await; match current { @@ -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, @@ -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(()); @@ -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) => { diff --git a/src/portable_cache.rs b/src/portable_cache.rs index 6bc109773..9404596d1 100644 --- a/src/portable_cache.rs +++ b/src/portable_cache.rs @@ -51,9 +51,10 @@ where } fn remove_key(&mut self, key: &K) -> Option> { - 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) } } diff --git a/src/send.rs b/src/send.rs index 987e68480..4965ffe05 100644 --- a/src/send.rs +++ b/src/send.rs @@ -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(); diff --git a/wacore/src/protocol/keepalive.rs b/wacore/src/protocol/keepalive.rs index 6945834d4..695e7d662 100644 --- a/wacore/src/protocol/keepalive.rs +++ b/wacore/src/protocol/keepalive.rs @@ -23,7 +23,7 @@ pub fn ms_since(timestamp_ms: u64) -> Option { 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)) } @@ -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), @@ -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)); } diff --git a/wacore/src/store/signal_cache.rs b/wacore/src/store/signal_cache.rs index 588a255ff..4549089c9 100644 --- a/wacore/src/store/signal_cache.rs +++ b/wacore/src/store/signal_cache.rs @@ -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<()> { - 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(); + 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, Option>)> = 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(()) } From b950542452db8eae65c1f73790a3cb23cbd6c203 Mon Sep 17 00:00:00 2001 From: jlucaso1 Date: Sat, 28 Mar 2026 23:26:43 +0000 Subject: [PATCH 2/4] fix: address review findings in flush durability, cache eviction, and TOCTOU comment MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Signal cache flush: snapshot dirty sets without draining, clear only after all writes succeed (preserves retry on partial failure) - PortableCache: revert lazy deletion — remove→reinsert broke FIFO eviction order; O(capacity) retain is correct and bounded - increment_retry_count: fix inaccurate comment about per-chat serialization (spawn_retry_receipt detaches) --- src/message.rs | 7 ++-- src/portable_cache.rs | 7 ++-- wacore/src/store/signal_cache.rs | 64 ++++++++++++++++++++++---------- 3 files changed, 52 insertions(+), 26 deletions(-) diff --git a/src/message.rs b/src/message.rs index ed4a84814..714bf4022 100644 --- a/src/message.rs +++ b/src/message.rs @@ -132,9 +132,10 @@ impl Client { /// Increments the retry count for a message and returns the new count. /// Returns `None` if max retries have been reached. /// - /// 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. + /// Note: get-then-insert has a theoretical TOCTOU window since + /// `spawn_retry_receipt` detaches. In practice, retries for the same + /// message are rare and a double-send is benign (recipients deduplicate + /// by message ID). async fn increment_retry_count(&self, cache_key: &str) -> Option { let current = self.message_retry_counts.get(&cache_key.to_string()).await; match current { diff --git a/src/portable_cache.rs b/src/portable_cache.rs index 9404596d1..6bc109773 100644 --- a/src/portable_cache.rs +++ b/src/portable_cache.rs @@ -51,10 +51,9 @@ where } fn remove_key(&mut self, key: &K) -> Option> { - // 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) + let entry = self.map.remove(key)?; + self.insertion_order.retain(|ik| ik != key); + Some(entry) } } diff --git a/wacore/src/store/signal_cache.rs b/wacore/src/store/signal_cache.rs index 4549089c9..c47e6bd70 100644 --- a/wacore/src/store/signal_cache.rs +++ b/wacore/src/store/signal_cache.rs @@ -300,15 +300,14 @@ impl SignalStoreCache { /// 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. + /// Dirty sets are only cleared after ALL writes succeed, preserving retry + /// semantics on partial failure. pub async fn flush(&self, backend: &dyn SignalStore) -> Result<()> { // 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(); + let (session_writes, session_delete_keys, session_dirty_keys) = { + let state = self.sessions.lock().await; + let dirty_keys: Vec<_> = state.dirty.iter().cloned().collect(); + let deleted_keys: Vec<_> = state.deleted.iter().cloned().collect(); 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()) { @@ -318,13 +317,13 @@ impl SignalStoreCache { writes.push((address.clone(), bytes)); } } - (writes, deleted_keys) + (writes, deleted_keys, dirty_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 (identity_writes, identity_delete_keys, identity_dirty_keys) = { + let state = self.identities.lock().await; + let dirty_keys: Vec<_> = state.dirty.iter().cloned().collect(); + let deleted_keys: Vec<_> = state.deleted.iter().cloned().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()) { @@ -337,12 +336,12 @@ impl SignalStoreCache { writes.push((address.clone(), key)); } } - (writes, deleted_keys) + (writes, deleted_keys, dirty_keys) }; - let sender_key_ops = { - let mut state = self.sender_keys.lock().await; - let dirty_keys: Vec<_> = state.dirty.drain().collect(); + let (sender_key_ops, sender_key_dirty_keys) = { + let state = self.sender_keys.lock().await; + let dirty_keys: Vec<_> = state.dirty.iter().cloned().collect(); let mut ops: Vec<(Arc, Option>)> = Vec::with_capacity(dirty_keys.len()); for name in &dirty_keys { match state.cache.get(name.as_ref()) { @@ -358,21 +357,21 @@ impl SignalStoreCache { None => {} } } - ops + (ops, dirty_keys) }; // 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_deletes { + for address in &session_delete_keys { backend.delete_session(address).await?; } for (address, key) in &identity_writes { backend.put_identity(address, *key).await?; } - for address in &identity_deletes { + for address in &identity_delete_keys { backend.delete_identity(address).await?; } @@ -383,6 +382,33 @@ impl SignalStoreCache { } } + // Phase 3: all writes succeeded — remove only the flushed keys from dirty sets. + // New mutations that occurred during Phase 2 remain in the dirty sets. + { + let mut state = self.sessions.lock().await; + for key in &session_dirty_keys { + state.dirty.remove(key); + } + for key in &session_delete_keys { + state.deleted.remove(key); + } + } + { + let mut state = self.identities.lock().await; + for key in &identity_dirty_keys { + state.dirty.remove(key); + } + for key in &identity_delete_keys { + state.deleted.remove(key); + } + } + { + let mut state = self.sender_keys.lock().await; + for key in &sender_key_dirty_keys { + state.dirty.remove(key); + } + } + Ok(()) } From 782515c36e1d439a053db9e11ad87ae754a9ccad Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jo=C3=A3o=20Lucas?= Date: Sat, 28 Mar 2026 23:18:30 -0300 Subject: [PATCH 3/4] fix: flush signal cache per-store to eliminate dirty-set race MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Replace the 3-phase snapshot-then-release pattern with per-store lock-through-IO: each store (sessions, identities, sender_keys) is flushed independently under its own lock. This eliminates the race where a mutation between Phase 1 (snapshot) and Phase 3 (clear) could lose a dirty marker for a newer value, since the lock is now held from snapshot through write through clear. Only one store is locked during its I/O — the other two remain free for concurrent encrypt/decrypt. This is strictly better than the original code (which held all 3 locks simultaneously) and correct unlike the 3-phase approach (which had the dirty-set race). --- wacore/src/store/signal_cache.rs | 119 ++++++++++++------------------- 1 file changed, 47 insertions(+), 72 deletions(-) diff --git a/wacore/src/store/signal_cache.rs b/wacore/src/store/signal_cache.rs index c47e6bd70..fa121b18c 100644 --- a/wacore/src/store/signal_cache.rs +++ b/wacore/src/store/signal_cache.rs @@ -294,37 +294,48 @@ impl SignalStoreCache { // === Flush === - /// Flush all dirty state to the backend in a single batch. + /// Flush all dirty state to the backend. /// - /// 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 only cleared after ALL writes succeed, preserving retry - /// semantics on partial failure. + /// Each store (sessions, identities, sender_keys) is flushed independently + /// under its own lock. This means: + /// - Only ONE store is locked during its I/O — the other two are free for + /// concurrent encrypt/decrypt operations. + /// - No race between snapshot and clear — the lock is held throughout, so + /// mutations to the same store are blocked until the flush completes. + /// - Dirty sets are cleared only after successful writes. pub async fn flush(&self, backend: &dyn SignalStore) -> Result<()> { - // Phase 1: snapshot + serialize under lock, then release. - let (session_writes, session_delete_keys, session_dirty_keys) = { - let state = self.sessions.lock().await; + // Flush sessions + { + let mut state = self.sessions.lock().await; let dirty_keys: Vec<_> = state.dirty.iter().cloned().collect(); let deleted_keys: Vec<_> = state.deleted.iter().cloned().collect(); - 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)); + backend.put_session(address, &bytes).await?; } } - (writes, deleted_keys, dirty_keys) - }; + for address in &deleted_keys { + backend.delete_session(address).await?; + } - let (identity_writes, identity_delete_keys, identity_dirty_keys) = { - let state = self.identities.lock().await; + for key in &dirty_keys { + state.dirty.remove(key); + } + for key in &deleted_keys { + state.deleted.remove(key); + } + } + + // Flush identities + { + let mut state = self.identities.lock().await; let dirty_keys: Vec<_> = state.dirty.iter().cloned().collect(); let deleted_keys: Vec<_> = state.deleted.iter().cloned().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(|_| { @@ -333,78 +344,42 @@ impl SignalStoreCache { data.len() ) })?; - writes.push((address.clone(), key)); + backend.put_identity(address, key).await?; } } - (writes, deleted_keys, dirty_keys) - }; + for address in &deleted_keys { + backend.delete_identity(address).await?; + } - let (sender_key_ops, sender_key_dirty_keys) = { - let state = self.sender_keys.lock().await; + for key in &dirty_keys { + state.dirty.remove(key); + } + for key in &deleted_keys { + state.deleted.remove(key); + } + } + + // Flush sender keys + { + let mut state = self.sender_keys.lock().await; let dirty_keys: Vec<_> = state.dirty.iter().cloned().collect(); - let mut ops: Vec<(Arc, Option>)> = 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))); + backend.put_sender_key(name, &bytes).await?; } Some(None) => { - ops.push((name.clone(), None)); + backend.delete_sender_key(name).await?; } None => {} } } - (ops, dirty_keys) - }; - - // 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_delete_keys { - backend.delete_session(address).await?; - } - - for (address, key) in &identity_writes { - backend.put_identity(address, *key).await?; - } - for address in &identity_delete_keys { - backend.delete_identity(address).await?; - } - 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?, - } - } - - // Phase 3: all writes succeeded — remove only the flushed keys from dirty sets. - // New mutations that occurred during Phase 2 remain in the dirty sets. - { - let mut state = self.sessions.lock().await; - for key in &session_dirty_keys { - state.dirty.remove(key); - } - for key in &session_delete_keys { - state.deleted.remove(key); - } - } - { - let mut state = self.identities.lock().await; - for key in &identity_dirty_keys { - state.dirty.remove(key); - } - for key in &identity_delete_keys { - state.deleted.remove(key); - } - } - { - let mut state = self.sender_keys.lock().await; - for key in &sender_key_dirty_keys { + for key in &dirty_keys { state.dirty.remove(key); } } From a65d796a9aa5298dc72df2b1ec9c6574412883d5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jo=C3=A3o=20Lucas?= Date: Sat, 28 Mar 2026 23:34:54 -0300 Subject: [PATCH 4/4] fix: dispatch Disconnected event after cleanup, not before The removal of cleanup_connection_state() from inside the message loop created a regression: Event::Disconnected was dispatched while is_connected was still true and transport handles still set. Handlers could observe stale state. Move the dispatch to run() after cleanup_connection_state() completes, matching the original ordering where cleanup ran before the event. --- src/client.rs | 16 +++++++++++++--- 1 file changed, 13 insertions(+), 3 deletions(-) diff --git a/src/client.rs b/src/client.rs index feff1a534..e175d7fb1 100644 --- a/src/client.rs +++ b/src/client.rs @@ -873,25 +873,36 @@ impl Client { if let Err(connect_err) = self.connect().await { error!("Failed to connect: {connect_err:#}. Will retry..."); } else { - if self.read_messages_loop().await.is_err() { + let unexpected_disconnect = if self.read_messages_loop().await.is_err() { // Check intentional_reconnect AFTER read loop exits — reconnect() // sets this flag while the loop is running, so it must be read here. if self.expected_disconnect.load(Ordering::Relaxed) || self.intentional_reconnect.swap(false, Ordering::Relaxed) { debug!("Message loop exited during expected disconnect."); + false } else { warn!( "Message loop exited with an error. Will attempt to reconnect if enabled." ); + true } } else if self.expected_disconnect.load(Ordering::Relaxed) { debug!("Message loop exited gracefully (expected disconnect)."); + false } else { info!("Message loop exited gracefully."); - } + false + }; self.cleanup_connection_state().await; + + // Dispatch after cleanup so handlers see cleared connection state. + if unexpected_disconnect { + self.core + .event_bus + .dispatch(&Event::Disconnected(crate::types::events::Disconnected)); + } } if !self.enable_auto_reconnect.load(Ordering::Relaxed) { @@ -1309,7 +1320,6 @@ impl Client { }, Ok(crate::transport::TransportEvent::Disconnected) | Err(_) => { if !self.expected_disconnect.load(Ordering::Relaxed) { - self.core.event_bus.dispatch(&Event::Disconnected(crate::types::events::Disconnected)); debug!("Transport disconnected unexpectedly."); return Err(anyhow::anyhow!("Transport disconnected unexpectedly")); } else {