Skip to main content

whatsapp_rust/
retry.rs

1use crate::client::Client;
2use crate::features::{MessageRetransmission, RetryRequestError};
3use crate::message::RetryReason;
4use crate::send::SendError;
5use crate::types::events::Receipt;
6use log::{debug, info, warn};
7use wacore::types::message::MessageCategory;
8
9use scopeguard;
10use std::sync::Arc;
11use wacore::iq::prekeys::{OneTimePreKeyNode, SignedPreKeyNode};
12use wacore::libsignal::protocol::{PreKeyBundle, PublicKey};
13use wacore::protocol::ProtocolNode;
14use wacore::protocol::retry::{MAX_RETRY_COUNT, MIN_RETRY_FOR_BASE_KEY_CHECK};
15use wacore::types::jid::JidExt;
16use wacore_binary::JidExt as _;
17#[cfg(test)]
18use wacore_binary::NodeContent;
19use wacore_binary::builder::NodeBuilder;
20use wacore_binary::{Jid, Node, OwnedNodeRef};
21use wacore_binary::{NodeContentRef, NodeRef};
22use waproto::whatsapp as wa;
23
24/// Helper to extract bytes content from a Node (used in tests).
25#[cfg(test)]
26fn get_bytes_content(node: &Node) -> Option<&[u8]> {
27    match &node.content {
28        Some(NodeContent::Bytes(b)) => Some(b.as_slice()),
29        _ => None,
30    }
31}
32
33/// Helper to extract bytes content from a NodeRef.
34fn get_bytes_content_ref<'a>(node: &'a NodeRef<'_>) -> Option<&'a [u8]> {
35    match node.content.as_ref() {
36        Some(NodeContentRef::Bytes(b)) => Some(b.as_ref()),
37        _ => None,
38    }
39}
40
41/// Throttle for the "no-keys + retry≥2" forced-recreate fallback. Mirrors
42/// whatsmeow's `recreateSessionTimeout` (`retry.go:156`).
43const RECREATE_SESSION_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(3600);
44
45#[derive(Clone, Copy)]
46enum RetransmissionRoute {
47    Direct,
48    Group,
49    Status,
50    BroadcastList,
51}
52
53impl RetransmissionRoute {
54    const fn uses_sender_key(self) -> bool {
55        matches!(self, Self::Group | Self::Status)
56    }
57}
58
59#[inline]
60fn is_own_account_jid(jid: &Jid, own_pn: Option<&Jid>, own_lid: Option<&Jid>) -> bool {
61    own_pn.is_some_and(|pn| jid.is_same_user_as(pn))
62        || own_lid.is_some_and(|lid| jid.is_same_user_as(lid))
63}
64
65struct PreparedRetransmission {
66    route: RetransmissionRoute,
67    chat: Jid,
68    wire_requester: Jid,
69    encryption_jid: Jid,
70    message: wa::Message,
71    message_id: String,
72    retry_count: u8,
73    recipient: Option<Jid>,
74    group_info: Option<Arc<wacore::client::context::GroupInfo>>,
75    /// Canonical unpadded protobuf bytes shared with the recent-message cache.
76    /// Public retransmissions provide them; the automatic path may fall back
77    /// to its already-decoded message when the cache bytes are unavailable.
78    pre_encoded: Option<Arc<Vec<u8>>>,
79}
80
81fn validate_retransmission(
82    chat: &Jid,
83    requester: &Jid,
84    message_id: &str,
85    retry_count: u8,
86    recipient: Option<&Jid>,
87) -> Result<RetransmissionRoute, SendError> {
88    if chat.is_empty() || requester.is_empty() {
89        return Err(SendError::InvalidRequest(
90            "retransmission JIDs must not be empty".into(),
91        ));
92    }
93    if message_id.is_empty() {
94        return Err(SendError::InvalidRequest(
95            "retransmission message ID must not be empty".into(),
96        ));
97    }
98    if !(1..MAX_RETRY_COUNT).contains(&retry_count) {
99        return Err(SendError::InvalidRequest(format!(
100            "retry count must be in 1..{MAX_RETRY_COUNT}"
101        )));
102    }
103
104    let requester_is_user = matches!(
105        requester.server,
106        wacore_binary::Server::Pn
107            | wacore_binary::Server::Lid
108            | wacore_binary::Server::Hosted
109            | wacore_binary::Server::HostedLid
110            | wacore_binary::Server::Bot
111    );
112    if !requester_is_user {
113        return Err(SendError::InvalidRequest(
114            "retransmission requester must be a user device JID".into(),
115        ));
116    }
117
118    let route = if chat.is_group() {
119        RetransmissionRoute::Group
120    } else if chat.is_status_broadcast() {
121        if !matches!(
122            requester.server,
123            wacore_binary::Server::Pn | wacore_binary::Server::Lid
124        ) {
125            return Err(SendError::InvalidRequest(
126                "status retransmission requester must be a PN or LID device".into(),
127            ));
128        }
129        RetransmissionRoute::Status
130    } else if chat.is_broadcast_list() {
131        if !matches!(
132            requester.server,
133            wacore_binary::Server::Pn | wacore_binary::Server::Lid
134        ) {
135            return Err(SendError::InvalidRequest(
136                "broadcast retransmission requester must be a PN or LID device".into(),
137            ));
138        }
139        RetransmissionRoute::BroadcastList
140    } else if matches!(
141        chat.server,
142        wacore_binary::Server::Pn
143            | wacore_binary::Server::Lid
144            | wacore_binary::Server::Hosted
145            | wacore_binary::Server::HostedLid
146            | wacore_binary::Server::Bot
147    ) {
148        RetransmissionRoute::Direct
149    } else {
150        return Err(SendError::InvalidRequest(
151            "unsupported retransmission chat class".into(),
152        ));
153    };
154
155    if recipient.is_some() && !matches!(route, RetransmissionRoute::Direct) {
156        return Err(SendError::InvalidRequest(
157            "recipient is only valid for direct retransmissions".into(),
158        ));
159    }
160    if recipient.is_some_and(|recipient| {
161        recipient.is_empty()
162            || !matches!(
163                recipient.server,
164                wacore_binary::Server::Pn
165                    | wacore_binary::Server::Lid
166                    | wacore_binary::Server::Hosted
167                    | wacore_binary::Server::HostedLid
168                    | wacore_binary::Server::Bot
169            )
170    }) {
171        return Err(SendError::InvalidRequest(
172            "retransmission recipient must be a user JID".into(),
173        ));
174    }
175
176    Ok(route)
177}
178
179pub(crate) enum RetryReceiptSendOutcome {
180    Sent { included_keys: bool },
181    Suppressed,
182}
183
184/// Separated chat and requester JIDs for retry receipt handling.
185/// Mirrors WAWebHandleRetryRequest `getActualChatInfo` + `getTargetChat`.
186struct RetryChatInfo {
187    /// Bare chat JID (no device suffix) for message lookup.
188    chat: Jid,
189    /// Device-specific JID of the requesting device, for session management.
190    requester: Jid,
191    /// Raw `from` JID from the receipt, for stanza `to` attribute.
192    /// WA Web preserves the original `from` (variable `m`) for the retry stanza.
193    original_from: Jid,
194    /// Receipt's `recipient` attribute, if present. WA Web's
195    /// `handleRetryRequest` propagates this verbatim into the retry resend
196    /// (only self-DM and bot receipts carry it).
197    recipient: Option<Jid>,
198    /// True if the requester is a bot JID (skip namespace normalization).
199    is_bot: bool,
200    /// WA Web's `bot_retry` parser path: only primary `@bot` JIDs, not legacy PN bots.
201    is_fbid_bot_retry: bool,
202}
203
204fn is_fbid_bot_retry_jid(jid: &Jid) -> bool {
205    jid.server == wacore_binary::Server::Bot && jid.device() == 0
206}
207
208/// Resolve the chat and requester JIDs from a retry receipt, separating
209/// message-lookup concerns from session-management concerns.
210/// Mirrors WAWebHandleRetryRequest `getActualChatInfo` + `getTargetChat`.
211fn resolve_retry_chat_info(
212    receipt: &Receipt,
213    node: &NodeRef<'_>,
214    own_pn: Option<&Jid>,
215    own_lid: Option<&Jid>,
216) -> Option<RetryChatInfo> {
217    let from = &receipt.source.chat;
218
219    if from.is_group() || from.is_status_broadcast() || from.is_broadcast_list() {
220        // Group-like chats: chat is already the group/broadcast JID.
221        // Requester is the participant attr (the actual retrying device).
222        let participant = node.attrs().optional_jid("participant");
223        let is_fbid_bot_retry =
224            from.is_group() && participant.as_ref().is_some_and(is_fbid_bot_retry_jid);
225        let requester = participant.unwrap_or_else(|| receipt.source.sender.clone());
226        let is_bot = requester.is_bot();
227        Some(RetryChatInfo {
228            chat: from.clone(),
229            requester,
230            original_from: from.clone(),
231            recipient: node.attrs().optional_jid("recipient"),
232            is_bot,
233            is_fbid_bot_retry,
234        })
235    } else {
236        // DM: resolve chat target via getTargetChat logic.
237        let recipient = node.attrs().optional_jid("recipient");
238        let is_bot = from.is_bot();
239
240        // WA Web getTargetChat (RetryRequest.js:339-371):
241        // 1. Bot + recipient → chat = recipient
242        // 2. Peer device + recipient → chat = recipient
243        // 3. Peer device without recipient → WA Web aborts (returns null).
244        // 4. Normal user → chat = asUserWidOrThrow(from) = from.to_non_ad()
245        let is_peer = is_own_account_jid(from, own_pn, own_lid);
246
247        let chat = if is_bot && let Some(r) = recipient.as_ref() {
248            r.to_non_ad()
249        } else if is_peer {
250            match recipient.as_ref() {
251                Some(r) => r.to_non_ad(),
252                None => {
253                    log::warn!("Ignoring peer device retry without recipient attr");
254                    return None;
255                }
256            }
257        } else {
258            from.to_non_ad()
259        };
260
261        let requester = if from.device() == 0 && from.agent == 0 {
262            chat.clone()
263        } else {
264            from.clone()
265        };
266
267        Some(RetryChatInfo {
268            chat,
269            requester,
270            original_from: from.clone(),
271            recipient,
272            is_bot,
273            is_fbid_bot_retry: is_fbid_bot_retry_jid(from),
274        })
275    }
276}
277
278fn validate_retry_prekey_presence(
279    keys_node: &NodeRef<'_>,
280    is_fbid_bot_retry: bool,
281) -> Result<(), anyhow::Error> {
282    if !is_fbid_bot_retry && keys_node.get_optional_child("key").is_none() {
283        anyhow::bail!("regular retry key bundle missing one-time prekey");
284    }
285    Ok(())
286}
287
288// No retry_count in the key: concurrent receipts for the same participant must
289// serialize, otherwise two update_local_signal_session calls race on session state.
290fn build_retry_processing_key(chat: &Jid, message_id: &str, participant_jid: &Jid) -> String {
291    let mut key = String::with_capacity(message_id.len() + 64);
292    chat.push_to(&mut key);
293    key.push(':');
294    key.push_str(message_id);
295    key.push(':');
296    participant_jid.push_to(&mut key);
297    key
298}
299
300impl Client {
301    async fn resolve_retransmission_encryption_jid(
302        &self,
303        route: RetransmissionRoute,
304        requester: &Jid,
305    ) -> Result<Jid, anyhow::Error> {
306        if matches!(route, RetransmissionRoute::Status) && requester.is_pn() {
307            return match self.get_lid_pn_entry(requester).await? {
308                Some(mapping) => Ok(Jid {
309                    user: wacore_binary::CompactString::new(&mapping.lid),
310                    server: wacore_binary::Server::Lid,
311                    device: requester.device,
312                    agent: requester.agent,
313                    integrator: requester.integrator,
314                }),
315                // WAWebResendStatusMsg explicitly falls back to the PN device
316                // when no LID mapping is available.
317                None => Ok(requester.clone()),
318            };
319        }
320        Ok(self.resolve_encryption_jid(requester).await)
321    }
322
323    /// Retransmit a message to one requesting device.
324    ///
325    /// The client derives the stanza from native protocol data and retains
326    /// ownership of routing, encryption, sender-key tracking, persistence, and
327    /// transport. The original message ID and retry count are preserved.
328    pub async fn retransmit_message(
329        &self,
330        request: MessageRetransmission,
331    ) -> Result<(), SendError> {
332        let route = validate_retransmission(
333            &request.chat,
334            &request.requester,
335            &request.message_id,
336            request.retry_count,
337            request.recipient.as_ref(),
338        )?;
339
340        if matches!(route, RetransmissionRoute::Direct) {
341            let snapshot = self.persistence_manager.get_device_snapshot();
342            let requester_is_local = is_own_account_jid(
343                &request.requester,
344                snapshot.pn.as_ref(),
345                snapshot.lid.as_ref(),
346            );
347            if request.recipient.is_some() {
348                if !requester_is_local && !request.requester.is_bot() {
349                    return Err(SendError::InvalidRequest(
350                        "a direct retransmission recipient is only valid for a local device or bot"
351                            .into(),
352                    ));
353                }
354            } else if requester_is_local {
355                return Err(SendError::InvalidRequest(
356                    "a direct retransmission to another local device requires a recipient".into(),
357                ));
358            }
359
360            let routing_chat = request.recipient.as_ref().unwrap_or(&request.requester);
361            if !self
362                .jids_share_user_identity(&request.chat, routing_chat)
363                .await
364                .map_err(SendError::from_anyhow)?
365            {
366                return Err(SendError::InvalidRequest(
367                    "direct retransmission chat does not match its routing identity".into(),
368                ));
369            }
370        }
371
372        let group_info = if matches!(route, RetransmissionRoute::Group) {
373            Some(
374                self.groups()
375                    .query_info_with_freshness(&request.chat, request.group_metadata_freshness)
376                    .await?,
377            )
378        } else {
379            None
380        };
381
382        let encryption_jid = self
383            .resolve_retransmission_encryption_jid(route, &request.requester)
384            .await
385            .map_err(SendError::from_anyhow)?;
386        if route.uses_sender_key() {
387            let chat_key = request.chat.to_string();
388            self.mark_forget_sender_key(&chat_key, std::slice::from_ref(&encryption_jid))
389                .await
390                .map_err(SendError::from_anyhow)?;
391        }
392
393        let MessageRetransmission {
394            chat,
395            requester: wire_requester,
396            message,
397            message_id,
398            retry_count,
399            recipient,
400            group_metadata_freshness: _,
401        } = request;
402        let pre_encoded = Arc::new(waproto::codec::message_to_vec(&message));
403        self.add_recent_message(&chat, &message_id, &message, Some(Arc::clone(&pre_encoded)))
404            .await;
405        self.retransmit_message_prepared(PreparedRetransmission {
406            route,
407            wire_requester,
408            encryption_jid,
409            chat,
410            message,
411            message_id,
412            retry_count,
413            recipient,
414            group_info,
415            pre_encoded: Some(pre_encoded),
416        })
417        .await
418        .map_err(SendError::from_anyhow)
419    }
420
421    /// Handle an inbound `<receipt type="retry">`.
422    ///
423    /// WA Web authorizes these through `isRetryEligible` (`WAWebApiMessageInfoStore`).
424    /// We enforce the reject reasons that need no per-recipient state:
425    /// `HIGH_RETRY_COUNT` (the `MAX_RETRY_COUNT` refusal), `MESSAGE_EXPIRED` /
426    /// `RECORD_MISSING` (the recent-message cache miss), and `DEVICE_NOT_IN_DATABASE`
427    /// (`should_drop_unknown_device_retry`); identity changes are handled during
428    /// repair (reg-id mismatch + base-key collision in `update_local_signal_session`).
429    /// `ALREADY_DELIVERED` and `DEVICE_NOT_RECIPIENT` need a per-(message, device)
430    /// receipt store we do not keep, so they are a known parity gap, not enforced here.
431    #[cfg_attr(feature = "tracing", tracing::instrument(name = "wa.retry.handle_receipt", level = "debug", skip_all, fields(chat = %receipt.source.chat.observe(), sender = %receipt.source.sender.observe(), count = tracing::field::Empty), err(Debug)))]
432    pub(crate) async fn handle_retry_receipt(
433        self: &Arc<Self>,
434        receipt: &Receipt,
435        node: &Arc<OwnedNodeRef>,
436    ) -> Result<(), anyhow::Error> {
437        let nr = node.get();
438        let retry_child = nr
439            .get_optional_child("retry")
440            .ok_or_else(|| anyhow::anyhow!("<retry> child missing from receipt"))?;
441
442        let message_id = retry_child
443            .get_attr("id")
444            .map(|v| v.as_str())
445            .ok_or_else(|| anyhow::anyhow!("<retry> missing 'id' attribute"))?
446            .into_owned();
447        let retry_count: u8 = retry_child
448            .get_attr("count")
449            .map(|v| v.as_str())
450            .and_then(|s| s.parse().ok())
451            .unwrap_or(1);
452        // Record the count on the span so retry-storm depth is aggregable per
453        // sender even when the cap refuses early below.
454        #[cfg(feature = "tracing")]
455        tracing::Span::current().record("count", retry_count);
456
457        // Refuse to handle retries that have exceeded the maximum attempts.
458        // This prevents infinite retry loops and matches WhatsApp Web's behavior.
459        // Logged at debug: remote-driven, expected and fully handled — WA Web
460        // emits this refusal via WALogger.LOG (informational), not WARN.
461        if retry_count >= MAX_RETRY_COUNT {
462            debug!(
463                "Refusing retry #{} for message {} from {}: exceeds max attempts ({})",
464                retry_count,
465                message_id,
466                receipt.source.sender.observe(),
467                MAX_RETRY_COUNT
468            );
469            wacore::telemetry::retry_refused();
470            return Ok(());
471        }
472
473        let device_snapshot = self.persistence_manager.get_device_snapshot();
474        let Some(mut info) = resolve_retry_chat_info(
475            receipt,
476            nr,
477            device_snapshot.pn.as_ref(),
478            device_snapshot.lid.as_ref(),
479        ) else {
480            return Ok(());
481        };
482        let route = match validate_retransmission(
483            &info.chat,
484            &info.requester,
485            &message_id,
486            retry_count,
487            info.recipient.as_ref(),
488        ) {
489            Ok(route) => route,
490            Err(error) => {
491                debug!("Ignoring malformed retry request: {error}");
492                return Ok(());
493            }
494        };
495        let uses_sender_key = route.uses_sender_key();
496
497        // WA Web doesn't dedupe receipts (Message/Queue.js just serializes per-chat);
498        // MAX_RETRY_COUNT covers loop prevention. This lock only guards against
499        // two concurrent receipts racing on session state.
500        let processing_key = build_retry_processing_key(&info.chat, &message_id, &info.requester);
501
502        if !self
503            .pending_retries
504            .lock()
505            .unwrap_or_else(|p| p.into_inner())
506            .insert(processing_key.clone())
507        {
508            log::debug!("Ignoring retry for {processing_key}: a retry is already in progress.");
509            return Ok(());
510        }
511        // processing_key isn't needed by name after this point — move it into
512        // the scopeguard instead of cloning again.
513        let pending = Arc::clone(&self.pending_retries);
514        let _guard = scopeguard::guard((), move |()| {
515            pending
516                .lock()
517                .unwrap_or_else(|p| p.into_inner())
518                .remove(&processing_key);
519        });
520
521        // A retry from a device missing from our registry signals a stale device
522        // list for this user, so refresh it (rate-limited, dedup'd) to learn the
523        // device for the next send. Done before the message-cache lookup so an
524        // evicted retry still triggers it.
525        let sender_device_id = info.requester.device() as u32;
526        let device_known = self
527            .has_device(&info.requester.user, sender_device_id)
528            .await;
529        if !device_known {
530            // Parity with WA Web's MdRetryFromUnknownDevice WAM (id 2178), which
531            // commits here only — not from the shared inbound device sync, which
532            // schedule_unknown_device_sync is also called from elsewhere.
533            wacore::telemetry::retry_unknown_device(if sender_device_id == 0 {
534                "primary"
535            } else {
536                "companion"
537            });
538            self.schedule_unknown_device_sync(info.requester.to_non_ad(), receipt.offline)
539                .await;
540        }
541
542        // Peek keeps the message in the cache, so we avoid the decode + re-encode
543        // and the background DB delete + re-store that take + re-add did on every
544        // retry (pure churn during retry storms). Fall back to the consuming take +
545        // re-add only on an L1 miss (DB-only mode, or after eviction), where peek
546        // can't serve it; that path still re-adds so other devices can retry.
547        let (original_msg, alt_chat) = match self.peek_recent_message(&info.chat, &message_id).await
548        {
549            Some(result) => result,
550            None => match self.take_recent_message(&info.chat, &message_id).await {
551                Some(result) => {
552                    self.add_recent_message(&info.chat, &message_id, &result.0, None)
553                        .await;
554                    result
555                }
556                None => {
557                    log::debug!(
558                        "Ignoring retry for message {message_id}: already handled or not found in cache."
559                    );
560                    return Ok(());
561                }
562            },
563        };
564
565        // When message was found via alternate PN<->LID key, the Signal session
566        // lives in the stored message's namespace (not the receipt's). Build the
567        // encryption JID from that namespace + requester's device, skipping
568        // resolve_encryption_jid (which would map back to the primary namespace).
569        // WA Web: `e.from.isBot() ? (p = e.from) : (p = d.isLid() ? toLid(e.from) : toPn(e.from))`
570        // Bots skip namespace normalization (WAWebHandleRetryRequest:311-312).
571        let resolved_jid = if let Some(alt_chat) = alt_chat
572            && !uses_sender_key
573            && !info.is_bot
574        {
575            let requester = &info.requester;
576            info.requester = Jid {
577                user: alt_chat.user,
578                server: alt_chat.server,
579                device: requester.device,
580                agent: requester.agent,
581                integrator: requester.integrator,
582            };
583            info.requester.clone()
584        } else {
585            self.resolve_retransmission_encryption_jid(route, &info.requester)
586                .await?
587        };
588
589        let keys_node_present = nr.get_optional_child("keys").is_some();
590        if wacore::protocol::retry::should_drop_unknown_device_retry(
591            keys_node_present,
592            device_known,
593        ) {
594            warn!(
595                "handle_retry_receipt: device not found for device={}, user={}",
596                sender_device_id, info.requester.user
597            );
598            return Ok(());
599        }
600
601        // Check if this is a retry from our own device (peer).
602        let is_peer = is_own_account_jid(
603            &info.requester,
604            device_snapshot.pn.as_ref(),
605            device_snapshot.lid.as_ref(),
606        );
607
608        // Volume-throttling inbound retries diverges from WA Web (which
609        // processes every receipt), so it is an operator opt-in, gated here
610        // before the expensive repair stages below. Own devices (`is_peer`) and
611        // DMs are never gated: dropping their retries has no safe SKDM fallback.
612        if uses_sender_key
613            && !is_peer
614            && let Some(policy) = self.retry_admission.get()
615            && !policy.admit(&info.chat, &info.requester, retry_count)
616        {
617            debug!(
618                "Retry receipt from {} in {} dropped by RetryAdmission policy",
619                info.requester.observe(),
620                info.chat.observe()
621            );
622            return Ok(());
623        }
624
625        // Fetch group info (cache-first, server on miss) — used for SKDM rotation + addressing_mode.
626        // Without this, a cold cache would silently default to PN semantics for LID groups.
627        let cached_group_info = if info.chat.is_group() {
628            match self.groups().query_info(&info.chat).await {
629                Ok(gi) => Some(gi),
630                Err(e) => {
631                    log::warn!(
632                        "Failed to fetch group info for retry of msg {} in {}: {e}",
633                        message_id,
634                        info.chat.observe()
635                    );
636                    None
637                }
638            }
639        } else {
640            None
641        };
642
643        // WA Web rotateKey: unknown device (not in participant list, not LID) →
644        // force full sender key rotation by clearing all sender key device tracking.
645        // This is separate from updateLocalSignalSession and specific to group retries.
646        let mut rotated_sender_key = false;
647        if matches!(route, RetransmissionRoute::Group) && !info.requester.is_lid() {
648            let group_jid = info.chat.to_string();
649            let is_known_participant = cached_group_info
650                .as_ref()
651                .is_some_and(|g| g.participants.iter().any(|p| p.user == info.requester.user));
652
653            if !is_known_participant {
654                log::warn!(
655                    "Unknown device {} in group {} — forcing full sender key rotation \
656                     (matches WA Web's rotateKey behavior)",
657                    info.requester.observe(),
658                    group_jid
659                );
660                let _distribution_guard = self.group_distribution_lock(&info.chat).await;
661
662                // WA Web: deleteGroupSenderKeyInfo(groupWid, ownWid) — delete our own
663                // sender key for forward secrecy. When addressing mode is known,
664                // delete only that namespace; otherwise both.
665                let addressing_mode = cached_group_info.as_ref().map(|g| g.addressing_mode);
666                let jids_to_delete: Vec<_> = match addressing_mode {
667                    Some(wacore::types::message::AddressingMode::Lid) => {
668                        device_snapshot.lid.as_ref().into_iter().collect()
669                    }
670                    Some(wacore::types::message::AddressingMode::Pn) => {
671                        device_snapshot.pn.as_ref().into_iter().collect()
672                    }
673                    None => device_snapshot
674                        .lid
675                        .as_ref()
676                        .into_iter()
677                        .chain(device_snapshot.pn.as_ref())
678                        .collect(),
679                };
680
681                for own_jid in jids_to_delete {
682                    use wacore::libsignal::store::sender_key_name::SenderKeyName;
683                    let sk_name = SenderKeyName::from_parts(
684                        &group_jid,
685                        own_jid.to_protocol_address().as_str(),
686                    );
687                    self.signal_cache
688                        .delete_sender_key(sk_name.cache_key())
689                        .await;
690                }
691
692                // DB first, then cache invalidate — prevents a concurrent
693                // resolve_skdm_targets from reviving stale cache entries.
694                if let Err(e) = self.reset_sender_key_device_tracking(&group_jid).await {
695                    log::warn!("Failed to clear sender key devices for rotation: {}", e);
696                }
697                rotated_sender_key = true;
698            }
699        }
700        if rotated_sender_key {
701            self.flush_signal_cache_batch_safe_logged("unknown-participant rotation", None)
702                .await;
703        }
704
705        // Mirror WAWebUpdateLocalSignalSession for all chat types: markForgetSenderKey
706        // (group/status) + processKeyBundle + regId-mismatch delete + base-key logic.
707        // Must run before ensureE2ESessions so any session deletion here is rebuilt there.
708        if !self
709            .update_local_signal_session(
710                &info,
711                &resolved_jid,
712                &message_id,
713                retry_count,
714                nr,
715                is_peer,
716            )
717            .await
718        {
719            return Ok(());
720        }
721
722        // Whatsmeow parity (`retry.go:284`). WA Web's regId/base-key check
723        // doesn't catch silently-diverged sessions; this fallback does.
724        if nr.get_optional_child("keys").is_none() {
725            // Hold the per-peer session lock across the throttle check+stamp AND
726            // the delete so the recreate decision is atomic per peer. The
727            // `session_recreate_history` get+insert is not atomic on its own,
728            // and retry receipts for different message_ids from the same peer
729            // dispatch concurrently (detached spawn in `handle_receipt`), so
730            // without this lock two of them could both pass the throttle and
731            // recreate. Mirrors whatsmeow holding `sessionRecreateHistoryLock`
732            // across its check+stamp (`retry.go:160`). This is the same per-peer
733            // lock the delete already used, so it adds no new lock.
734            let signal_address = resolved_jid.to_protocol_address();
735            let lock = self.session_lock_for(signal_address.as_str()).await;
736            let guard = lock.lock().await;
737            if let Some(reason) = self
738                .should_recreate_session(retry_count, &resolved_jid)
739                .await
740            {
741                info!(
742                    "Recreating session with {} for retry of {message_id}: {reason}",
743                    resolved_jid.observe()
744                );
745                self.signal_cache.delete_session(&signal_address).await;
746                drop(guard);
747                self.flush_signal_cache_batch_safe_logged(
748                    "should_recreate_session",
749                    Some(&message_id),
750                )
751                .await;
752            }
753        }
754
755        // Bound the aggregate resend rate per group (the anti-abuse signal): a
756        // PN to LID fan-out has many distinct devices retry the same messages,
757        // which per-device/per-message caps miss. Group-only: the requester was
758        // marked for fresh SKDM above so future messages recover, and it
759        // re-requests this one on its own timer once the bucket refills. DMs have
760        // no SKDM fallback, so they keep the unconditional resend (bounded by
761        // MAX_RETRY_COUNT) rather than risk dropping a delivery.
762        if info.chat.is_group() && !self.resend_rate_limiter.try_acquire(&info.chat).await {
763            debug!(
764                "Throttling resend of {} to {}: per-chat resend rate cap reached",
765                message_id,
766                info.chat.observe()
767            );
768            return Ok(());
769        }
770
771        info!(
772            "Resending message {} to {} (retry #{})",
773            message_id,
774            info.chat.observe(),
775            retry_count
776        );
777
778        let wire_requester = if matches!(route, RetransmissionRoute::Direct) {
779            info.original_from
780        } else {
781            info.requester
782        };
783        self.retransmit_message_prepared(PreparedRetransmission {
784            route,
785            chat: info.chat,
786            wire_requester,
787            encryption_jid: resolved_jid,
788            message: original_msg,
789            message_id,
790            retry_count,
791            recipient: info.recipient,
792            group_info: cached_group_info,
793            pre_encoded: None,
794        })
795        .await?;
796
797        Ok(())
798    }
799
800    async fn send_retry_stanza(&self, stanza: Node) -> Result<(), anyhow::Error> {
801        self.persist_signal_state_pre_wire().await?;
802        self.send_node(stanza).await?;
803        Ok(())
804    }
805
806    async fn retransmit_message_prepared(
807        &self,
808        request: PreparedRetransmission,
809    ) -> Result<(), anyhow::Error> {
810        let PreparedRetransmission {
811            route,
812            chat,
813            wire_requester,
814            encryption_jid,
815            message,
816            message_id,
817            retry_count,
818            recipient,
819            group_info,
820            pre_encoded,
821        } = request;
822
823        if matches!(route, RetransmissionRoute::Status) {
824            return self
825                .retransmit_status_message(
826                    chat,
827                    encryption_jid,
828                    message,
829                    message_id,
830                    pre_encoded.as_deref().map(Vec::as_slice),
831                )
832                .await;
833        }
834
835        // Every remaining route is pairwise, including broadcast-list
836        // participants, and shares the normal session recovery path.
837        self.ensure_e2e_sessions_resolved(std::slice::from_ref(&encryption_jid))
838            .await?;
839        let signal_address = encryption_jid.to_protocol_address();
840        let session_mutex = self.session_lock_for(signal_address.as_str()).await;
841        let session_guard = session_mutex.lock().await;
842        let mut store_adapter = self.signal_adapter().await;
843        let device_snapshot = self.persistence_manager.get_device_snapshot();
844        let edit = wacore::types::message::EditAttribute::infer_from_message(&message);
845
846        let destination = match route {
847            RetransmissionRoute::Direct => wacore::send::PairwiseRetryDestination::Direct {
848                to: wire_requester,
849                recipient,
850            },
851            RetransmissionRoute::Group => {
852                let addressing_mode = group_info
853                    .as_ref()
854                    .map(|info| info.addressing_mode)
855                    .unwrap_or_default();
856                wacore::send::PairwiseRetryDestination::Participant {
857                    to: chat,
858                    participant: wire_requester,
859                    addressing_mode: Some(addressing_mode),
860                }
861            }
862            RetransmissionRoute::BroadcastList => {
863                wacore::send::PairwiseRetryDestination::Participant {
864                    to: chat,
865                    participant: wire_requester,
866                    addressing_mode: None,
867                }
868            }
869            RetransmissionRoute::Status => unreachable!("status handled above"),
870        };
871        let stanza = wacore::send::prepare_pairwise_retry_stanza(
872            &mut store_adapter.session_store,
873            &mut store_adapter.identity_store,
874            wacore::send::PairwiseRetryRequest {
875                destination,
876                encryption_jid,
877                message: &message,
878                message_id,
879                retry_count,
880                account: device_snapshot.account.as_deref(),
881                edit,
882                pre_encoded: pre_encoded.as_deref().map(Vec::as_slice),
883            },
884        )
885        .await?;
886
887        // Persistence may need the processing permit, whose holder may in turn
888        // need this session lock. Release it before the durability gate.
889        drop(session_guard);
890        self.send_retry_stanza(stanza).await
891    }
892
893    /// Rebuild a status message for exactly the requesting device. The retry
894    /// count remains an operation-level guard; the captured status wire does not
895    /// encode it on either the skmsg or SKDM `<enc>` node.
896    async fn retransmit_status_message(
897        &self,
898        chat: Jid,
899        requester: Jid,
900        message: wa::Message,
901        message_id: String,
902        pre_encoded: Option<&[u8]>,
903    ) -> Result<(), anyhow::Error> {
904        let snapshot = self.persistence_manager.get_device_snapshot();
905        let own_pn = snapshot
906            .pn
907            .as_ref()
908            .ok_or(crate::client::ClientError::NotLoggedIn)?;
909        let own_lid = snapshot
910            .lid
911            .as_ref()
912            .ok_or_else(|| anyhow::anyhow!("cannot retransmit status without a device LID"))?;
913        let is_sending_device = (requester.is_same_user_as(own_pn)
914            && requester.device == own_pn.device)
915            || (requester.is_same_user_as(own_lid) && requester.device == own_lid.device);
916        if is_sending_device {
917            anyhow::bail!("cannot retransmit a status to the sending device itself");
918        }
919
920        let chat_key = chat.to_string();
921        let distribution_guard = self.group_distribution_lock(&chat).await;
922        let group_info = wacore::client::context::GroupInfo::new(
923            Vec::new(),
924            wacore::types::message::AddressingMode::Lid,
925        );
926
927        let can_reuse_encoding = message.message_context_info.is_unset();
928        let encoded_fallback = (pre_encoded.is_none() && can_reuse_encoding)
929            .then(|| waproto::codec::message_to_vec(&message));
930        let encoded = pre_encoded
931            .filter(|_| can_reuse_encoding)
932            .or(encoded_fallback.as_deref());
933        let device_store = self.persistence_manager.get_device_arc().await;
934        let mut store_adapter = self.signal_adapter_from(device_store);
935        let mut stores = store_adapter.as_signal_stores();
936        let edit = wacore::types::message::EditAttribute::infer_from_message(&message);
937        let prepared = match wacore::send::prepare_group_stanza(
938            &*self.runtime,
939            &mut stores,
940            self,
941            wacore::send::GroupStanzaRequest {
942                group: &group_info,
943                own_jid: own_pn,
944                own_lid,
945                account: snapshot.account.as_deref(),
946                to: &chat,
947                message: &message,
948                message_id: &message_id,
949                force_distribution: false,
950                distribution_targets: Some(vec![requester]),
951                distribution_policy: wacore::send::SenderKeyDistributionPolicy::Required,
952                phash_devices: None,
953                edit: edit.as_ref(),
954                extra_nodes: &[],
955                pre_encoded: encoded,
956            },
957        )
958        .await
959        {
960            Ok(prepared) => prepared,
961            Err(error) => {
962                // Do not hold the sender-key distribution lane across registry
963                // I/O. The typed failure retains the original source chain and
964                // identifies only users whose pre-key lookup returned 406.
965                drop(distribution_guard);
966                if let Some(failure) =
967                    error.downcast_ref::<wacore::send::RequiredSenderKeyDistributionError>()
968                {
969                    for user in failure.stale_device_users() {
970                        self.invalidate_device_cache(user).await;
971                    }
972                }
973                return Err(error);
974            }
975        };
976        self.send_retry_stanza(prepared.node).await?;
977        self.update_sender_key_devices(&chat_key, &prepared.skdm_devices)
978            .await;
979        drop(distribution_guard);
980        for user in &prepared.stale_device_users {
981            self.invalidate_device_cache(user).await;
982        }
983        Ok(())
984    }
985
986    /// Mirrors WAWebUpdateLocalSignalSession (`WAWeb/Update/LocalSignalSession.js`).
987    /// Runs before ensureE2ESessions + sendRetry for all chat types (DM, group,
988    /// status). Order and semantics match the WA Web implementation:
989    ///   1. markForgetSenderKey for group/status (participant needs fresh SKDM)
990    ///   2. processKeyBundle if `<keys>` present
991    ///   3. If no bundle AND stored regId differs → delete session
992    ///   4. retry == 2 → save current base key, return (no delete)
993    ///   5. retry > 2 AND same base key → delete session (force re-establish)
994    ///
995    /// Unlike the previous DM-only path, this does NOT unconditionally delete
996    /// the session on every retry — WA Web preserves it on retry==1 and on
997    /// retry>2 when the base key already changed (session was regenerated
998    /// legitimately). The subsequent `ensure_e2e_sessions_resolved` call in
999    /// `handle_retry_receipt` rebuilds any session this function deleted.
1000    #[cfg_attr(feature = "tracing", tracing::instrument(name = "wa.retry.update_local_session", level = "debug", skip_all, fields(chat = %info.chat.observe(), peer = %resolved_jid.observe(), retry = retry_count)))]
1001    async fn update_local_signal_session(
1002        &self,
1003        info: &RetryChatInfo,
1004        resolved_jid: &Jid,
1005        message_id: &str,
1006        retry_count: u8,
1007        node: &NodeRef<'_>,
1008        is_peer: bool,
1009    ) -> bool {
1010        // 1. markForgetSenderKey (WA Web L33-38). Rust unifies group and status
1011        //    under a single storage (chat JID as the key) — markForgetSenderKey
1012        //    handles both `@g.us` and `status@broadcast` as opaque group_jid.
1013        if info.chat.is_group() || info.chat.is_status_broadcast() {
1014            let group_jid = info.chat.to_string();
1015            match self
1016                .mark_forget_sender_key(&group_jid, std::slice::from_ref(resolved_jid))
1017                .await
1018            {
1019                Ok(()) => {
1020                    let chat_type = if info.chat.is_status_broadcast() {
1021                        "status broadcast"
1022                    } else {
1023                        "group"
1024                    };
1025                    // debug, not info: one line per retry receipt, and a broken
1026                    // cohort in a large group emits tens of thousands per day
1027                    // (WA Web logs the same event at its verbose LOG level).
1028                    debug!(
1029                        "Marked {} for fresh SKDM in {} {} due to retry receipt",
1030                        resolved_jid.observe(),
1031                        chat_type,
1032                        group_jid
1033                    );
1034                }
1035                Err(e) => log::warn!(
1036                    "Failed to mark sender key forget for {} in {}: {}",
1037                    info.requester.observe(),
1038                    group_jid,
1039                    e
1040                ),
1041            }
1042        }
1043
1044        // 2. processKeyBundle (WA Web L51). Previously gated behind
1045        //    `!is_status_broadcast()`; WA Web runs it unconditionally.
1046        let keys_node_present = node.get_optional_child("keys").is_some();
1047        let key_bundle_result = self
1048            .process_retry_key_bundle(node, resolved_jid, is_peer, info.is_fbid_bot_retry)
1049            .await;
1050        let key_bundle_processed = key_bundle_result.is_ok();
1051
1052        // 3. No bundle + regId mismatch → delete session (WA Web L52-65).
1053        //    Gate on `!keys_node_present` so a rejected bundle (security
1054        //    refusal for peer reg-ID change, parse errors, invalid reg ID)
1055        //    doesn't trigger destructive session deletion as a side effect.
1056        if !key_bundle_processed && keys_node_present {
1057            log::warn!(
1058                "Key bundle present but rejected for {}: {:?} — aborting retry resend",
1059                resolved_jid.observe(),
1060                key_bundle_result.as_ref().err()
1061            );
1062            return false;
1063        }
1064        if !key_bundle_processed && !keys_node_present {
1065            if let Err(ref e) = key_bundle_result {
1066                // Demoted to debug on the happy path (peer retry without re-key):
1067                // only warn when a regId mismatch triggers a delete below.
1068                log::debug!(
1069                    "No key bundle in retry receipt for {}: {}. Checking for reg ID mismatch.",
1070                    resolved_jid.observe(),
1071                    e
1072                );
1073            }
1074
1075            if let Some(received_reg_id) =
1076                wacore::protocol::retry::extract_registration_id_from_node_ref(node)
1077            {
1078                let signal_address = resolved_jid.to_protocol_address();
1079                let device_snapshot = self.persistence_manager.get_device_snapshot();
1080                let session = self
1081                    .signal_cache
1082                    .peek_session(&signal_address, &*device_snapshot.backend)
1083                    .await
1084                    .ok()
1085                    .flatten();
1086
1087                if let Some(session) = session
1088                    && let Ok(stored_reg_id) = session.remote_registration_id()
1089                    && stored_reg_id != 0
1090                    && stored_reg_id != received_reg_id
1091                {
1092                    info!(
1093                        "Registration ID mismatch for {} (stored: {}, received: {}). \
1094                         Deleting session since no key bundle provided.",
1095                        wacore::types::jid::observe_protocol_address(&signal_address),
1096                        stored_reg_id,
1097                        received_reg_id
1098                    );
1099                    let lock = self.session_lock_for(signal_address.as_str()).await;
1100                    let _guard = lock.lock().await;
1101                    self.signal_cache.delete_session(&signal_address).await;
1102                    drop(_guard);
1103                    self.flush_signal_cache_batch_safe_logged(
1104                        "reg ID mismatch session deletion",
1105                        None,
1106                    )
1107                    .await;
1108                }
1109            }
1110        }
1111
1112        // 4-5. Base-key collision logic (WA Web L66-80). Applied to ALL chat
1113        //      types now — previously only ran in the DM branch.
1114        let signal_address = resolved_jid.to_protocol_address();
1115        let device_snapshot = self.persistence_manager.get_device_snapshot();
1116        let session = self
1117            .signal_cache
1118            .peek_session(&signal_address, &*device_snapshot.backend)
1119            .await
1120            .ok()
1121            .flatten();
1122
1123        let Some(session) = session else {
1124            return true;
1125        };
1126        let Ok(current_base_key) = session.alice_base_key() else {
1127            return true;
1128        };
1129
1130        let addr_str = signal_address.as_str();
1131        if retry_count == MIN_RETRY_FOR_BASE_KEY_CHECK {
1132            // retry == 2: save base key, do NOT delete (WA Web L66-67).
1133            match device_snapshot
1134                .backend
1135                .save_base_key(addr_str, message_id, current_base_key)
1136                .await
1137            {
1138                Ok(()) => info!(
1139                    "Saved base key for {} at retry #{} for collision detection",
1140                    wacore::types::jid::observe_protocol_address(&signal_address),
1141                    retry_count
1142                ),
1143                Err(e) => warn!(
1144                    "Failed to save base key for {}: {}",
1145                    wacore::types::jid::observe_protocol_address(&signal_address),
1146                    e
1147                ),
1148            }
1149            return true;
1150        }
1151
1152        if retry_count > MIN_RETRY_FOR_BASE_KEY_CHECK {
1153            match device_snapshot
1154                .backend
1155                .has_same_base_key(addr_str, message_id, current_base_key)
1156                .await
1157            {
1158                Ok(true) => {
1159                    // Informational, not WARN: this is the corrective action WA
1160                    // Web takes here too (WAWebUpdateLocalSignalSession logs the
1161                    // same-base-key delete via WALogger.LOG), and the three
1162                    // sibling branches of this routine already log at info.
1163                    info!(
1164                        "Base key collision detected for {} (msg {}) at retry #{}. \
1165                         Session hasn't been regenerated. Forcing fresh session.",
1166                        wacore::types::jid::observe_protocol_address(&signal_address),
1167                        message_id,
1168                        retry_count
1169                    );
1170                    wacore::telemetry::base_key_collision();
1171                    let _ = device_snapshot
1172                        .backend
1173                        .delete_base_key(addr_str, message_id)
1174                        .await;
1175                    let lock = self.session_lock_for(signal_address.as_str()).await;
1176                    let _guard = lock.lock().await;
1177                    self.signal_cache.delete_session(&signal_address).await;
1178                    drop(_guard);
1179                    self.flush_signal_cache_batch_safe_logged(
1180                        "base key collision — forcing fresh session",
1181                        None,
1182                    )
1183                    .await;
1184                }
1185                Ok(false) => {
1186                    info!(
1187                        "Base key changed for {} (msg {}) at retry #{} - session regenerated",
1188                        wacore::types::jid::observe_protocol_address(&signal_address),
1189                        message_id,
1190                        retry_count
1191                    );
1192                    let _ = device_snapshot
1193                        .backend
1194                        .delete_base_key(addr_str, message_id)
1195                        .await;
1196                }
1197                Err(e) => {
1198                    warn!(
1199                        "Failed to check base key for {}: {}",
1200                        wacore::types::jid::observe_protocol_address(&signal_address),
1201                        e
1202                    );
1203                }
1204            }
1205        }
1206        true
1207    }
1208
1209    /// Mirrors whatsmeow's `shouldRecreateSession`. Returns `Some(reason)`
1210    /// and bumps the history clock if we should drop the local session for
1211    /// `jid`; `None` otherwise. Two conditions trigger:
1212    ///   1. No session present locally.
1213    ///   2. `retry_count >= 2` and >`RECREATE_SESSION_TIMEOUT` since the
1214    ///      last recreate for this JID.
1215    ///
1216    /// Callers pair this with `signal_cache.delete_session` so the next
1217    /// `ensure_e2e_sessions_resolved` does the prekey fetch + rebuild.
1218    async fn should_recreate_session(&self, retry_count: u8, jid: &Jid) -> Option<&'static str> {
1219        self.should_recreate_session_at(retry_count, jid, wacore::time::Instant::now())
1220            .await
1221    }
1222
1223    /// Injectable-clock variant for testing the throttle expiry path.
1224    /// wacore::time::Instant is std::time::Instant-backed so subtracting a
1225    /// Duration to fabricate a "past" stamp saturates to 0 in young test
1226    /// runtimes; passing a future `now` instead exercises the same branch.
1227    async fn should_recreate_session_at(
1228        &self,
1229        retry_count: u8,
1230        jid: &Jid,
1231        now: wacore::time::Instant,
1232    ) -> Option<&'static str> {
1233        let signal_address = jid.to_protocol_address();
1234        let device_snapshot = self.persistence_manager.get_device_snapshot();
1235        // Whatsmeow returns `false` on `ContainsSession` errors so a transient
1236        // backend read failure doesn't masquerade as "no session" and trigger
1237        // an unnecessary delete + prekey fetch (`retry.go:161-163`).
1238        let has_session = match self
1239            .signal_cache
1240            .has_session(&signal_address, &*device_snapshot.backend)
1241            .await
1242        {
1243            Ok(present) => present,
1244            Err(e) => {
1245                warn!(
1246                    "should_recreate_session: has_session failed for {}: {} — skipping recreate",
1247                    signal_address, e
1248                );
1249                return None;
1250            }
1251        };
1252
1253        let history = &self.session_recreate_history;
1254
1255        if !has_session {
1256            history.insert(jid.clone(), now).await;
1257            return Some("we don't have a Signal session with them");
1258        }
1259
1260        if retry_count < MIN_RETRY_FOR_BASE_KEY_CHECK {
1261            return None;
1262        }
1263
1264        // Throttle: skip if this peer was recreated within the timeout. This
1265        // explicit age check against the injectable `now` is the authoritative,
1266        // deterministic gate. The cache's 1h TTL on `session_recreate_history`
1267        // is only a memory backstop (lazy eviction independent of the stored
1268        // `now`, so it can't drive the throttle decision).
1269        // Do NOT drop this check as "redundant with the TTL".
1270        if let Some(prev) = history.get(jid).await
1271            && now.saturating_duration_since(prev) < RECREATE_SESSION_TIMEOUT
1272        {
1273            return None;
1274        }
1275
1276        history.insert(jid.clone(), now).await;
1277        Some("retry count > 1 and over an hour since last recreation")
1278    }
1279
1280    /// Extracts and processes the key bundle from a retry receipt.
1281    /// This allows us to establish a new session with the requester using their fresh prekeys.
1282    ///
1283    /// # Arguments
1284    /// * `node` - The retry receipt node containing the key bundle
1285    /// * `requester_jid` - The JID of the device requesting the retry
1286    /// * `is_peer` - Whether this is a peer device (our own device)
1287    #[cfg_attr(feature = "tracing", tracing::instrument(name = "wa.retry.process_key_bundle", level = "debug", skip_all, fields(peer = %requester_jid.observe(), is_peer, is_fbid_bot_retry), err(Debug)))]
1288    async fn process_retry_key_bundle(
1289        &self,
1290        node: &NodeRef<'_>,
1291        requester_jid: &Jid,
1292        is_peer: bool,
1293        is_fbid_bot_retry: bool,
1294    ) -> Result<(), anyhow::Error> {
1295        let keys_node = node
1296            .get_optional_child("keys")
1297            .ok_or_else(|| anyhow::anyhow!("<keys> child missing from retry receipt"))?;
1298        validate_retry_prekey_presence(keys_node, is_fbid_bot_retry)?;
1299
1300        // Use the centralized extractor so the >4-byte rejection rule applies
1301        // here too, not just on the no-keys retry path.
1302        let registration_id =
1303            wacore::protocol::retry::extract_registration_id_from_node_ref(node).unwrap_or(0);
1304
1305        if registration_id == 0 {
1306            return Err(anyhow::anyhow!("Invalid registration ID in retry receipt"));
1307        }
1308
1309        // Use requester_jid directly — the caller already resolved the correct
1310        // namespace (including alternate PN/LID normalization). Re-resolving
1311        // here would undo that normalization.
1312        let signal_address = requester_jid.to_protocol_address();
1313
1314        // Check if the registration ID changed (indicates device reinstall).
1315        // Read session through cache for consistent state.
1316        {
1317            let device_snapshot = self.persistence_manager.get_device_snapshot();
1318            let session = self
1319                .signal_cache
1320                .peek_session(&signal_address, &*device_snapshot.backend)
1321                .await
1322                .ok()
1323                .flatten();
1324
1325            if let Some(session) = session {
1326                let existing_reg_id = session.remote_registration_id()?;
1327                if existing_reg_id != 0 && existing_reg_id != registration_id {
1328                    // WhatsApp Web throws an error for peer device registration ID changes.
1329                    // This is a security measure - peer devices should maintain consistent identity.
1330                    if is_peer {
1331                        return Err(anyhow::anyhow!(
1332                            "Registration ID changed for peer device {} (was {}, now {}). \
1333                             This may indicate the device was reinstalled.",
1334                            signal_address,
1335                            existing_reg_id,
1336                            registration_id
1337                        ));
1338                    }
1339                    info!(
1340                        "Registration ID changed for {} (was {}, now {}). Session will be replaced.",
1341                        signal_address, existing_reg_id, registration_id
1342                    );
1343                }
1344            }
1345        }
1346
1347        // Extract identity key.
1348        let identity_bytes = keys_node
1349            .get_optional_child("identity")
1350            .and_then(get_bytes_content_ref)
1351            .ok_or_else(|| anyhow::anyhow!("Missing identity key in retry receipt"))?;
1352        let identity_key = PublicKey::from_djb_public_key_bytes(identity_bytes)?;
1353
1354        // Companion devices ADV-bind the fetched identity via <device-identity>;
1355        // reject a present-but-invalid one so a relay can't swap in a forged key.
1356        // Mirrors the prekey-fetch path. The account key is the in-blob
1357        // `account_signature_key` or, when the server omits it, the contact's
1358        // primary (device 0) identity from the store. An unverifiable-for-lack-of-key
1359        // chain or a missing device-identity is logged, not fatal.
1360        if requester_jid.device != 0
1361            && let Some(device_identity) = keys_node
1362                .get_optional_child("device-identity")
1363                .and_then(get_bytes_content_ref)
1364        {
1365            let fetched_identity: [u8; 32] = identity_bytes
1366                .try_into()
1367                .map_err(|_| anyhow::anyhow!("identity key in retry receipt is not 32 bytes"))?;
1368            let account_identity = self.load_account_identity(requester_jid).await;
1369            match wacore::adv::validate_adv_with_identity_key(
1370                device_identity,
1371                &fetched_identity,
1372                account_identity.as_ref(),
1373            ) {
1374                wacore::adv::AdvValidation::Valid => {}
1375                wacore::adv::AdvValidation::Invalid => {
1376                    return Err(anyhow::anyhow!(
1377                        "device-identity ADV validation failed for companion {requester_jid}"
1378                    ));
1379                }
1380                wacore::adv::AdvValidation::NoAccountKey => log::debug!(
1381                    "retry key bundle for companion {requester_jid} omits account_signature_key and no stored account identity; proceeding without ADV validation"
1382                ),
1383            }
1384        } else if requester_jid.device != 0 {
1385            log::warn!(
1386                "retry key bundle for companion {requester_jid} omits <device-identity>; proceeding without ADV validation"
1387            );
1388        }
1389
1390        // Extract prekey (optional in some cases).
1391        let prekey_data = if let Some(key_ref) = keys_node.get_optional_child("key") {
1392            let prekey_node = OneTimePreKeyNode::try_from_node_ref(key_ref)?;
1393            let prekey_public = PublicKey::from_djb_public_key_bytes(&prekey_node.public_bytes)?;
1394            Some((prekey_node.id.into(), prekey_public))
1395        } else {
1396            None
1397        };
1398
1399        // Extract signed prekey.
1400        let skey_ref = keys_node
1401            .get_optional_child("skey")
1402            .ok_or_else(|| anyhow::anyhow!("Missing signed prekey in retry receipt"))?;
1403
1404        let signed_prekey = SignedPreKeyNode::try_from_node_ref(skey_ref)?;
1405        let skey_public = PublicKey::from_djb_public_key_bytes(&signed_prekey.public_bytes)?;
1406        let skey_signature: [u8; 64] = signed_prekey
1407            .signature
1408            .as_slice()
1409            .try_into()
1410            .map_err(|_| anyhow::anyhow!("Invalid signature length"))?;
1411
1412        // Build and process the prekey bundle.
1413        let bundle = PreKeyBundle::new(
1414            registration_id,
1415            u32::from(requester_jid.device).into(),
1416            prekey_data,
1417            signed_prekey.id.into(),
1418            skey_public,
1419            skey_signature.into(),
1420            identity_key.into(),
1421        )?;
1422
1423        let mut adapter = self.signal_adapter().await;
1424        let mut rng = rand::make_rng::<rand::rngs::StdRng>();
1425        self.install_prekey_bundle_cached(requester_jid, &bundle, &mut adapter, &mut rng)
1426            .await?;
1427
1428        self.flush_signal_cache_batch_safe().await?;
1429
1430        info!(
1431            "Processed key bundle from retry receipt for {}",
1432            signal_address
1433        );
1434
1435        Ok(())
1436    }
1437
1438    /// Sends a retry receipt to request the sender to resend a message.
1439    ///
1440    /// # Arguments
1441    /// * `info` - The message info for the failed message
1442    /// * `retry_count` - The retry attempt number (1-5). This is sent to the sender so they
1443    ///   know which attempt this is. The sender may use this to decide whether to resend.
1444    /// * `reason` - The retry reason code (matches WhatsApp Web's RetryReason enum). This helps
1445    ///   the sender understand why the message couldn't be decrypted.
1446    #[cfg_attr(feature = "tracing", tracing::instrument(name = "wa.retry.send_receipt", level = "debug", skip_all, fields(chat = %info.source.chat.observe(), sender = %info.source.sender.observe(), retry = retry_count), err(Debug)))]
1447    pub(crate) async fn send_retry_receipt(
1448        &self,
1449        info: &crate::types::message::MessageInfo,
1450        retry_count: u8,
1451        reason: RetryReason,
1452        force_include_keys: bool,
1453    ) -> Result<RetryReceiptSendOutcome, RetryRequestError> {
1454        let device_snapshot = self.persistence_manager.get_device_snapshot();
1455
1456        // WA Web's sendRetryReceipt aborts only when `!to.isBot() && participant.isBot()`,
1457        // with participant null for DMs. A bot DM is chat == sender == bot, so it is NOT
1458        // suppressed and the retry is sent; only a bot reply in a non-bot group is dropped.
1459        // Same helper the ack and self-fanout paths already use.
1460        if info.source.is_bot_authored_non_bot_chat() {
1461            log::debug!(
1462                "Skipping retry receipt for message {} from bot {} in non-bot chat {}",
1463                info.id,
1464                info.source.sender.observe(),
1465                info.source.chat.observe()
1466            );
1467            return Ok(RetryReceiptSendOutcome::Suppressed);
1468        }
1469
1470        debug!(
1471            "Sending retry receipt #{} for message {} in chat {} from {} (reason: {:?})",
1472            retry_count,
1473            info.id,
1474            info.source.chat.observe(),
1475            info.source.sender.observe(),
1476            reason
1477        );
1478
1479        // Build the retry element with the error code (matches WhatsApp Web's format)
1480        let mut retry_builder = NodeBuilder::new("retry")
1481            .attr("v", "1")
1482            .attr("id", info.id.clone())
1483            .attr("t", info.timestamp.timestamp())
1484            .attr("count", retry_count);
1485
1486        // Include the error code if it's not UnknownError (matches WhatsApp Web's behavior
1487        // where error is only included when there's a specific reason)
1488        if reason != RetryReason::UnknownError {
1489            retry_builder = retry_builder.attr("error", reason as u8);
1490        }
1491
1492        let retry_node = retry_builder.build();
1493
1494        let registration_id_bytes = device_snapshot.registration_id.to_be_bytes().to_vec();
1495        let registration_node = NodeBuilder::new("registration")
1496            .bytes(registration_id_bytes)
1497            .build();
1498
1499        let receipt_to = if info.source.is_group {
1500            &info.source.chat
1501        } else {
1502            &info.source.sender
1503        };
1504        let include_keys = wacore::protocol::retry::should_include_keys_with_policy(
1505            retry_count,
1506            force_include_keys,
1507            receipt_to.is_hosted(),
1508        );
1509
1510        let keys_node = if include_keys {
1511            // Validate the account BEFORE reserving/marking the prekey: a missing
1512            // account bails here, and marking after would abandon a one-time
1513            // prekey from the upload window without any receipt going out.
1514            let device_identity_bytes = waproto::codec::adv_signed_device_identity_to_vec(
1515                device_snapshot.account.as_deref().ok_or_else(|| {
1516                    anyhow::anyhow!("Missing device account info for retry receipt")
1517                })?,
1518            );
1519
1520            // markKeyAsUploaded: the retry prekey goes directly to the peer, so
1521            // it must not also be re-offered to the server pool (a third party
1522            // could consume the same one-time id and fail to decrypt). Hold
1523            // prekey_upload_lock so get-or-gen and the mark are one atomic step
1524            // against the batch upload path.
1525            let prekey_guard = self.prekey_upload_lock.lock().await;
1526            let (new_prekey_id, new_prekey_public) = self.get_or_gen_single_pre_key().await?;
1527            self.mark_single_prekey_uploaded(&prekey_guard, new_prekey_id)
1528                .await?;
1529            drop(prekey_guard);
1530
1531            Some(wacore::protocol::retry::build_retry_keys_node(
1532                &device_snapshot.identity_key.public_key,
1533                new_prekey_id,
1534                &new_prekey_public,
1535                device_snapshot.signed_pre_key_id,
1536                &device_snapshot.signed_pre_key.public_key,
1537                device_snapshot.signed_pre_key_signature.to_vec(),
1538                device_identity_bytes,
1539            ))
1540        } else {
1541            None
1542        };
1543
1544        // Build the receipt node. For group messages, include the participant attribute
1545        // to identify which group member should resend. For DMs, omit it since the
1546        // "to" address already identifies the sender.
1547        let mut builder = NodeBuilder::new("receipt")
1548            .attr("to", receipt_to)
1549            .attr("id", info.id.clone())
1550            .attr("type", "retry");
1551
1552        if info.source.is_group {
1553            builder = builder.attr("participant", &info.source.sender);
1554        }
1555
1556        // Handle peer vs device sync messages (matches WhatsApp Web's sendRetryReceipt):
1557        // WhatsApp Web checks: if (to.isUser()) { if (isMeAccount(to)) { ... } }
1558        // This means the category/recipient logic ONLY applies to DMs (not groups).
1559        // For groups, only the participant attribute is set (handled above).
1560        if !info.source.is_group {
1561            let is_from_own_account = device_snapshot
1562                .pn
1563                .as_ref()
1564                .is_some_and(|pn| info.source.sender.is_same_user_as(pn))
1565                || device_snapshot
1566                    .lid
1567                    .as_ref()
1568                    .is_some_and(|lid| info.source.sender.is_same_user_as(lid));
1569
1570            if is_from_own_account {
1571                if info.category == MessageCategory::Peer {
1572                    builder = builder.attr("category", MessageCategory::Peer.as_str());
1573                } else {
1574                    // Include recipient so the sender can look up the original message.
1575                    // Without this, the retry fails silently (getTargetChat returns null).
1576                    let recipient = info.source.recipient.as_ref().unwrap_or(&info.source.chat);
1577                    builder = builder.attr("recipient", recipient);
1578                }
1579            }
1580        }
1581
1582        // Build the final child list after the policy has decided whether this
1583        // request carries key material.
1584        let receipt_node = if let Some(keys) = keys_node {
1585            builder
1586                .children([retry_node, registration_node, keys])
1587                .build()
1588        } else {
1589            builder.children([retry_node, registration_node]).build()
1590        };
1591
1592        drop(device_snapshot);
1593        self.send_node(receipt_node).await?;
1594        Ok(RetryReceiptSendOutcome::Sent {
1595            included_keys: include_keys,
1596        })
1597    }
1598
1599    /// Sends an `enc_rekey_retry` receipt for VoIP call encryption re-keying.
1600    ///
1601    /// WA Web: When a peer fails to decrypt VoIP call encryption data (e.g.,
1602    /// `<enc>` within a `<call>` stanza), the receiver sends this receipt asking
1603    /// the sender to re-key.  The receipt uses `<enc_rekey>` child instead of
1604    /// `<retry>`, carrying VoIP call context (`call-id`, `call-creator`).
1605    ///
1606    /// WA Web reference: `ENC_RETRY_RECEIPT_ATTRS.GROUP_CALL = "enc_rekey_retry"`,
1607    /// constructed in `WAWebVoipSignalingEnums` module.
1608    #[allow(dead_code)] // Will be used when call handling is implemented (#345)
1609    #[cfg_attr(feature = "tracing", tracing::instrument(name = "wa.retry.send_enc_rekey_receipt", level = "debug", skip_all, fields(peer = %peer_jid.observe(), retry = retry_count), err(Debug)))]
1610    pub(crate) async fn send_enc_rekey_retry_receipt(
1611        &self,
1612        stanza_id: &str,
1613        peer_jid: &Jid,
1614        call_id: &str,
1615        call_creator: &Jid,
1616        retry_count: u8,
1617    ) -> Result<(), anyhow::Error> {
1618        let device_snapshot = self.persistence_manager.get_device_snapshot();
1619
1620        let registration_id_bytes = device_snapshot.registration_id.to_be_bytes().to_vec();
1621
1622        // WA Web: <enc_rekey call-creator="JID" call-id="..." count="N"/>
1623        let enc_rekey_node = NodeBuilder::new("enc_rekey")
1624            .attr("call-creator", call_creator)
1625            .attr("call-id", call_id)
1626            .attr("count", retry_count)
1627            .build();
1628
1629        let registration_node = NodeBuilder::new("registration")
1630            .bytes(registration_id_bytes)
1631            .build();
1632
1633        let receipt_node = NodeBuilder::new("receipt")
1634            .attr("to", peer_jid)
1635            .attr("id", stanza_id)
1636            .attr("type", "enc_rekey_retry")
1637            .children([enc_rekey_node, registration_node])
1638            .build();
1639
1640        info!(
1641            "Sending enc_rekey_retry receipt for call-id={} to {} (count={})",
1642            call_id,
1643            peer_jid.observe(),
1644            retry_count
1645        );
1646
1647        self.send_node(receipt_node).await?;
1648        Ok(())
1649    }
1650}
1651
1652#[cfg(test)]
1653mod tests {
1654    use super::*;
1655    use crate::store::persistence_manager::PersistenceManager;
1656    use crate::test_utils::MockHttpClient;
1657    use std::borrow::Cow;
1658    use std::sync::Arc;
1659    use wacore::libsignal::protocol::{IdentityKeyPair, KeyPair};
1660    use wacore::types::jid::JidExt as _;
1661    use wacore_binary::{Jid, JidExt};
1662    use waproto::whatsapp as wa;
1663
1664    fn resolve_retry_chat_info(
1665        receipt: &Receipt,
1666        node: &NodeRef<'_>,
1667        own_pn: Option<&Jid>,
1668        own_lid: Option<&Jid>,
1669    ) -> RetryChatInfo {
1670        super::resolve_retry_chat_info(receipt, node, own_pn, own_lid)
1671            .expect("retry should resolve a target chat")
1672    }
1673
1674    fn maybe_resolve_retry_chat_info(
1675        receipt: &Receipt,
1676        node: &NodeRef<'_>,
1677        own_pn: Option<&Jid>,
1678        own_lid: Option<&Jid>,
1679    ) -> Option<RetryChatInfo> {
1680        super::resolve_retry_chat_info(receipt, node, own_pn, own_lid)
1681    }
1682
1683    async fn attach_mock_noise_socket(client: &Client) {
1684        use crate::socket::NoiseSocket;
1685        use crate::transport::mock::MockTransport;
1686        use wacore::handshake::NoiseCipher;
1687
1688        let key = [0u8; 32];
1689        let socket = NoiseSocket::new(
1690            Arc::new(crate::runtime_impl::TokioRuntime),
1691            Arc::new(MockTransport),
1692            NoiseCipher::new(&key).expect("write cipher"),
1693            NoiseCipher::new(&key).expect("read cipher"),
1694        );
1695        *client.noise_socket.lock().await = Some(Arc::new(socket));
1696    }
1697
1698    async fn seed_retry_lease(
1699        client: &Client,
1700        address: &wacore::libsignal::protocol::ProtocolAddress,
1701        durable: bool,
1702    ) {
1703        use wacore::libsignal::protocol::SessionRecord;
1704
1705        let mut record = SessionRecord::new_fresh();
1706        record.reserve_sender_chain_counters(0);
1707        client.signal_cache.put_session(address, record).await;
1708        if !durable {
1709            return;
1710        }
1711
1712        client.flush_signal_cache().await.expect("durable lease");
1713        let snapshot = client.persistence_manager.get_device_snapshot();
1714        let record = client
1715            .signal_cache
1716            .get_session(address, &*snapshot.backend)
1717            .await
1718            .expect("session read")
1719            .expect("leased session");
1720        assert!(record.reserved_sender_chain_index() > 0);
1721        client.signal_cache.put_session(address, record).await;
1722    }
1723
1724    #[tokio::test]
1725    async fn retry_pre_wire_flush_failure_never_reaches_send_node() {
1726        use std::sync::atomic::Ordering;
1727
1728        let client =
1729            crate::test_utils::create_test_client_with_name("retry_pre_wire_failure").await;
1730        attach_mock_noise_socket(&client).await;
1731        let address = Jid::lid_device("100000000001035".to_string(), 7).to_protocol_address();
1732        seed_retry_lease(&client, &address, false).await;
1733        assert!(client.signal_cache.needs_pre_wire_flush().await);
1734
1735        client.inbound_commit_batch.reset();
1736        client
1737            .inbound_commit_batch
1738            .fail_flushes
1739            .store(true, Ordering::Release);
1740
1741        let id = "RETRY_PRE_WIRE_FAILURE";
1742        let mut waiter =
1743            client.wait_for_sent_node(crate::client::NodeFilter::tag("message").attr("id", id));
1744        let result = client
1745            .send_retry_stanza(NodeBuilder::new("message").attr("id", id).build())
1746            .await;
1747
1748        client
1749            .inbound_commit_batch
1750            .fail_flushes
1751            .store(false, Ordering::Release);
1752        assert!(result.is_err(), "the failed durability gate must abort");
1753        assert!(
1754            waiter.try_recv().expect("waiter stays live").is_none(),
1755            "send_node must not observe a stanza before durability"
1756        );
1757        assert!(
1758            client.signal_cache.needs_pre_wire_flush().await,
1759            "the failed reservation must remain gated"
1760        );
1761    }
1762
1763    #[tokio::test]
1764    async fn retry_inside_durable_lease_skips_synchronous_full_flush() {
1765        use std::sync::atomic::Ordering;
1766
1767        let client =
1768            crate::test_utils::create_test_client_with_name("retry_covered_by_lease").await;
1769        attach_mock_noise_socket(&client).await;
1770        let address = Jid::lid_device("100000000001036".to_string(), 8).to_protocol_address();
1771        seed_retry_lease(&client, &address, true).await;
1772        assert!(!client.signal_cache.needs_pre_wire_flush().await);
1773
1774        client.inbound_commit_batch.reset();
1775        client
1776            .inbound_commit_batch
1777            .fail_flushes
1778            .store(true, Ordering::Release);
1779
1780        let id = "RETRY_COVERED_BY_LEASE";
1781        let waiter =
1782            client.wait_for_sent_node(crate::client::NodeFilter::tag("message").attr("id", id));
1783        let result = client
1784            .send_retry_stanza(NodeBuilder::new("message").attr("id", id).build())
1785            .await;
1786
1787        client
1788            .inbound_commit_batch
1789            .fail_flushes
1790            .store(false, Ordering::Release);
1791        result.expect("an existing durable lease must not synchronously flush");
1792        let sent = waiter.await.expect("retry stanza reached send_node");
1793        assert_eq!(sent.attrs().required_string("id").unwrap(), id);
1794    }
1795
1796    #[tokio::test]
1797    async fn recent_message_cache_insert_and_take() {
1798        let _ = env_logger::builder().is_test(true).try_init();
1799
1800        let backend = crate::test_utils::create_test_backend().await;
1801        let pm = Arc::new(
1802            PersistenceManager::new(backend)
1803                .await
1804                .expect("persistence manager should initialize"),
1805        );
1806        // Enable L1 cache so MockBackend (which doesn't persist) works for this test
1807        let mut config = crate::cache_config::CacheConfig::default();
1808        config.recent_messages.capacity = 1_000;
1809        let (client, _sync_rx) = Client::new_with_cache_config(
1810            Arc::new(crate::runtime_impl::TokioRuntime),
1811            pm.clone(),
1812            Arc::new(crate::transport::mock::MockTransportFactory::new()),
1813            Arc::new(MockHttpClient),
1814            None,
1815            config,
1816        )
1817        .await;
1818
1819        let chat: Jid = "120363021033254949@g.us"
1820            .parse()
1821            .expect("test JID should be valid");
1822        let msg_id = "ABC123".to_string();
1823        let msg = wa::Message {
1824            conversation: Some("hello".into()),
1825            ..Default::default()
1826        };
1827
1828        // Insert via the new async API
1829        client.add_recent_message(&chat, &msg_id, &msg, None).await;
1830
1831        // First take should return and remove it from cache
1832        let taken = client.take_recent_message(&chat, &msg_id).await;
1833        assert!(taken.is_some());
1834        let (msg, alt_chat) = taken.unwrap();
1835        assert!(alt_chat.is_none(), "primary key should match");
1836        assert_eq!(msg.conversation.as_deref(), Some("hello"));
1837
1838        // Second take should return None
1839        let taken_again = client.take_recent_message(&chat, &msg_id).await;
1840        assert!(taken_again.is_none());
1841    }
1842
1843    /// DB-only path (no L1 cache, capacity 0 -- the harness/default): the wave
1844    /// that resolves the chat directly and stores the caller's borrowed id must
1845    /// still round-trip through the backend, so take_recent_message finds it.
1846    #[tokio::test]
1847    async fn recent_message_db_only_round_trip() {
1848        let _ = env_logger::builder().is_test(true).try_init();
1849
1850        let backend = crate::test_utils::create_test_backend().await;
1851        let pm = Arc::new(
1852            PersistenceManager::new(backend)
1853                .await
1854                .expect("persistence manager should initialize"),
1855        );
1856        // Capacity 0 keeps the L1 cache off, so the store + retrieve goes through
1857        // the backend -- exactly the DB-only branch add_recent_message took.
1858        let config = crate::cache_config::CacheConfig::default();
1859        assert_eq!(
1860            config.recent_messages.capacity, 0,
1861            "this test asserts the DB-only (capacity 0) path"
1862        );
1863        let (client, _sync_rx) = Client::new_with_cache_config(
1864            Arc::new(crate::runtime_impl::TokioRuntime),
1865            pm.clone(),
1866            Arc::new(crate::transport::mock::MockTransportFactory::new()),
1867            Arc::new(MockHttpClient),
1868            None,
1869            config,
1870        )
1871        .await;
1872
1873        let chat: Jid = "120363021033254949@g.us"
1874            .parse()
1875            .expect("test JID should be valid");
1876        let msg_id = "DBONLY1".to_string();
1877        let msg = wa::Message {
1878            conversation: Some("db-only".into()),
1879            ..Default::default()
1880        };
1881
1882        client.add_recent_message(&chat, &msg_id, &msg, None).await;
1883
1884        let taken = client.take_recent_message(&chat, &msg_id).await;
1885        assert!(
1886            taken.is_some(),
1887            "a DB-only stored message must be retrievable from the backend"
1888        );
1889        let (got, _alt) = taken.unwrap();
1890        assert_eq!(got.conversation.as_deref(), Some("db-only"));
1891
1892        let again = client.take_recent_message(&chat, &msg_id).await;
1893        assert!(again.is_none(), "take consumes the DB-only message");
1894    }
1895
1896    #[tokio::test]
1897    async fn peek_recent_message_does_not_consume() {
1898        let _ = env_logger::builder().is_test(true).try_init();
1899
1900        let backend = crate::test_utils::create_test_backend().await;
1901        let pm = Arc::new(
1902            PersistenceManager::new(backend)
1903                .await
1904                .expect("persistence manager should initialize"),
1905        );
1906        let mut config = crate::cache_config::CacheConfig::default();
1907        config.recent_messages.capacity = 1_000;
1908        let (client, _sync_rx) = Client::new_with_cache_config(
1909            Arc::new(crate::runtime_impl::TokioRuntime),
1910            pm.clone(),
1911            Arc::new(crate::transport::mock::MockTransportFactory::new()),
1912            Arc::new(MockHttpClient),
1913            None,
1914            config,
1915        )
1916        .await;
1917
1918        let chat: Jid = "120363021033254949@g.us".parse().unwrap();
1919        let msg_id = "PEEK1".to_string();
1920        let msg = wa::Message {
1921            conversation: Some("hi".into()),
1922            ..Default::default()
1923        };
1924        client.add_recent_message(&chat, &msg_id, &msg, None).await;
1925
1926        // Peeking twice both return the message and leave it in the cache...
1927        for _ in 0..2 {
1928            let peeked = client.peek_recent_message(&chat, &msg_id).await;
1929            let (m, alt) = peeked.expect("peek should find the cached message");
1930            assert!(alt.is_none());
1931            assert_eq!(m.conversation.as_deref(), Some("hi"));
1932        }
1933        // ...so a subsequent take still finds it (peek didn't remove it).
1934        assert!(client.take_recent_message(&chat, &msg_id).await.is_some());
1935    }
1936
1937    #[test]
1938    fn get_bytes_content_extracts_bytes() {
1939        use wacore_binary::{Attrs, Node};
1940
1941        // Test with bytes content
1942        let node = Node {
1943            tag: Cow::Borrowed("test"),
1944            attrs: Attrs::new(),
1945            content: Some(NodeContent::Bytes(vec![1, 2, 3, 4])),
1946        };
1947        assert_eq!(get_bytes_content(&node), Some(&[1, 2, 3, 4][..]));
1948
1949        // Test with string content (should return None)
1950        let node_str = Node {
1951            tag: Cow::Borrowed("test"),
1952            attrs: Attrs::new(),
1953            content: Some(NodeContent::String("hello".into())),
1954        };
1955        assert_eq!(get_bytes_content(&node_str), None);
1956
1957        // Test with no content
1958        let node_empty = Node {
1959            tag: Cow::Borrowed("test"),
1960            attrs: Attrs::new(),
1961            content: None,
1962        };
1963        assert_eq!(get_bytes_content(&node_empty), None);
1964    }
1965
1966    #[test]
1967    fn peer_detection_logic() {
1968        let our_jid = Jid::pn("559911112222");
1969        let peer_jid = Jid::pn_device("559911112222", 1);
1970        let other_jid = Jid::pn("559933334444");
1971
1972        assert_eq!(our_jid.user, peer_jid.user);
1973        assert_ne!(our_jid.user, other_jid.user);
1974    }
1975
1976    /// Integration test for retry receipt attribute logic.
1977    /// Tests the fix for lost device sync messages (AC7B18EBD4445BFC55C0EA3CF9F913F8 case).
1978    /// Matches WhatsApp Web's sendRetryReceipt: if (to.isUser()) { if (isMeAccount(to)) { ... } }
1979    #[test]
1980    fn retry_receipt_attributes_for_device_sync_vs_peer_vs_group() {
1981        use wacore::types::message::{MessageCategory, MessageInfo, MessageSource};
1982        use wacore_binary::builder::NodeBuilder;
1983
1984        let our_pn = Jid::pn("559999999999");
1985        let our_lid = Jid::lid("100000000000001");
1986
1987        fn build_retry_receipt(info: &MessageInfo, our_pn: &Jid, our_lid: &Jid) -> Node {
1988            // Mirror production routing: groups → chat JID, DMs → sender JID
1989            let receipt_to = if info.source.is_group {
1990                &info.source.chat
1991            } else {
1992                &info.source.sender
1993            };
1994            let mut builder = NodeBuilder::new("receipt")
1995                .attr("to", receipt_to)
1996                .attr("id", info.id.clone())
1997                .attr("type", "retry");
1998
1999            if info.source.is_group {
2000                builder = builder.attr("participant", &info.source.sender);
2001            }
2002
2003            if !info.source.is_group {
2004                let is_from_own_account = info.source.sender.is_same_user_as(our_pn)
2005                    || info.source.sender.is_same_user_as(our_lid);
2006
2007                if is_from_own_account {
2008                    if info.category == MessageCategory::Peer {
2009                        builder = builder.attr("category", MessageCategory::Peer.as_str());
2010                    } else {
2011                        let recipient = info.source.recipient.as_ref().unwrap_or(&info.source.chat);
2012                        builder = builder.attr("recipient", recipient);
2013                    }
2014                }
2015            }
2016
2017            builder.build()
2018        }
2019
2020        // Case 1: Device sync DM
2021        let recipient_lid = Jid::lid("200000000000002");
2022        let device_sync_info = MessageInfo {
2023            id: "DEVICE_SYNC_MSG_001".to_string(),
2024            source: MessageSource {
2025                chat: recipient_lid.clone(),
2026                sender: our_lid.clone(),
2027                is_from_me: true,
2028                is_group: false,
2029                recipient: Some(recipient_lid.clone()),
2030                ..Default::default()
2031            },
2032            category: MessageCategory::default(),
2033            ..Default::default()
2034        };
2035
2036        let node = build_retry_receipt(&device_sync_info, &our_pn, &our_lid);
2037        assert_eq!(
2038            node.attrs
2039                .get("recipient")
2040                .map(|v| v == "200000000000002@lid"),
2041            Some(true),
2042            "Device sync DM should include recipient"
2043        );
2044        assert!(
2045            node.attrs.get("category").is_none(),
2046            "Device sync DM should NOT have category=peer"
2047        );
2048        assert!(
2049            node.attrs.get("participant").is_none(),
2050            "DM should NOT have participant"
2051        );
2052
2053        // Case 2: Peer DM with category="peer"
2054        let other_pn = Jid::pn("551188888888");
2055        let peer_info = MessageInfo {
2056            id: "PEER123".to_string(),
2057            source: MessageSource {
2058                chat: other_pn.clone(),
2059                sender: our_pn.clone(),
2060                is_from_me: true,
2061                is_group: false,
2062                recipient: None,
2063                ..Default::default()
2064            },
2065            category: MessageCategory::Peer,
2066            ..Default::default()
2067        };
2068
2069        let node = build_retry_receipt(&peer_info, &our_pn, &our_lid);
2070        assert_eq!(
2071            node.attrs.get("category").map(|v| v == "peer"),
2072            Some(true),
2073            "Peer DM should have category=peer"
2074        );
2075        assert!(
2076            node.attrs.get("recipient").is_none(),
2077            "Peer DM should NOT have recipient"
2078        );
2079
2080        // Case 3: Group message from our own account
2081        let group_info = MessageInfo {
2082            id: "GROUP123".to_string(),
2083            source: MessageSource {
2084                chat: "123456789@g.us".parse().unwrap(),
2085                sender: our_lid.clone(),
2086                is_from_me: true,
2087                is_group: true,
2088                recipient: None,
2089                ..Default::default()
2090            },
2091            category: MessageCategory::default(),
2092            ..Default::default()
2093        };
2094
2095        let node = build_retry_receipt(&group_info, &our_pn, &our_lid);
2096        assert!(
2097            node.attrs.get("participant").is_some(),
2098            "Group should have participant"
2099        );
2100        assert!(
2101            node.attrs.get("category").is_none(),
2102            "Group should NOT have category"
2103        );
2104        assert!(
2105            node.attrs.get("recipient").is_none(),
2106            "Group should NOT have recipient"
2107        );
2108
2109        // Case 4: DM from someone else
2110        let other_dm_info = MessageInfo {
2111            id: "OTHER123".to_string(),
2112            source: MessageSource {
2113                chat: other_pn.clone(),
2114                sender: other_pn.clone(),
2115                is_from_me: false,
2116                is_group: false,
2117                recipient: None,
2118                ..Default::default()
2119            },
2120            category: MessageCategory::default(),
2121            ..Default::default()
2122        };
2123
2124        let node = build_retry_receipt(&other_dm_info, &our_pn, &our_lid);
2125        assert!(
2126            node.attrs.get("category").is_none(),
2127            "DM from other should NOT have category"
2128        );
2129        assert!(
2130            node.attrs.get("recipient").is_none(),
2131            "DM from other should NOT have recipient"
2132        );
2133    }
2134
2135    /// Verify enc_rekey_retry receipt node structure matches WhatsApp Web:
2136    /// <receipt to="peer" id="stanza_id" type="enc_rekey_retry">
2137    ///   <enc_rekey call-creator="creator_jid" call-id="..." count="N"/>
2138    ///   <registration>{4-byte big-endian reg id}</registration>
2139    /// </receipt>
2140    #[test]
2141    fn enc_rekey_retry_receipt_node_structure() {
2142        use wacore_binary::builder::NodeBuilder;
2143
2144        let peer_jid: Jid = "5511999999999@s.whatsapp.net".parse().expect("peer JID");
2145        let call_creator: Jid = "5511888888888@s.whatsapp.net".parse().expect("creator JID");
2146        let call_id = "CALL-ABC-123";
2147        let stanza_id = "3EB0AABBCCDD";
2148        let retry_count: u8 = 2;
2149        let registration_id: u32 = 12345;
2150
2151        // Build the receipt exactly as send_enc_rekey_retry_receipt does
2152        let enc_rekey_node = NodeBuilder::new("enc_rekey")
2153            .attr("call-creator", call_creator)
2154            .attr("call-id", call_id)
2155            .attr("count", retry_count)
2156            .build();
2157
2158        let registration_node = NodeBuilder::new("registration")
2159            .bytes(registration_id.to_be_bytes().to_vec())
2160            .build();
2161
2162        let receipt_node = NodeBuilder::new("receipt")
2163            .attr("to", peer_jid)
2164            .attr("id", stanza_id)
2165            .attr("type", "enc_rekey_retry")
2166            .children([enc_rekey_node, registration_node])
2167            .build();
2168
2169        // Verify top-level receipt attributes
2170        assert_eq!(
2171            receipt_node.attrs().optional_string("type").as_deref(),
2172            Some("enc_rekey_retry"),
2173            "receipt type must be enc_rekey_retry"
2174        );
2175        assert!(
2176            receipt_node
2177                .attrs
2178                .get("to")
2179                .is_some_and(|v| *v == "5511999999999@s.whatsapp.net"),
2180            "receipt 'to' must be peer JID"
2181        );
2182        assert_eq!(
2183            receipt_node.attrs().optional_string("id").as_deref(),
2184            Some("3EB0AABBCCDD")
2185        );
2186
2187        // Verify <enc_rekey> child (NOT <retry>)
2188        assert!(
2189            receipt_node.get_optional_child("retry").is_none(),
2190            "enc_rekey_retry must NOT contain <retry> child"
2191        );
2192        let enc_rekey = receipt_node
2193            .get_optional_child("enc_rekey")
2194            .expect("<enc_rekey> child must exist");
2195        assert_eq!(
2196            enc_rekey.attrs().optional_string("call-id").as_deref(),
2197            Some("CALL-ABC-123")
2198        );
2199        assert!(
2200            enc_rekey
2201                .attrs
2202                .get("call-creator")
2203                .is_some_and(|v| *v == "5511888888888@s.whatsapp.net"),
2204            "enc_rekey 'call-creator' must be creator JID"
2205        );
2206        assert_eq!(
2207            enc_rekey.attrs().optional_string("count").as_deref(),
2208            Some("2")
2209        );
2210
2211        // Verify <registration> child
2212        let registration = receipt_node
2213            .get_optional_child("registration")
2214            .expect("<registration> child must exist");
2215        let reg_bytes = match &registration.content {
2216            Some(NodeContent::Bytes(b)) => b.clone(),
2217            _ => panic!("registration must contain bytes"),
2218        };
2219        assert_eq!(
2220            u32::from_be_bytes(reg_bytes.try_into().unwrap()),
2221            12345,
2222            "registration ID must be 4-byte big-endian"
2223        );
2224    }
2225
2226    #[test]
2227    fn prekey_id_parsing() {
2228        // PreKey IDs are 3 bytes big-endian
2229        let id_bytes = [0x01, 0x02, 0x03];
2230        let prekey_id = u32::from_be_bytes([0, id_bytes[0], id_bytes[1], id_bytes[2]]);
2231        assert_eq!(prekey_id, 0x00010203);
2232
2233        // Signed prekey IDs follow the same format
2234        let skey_id_bytes = [0xFF, 0xFE, 0xFD];
2235        let skey_id = u32::from_be_bytes([0, skey_id_bytes[0], skey_id_bytes[1], skey_id_bytes[2]]);
2236        assert_eq!(skey_id, 0x00FFFEFD);
2237    }
2238
2239    #[tokio::test]
2240    async fn base_key_store_operations() {
2241        let _ = env_logger::builder().is_test(true).try_init();
2242
2243        let backend = crate::test_utils::create_test_backend().await;
2244
2245        let address = "12345.0:1";
2246        let msg_id = "ABC123";
2247        let base_key = vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
2248
2249        // Initially, has_same_base_key should return false (no saved key)
2250        let result = backend.has_same_base_key(address, msg_id, &base_key).await;
2251        assert!(result.is_ok());
2252        assert!(!result.unwrap());
2253
2254        // Save the base key
2255        let save_result = backend.save_base_key(address, msg_id, &base_key).await;
2256        assert!(save_result.is_ok());
2257
2258        // Same key should now match (collision detected)
2259        let result = backend.has_same_base_key(address, msg_id, &base_key).await;
2260        assert!(result.is_ok());
2261        assert!(result.unwrap());
2262
2263        // Different key should NOT match (no collision)
2264        let different_key = vec![10, 9, 8, 7, 6, 5, 4, 3, 2, 1];
2265        let result = backend
2266            .has_same_base_key(address, msg_id, &different_key)
2267            .await;
2268        assert!(result.is_ok());
2269        assert!(!result.unwrap());
2270
2271        // Delete the base key
2272        let delete_result = backend.delete_base_key(address, msg_id).await;
2273        assert!(delete_result.is_ok());
2274
2275        // After deletion, has_same_base_key should return false
2276        let result = backend.has_same_base_key(address, msg_id, &base_key).await;
2277        assert!(result.is_ok());
2278        assert!(!result.unwrap());
2279    }
2280
2281    #[tokio::test]
2282    async fn base_key_store_upsert() {
2283        let _ = env_logger::builder().is_test(true).try_init();
2284
2285        let backend = crate::test_utils::create_test_backend().await;
2286
2287        let address = "12345.0:1";
2288        let msg_id = "MSG001";
2289        let first_key = vec![1, 2, 3];
2290        let second_key = vec![4, 5, 6];
2291
2292        // Save first key
2293        backend
2294            .save_base_key(address, msg_id, &first_key)
2295            .await
2296            .unwrap();
2297        assert!(
2298            backend
2299                .has_same_base_key(address, msg_id, &first_key)
2300                .await
2301                .unwrap()
2302        );
2303        assert!(
2304            !backend
2305                .has_same_base_key(address, msg_id, &second_key)
2306                .await
2307                .unwrap()
2308        );
2309
2310        // Save second key (upsert should replace)
2311        backend
2312            .save_base_key(address, msg_id, &second_key)
2313            .await
2314            .unwrap();
2315        assert!(
2316            !backend
2317                .has_same_base_key(address, msg_id, &first_key)
2318                .await
2319                .unwrap()
2320        );
2321        assert!(
2322            backend
2323                .has_same_base_key(address, msg_id, &second_key)
2324                .await
2325                .unwrap()
2326        );
2327    }
2328
2329    #[tokio::test]
2330    async fn base_key_store_multiple_messages() {
2331        let _ = env_logger::builder().is_test(true).try_init();
2332
2333        let backend = crate::test_utils::create_test_backend().await;
2334
2335        let address = "12345.0:1";
2336        let msg_id_1 = "MSG001";
2337        let msg_id_2 = "MSG002";
2338        let key_1 = vec![1, 2, 3];
2339        let key_2 = vec![4, 5, 6];
2340
2341        // Save keys for different messages
2342        backend
2343            .save_base_key(address, msg_id_1, &key_1)
2344            .await
2345            .unwrap();
2346        backend
2347            .save_base_key(address, msg_id_2, &key_2)
2348            .await
2349            .unwrap();
2350
2351        // Each message should have its own key
2352        assert!(
2353            backend
2354                .has_same_base_key(address, msg_id_1, &key_1)
2355                .await
2356                .unwrap()
2357        );
2358        assert!(
2359            !backend
2360                .has_same_base_key(address, msg_id_1, &key_2)
2361                .await
2362                .unwrap()
2363        );
2364        assert!(
2365            !backend
2366                .has_same_base_key(address, msg_id_2, &key_1)
2367                .await
2368                .unwrap()
2369        );
2370        assert!(
2371            backend
2372                .has_same_base_key(address, msg_id_2, &key_2)
2373                .await
2374                .unwrap()
2375        );
2376
2377        // Delete one message's key, other should remain
2378        backend.delete_base_key(address, msg_id_1).await.unwrap();
2379        assert!(
2380            !backend
2381                .has_same_base_key(address, msg_id_1, &key_1)
2382                .await
2383                .unwrap()
2384        );
2385        assert!(
2386            backend
2387                .has_same_base_key(address, msg_id_2, &key_2)
2388                .await
2389                .unwrap()
2390        );
2391    }
2392
2393    /// Build a minimal `<receipt>` Node representing an incoming retry receipt
2394    /// without `<keys>`. Used by tests that exercise the no-bundle path of
2395    /// `update_local_signal_session`.
2396    fn build_retry_receipt_without_keys() -> Node {
2397        use wacore_binary::builder::NodeBuilder;
2398        NodeBuilder::new("receipt").build()
2399    }
2400
2401    /// Build a `<receipt>` with a `<registration>` child carrying `reg_id` (big
2402    /// endian). Used to exercise the reg-ID-mismatch branch without a full
2403    /// `<keys>` bundle.
2404    fn build_retry_receipt_with_registration(reg_id: u32) -> Node {
2405        use wacore_binary::builder::NodeBuilder;
2406        NodeBuilder::new("receipt")
2407            .children([NodeBuilder::new("registration")
2408                .bytes(reg_id.to_be_bytes().to_vec())
2409                .build()])
2410            .build()
2411    }
2412
2413    fn dm_retry_info(resolved_jid: &Jid) -> RetryChatInfo {
2414        RetryChatInfo {
2415            chat: resolved_jid.to_non_ad(),
2416            requester: resolved_jid.clone(),
2417            original_from: resolved_jid.clone(),
2418            recipient: None,
2419            is_bot: false,
2420            is_fbid_bot_retry: false,
2421        }
2422    }
2423
2424    // Produces a parseable SessionRecord so peek_session succeeds and
2425    // alice_base_key/remote_registration_id return meaningful values.
2426    fn valid_serialized_session(remote_regid: u32, base_key: Vec<u8>) -> Vec<u8> {
2427        use wacore::libsignal::protocol::{SessionRecord, SessionState};
2428        use waproto::whatsapp::SessionStructure;
2429
2430        let state = SessionState::from_session_structure(SessionStructure {
2431            session_version: Some(3),
2432            local_identity_public: None,
2433            remote_identity_public: None,
2434            root_key: None,
2435            previous_counter: Some(0),
2436            sender_chain: buffa::MessageField::default(),
2437            receiver_chains: vec![],
2438            pending_pre_key: buffa::MessageField::default(),
2439            remote_registration_id: Some(remote_regid),
2440            local_registration_id: Some(0),
2441            alice_base_key: Some(base_key),
2442            needs_refresh: None,
2443            pending_key_exchange: buffa::MessageField::default(),
2444        });
2445        SessionRecord::new(state)
2446            .serialize()
2447            .expect("serialize session record")
2448    }
2449
2450    /// WA Web compliance: at retry #1 with no `<keys>`, `updateLocalSignalSession`
2451    /// does NOT delete the session. Previously the Rust DM path unconditionally
2452    /// deleted on every retry — this regressed legitimate sessions and forced
2453    /// unnecessary prekey bundle fetches.
2454    /// Ref: `WAWeb/Update/LocalSignalSession.js` (no delete on retry==1)
2455    #[tokio::test]
2456    async fn update_local_signal_session_preserves_dm_session_at_retry_1() {
2457        let client =
2458            crate::test_utils::create_test_client_with_failing_http("retry_preserve_retry_1").await;
2459        let user = "100000000000088".to_string();
2460        let resolved_jid = Jid::lid_device(user.clone(), 33);
2461
2462        let backend = client.persistence_manager.backend();
2463        let device_0 = Jid::lid_device(user.clone(), 0).to_protocol_address();
2464        let device_33 = Jid::lid_device(user, 33).to_protocol_address();
2465
2466        // Real serializable SessionRecords — peek_session must return Some(...)
2467        // so the function reaches the base-key branch at retry==1 and exercises
2468        // the "no delete" rule. Invalid bytes would short-circuit via .ok().flatten().
2469        let session_bytes_33 = valid_serialized_session(4242, vec![0xAA; 32]);
2470        let session_bytes_0 = valid_serialized_session(4243, vec![0xBB; 32]);
2471        backend
2472            .put_session(device_0.as_str(), &session_bytes_0)
2473            .await
2474            .unwrap();
2475        backend
2476            .put_session(device_33.as_str(), &session_bytes_33)
2477            .await
2478            .unwrap();
2479
2480        let node = build_retry_receipt_without_keys();
2481        let node_ref = node.as_node_ref();
2482        client
2483            .update_local_signal_session(
2484                &dm_retry_info(&resolved_jid),
2485                &resolved_jid,
2486                "MSG-RETRY-1",
2487                1,
2488                &node_ref,
2489                false,
2490            )
2491            .await;
2492        client.flush_signal_cache().await.unwrap();
2493
2494        assert!(
2495            backend
2496                .get_session(device_0.as_str())
2497                .await
2498                .unwrap()
2499                .is_some(),
2500            "non-requesting device session must be preserved"
2501        );
2502        assert!(
2503            backend
2504                .get_session(device_33.as_str())
2505                .await
2506                .unwrap()
2507                .is_some(),
2508            "requesting device session with valid record must be preserved at retry #1"
2509        );
2510    }
2511
2512    /// Production scenario from debug-1776271138: peer sends retry receipt
2513    /// without `<keys>` but with `<registration>` whose reg_id differs from
2514    /// our stored session. WA Web deletes the session (LocalSignalSession.js
2515    /// L52-65) so the next ensureE2ESessions fetches a fresh bundle.
2516    #[tokio::test]
2517    async fn update_local_signal_session_deletes_on_regid_mismatch() {
2518        let client =
2519            crate::test_utils::create_test_client_with_failing_http("retry_regid_mismatch").await;
2520        let resolved_jid = Jid::lid_device("100000000000099".to_string(), 17);
2521        let signal_address = resolved_jid.to_protocol_address();
2522        let backend = client.persistence_manager.backend();
2523
2524        let stored_regid = 4242u32;
2525        let session_bytes = valid_serialized_session(stored_regid, vec![0xAA; 32]);
2526        backend
2527            .put_session(signal_address.as_str(), &session_bytes)
2528            .await
2529            .unwrap();
2530
2531        let received_regid = 0xDEAD_BEEFu32;
2532        assert_ne!(stored_regid, received_regid);
2533        let node = build_retry_receipt_with_registration(received_regid);
2534        let node_ref = node.as_node_ref();
2535        client
2536            .update_local_signal_session(
2537                &dm_retry_info(&resolved_jid),
2538                &resolved_jid,
2539                "MSG-REGID",
2540                1,
2541                &node_ref,
2542                false,
2543            )
2544            .await;
2545        client.flush_signal_cache().await.unwrap();
2546
2547        assert!(
2548            backend
2549                .get_session(signal_address.as_str())
2550                .await
2551                .unwrap()
2552                .is_none(),
2553            "session must be deleted when retry has no keys and reg IDs differ"
2554        );
2555    }
2556
2557    /// Unparseable session bytes: peek_session returns None via .ok().flatten(),
2558    /// so every branch that dereferences a session is skipped. Verifies we
2559    /// don't panic or re-process stale bytes when the record can't decode.
2560    #[tokio::test]
2561    async fn update_local_signal_session_handles_unparseable_session_gracefully() {
2562        let client =
2563            crate::test_utils::create_test_client_with_failing_http("retry_unparseable_session")
2564                .await;
2565        let resolved_jid = Jid::lid_device("100000000000099".to_string(), 17);
2566        let signal_address = resolved_jid.to_protocol_address();
2567        let backend = client.persistence_manager.backend();
2568
2569        backend
2570            .put_session(signal_address.as_str(), b"invalid-session")
2571            .await
2572            .unwrap();
2573
2574        let node = build_retry_receipt_with_registration(0xDEAD_BEEF);
2575        let node_ref = node.as_node_ref();
2576        client
2577            .update_local_signal_session(
2578                &dm_retry_info(&resolved_jid),
2579                &resolved_jid,
2580                "MSG-REGID",
2581                1,
2582                &node_ref,
2583                false,
2584            )
2585            .await;
2586        client.flush_signal_cache().await.unwrap();
2587
2588        assert!(
2589            backend
2590                .get_session(signal_address.as_str())
2591                .await
2592                .unwrap()
2593                .is_some(),
2594            "unparseable bytes skip every branch; nothing should delete them"
2595        );
2596    }
2597
2598    /// Verify the function is a safe no-op when there is no session at all.
2599    /// This is the common case for retries from devices we haven't messaged
2600    /// yet (e.g., a new companion device).
2601    #[tokio::test]
2602    async fn update_local_signal_session_no_session_is_noop() {
2603        let client =
2604            crate::test_utils::create_test_client_with_failing_http("retry_no_session").await;
2605        let resolved_jid = Jid::lid_device("100000000000199".to_string(), 42);
2606        let node = build_retry_receipt_without_keys();
2607        let node_ref = node.as_node_ref();
2608        client
2609            .update_local_signal_session(
2610                &dm_retry_info(&resolved_jid),
2611                &resolved_jid,
2612                "MSG-NOSESS",
2613                1,
2614                &node_ref,
2615                false,
2616            )
2617            .await;
2618    }
2619
2620    /// Group/status at retry #1 must not delete any session. Group/status
2621    /// previously skipped the base-key path entirely; now it runs but the
2622    /// retry==1 short-circuit still prevents deletion.
2623    #[tokio::test]
2624    async fn update_local_signal_session_preserves_group_session_at_retry_1() {
2625        let client =
2626            crate::test_utils::create_test_client_with_failing_http("retry_group_preserve").await;
2627        let resolved_jid = Jid::lid_device("100000000000088".to_string(), 33);
2628        let signal_address = resolved_jid.to_protocol_address();
2629        let backend = client.persistence_manager.backend();
2630
2631        let session_bytes = valid_serialized_session(9999, vec![0xCC; 32]);
2632        backend
2633            .put_session(signal_address.as_str(), &session_bytes)
2634            .await
2635            .unwrap();
2636
2637        let group_chat: Jid = "120363042537531116@g.us".parse().unwrap();
2638        let info = RetryChatInfo {
2639            chat: group_chat.clone(),
2640            requester: resolved_jid.clone(),
2641            original_from: group_chat,
2642            recipient: None,
2643            is_bot: false,
2644            is_fbid_bot_retry: false,
2645        };
2646
2647        let node = build_retry_receipt_without_keys();
2648        let node_ref = node.as_node_ref();
2649        client
2650            .update_local_signal_session(&info, &resolved_jid, "MSG-GRP-1", 1, &node_ref, false)
2651            .await;
2652        client.flush_signal_cache().await.unwrap();
2653
2654        assert!(
2655            backend
2656                .get_session(signal_address.as_str())
2657                .await
2658                .unwrap()
2659                .is_some(),
2660            "group retry at #1 should not delete the session"
2661        );
2662    }
2663
2664    #[tokio::test]
2665    async fn update_local_signal_session_cools_resolved_sender_key_namespace() {
2666        let client = crate::test_utils::create_test_client_with_failing_http(
2667            "retry_sender_key_resolved_namespace",
2668        )
2669        .await;
2670        let group = "120363000000000006@g.us";
2671        let requester_pn: Jid = "12025550108:33@s.whatsapp.net".parse().unwrap();
2672        let resolved_lid: Jid = "100000000000088:33@lid".parse().unwrap();
2673        client
2674            .persistence_manager
2675            .set_sender_key_status(
2676                group,
2677                &[
2678                    ("12025550108:33@s.whatsapp.net", true),
2679                    ("100000000000088:33@lid", true),
2680                ],
2681            )
2682            .await
2683            .unwrap();
2684
2685        let rows = client
2686            .persistence_manager
2687            .get_sender_key_devices(group)
2688            .await
2689            .unwrap();
2690        let cached = client
2691            .sender_key_device_cache
2692            .get_or_init(group, async {
2693                Arc::new(crate::sender_key_device_cache::SenderKeyDeviceMap::from_db_rows(&rows))
2694            })
2695            .await;
2696
2697        let info = RetryChatInfo {
2698            chat: group.parse().unwrap(),
2699            requester: requester_pn,
2700            original_from: group.parse().unwrap(),
2701            recipient: None,
2702            is_bot: false,
2703            is_fbid_bot_retry: false,
2704        };
2705        let node = build_retry_receipt_without_keys();
2706        assert!(
2707            client
2708                .update_local_signal_session(
2709                    &info,
2710                    &resolved_lid,
2711                    "MSG-GRP-NAMESPACE",
2712                    1,
2713                    &node.as_node_ref(),
2714                    false,
2715                )
2716                .await
2717        );
2718
2719        assert_eq!(cached.device_has_key("100000000000088", 33), Some(false));
2720        assert_eq!(cached.device_has_key("12025550108", 33), Some(true));
2721        let persisted = crate::sender_key_device_cache::SenderKeyDeviceMap::from_db_rows(
2722            &client
2723                .persistence_manager
2724                .get_sender_key_devices(group)
2725                .await
2726                .unwrap(),
2727        );
2728        assert_eq!(persisted.device_has_key("100000000000088", 33), Some(false));
2729        assert_eq!(persisted.device_has_key("12025550108", 33), Some(true));
2730    }
2731
2732    #[tokio::test]
2733    async fn status_retransmission_resolution_is_cache_aside_with_pn_fallback() {
2734        let client = crate::test_utils::create_test_client_with_failing_http(
2735            "retry_status_requester_resolution",
2736        )
2737        .await;
2738        client
2739            .add_lid_pn_mapping(
2740                "100000000000089",
2741                "12025550109",
2742                crate::lid_pn_cache::LearningSource::Usync,
2743            )
2744            .await
2745            .unwrap();
2746        client.lid_pn_cache.clear().await;
2747
2748        let mapped_pn: Jid = "12025550109:19@s.whatsapp.net".parse().unwrap();
2749        let mapped = client
2750            .resolve_retransmission_encryption_jid(RetransmissionRoute::Status, &mapped_pn)
2751            .await
2752            .unwrap();
2753        assert_eq!(mapped, "100000000000089:19@lid".parse::<Jid>().unwrap());
2754
2755        let unmapped_pn: Jid = "12025550110:20@s.whatsapp.net".parse().unwrap();
2756        let fallback = client
2757            .resolve_retransmission_encryption_jid(RetransmissionRoute::Status, &unmapped_pn)
2758            .await
2759            .unwrap();
2760        assert_eq!(fallback, unmapped_pn);
2761    }
2762
2763    /// `should_recreate_session` mirrors whatsmeow `shouldRecreateSession`:
2764    /// 1) no session → always recreate;
2765    /// 2) session exists + retry<2 → never recreate;
2766    /// 3) session exists + retry≥2 + first time (or >1h since last) → recreate.
2767    /// 4) session exists + retry≥2 + recreated <1h ago → throttled, do not recreate.
2768    #[tokio::test]
2769    async fn should_recreate_session_matrix() {
2770        let client =
2771            crate::test_utils::create_test_client_with_failing_http("should_recreate_session")
2772                .await;
2773
2774        // Use disjoint JIDs per scenario so the negative-cache populated by
2775        // `has_session` on the "no session" branch can't shadow the later
2776        // backend put for the "session present" branches.
2777        let jid_with = Jid::lid_device("999999999999991".to_string(), 3);
2778        let jid_without = Jid::lid_device("999999999999992".to_string(), 3);
2779
2780        // Seed a session for jid_with BEFORE the first has_session lookup so
2781        // the cache caches the hit, not the miss.
2782        let session_bytes = valid_serialized_session(7777, vec![0xEE; 32]);
2783        client
2784            .persistence_manager
2785            .backend()
2786            .put_session(jid_with.to_protocol_address().as_str(), &session_bytes)
2787            .await
2788            .unwrap();
2789
2790        // 1) session present + retry<2 → never recreate, no history stamp.
2791        assert!(
2792            client.should_recreate_session(1, &jid_with).await.is_none(),
2793            "retry<2 with session present should not recreate"
2794        );
2795        assert!(
2796            client
2797                .session_recreate_history
2798                .get(&jid_with)
2799                .await
2800                .is_none(),
2801            "no-op path must not stamp the history"
2802        );
2803
2804        // 2) session present + retry≥2 + cold history → recreate, stamp history.
2805        assert!(
2806            client
2807                .should_recreate_session(2, &jid_with)
2808                .await
2809                .is_some_and(|r| r.contains("retry count > 1")),
2810            "retry≥2 with cold history should recreate"
2811        );
2812        let after_first = client.session_recreate_history.get(&jid_with).await;
2813        assert!(after_first.is_some(), "first recreate must stamp history");
2814
2815        // 3) session present + retry≥2 + recent history → throttled.
2816        assert!(
2817            client.should_recreate_session(3, &jid_with).await.is_none(),
2818            "retry≥2 within {}s should be throttled",
2819            RECREATE_SESSION_TIMEOUT.as_secs()
2820        );
2821        let after_second = client.session_recreate_history.get(&jid_with).await;
2822        assert_eq!(
2823            after_first, after_second,
2824            "throttled path must not re-stamp the history"
2825        );
2826
2827        // 4) Past the throttle window → fresh recreate. Use a future `now`
2828        // (subtracting from a young runtime's Instant would saturate to zero).
2829        let stamp_then = after_first.expect("first recreate stamped history");
2830        let well_past = stamp_then + RECREATE_SESSION_TIMEOUT + std::time::Duration::from_secs(1);
2831        assert!(
2832            client
2833                .should_recreate_session_at(3, &jid_with, well_past)
2834                .await
2835                .is_some_and(|r| r.contains("over an hour")),
2836            "entry past the throttle window must allow a fresh recreate"
2837        );
2838
2839        // 5) no session → recreate regardless of retry count.
2840        assert!(
2841            client
2842                .should_recreate_session(0, &jid_without)
2843                .await
2844                .is_some_and(|r| r.contains("don't have a Signal session")),
2845            "missing session should recreate"
2846        );
2847    }
2848
2849    /// The `session_recreate_history` is capacity-bounded (256), unlike the
2850    /// old age-only prune which never evicted a still-recent entry. Under more
2851    /// than that many distinct peers retrying within the window, the cache can evict
2852    /// a recent entry, costing at most one extra recreate for that peer
2853    /// (bounded and self-healing: re-stamped on the next receipt), never the
2854    /// unbounded prekey loop the throttle prevents. Documents that trade-off.
2855    #[tokio::test]
2856    async fn session_recreate_history_is_capacity_bounded() {
2857        let client =
2858            crate::test_utils::create_test_client_with_failing_http("session_recreate_history_cap")
2859                .await;
2860        let now = wacore::time::Instant::now();
2861        let cap: u64 = 256;
2862
2863        // Insert well over the cap of distinct, all-recent peers.
2864        for i in 0..(cap * 2) {
2865            let jid = Jid::lid_device(format!("{}", 900_000_000_000_000u64 + i), 3);
2866            client.session_recreate_history.insert(jid, now).await;
2867        }
2868        client.session_recreate_history.run_pending_tasks().await;
2869
2870        let count = client.session_recreate_history.entry_count();
2871        assert!(
2872            count <= cap,
2873            "capacity must bound the throttle history (got {count}, cap {cap}); \
2874             a still-recent entry can be evicted under heavy peer load"
2875        );
2876    }
2877
2878    /// The resend rate limiter is reachable and tunable through the public
2879    /// `Client` API, and its drops surface on `stats().resends_throttled`. Covers
2880    /// the wiring the `handle_retry_receipt` hook relies on; the bucket logic
2881    /// itself is unit-tested in `resend_rate_limiter`.
2882    #[tokio::test]
2883    async fn client_resend_rate_limiter_is_wired_and_tunable() {
2884        let client =
2885            crate::test_utils::create_test_client_with_failing_http("resend_rl_wired").await;
2886        let chat: Jid = "120363021033254949@g.us".parse().unwrap();
2887
2888        // Tight ceiling, no refill: the bucket holds exactly `burst` tokens.
2889        client.set_resend_rate_limit(3, 0);
2890        let mut allowed = 0;
2891        for _ in 0..10 {
2892            if client.resend_rate_limiter.try_acquire(&chat).await {
2893                allowed += 1;
2894            }
2895        }
2896        assert_eq!(allowed, 3, "client honors the configured per-chat burst");
2897        assert_eq!(
2898            client.stats().resends_throttled,
2899            7,
2900            "public counter tracks dropped resends"
2901        );
2902
2903        // Disabling restores unthrottled behavior.
2904        client.set_resend_rate_limit(0, 0);
2905        let other: Jid = "120363000000000001@g.us".parse().unwrap();
2906        for _ in 0..50 {
2907            assert!(client.resend_rate_limiter.try_acquire(&other).await);
2908        }
2909    }
2910
2911    /// End-to-end: a throttled group retry drops the resend (returns Ok, sends
2912    /// nothing) while the path up to the limiter still runs, and the cached
2913    /// message is retained for the device's later re-request. Exercises the hook
2914    /// placement and the no-resend-on-refusal semantics the unit tests cannot.
2915    #[tokio::test]
2916    async fn handle_retry_receipt_drops_throttled_group_resend() {
2917        use wacore_binary::builder::NodeBuilder;
2918
2919        let backend = crate::test_utils::create_test_backend().await;
2920        let pm = Arc::new(PersistenceManager::new(backend).await.unwrap());
2921        let mut config = crate::cache_config::CacheConfig::default();
2922        config.recent_messages.capacity = 1_000;
2923        let (client, _rx) = Client::new_with_cache_config(
2924            Arc::new(crate::runtime_impl::TokioRuntime),
2925            pm,
2926            Arc::new(crate::transport::mock::MockTransportFactory::new()),
2927            Arc::new(MockHttpClient),
2928            None,
2929            config,
2930        )
2931        .await;
2932
2933        let group: Jid = "120363021033254949@g.us".parse().unwrap();
2934        let msg_id = "RLMSG001";
2935        client
2936            .add_recent_message(
2937                &group,
2938                msg_id,
2939                &wa::Message {
2940                    conversation: Some("hi".into()),
2941                    ..Default::default()
2942                },
2943                None,
2944            )
2945            .await;
2946
2947        // Drain the single token so the incoming retry must be throttled; the
2948        // throttle returns before any network resend, keeping the test offline.
2949        client.set_resend_rate_limit(1, 0);
2950        assert!(client.resend_rate_limiter.try_acquire(&group).await);
2951
2952        // Inbound group retry from a device-0 LID participant: has_device's
2953        // device-0 fast path makes it known, LID skips rotateKey, and no <keys>
2954        // leaves update_local_signal_session a noop on a missing session.
2955        let node = NodeBuilder::new("receipt")
2956            .attr("participant", "555000111@lid")
2957            .children([NodeBuilder::new("retry")
2958                .attr("id", msg_id)
2959                .attr("count", "1")
2960                .build()])
2961            .build();
2962        let node_ref = crate::test_utils::node_to_owned_ref(&node);
2963        let receipt = Receipt::builder()
2964            .source(crate::types::message::MessageSource {
2965                chat: group.clone(),
2966                sender: "555000111@lid".parse().unwrap(),
2967                is_group: true,
2968                ..Default::default()
2969            })
2970            .message_ids(vec![msg_id.to_string()])
2971            .timestamp(wacore::time::now_utc())
2972            .r#type(crate::types::presence::ReceiptType::Retry)
2973            .offline(false)
2974            .build();
2975
2976        let result = client.handle_retry_receipt(&receipt, &node_ref).await;
2977        assert!(
2978            result.is_ok(),
2979            "a throttled retry returns Ok(()), not an error"
2980        );
2981        assert_eq!(
2982            client.stats().resends_throttled,
2983            1,
2984            "the resend was dropped by the limiter"
2985        );
2986        assert!(
2987            client.peek_recent_message(&group, msg_id).await.is_some(),
2988            "throttling keeps the message cached for the device's re-request"
2989        );
2990        assert_eq!(
2991            client.pending_retries.lock().unwrap().len(),
2992            0,
2993            "the in-progress marker is cleared after the throttled return"
2994        );
2995    }
2996
2997    #[tokio::test]
2998    async fn unknown_participant_rotation_is_durable_before_throttled_return() {
2999        use wacore::libsignal::protocol::{SENDERKEY_MESSAGE_CURRENT_VERSION, SenderKeyRecord};
3000        use wacore::libsignal::store::sender_key_name::SenderKeyName;
3001        use wacore_binary::builder::NodeBuilder;
3002
3003        let backend = crate::test_utils::create_test_backend().await;
3004        let pm = Arc::new(PersistenceManager::new(backend.clone()).await.unwrap());
3005        let mut config = crate::cache_config::CacheConfig::default();
3006        config.recent_messages.capacity = 1_000;
3007        let (client, _rx) = Client::new_with_cache_config(
3008            Arc::new(crate::runtime_impl::TokioRuntime),
3009            pm,
3010            Arc::new(crate::transport::mock::MockTransportFactory::new()),
3011            Arc::new(MockHttpClient),
3012            None,
3013            config,
3014        )
3015        .await;
3016
3017        let own_lid: Jid = "100000000001040:13@lid".parse().unwrap();
3018        client
3019            .persistence_manager
3020            .process_command(crate::store::commands::DeviceCommand::SetLid(Some(
3021                own_lid.clone(),
3022            )))
3023            .await;
3024        let group: Jid = "120363021033254950@g.us".parse().unwrap();
3025        let group_id = group.to_string();
3026        let sender_key_name =
3027            SenderKeyName::from_parts(&group_id, own_lid.to_protocol_address().as_str());
3028        let mut rng = rand::make_rng::<rand::rngs::StdRng>();
3029        let key_pair = KeyPair::generate(&mut rng);
3030        let mut record = SenderKeyRecord::new_empty();
3031        record
3032            .add_sender_key_state(
3033                SENDERKEY_MESSAGE_CURRENT_VERSION,
3034                9,
3035                0,
3036                &[7; 32],
3037                key_pair.public_key,
3038                Some(key_pair.private_key),
3039            )
3040            .unwrap();
3041        client
3042            .signal_cache
3043            .put_sender_key(&sender_key_name, record)
3044            .await;
3045        client.flush_signal_cache().await.unwrap();
3046        assert!(
3047            backend
3048                .get_sender_key(sender_key_name.cache_key())
3049                .await
3050                .unwrap()
3051                .is_some()
3052        );
3053
3054        let msg_id = "ROTATEFLUSH001";
3055        client
3056            .add_recent_message(
3057                &group,
3058                msg_id,
3059                &wa::Message {
3060                    conversation: Some("hi".into()),
3061                    ..Default::default()
3062                },
3063                None,
3064            )
3065            .await;
3066        client.set_resend_rate_limit(1, 0);
3067        assert!(client.resend_rate_limiter.try_acquire(&group).await);
3068
3069        let requester: Jid = "15551234002@s.whatsapp.net".parse().unwrap();
3070        let node = NodeBuilder::new("receipt")
3071            .attr("participant", &requester)
3072            .children([NodeBuilder::new("retry")
3073                .attr("id", msg_id)
3074                .attr("count", "1")
3075                .build()])
3076            .build();
3077        let node_ref = crate::test_utils::node_to_owned_ref(&node);
3078        let receipt = Receipt::builder()
3079            .source(crate::types::message::MessageSource {
3080                chat: group.clone(),
3081                sender: requester,
3082                is_group: true,
3083                ..Default::default()
3084            })
3085            .message_ids(vec![msg_id.to_string()])
3086            .timestamp(wacore::time::now_utc())
3087            .r#type(crate::types::presence::ReceiptType::Retry)
3088            .offline(false)
3089            .build();
3090
3091        client
3092            .handle_retry_receipt(&receipt, &node_ref)
3093            .await
3094            .unwrap();
3095        assert!(
3096            backend
3097                .get_sender_key(sender_key_name.cache_key())
3098                .await
3099                .unwrap()
3100                .is_none(),
3101            "early retry return must not leave the retired key durable"
3102        );
3103    }
3104
3105    /// Atomicity guard for the per-peer session lock the retry caller wraps
3106    /// around the recreate check+stamp. The cache's get+insert is not atomic, and
3107    /// same-peer retries for different message_ids dispatch concurrently, so
3108    /// without the lock both could observe a cold history and recreate. Holding
3109    /// `session_lock_for` serializes the decision: exactly one recreate fires.
3110    /// (Mirrors the caller's lock; the matrix test covers the sequential logic.)
3111    #[tokio::test]
3112    async fn concurrent_same_peer_recreate_check_is_serialized() {
3113        let client =
3114            crate::test_utils::create_test_client_with_failing_http("concurrent_recreate").await;
3115        let jid = Jid::lid_device("999999999999993".to_string(), 3);
3116
3117        // Seed a session so the retry>=2 throttle branch is exercised (the
3118        // no-session branch always stamps and would not show serialization).
3119        let session_bytes = valid_serialized_session(8888, vec![0xCC; 32]);
3120        client
3121            .persistence_manager
3122            .backend()
3123            .put_session(jid.to_protocol_address().as_str(), &session_bytes)
3124            .await
3125            .unwrap();
3126
3127        let c1 = client.clone();
3128        let j1 = jid.clone();
3129        let task1 = async move {
3130            let addr = j1.to_protocol_address();
3131            let lock = c1.session_lock_for(addr.as_str()).await;
3132            let _g = lock.lock().await;
3133            c1.should_recreate_session(2, &j1).await.is_some()
3134        };
3135        let c2 = client.clone();
3136        let j2 = jid.clone();
3137        let task2 = async move {
3138            let addr = j2.to_protocol_address();
3139            let lock = c2.session_lock_for(addr.as_str()).await;
3140            let _g = lock.lock().await;
3141            c2.should_recreate_session(2, &j2).await.is_some()
3142        };
3143        let (a, b) = tokio::join!(task1, task2);
3144
3145        assert_eq!(
3146            usize::from(a) + usize::from(b),
3147            1,
3148            "exactly one of two concurrent same-peer recreate checks may fire; \
3149             the per-peer session lock serializes the non-atomic get+insert"
3150        );
3151    }
3152
3153    /// WA Web calls `ensureE2ESessions([g])` before resending for all chat types
3154    /// (RetryRequest.js:200). When the session already exists, this MUST be a
3155    /// fast no-op — otherwise group/status retries would hit the network on
3156    /// every receipt, defeating the cache. Regression guard for the group-branch
3157    /// call added alongside this test.
3158    #[tokio::test]
3159    async fn ensure_e2e_sessions_resolved_is_noop_when_session_exists() {
3160        use std::sync::atomic::Ordering;
3161
3162        let client = crate::test_utils::create_test_client_with_failing_http(
3163            "group_retry_ensure_sessions_noop",
3164        )
3165        .await;
3166
3167        // Bypass the offline-delivery wait that ensureE2ESessions does first.
3168        client.offline_sync_completed.store(true, Ordering::Relaxed);
3169
3170        let resolved_jid = Jid::lid_device("100000000000199".to_string(), 17);
3171        let signal_address = resolved_jid.to_protocol_address();
3172
3173        let session_bytes = valid_serialized_session(5555, vec![0xDD; 32]);
3174        client
3175            .persistence_manager
3176            .backend()
3177            .put_session(signal_address.as_str(), &session_bytes)
3178            .await
3179            .unwrap();
3180
3181        // With a session present, no prekey fetch should happen (the test
3182        // client has no wired IQ responder, so a fetch would hang/error).
3183        client
3184            .ensure_e2e_sessions_resolved(std::slice::from_ref(&resolved_jid))
3185            .await
3186            .expect("no-op when session exists");
3187    }
3188
3189    #[tokio::test]
3190    async fn retry_key_bundle_requires_one_time_prekey_except_fbid_bot() {
3191        let backend = crate::test_utils::create_test_backend().await;
3192        let pm = Arc::new(
3193            PersistenceManager::new(backend)
3194                .await
3195                .expect("persistence manager should initialize"),
3196        );
3197        let (client, _sync_rx) = Client::new(
3198            Arc::new(crate::runtime_impl::TokioRuntime),
3199            pm,
3200            Arc::new(crate::transport::mock::MockTransportFactory::new()),
3201            Arc::new(MockHttpClient),
3202            None,
3203        )
3204        .await;
3205
3206        let mut rng = rand::make_rng::<rand::rngs::StdRng>();
3207        let remote_identity = IdentityKeyPair::generate(&mut rng);
3208        let signed_prekey = KeyPair::generate(&mut rng);
3209        let signed_prekey_signature = remote_identity
3210            .private_key()
3211            .calculate_signature(&signed_prekey.public_key.serialize(), &mut rng)
3212            .expect("signed prekey signature should be valid");
3213
3214        let regular_requester = Jid::pn_device("559922223333", 1);
3215        let keys = NodeBuilder::new("keys")
3216            .children([
3217                NodeBuilder::new("type").bytes(vec![5]).build(),
3218                NodeBuilder::new("identity")
3219                    .bytes(
3220                        remote_identity
3221                            .identity_key()
3222                            .public_key()
3223                            .public_key_bytes()
3224                            .to_vec(),
3225                    )
3226                    .build(),
3227                SignedPreKeyNode::new(
3228                    100,
3229                    signed_prekey.public_key.public_key_bytes().to_vec(),
3230                    signed_prekey_signature.to_vec(),
3231                )
3232                .into_node(),
3233            ])
3234            .build();
3235        let receipt = NodeBuilder::new("receipt")
3236            .children([
3237                NodeBuilder::new("registration")
3238                    .bytes(12345u32.to_be_bytes().to_vec())
3239                    .build(),
3240                keys,
3241            ])
3242            .build();
3243
3244        let err = client
3245            .process_retry_key_bundle(&receipt.as_node_ref(), &regular_requester, false, false)
3246            .await
3247            .expect_err("regular retry without one-time prekey must be rejected");
3248        assert!(
3249            err.to_string()
3250                .contains("regular retry key bundle missing one-time prekey")
3251        );
3252
3253        let fbid_bot_requester = Jid::new("200000000000002", wacore_binary::Server::Bot);
3254        client
3255            .process_retry_key_bundle(&receipt.as_node_ref(), &fbid_bot_requester, false, true)
3256            .await
3257            .expect("fbid bot retry without one-time prekey should establish a session");
3258
3259        let snapshot = client.persistence_manager.get_device_snapshot();
3260        let session = client
3261            .signal_cache
3262            .peek_session(
3263                &fbid_bot_requester.to_protocol_address(),
3264                &*snapshot.backend,
3265            )
3266            .await
3267            .expect("session lookup should succeed");
3268        assert!(session.is_some());
3269    }
3270
3271    #[test]
3272    fn bot_jid_detection() {
3273        // Test bot JID detection for bot message filtering
3274        use wacore_binary::JidExt as _;
3275
3276        // Regular user JID - not a bot
3277        let regular_user: Jid = "1234567890@s.whatsapp.net".parse().unwrap();
3278        assert!(!regular_user.is_bot());
3279
3280        // Bot JID with bot server
3281        let bot_server: Jid = "somebot@bot".parse().unwrap();
3282        assert!(bot_server.is_bot());
3283
3284        // Legacy bot JID pattern (1313555...)
3285        let legacy_bot: Jid = "1313555123456@s.whatsapp.net".parse().unwrap();
3286        assert!(legacy_bot.is_bot());
3287
3288        // Legacy bot JID pattern (131655500...)
3289        let legacy_bot2: Jid = "131655500123456@s.whatsapp.net".parse().unwrap();
3290        assert!(legacy_bot2.is_bot());
3291
3292        // Similar but not bot (doesn't start with exact prefix)
3293        let not_bot: Jid = "1313556123456@s.whatsapp.net".parse().unwrap();
3294        assert!(!not_bot.is_bot());
3295    }
3296
3297    #[test]
3298    fn extract_registration_id_from_node_test() {
3299        use wacore::protocol::retry::{
3300            extract_registration_id_from_node, extract_registration_id_from_node_ref,
3301        };
3302        use wacore_binary::{Attrs, Node};
3303
3304        let reg_receipt = |bytes: Vec<u8>| Node {
3305            tag: Cow::Borrowed("receipt"),
3306            attrs: Attrs::new(),
3307            content: Some(NodeContent::Nodes(vec![Node {
3308                tag: Cow::Borrowed("registration"),
3309                attrs: Attrs::new(),
3310                content: Some(NodeContent::Bytes(bytes)),
3311            }])),
3312        };
3313
3314        // 4-byte registration ID.
3315        let parent = reg_receipt(vec![0x00, 0x01, 0x02, 0x03]);
3316        assert_eq!(extract_registration_id_from_node(&parent), Some(0x00010203));
3317        assert_eq!(
3318            extract_registration_id_from_node_ref(&parent.as_node_ref()),
3319            Some(0x00010203)
3320        );
3321
3322        // 3-byte registration ID (variable length, left zero-padded).
3323        let parent_short = reg_receipt(vec![0x01, 0x02, 0x03]);
3324        assert_eq!(
3325            extract_registration_id_from_node(&parent_short),
3326            Some(0x00010203)
3327        );
3328        assert_eq!(
3329            extract_registration_id_from_node_ref(&parent_short.as_node_ref()),
3330            Some(0x00010203)
3331        );
3332
3333        // Oversized (>4 byte) payload: rejected, not truncated, on both paths.
3334        let parent_oversized = reg_receipt(vec![0x01, 0x02, 0x03, 0x04, 0x05]);
3335        assert_eq!(extract_registration_id_from_node(&parent_oversized), None);
3336        assert_eq!(
3337            extract_registration_id_from_node_ref(&parent_oversized.as_node_ref()),
3338            None
3339        );
3340
3341        // No registration node.
3342        let parent_no_reg = Node {
3343            tag: Cow::Borrowed("receipt"),
3344            attrs: Attrs::new(),
3345            content: Some(NodeContent::Nodes(vec![])),
3346        };
3347        assert_eq!(extract_registration_id_from_node(&parent_no_reg), None);
3348        assert_eq!(
3349            extract_registration_id_from_node_ref(&parent_no_reg.as_node_ref()),
3350            None
3351        );
3352
3353        // Empty bytes.
3354        let parent_empty = reg_receipt(vec![]);
3355        assert_eq!(extract_registration_id_from_node(&parent_empty), None);
3356        assert_eq!(
3357            extract_registration_id_from_node_ref(&parent_empty.as_node_ref()),
3358            None
3359        );
3360    }
3361
3362    #[test]
3363    fn group_or_status_detection_for_sender_key_handling() {
3364        // Test that both groups and status broadcasts trigger sender key handling
3365        use wacore_binary::JidExt as _;
3366
3367        let group: Jid = "120363021033254949@g.us".parse().unwrap();
3368        let status: Jid = "status@broadcast".parse().unwrap();
3369        let dm: Jid = "1234567890@s.whatsapp.net".parse().unwrap();
3370
3371        // Both group and status should trigger sender key deletion
3372        assert!(group.is_group() || group.is_status_broadcast());
3373        assert!(status.is_group() || status.is_status_broadcast());
3374
3375        // DM should NOT trigger sender key deletion
3376        assert!(!(dm.is_group() || dm.is_status_broadcast()));
3377    }
3378
3379    #[test]
3380    fn retransmission_route_validation_is_strict_and_typed() {
3381        let direct: Jid = "12025550100@s.whatsapp.net".parse().unwrap();
3382        let requester: Jid = "12025550100:7@s.whatsapp.net".parse().unwrap();
3383        let group: Jid = "120363000000000001@g.us".parse().unwrap();
3384        let status = Jid::status_broadcast();
3385        let broadcast: Jid = "1234567890@broadcast".parse().unwrap();
3386
3387        assert!(matches!(
3388            validate_retransmission(&direct, &requester, "DM1", 1, Some(&direct)),
3389            Ok(RetransmissionRoute::Direct)
3390        ));
3391        assert!(matches!(
3392            validate_retransmission(&group, &requester, "GROUP1", 1, None),
3393            Ok(RetransmissionRoute::Group)
3394        ));
3395        assert!(matches!(
3396            validate_retransmission(&status, &requester, "STATUS1", 1, None),
3397            Ok(RetransmissionRoute::Status)
3398        ));
3399        assert!(matches!(
3400            validate_retransmission(&broadcast, &requester, "BROADCAST1", 1, None),
3401            Ok(RetransmissionRoute::BroadcastList)
3402        ));
3403
3404        for (id, count) in [("ZERO", 0), ("", 1), ("MAX", MAX_RETRY_COUNT)] {
3405            assert!(
3406                validate_retransmission(&direct, &requester, id, count, None).is_err(),
3407                "invalid id/count pair must fail: {id:?}/{count}"
3408            );
3409        }
3410        assert!(
3411            validate_retransmission(&group, &requester, "GROUP2", 1, Some(&direct)).is_err(),
3412            "recipient is only meaningful on a direct retry"
3413        );
3414        assert!(
3415            validate_retransmission(&status, &group, "STATUS2", 1, None).is_err(),
3416            "a group JID cannot be a requesting status device"
3417        );
3418    }
3419
3420    #[tokio::test]
3421    async fn public_peer_retransmission_requires_a_recipient() {
3422        let client = crate::test_utils::create_test_client().await;
3423        let own_pn: Jid = "12025550100:13@s.whatsapp.net".parse().unwrap();
3424        client
3425            .persistence_manager
3426            .process_command(crate::store::commands::DeviceCommand::SetId(Some(
3427                own_pn.clone(),
3428            )))
3429            .await;
3430
3431        let chat: Jid = "12025550101@s.whatsapp.net".parse().unwrap();
3432        let requester = own_pn.with_device(7);
3433        let request = MessageRetransmission::new(
3434            chat,
3435            requester,
3436            wa::Message::default(),
3437            "PEER-RETRY-1".to_string(),
3438            1,
3439        );
3440
3441        let error = client
3442            .retransmit_message(request)
3443            .await
3444            .expect_err("a peer route without its actual chat cannot be sent");
3445        assert!(matches!(error, SendError::InvalidRequest(_)));
3446        assert!(error.to_string().contains("requires a recipient"));
3447    }
3448
3449    #[tokio::test]
3450    async fn public_direct_retransmission_binds_chat_to_routing_identity() {
3451        let client = crate::test_utils::create_test_client().await;
3452        let chat = Jid::pn("12025550104");
3453        let requester = Jid::pn_device("12025550105", 7);
3454        let bot_requester: Jid = "200000000000002@bot".parse().unwrap();
3455
3456        for request in [
3457            MessageRetransmission::new(
3458                chat.clone(),
3459                requester,
3460                wa::Message::default(),
3461                "DIRECT-CHAT-MISMATCH-1".to_string(),
3462                1,
3463            ),
3464            MessageRetransmission::new(
3465                chat.clone(),
3466                bot_requester,
3467                wa::Message::default(),
3468                "DIRECT-RECIPIENT-MISMATCH-1".to_string(),
3469                1,
3470            )
3471            .with_recipient(Jid::pn("12025550106")),
3472        ] {
3473            let error = client
3474                .retransmit_message(request)
3475                .await
3476                .expect_err("an unrelated routing identity must be rejected");
3477            assert!(matches!(error, SendError::InvalidRequest(_)));
3478            assert!(error.to_string().contains("routing identity"));
3479        }
3480    }
3481
3482    #[tokio::test]
3483    async fn public_direct_recipient_rejects_an_unrelated_requester() {
3484        let client = crate::test_utils::create_test_client().await;
3485        let chat = Jid::pn("12025550108");
3486        let request = MessageRetransmission::new(
3487            chat.clone(),
3488            Jid::pn_device("12025550109", 7),
3489            wa::Message::default(),
3490            "DIRECT-RECIPIENT-SOURCE-1".to_string(),
3491            1,
3492        )
3493        .with_recipient(chat);
3494
3495        let error = client
3496            .retransmit_message(request)
3497            .await
3498            .expect_err("a normal remote user cannot declare a recipient route");
3499        assert!(matches!(error, SendError::InvalidRequest(_)));
3500        assert!(error.to_string().contains("local device or bot"));
3501    }
3502
3503    #[tokio::test]
3504    async fn direct_retransmission_chat_accepts_known_pn_lid_alias() {
3505        let client = crate::test_utils::create_test_client().await;
3506        let pn = Jid::pn("12025550107");
3507        let lid = Jid::lid("100000000000107");
3508        client
3509            .lid_pn_cache
3510            .add(&wacore::types::lid_pn::LidPnEntry {
3511                lid: lid.user.as_str().into(),
3512                phone_number: pn.user.as_str().into(),
3513                created_at: 1,
3514                learning_source: wacore::types::lid_pn::LearningSource::Usync,
3515            })
3516            .await;
3517
3518        assert!(client.jids_share_user_identity(&pn, &lid).await.unwrap());
3519        assert!(client.jids_share_user_identity(&lid, &pn).await.unwrap());
3520    }
3521
3522    #[tokio::test]
3523    async fn public_retransmission_recaches_the_supplied_message() {
3524        let mut config = crate::cache_config::CacheConfig::default();
3525        config.recent_messages.capacity = 16;
3526        let client = crate::test_utils::create_test_client_with_config(
3527            "public_retransmission_cache",
3528            Arc::new(MockHttpClient),
3529            config,
3530        )
3531        .await;
3532        let chat = Jid::pn("12025550103");
3533        let requester = chat.with_device(7);
3534        crate::test_utils::seed_peer_session(&client, &requester).await;
3535        let message = wa::Message {
3536            conversation: Some("retry me".into()),
3537            ..Default::default()
3538        };
3539        let message_id = "PUBLIC-RETRY-CACHE-1";
3540
3541        // The fresh test session emits pkmsg and this client intentionally has
3542        // no device identity, so the wire attempt fails after the public API has
3543        // accepted and cached the supplied message.
3544        let result = client
3545            .retransmit_message(MessageRetransmission::new(
3546                chat.clone(),
3547                requester,
3548                message,
3549                message_id.to_string(),
3550                1,
3551            ))
3552            .await;
3553        assert!(result.is_err());
3554
3555        let (cached, alternate) = client
3556            .peek_recent_message(&chat, message_id)
3557            .await
3558            .expect("a later retry count must find the retransmitted message");
3559        assert!(alternate.is_none());
3560        assert_eq!(cached.conversation.as_deref(), Some("retry me"));
3561    }
3562
3563    #[test]
3564    fn resolve_retry_chat_info_broadcast_uses_participant_device() {
3565        let broadcast = "1234567890@broadcast";
3566        let participant = "12025550101:9@s.whatsapp.net";
3567        let node = NodeBuilder::new("receipt")
3568            .attr("participant", participant)
3569            .build();
3570        let receipt = make_test_receipt(broadcast);
3571        let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
3572
3573        assert!(info.chat.is_broadcast_list());
3574        assert_eq!(info.requester, participant.parse::<Jid>().unwrap());
3575    }
3576
3577    /// The key-bundle policy is driven only by explicit force, stateless routing,
3578    /// and the retry threshold. The diagnostic reason must not change the wire
3579    /// shape of a first retry.
3580    #[test]
3581    fn retry_key_inclusion_matches_canonical_policy() {
3582        use wacore::protocol::retry::{should_include_keys, should_include_keys_with_policy};
3583
3584        assert!(!should_include_keys(1, RetryReason::NoSession));
3585        assert!(!should_include_keys(
3586            1,
3587            RetryReason::UnknownCompanionNoPrekey
3588        ));
3589        assert!(should_include_keys_with_policy(1, true, false));
3590        assert!(should_include_keys_with_policy(1, false, true));
3591        assert!(should_include_keys(2, RetryReason::InvalidMessage));
3592        assert!(should_include_keys(3, RetryReason::BadMac));
3593    }
3594
3595    /// Helper to build a DM Receipt for testing resolve_retry_chat_info.
3596    fn make_test_receipt(from: &str) -> Receipt {
3597        Receipt::builder()
3598            .source(crate::types::message::MessageSource {
3599                chat: from.parse().unwrap(),
3600                sender: from.parse().unwrap(),
3601                ..Default::default()
3602            })
3603            .message_ids(vec!["MSG001".to_string()])
3604            .timestamp(wacore::time::now_utc())
3605            .r#type(crate::types::presence::ReceiptType::Retry)
3606            .offline(false)
3607            .build()
3608    }
3609
3610    #[test]
3611    fn resolve_retry_chat_info_dm_with_device() {
3612        use wacore_binary::builder::NodeBuilder;
3613
3614        // Node attrs are unused in the DM branch (no participant lookup)
3615        let node = NodeBuilder::new("receipt").build();
3616        let receipt = make_test_receipt("5511999999999:33@s.whatsapp.net");
3617        let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
3618
3619        // chat should be bare (device stripped)
3620        assert_eq!(info.chat.device(), 0);
3621        assert_eq!(info.chat.user, "5511999999999");
3622        assert!(info.chat.is_pn());
3623
3624        // requester should preserve device 33
3625        assert_eq!(info.requester.device(), 33);
3626        assert_eq!(info.requester.user, "5511999999999");
3627    }
3628
3629    #[test]
3630    fn resolve_retry_chat_info_lid_dm_with_device() {
3631        use wacore_binary::builder::NodeBuilder;
3632
3633        let node = NodeBuilder::new("receipt").build();
3634        let receipt = make_test_receipt("236395184570386:5@lid");
3635        let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
3636
3637        // chat should be bare LID (device stripped)
3638        assert_eq!(info.chat.device(), 0);
3639        assert_eq!(info.chat.user, "236395184570386");
3640        assert!(info.chat.is_lid());
3641
3642        // requester should preserve device 5
3643        assert_eq!(info.requester.device(), 5);
3644        assert_eq!(info.requester.user, "236395184570386");
3645        assert!(info.requester.is_lid());
3646    }
3647
3648    /// `info.recipient` must come from the receipt's `recipient` attribute,
3649    /// not derived from `info.chat`. Pre-fix, the DM resend used
3650    /// `info.chat.clone()` for the stanza's `recipient` — fine on the primary
3651    /// namespace but wrong whenever `take_recent_message` hit `alt_chat` (the
3652    /// original was sent under PN while the receipt arrived under LID, or
3653    /// vice-versa). WA Web's `WAWebHandleRetryRequest` forwards the receipt
3654    /// attr verbatim (`f && (k.recipient = f)`), so the resend's `recipient`
3655    /// matches the original outbound's namespace regardless of how the
3656    /// receipt's `from` was addressed.
3657    #[test]
3658    fn resolve_retry_chat_info_forwards_recipient_attribute_verbatim() {
3659        use wacore_binary::builder::NodeBuilder;
3660
3661        // Cross-namespace shape: receipt `from` is LID, `recipient` is PN.
3662        let node = NodeBuilder::new("receipt")
3663            .attr("recipient", "5500000000123@s.whatsapp.net")
3664            .build();
3665        let receipt = make_test_receipt("100000000000456:5@lid");
3666        let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
3667
3668        let recipient = info
3669            .recipient
3670            .as_ref()
3671            .expect("recipient must be populated from the node attr");
3672        assert_eq!(recipient.user, "5500000000123");
3673        assert!(recipient.is_pn(), "recipient namespace must be PN");
3674        assert_ne!(
3675            recipient.user, info.chat.user,
3676            "recipient must come from the node attr, not info.chat"
3677        );
3678
3679        // Inverse: absent attr → None (drops `recipient` from the resend
3680        // stanza, mirroring WA Web's `f && (k.recipient = f)`).
3681        let node_no_recipient = NodeBuilder::new("receipt").build();
3682        let info_no_recipient =
3683            resolve_retry_chat_info(&receipt, &node_no_recipient.as_node_ref(), None, None);
3684        assert!(
3685            info_no_recipient.recipient.is_none(),
3686            "missing `recipient` attr must propagate as None"
3687        );
3688    }
3689
3690    #[test]
3691    fn resolve_retry_chat_info_dm_bare() {
3692        use wacore_binary::builder::NodeBuilder;
3693
3694        let node = NodeBuilder::new("receipt").build();
3695        let receipt = make_test_receipt("5511999999999@s.whatsapp.net");
3696        let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
3697
3698        assert_eq!(info.chat.device(), 0);
3699        assert_eq!(info.requester.device(), 0);
3700        assert_eq!(info.chat, info.requester);
3701    }
3702
3703    #[test]
3704    fn resolve_retry_chat_info_group() {
3705        use wacore_binary::builder::NodeBuilder;
3706
3707        let node = NodeBuilder::new("receipt")
3708            .attr("from", "120363021033254949@g.us")
3709            .attr("id", "MSG001")
3710            .attr("participant", "236395184570386:33@lid")
3711            .attr("type", "retry")
3712            .build();
3713        let receipt = Receipt::builder()
3714            .source(crate::types::message::MessageSource {
3715                chat: "120363021033254949@g.us".parse().unwrap(),
3716                sender: "236395184570386:33@lid".parse().unwrap(),
3717                ..Default::default()
3718            })
3719            .message_ids(vec!["MSG001".to_string()])
3720            .timestamp(wacore::time::now_utc())
3721            .r#type(crate::types::presence::ReceiptType::Retry)
3722            .offline(false)
3723            .build();
3724        let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
3725
3726        assert!(info.chat.is_group());
3727        assert_eq!(info.chat.user, "120363021033254949");
3728        assert!(info.requester.is_lid());
3729        assert_eq!(info.requester.device(), 33);
3730    }
3731
3732    #[test]
3733    fn resolve_retry_chat_info_group_bot_device_marks_bot_namespace_only() {
3734        use wacore_binary::builder::NodeBuilder;
3735
3736        let node = NodeBuilder::new("receipt")
3737            .attr("participant", "somebot:4@bot")
3738            .build();
3739        let receipt = Receipt::builder()
3740            .source(crate::types::message::MessageSource {
3741                chat: "120363021033254949@g.us".parse().unwrap(),
3742                sender: "somebot:4@bot".parse().unwrap(),
3743                ..Default::default()
3744            })
3745            .message_ids(vec!["MSG001".to_string()])
3746            .timestamp(wacore::time::now_utc())
3747            .r#type(crate::types::presence::ReceiptType::Retry)
3748            .offline(false)
3749            .build();
3750
3751        let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
3752
3753        assert!(info.chat.is_group());
3754        assert!(info.is_bot);
3755        assert!(!info.is_fbid_bot_retry);
3756    }
3757
3758    #[test]
3759    fn resolve_retry_chat_info_group_primary_fbid_bot_marks_bot_retry() {
3760        use wacore_binary::builder::NodeBuilder;
3761
3762        let node = NodeBuilder::new("receipt")
3763            .attr("participant", "somebot@bot")
3764            .build();
3765        let receipt = Receipt::builder()
3766            .source(crate::types::message::MessageSource {
3767                chat: "120363021033254949@g.us".parse().unwrap(),
3768                sender: "somebot@bot".parse().unwrap(),
3769                ..Default::default()
3770            })
3771            .message_ids(vec!["MSG001".to_string()])
3772            .timestamp(wacore::time::now_utc())
3773            .r#type(crate::types::presence::ReceiptType::Retry)
3774            .offline(false)
3775            .build();
3776
3777        let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
3778
3779        assert!(info.chat.is_group());
3780        assert!(info.is_bot);
3781        assert!(info.is_fbid_bot_retry);
3782    }
3783
3784    #[test]
3785    fn resolve_retry_chat_info_status_broadcast() {
3786        use wacore_binary::builder::NodeBuilder;
3787
3788        let node = NodeBuilder::new("receipt")
3789            .attr("from", "status@broadcast")
3790            .attr("id", "3EB06D00CAB92340790621")
3791            .attr("participant", "236395184570386@lid")
3792            .attr("type", "retry")
3793            .build();
3794        let receipt = make_test_receipt("status@broadcast");
3795        let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
3796
3797        assert!(info.chat.is_status_broadcast());
3798        // requester should be the participant, not status@broadcast
3799        assert!(info.requester.is_lid());
3800        assert_eq!(info.requester.user, "236395184570386");
3801    }
3802
3803    #[test]
3804    fn resolve_retry_chat_info_status_broadcast_no_participant() {
3805        use wacore_binary::builder::NodeBuilder;
3806
3807        // Missing participant attr (edge case) — falls back to sender
3808        let node = NodeBuilder::new("receipt")
3809            .attr("from", "status@broadcast")
3810            .attr("id", "MSG001")
3811            .attr("type", "retry")
3812            .build();
3813        let receipt = make_test_receipt("status@broadcast");
3814        let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
3815
3816        assert!(info.chat.is_status_broadcast());
3817        assert!(info.requester.is_status_broadcast());
3818    }
3819
3820    // Different participants get different keys; same participant keeps the same
3821    // key across retry counts so pending_retries serializes concurrent receipts.
3822    #[test]
3823    fn retry_processing_key_per_participant() {
3824        let msg_id = "3EB06D00CAB92340790621";
3825
3826        let status_chat = Jid::status_broadcast();
3827        let status_participant_a: Jid = "236395184570386@lid".parse().unwrap();
3828        let status_participant_b: Jid = "559985213786@s.whatsapp.net".parse().unwrap();
3829        let status_key_a = build_retry_processing_key(&status_chat, msg_id, &status_participant_a);
3830        let status_key_b = build_retry_processing_key(&status_chat, msg_id, &status_participant_b);
3831        assert_ne!(
3832            status_key_a, status_key_b,
3833            "Different status participants must have different processing keys"
3834        );
3835        assert_eq!(
3836            status_key_a,
3837            build_retry_processing_key(&status_chat, msg_id, &status_participant_a),
3838            "Same participant must produce the same key — any retry count for that \
3839             participant serializes through pending_retries"
3840        );
3841
3842        let dm_chat = Jid::pn("559911112222");
3843        let dm_device_a = Jid::pn_device("559922223333", 1);
3844        let dm_device_b = Jid::pn_device("559922223333", 2);
3845        let dm_key_a = build_retry_processing_key(&dm_chat, msg_id, &dm_device_a);
3846        let dm_key_b = build_retry_processing_key(&dm_chat, msg_id, &dm_device_b);
3847        assert_ne!(
3848            dm_key_a, dm_key_b,
3849            "Different DM requester devices must have different processing keys"
3850        );
3851        assert_eq!(
3852            dm_key_a,
3853            build_retry_processing_key(&dm_chat, msg_id, &dm_device_a),
3854            "Same DM requester device must produce the same processing key"
3855        );
3856    }
3857
3858    /// Test that the recent message cache supports re-addition after take.
3859    /// This is critical for multi-device retries where another device can
3860    /// ask for the same message after the first retry already consumed it.
3861    #[tokio::test]
3862    async fn recent_message_cache_readd_after_take() {
3863        let _ = env_logger::builder().is_test(true).try_init();
3864
3865        let backend = crate::test_utils::create_test_backend().await;
3866        let pm = Arc::new(
3867            PersistenceManager::new(backend)
3868                .await
3869                .expect("persistence manager should initialize"),
3870        );
3871        // Enable L1 cache so MockBackend (which doesn't persist) works for this test
3872        let mut config = crate::cache_config::CacheConfig::default();
3873        config.recent_messages.capacity = 1_000;
3874        let (client, _sync_rx) = Client::new_with_cache_config(
3875            Arc::new(crate::runtime_impl::TokioRuntime),
3876            pm.clone(),
3877            Arc::new(crate::transport::mock::MockTransportFactory::new()),
3878            Arc::new(MockHttpClient),
3879            None,
3880            config,
3881        )
3882        .await;
3883
3884        let msg = wa::Message {
3885            extended_text_message: buffa::MessageField::some(wa::message::ExtendedTextMessage {
3886                text: Some("status text".to_string()),
3887                ..Default::default()
3888            }),
3889            ..Default::default()
3890        };
3891
3892        for (chat, msg_id) in [
3893            (Jid::status_broadcast(), "STATUS_MSG_001".to_string()),
3894            (Jid::pn("559911112222"), "DM_MSG_001".to_string()),
3895        ] {
3896            client.add_recent_message(&chat, &msg_id, &msg, None).await;
3897
3898            let taken = client.take_recent_message(&chat, &msg_id).await;
3899            assert!(taken.is_some(), "First take should succeed for {chat}");
3900
3901            let (taken_msg, _) = taken.unwrap();
3902            client
3903                .add_recent_message(&chat, &msg_id, &taken_msg, None)
3904                .await;
3905
3906            let taken2 = client.take_recent_message(&chat, &msg_id).await;
3907            assert!(
3908                taken2.is_some(),
3909                "Second take should succeed after re-add for {chat}"
3910            );
3911            assert_eq!(
3912                taken2
3913                    .unwrap()
3914                    .0
3915                    .extended_text_message
3916                    .as_option()
3917                    .unwrap()
3918                    .text
3919                    .as_deref(),
3920                Some("status text")
3921            );
3922        }
3923    }
3924
3925    /// Message stored under bare JID should be found when looking up via bare
3926    /// JID (the path resolve_retry_chat_info now provides for DMs).
3927    #[tokio::test]
3928    async fn dm_retry_message_lookup_uses_bare_jid() {
3929        let _ = env_logger::builder().is_test(true).try_init();
3930
3931        let backend = crate::test_utils::create_test_backend().await;
3932        let pm = Arc::new(
3933            PersistenceManager::new(backend)
3934                .await
3935                .expect("persistence manager should initialize"),
3936        );
3937        let mut config = crate::cache_config::CacheConfig::default();
3938        config.recent_messages.capacity = 1_000;
3939        let (client, _sync_rx) = Client::new_with_cache_config(
3940            Arc::new(crate::runtime_impl::TokioRuntime),
3941            pm.clone(),
3942            Arc::new(crate::transport::mock::MockTransportFactory::new()),
3943            Arc::new(MockHttpClient),
3944            None,
3945            config,
3946        )
3947        .await;
3948
3949        let bare_jid: Jid = "5511999999999@s.whatsapp.net".parse().unwrap();
3950        let msg_id = "RETRY_MSG_001";
3951        let msg = wa::Message {
3952            conversation: Some("test dm".into()),
3953            ..Default::default()
3954        };
3955
3956        // Store under bare JID (how send_message stores it)
3957        client
3958            .add_recent_message(&bare_jid, msg_id, &msg, None)
3959            .await;
3960
3961        // Lookup via bare JID should succeed (this is what info.chat provides)
3962        let taken = client.take_recent_message(&bare_jid, msg_id).await;
3963        assert!(taken.is_some(), "Lookup via bare JID should succeed");
3964        let (msg_out, alt_chat) = taken.unwrap();
3965        assert!(alt_chat.is_none(), "primary key should match for bare JID");
3966
3967        // Re-add under bare JID
3968        client
3969            .add_recent_message(&bare_jid, msg_id, &msg_out, None)
3970            .await;
3971
3972        // Second take should also work
3973        let taken2 = client.take_recent_message(&bare_jid, msg_id).await;
3974        assert!(
3975            taken2.is_some(),
3976            "Second lookup via bare JID should succeed after re-add"
3977        );
3978    }
3979
3980    /// Alternate PN/LID key lookup: a message stored under PN should be found
3981    /// when the primary lookup resolves to LID (because a mapping was added
3982    /// between send time and retry time).
3983    #[tokio::test]
3984    async fn alternate_key_lookup_pn_to_lid() {
3985        let _ = env_logger::builder().is_test(true).try_init();
3986
3987        let backend = crate::test_utils::create_test_backend().await;
3988        let pm = Arc::new(
3989            PersistenceManager::new(backend)
3990                .await
3991                .expect("persistence manager should initialize"),
3992        );
3993        let mut config = crate::cache_config::CacheConfig::default();
3994        config.recent_messages.capacity = 1_000;
3995        let (client, _sync_rx) = Client::new_with_cache_config(
3996            Arc::new(crate::runtime_impl::TokioRuntime),
3997            pm.clone(),
3998            Arc::new(crate::transport::mock::MockTransportFactory::new()),
3999            Arc::new(MockHttpClient),
4000            None,
4001            config,
4002        )
4003        .await;
4004
4005        let pn_jid: Jid = "5511999999999@s.whatsapp.net".parse().unwrap();
4006        let lid_jid: Jid = "236395184570386@lid".parse().unwrap();
4007        let msg_id = "RETRY_ALT_001";
4008        let msg = wa::Message {
4009            conversation: Some("alternate key test".into()),
4010            ..Default::default()
4011        };
4012
4013        // Store under PN (no LID mapping existed at send time)
4014        client.add_recent_message(&pn_jid, msg_id, &msg, None).await;
4015
4016        // Now add a LID mapping (simulates mapping arriving between send and retry)
4017        client
4018            .lid_pn_cache
4019            .add(&wacore::types::lid_pn::LidPnEntry {
4020                lid: lid_jid.user.as_str().into(),
4021                phone_number: pn_jid.user.as_str().into(),
4022                created_at: 0,
4023                learning_source: wacore::types::lid_pn::LearningSource::Usync,
4024            })
4025            .await;
4026
4027        // Lookup via LID: primary key resolves to LID (miss),
4028        // alternate key falls back to PN (hit)
4029        let taken = client.take_recent_message(&lid_jid, msg_id).await;
4030        assert!(
4031            taken.is_some(),
4032            "Alternate PN key lookup should find message stored under PN"
4033        );
4034        let (msg_out, alt_chat) = taken.unwrap();
4035        let alt_chat = alt_chat.expect("should be found via alternate key");
4036        assert!(alt_chat.is_pn(), "alternate chat should be PN");
4037        assert_eq!(alt_chat.user, pn_jid.user);
4038        assert_eq!(msg_out.conversation.as_deref(), Some("alternate key test"));
4039    }
4040
4041    /// swap_pn_lid_namespace should swap between PN and LID while preserving
4042    /// device/agent — this is the shared helper used for both alternate key
4043    /// computation and requester normalization after an alternate hit.
4044    #[tokio::test]
4045    async fn swap_pn_lid_namespace_preserves_device() {
4046        let _ = env_logger::builder().is_test(true).try_init();
4047
4048        let backend = crate::test_utils::create_test_backend().await;
4049        let pm = Arc::new(
4050            PersistenceManager::new(backend)
4051                .await
4052                .expect("persistence manager should initialize"),
4053        );
4054        let (client, _sync_rx) = Client::new(
4055            Arc::new(crate::runtime_impl::TokioRuntime),
4056            pm.clone(),
4057            Arc::new(crate::transport::mock::MockTransportFactory::new()),
4058            Arc::new(MockHttpClient),
4059            None,
4060        )
4061        .await;
4062
4063        let pn_jid: Jid = "5511999999999@s.whatsapp.net".parse().unwrap();
4064        let lid_jid: Jid = "236395184570386@lid".parse().unwrap();
4065
4066        client
4067            .lid_pn_cache
4068            .add(&wacore::types::lid_pn::LidPnEntry {
4069                lid: lid_jid.user.as_str().into(),
4070                phone_number: pn_jid.user.as_str().into(),
4071                created_at: 0,
4072                learning_source: wacore::types::lid_pn::LearningSource::Usync,
4073            })
4074            .await;
4075
4076        // LID:5 → PN:5
4077        let lid_with_device: Jid = "236395184570386:5@lid".parse().unwrap();
4078        let swapped = client.swap_pn_lid_namespace(&lid_with_device).await;
4079        let swapped = swapped.expect("should resolve LID→PN");
4080        assert!(swapped.is_pn());
4081        assert_eq!(swapped.user, "5511999999999");
4082        assert_eq!(swapped.device(), 5);
4083
4084        // PN:3 → LID:3
4085        let pn_with_device: Jid = "5511999999999:3@s.whatsapp.net".parse().unwrap();
4086        let swapped = client.swap_pn_lid_namespace(&pn_with_device).await;
4087        let swapped = swapped.expect("should resolve PN→LID");
4088        assert!(swapped.is_lid());
4089        assert_eq!(swapped.user, "236395184570386");
4090        assert_eq!(swapped.device(), 3);
4091
4092        // Group JID → None
4093        let group: Jid = "120363021033254949@g.us".parse().unwrap();
4094        assert!(client.swap_pn_lid_namespace(&group).await.is_none());
4095    }
4096
4097    /// Alternate key lookup via PN input: message stored under PN, LID mapping
4098    /// added later, lookup via PN. Exercises the `server != server` optimization
4099    /// where `to` is used directly as alternate (no cache round-trip).
4100    #[tokio::test]
4101    async fn alternate_key_lookup_pn_input_server_changed() {
4102        let _ = env_logger::builder().is_test(true).try_init();
4103
4104        let backend = crate::test_utils::create_test_backend().await;
4105        let pm = Arc::new(
4106            PersistenceManager::new(backend)
4107                .await
4108                .expect("persistence manager should initialize"),
4109        );
4110        let mut config = crate::cache_config::CacheConfig::default();
4111        config.recent_messages.capacity = 1_000;
4112        let (client, _sync_rx) = Client::new_with_cache_config(
4113            Arc::new(crate::runtime_impl::TokioRuntime),
4114            pm.clone(),
4115            Arc::new(crate::transport::mock::MockTransportFactory::new()),
4116            Arc::new(MockHttpClient),
4117            None,
4118            config,
4119        )
4120        .await;
4121
4122        let pn_jid: Jid = "5511999999999@s.whatsapp.net".parse().unwrap();
4123        let lid_jid: Jid = "236395184570386@lid".parse().unwrap();
4124        let msg_id = "RETRY_ALT_PN";
4125        let msg = wa::Message {
4126            conversation: Some("pn input alternate".into()),
4127            ..Default::default()
4128        };
4129
4130        // Store under PN (no mapping at send time)
4131        client.add_recent_message(&pn_jid, msg_id, &msg, None).await;
4132
4133        // Add LID mapping
4134        client
4135            .lid_pn_cache
4136            .add(&wacore::types::lid_pn::LidPnEntry {
4137                lid: lid_jid.user.as_str().into(),
4138                phone_number: pn_jid.user.as_str().into(),
4139                created_at: 0,
4140                learning_source: wacore::types::lid_pn::LearningSource::Usync,
4141            })
4142            .await;
4143
4144        // Lookup via PN: resolve_encryption_jid maps to LID (primary),
4145        // primary misses, server changed (Lid != Pn) → uses `to` directly
4146        let taken = client.take_recent_message(&pn_jid, msg_id).await;
4147        assert!(
4148            taken.is_some(),
4149            "Should find message via server-changed path"
4150        );
4151        let (msg_out, alt_chat) = taken.unwrap();
4152        let alt_chat = alt_chat.expect("should be alternate hit");
4153        assert!(
4154            alt_chat.is_pn(),
4155            "alternate chat should be PN (the original input)"
4156        );
4157        assert_eq!(alt_chat.user, pn_jid.user);
4158        assert_eq!(msg_out.conversation.as_deref(), Some("pn input alternate"));
4159    }
4160
4161    /// When no PN/LID mapping exists, no alternate is tried and take returns None.
4162    #[tokio::test]
4163    async fn no_alternate_without_mapping() {
4164        let _ = env_logger::builder().is_test(true).try_init();
4165
4166        let backend = crate::test_utils::create_test_backend().await;
4167        let pm = Arc::new(
4168            PersistenceManager::new(backend)
4169                .await
4170                .expect("persistence manager should initialize"),
4171        );
4172        let mut config = crate::cache_config::CacheConfig::default();
4173        config.recent_messages.capacity = 1_000;
4174        let (client, _sync_rx) = Client::new_with_cache_config(
4175            Arc::new(crate::runtime_impl::TokioRuntime),
4176            pm.clone(),
4177            Arc::new(crate::transport::mock::MockTransportFactory::new()),
4178            Arc::new(MockHttpClient),
4179            None,
4180            config,
4181        )
4182        .await;
4183
4184        let lid_jid: Jid = "236395184570386@lid".parse().unwrap();
4185        let msg_id = "RETRY_NO_ALT";
4186        let msg = wa::Message {
4187            conversation: Some("no alternate".into()),
4188            ..Default::default()
4189        };
4190
4191        // Store under LID, no PN mapping exists
4192        client
4193            .add_recent_message(&lid_jid, msg_id, &msg, None)
4194            .await;
4195
4196        // Lookup via LID: primary hits directly (same namespace)
4197        let taken = client.take_recent_message(&lid_jid, msg_id).await;
4198        assert!(taken.is_some());
4199        let (_, alt_chat) = taken.unwrap();
4200        assert!(alt_chat.is_none(), "primary hit should have no alt_chat");
4201
4202        // Now try looking up a message that doesn't exist at all
4203        let missing = client.take_recent_message(&lid_jid, "NONEXISTENT").await;
4204        assert!(missing.is_none(), "non-existent message should return None");
4205    }
4206
4207    /// When both primary and alternate miss, take returns None.
4208    #[tokio::test]
4209    async fn alternate_key_both_miss() {
4210        let _ = env_logger::builder().is_test(true).try_init();
4211
4212        let backend = crate::test_utils::create_test_backend().await;
4213        let pm = Arc::new(
4214            PersistenceManager::new(backend)
4215                .await
4216                .expect("persistence manager should initialize"),
4217        );
4218        let mut config = crate::cache_config::CacheConfig::default();
4219        config.recent_messages.capacity = 1_000;
4220        let (client, _sync_rx) = Client::new_with_cache_config(
4221            Arc::new(crate::runtime_impl::TokioRuntime),
4222            pm.clone(),
4223            Arc::new(crate::transport::mock::MockTransportFactory::new()),
4224            Arc::new(MockHttpClient),
4225            None,
4226            config,
4227        )
4228        .await;
4229
4230        let pn_jid: Jid = "5511999999999@s.whatsapp.net".parse().unwrap();
4231        let lid_jid: Jid = "236395184570386@lid".parse().unwrap();
4232
4233        // Add mapping but don't store any message
4234        client
4235            .lid_pn_cache
4236            .add(&wacore::types::lid_pn::LidPnEntry {
4237                lid: lid_jid.user.as_str().into(),
4238                phone_number: pn_jid.user.as_str().into(),
4239                created_at: 0,
4240                learning_source: wacore::types::lid_pn::LearningSource::Usync,
4241            })
4242            .await;
4243
4244        // Lookup via PN: primary (LID) misses, alternate (PN) also misses
4245        let taken = client.take_recent_message(&pn_jid, "MISSING").await;
4246        assert!(taken.is_none(), "both primary and alternate miss → None");
4247    }
4248
4249    // --- Peer device / bot / original_from tests ---
4250
4251    #[test]
4252    fn resolve_retry_chat_info_peer_device_with_recipient() {
4253        use wacore_binary::builder::NodeBuilder;
4254
4255        // Peer retry: from=our own JID, recipient=the actual chat partner
4256        let our_pn: Jid = "5511999999999@s.whatsapp.net".parse().unwrap();
4257        let recipient: Jid = "5522888888888@s.whatsapp.net".parse().unwrap();
4258
4259        let node = NodeBuilder::new("receipt")
4260            .attr("recipient", "5522888888888@s.whatsapp.net")
4261            .build();
4262        let receipt = make_test_receipt("5511999999999:2@s.whatsapp.net");
4263
4264        let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), Some(&our_pn), None);
4265
4266        // Chat should be the recipient (the actual conversation partner)
4267        assert_eq!(info.chat.user, recipient.user);
4268        assert_eq!(info.chat.device(), 0, "chat should be bare");
4269        // Requester is still our device
4270        assert_eq!(info.requester.user, our_pn.user);
4271        assert_eq!(info.requester.device(), 2);
4272    }
4273
4274    #[test]
4275    fn resolve_retry_chat_info_peer_device_without_recipient() {
4276        use wacore_binary::builder::NodeBuilder;
4277
4278        // Peer retry without recipient attr has no target chat in WA Web.
4279        let our_pn: Jid = "5511999999999@s.whatsapp.net".parse().unwrap();
4280        let node = NodeBuilder::new("receipt").build();
4281        let receipt = make_test_receipt("5511999999999:2@s.whatsapp.net");
4282
4283        let info =
4284            maybe_resolve_retry_chat_info(&receipt, &node.as_node_ref(), Some(&our_pn), None);
4285
4286        assert!(info.is_none());
4287    }
4288
4289    #[test]
4290    fn resolve_retry_chat_info_bot_with_recipient() {
4291        use wacore_binary::builder::NodeBuilder;
4292
4293        // Bot retry: from=bot JID, recipient=actual chat
4294        let node = NodeBuilder::new("receipt")
4295            .attr("recipient", "5522888888888@s.whatsapp.net")
4296            .build();
4297        let receipt = make_test_receipt("131355500001@s.whatsapp.net");
4298
4299        let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
4300
4301        assert!(info.is_bot, "bot JID should be detected");
4302        assert!(
4303            !info.is_fbid_bot_retry,
4304            "legacy PN bots use the regular retry parser"
4305        );
4306        // Chat should be the recipient
4307        assert_eq!(info.chat.user, "5522888888888");
4308        assert_eq!(info.chat.device(), 0);
4309    }
4310
4311    #[test]
4312    fn resolve_retry_chat_info_bot_without_recipient() {
4313        use wacore_binary::builder::NodeBuilder;
4314
4315        // Bot retry without recipient — falls through to normal DM path
4316        let node = NodeBuilder::new("receipt").build();
4317        let receipt = make_test_receipt("131355500001@s.whatsapp.net");
4318
4319        let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
4320
4321        assert!(info.is_bot);
4322        assert!(!info.is_fbid_bot_retry);
4323        // Without recipient, falls to from.to_non_ad()
4324        assert_eq!(info.chat.user, "131355500001");
4325    }
4326
4327    #[test]
4328    fn resolve_retry_chat_info_fbid_bot_dm_marks_bot_retry() {
4329        use wacore_binary::builder::NodeBuilder;
4330
4331        let node = NodeBuilder::new("receipt")
4332            .attr("recipient", "5522888888888@s.whatsapp.net")
4333            .build();
4334        let receipt = make_test_receipt("200000000000002@bot");
4335
4336        let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
4337
4338        assert!(info.is_bot);
4339        assert!(info.is_fbid_bot_retry);
4340        assert_eq!(info.chat.user, "5522888888888");
4341    }
4342
4343    #[test]
4344    fn resolve_retry_chat_info_preserves_original_from() {
4345        use wacore_binary::builder::NodeBuilder;
4346
4347        // DM with device suffix — original_from preserves the raw receipt from
4348        // (WA Web: variable m = e.from, used as-is for stanza to)
4349        let node = NodeBuilder::new("receipt").build();
4350        let receipt = make_test_receipt("5511999999999:33@s.whatsapp.net");
4351
4352        let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
4353
4354        // original_from keeps the full JID including device
4355        assert_eq!(info.original_from.device(), 33);
4356        assert_eq!(info.original_from.user, "5511999999999");
4357
4358        // chat is bare
4359        assert_eq!(info.chat.device(), 0);
4360        assert_eq!(info.chat.user, "5511999999999");
4361    }
4362
4363    #[test]
4364    fn resolve_retry_chat_info_peer_via_lid() {
4365        use wacore_binary::builder::NodeBuilder;
4366
4367        // Peer retry detected via LID (not PN)
4368        let our_lid: Jid = "236395184570386@lid".parse().unwrap();
4369        let recipient: Jid = "5522888888888@s.whatsapp.net".parse().unwrap();
4370
4371        let node = NodeBuilder::new("receipt")
4372            .attr("recipient", "5522888888888@s.whatsapp.net")
4373            .build();
4374        let receipt = make_test_receipt("236395184570386:5@lid");
4375
4376        let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, Some(&our_lid));
4377
4378        assert_eq!(info.chat.user, recipient.user);
4379        assert_eq!(info.chat.device(), 0);
4380        assert_eq!(info.requester.device(), 5);
4381    }
4382}