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#[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
33fn 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
41const 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 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
184struct RetryChatInfo {
187 chat: Jid,
189 requester: Jid,
191 original_from: Jid,
194 recipient: Option<Jid>,
198 is_bot: bool,
200 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
208fn 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 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 let recipient = node.attrs().optional_jid("recipient");
238 let is_bot = from.is_bot();
239
240 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
288fn 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 None => Ok(requester.clone()),
318 };
319 }
320 Ok(self.resolve_encryption_jid(requester).await)
321 }
322
323 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 #[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 #[cfg(feature = "tracing")]
455 tracing::Span::current().record("count", retry_count);
456
457 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 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 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 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 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 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 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 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 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 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 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 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 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 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 if nr.get_optional_child("keys").is_none() {
725 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 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 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 drop(session_guard);
890 self.send_retry_stanza(stanza).await
891 }
892
893 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 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 #[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 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!(
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 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 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 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 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 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 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 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 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 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 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 #[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 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 let signal_address = requester_jid.to_protocol_address();
1313
1314 {
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 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 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 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 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 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 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 #[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 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 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 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 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 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 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 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 let recipient = info.source.recipient.as_ref().unwrap_or(&info.source.chat);
1577 builder = builder.attr("recipient", recipient);
1578 }
1579 }
1580 }
1581
1582 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 #[allow(dead_code)] #[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 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 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 client.add_recent_message(&chat, &msg_id, &msg, None).await;
1830
1831 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 let taken_again = client.take_recent_message(&chat, &msg_id).await;
1840 assert!(taken_again.is_none());
1841 }
1842
1843 #[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 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 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 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 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 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 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 #[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 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 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 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 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 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 #[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 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 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 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 let registration = receipt_node
2213 .get_optional_child("registration")
2214 .expect("<registration> child must exist");
2215 let reg_bytes = match ®istration.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 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 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 let result = backend.has_same_base_key(address, msg_id, &base_key).await;
2251 assert!(result.is_ok());
2252 assert!(!result.unwrap());
2253
2254 let save_result = backend.save_base_key(address, msg_id, &base_key).await;
2256 assert!(save_result.is_ok());
2257
2258 let result = backend.has_same_base_key(address, msg_id, &base_key).await;
2260 assert!(result.is_ok());
2261 assert!(result.unwrap());
2262
2263 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 let delete_result = backend.delete_base_key(address, msg_id).await;
2273 assert!(delete_result.is_ok());
2274
2275 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 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 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 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 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 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 fn build_retry_receipt_without_keys() -> Node {
2397 use wacore_binary::builder::NodeBuilder;
2398 NodeBuilder::new("receipt").build()
2399 }
2400
2401 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 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 #[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 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 #[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 #[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 #[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 #[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 #[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 let jid_with = Jid::lid_device("999999999999991".to_string(), 3);
2778 let jid_without = Jid::lid_device("999999999999992".to_string(), 3);
2779
2780 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 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 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 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 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 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 #[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 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 #[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 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 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 #[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 client.set_resend_rate_limit(1, 0);
2950 assert!(client.resend_rate_limiter.try_acquire(&group).await);
2951
2952 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 #[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 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 #[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 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 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(), ®ular_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 use wacore_binary::JidExt as _;
3275
3276 let regular_user: Jid = "1234567890@s.whatsapp.net".parse().unwrap();
3278 assert!(!regular_user.is_bot());
3279
3280 let bot_server: Jid = "somebot@bot".parse().unwrap();
3282 assert!(bot_server.is_bot());
3283
3284 let legacy_bot: Jid = "1313555123456@s.whatsapp.net".parse().unwrap();
3286 assert!(legacy_bot.is_bot());
3287
3288 let legacy_bot2: Jid = "131655500123456@s.whatsapp.net".parse().unwrap();
3290 assert!(legacy_bot2.is_bot());
3291
3292 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 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 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 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 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 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 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 assert!(group.is_group() || group.is_status_broadcast());
3373 assert!(status.is_group() || status.is_status_broadcast());
3374
3375 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 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 #[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 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 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 assert_eq!(info.chat.device(), 0);
3621 assert_eq!(info.chat.user, "5511999999999");
3622 assert!(info.chat.is_pn());
3623
3624 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 assert_eq!(info.chat.device(), 0);
3639 assert_eq!(info.chat.user, "236395184570386");
3640 assert!(info.chat.is_lid());
3641
3642 assert_eq!(info.requester.device(), 5);
3644 assert_eq!(info.requester.user, "236395184570386");
3645 assert!(info.requester.is_lid());
3646 }
3647
3648 #[test]
3658 fn resolve_retry_chat_info_forwards_recipient_attribute_verbatim() {
3659 use wacore_binary::builder::NodeBuilder;
3660
3661 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 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 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 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 #[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 #[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 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 #[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 client
3958 .add_recent_message(&bare_jid, msg_id, &msg, None)
3959 .await;
3960
3961 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 client
3969 .add_recent_message(&bare_jid, msg_id, &msg_out, None)
3970 .await;
3971
3972 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 #[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 client.add_recent_message(&pn_jid, msg_id, &msg, None).await;
4015
4016 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 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 #[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 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 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 let group: Jid = "120363021033254949@g.us".parse().unwrap();
4094 assert!(client.swap_pn_lid_namespace(&group).await.is_none());
4095 }
4096
4097 #[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 client.add_recent_message(&pn_jid, msg_id, &msg, None).await;
4132
4133 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 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 #[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 client
4193 .add_recent_message(&lid_jid, msg_id, &msg, None)
4194 .await;
4195
4196 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 let missing = client.take_recent_message(&lid_jid, "NONEXISTENT").await;
4204 assert!(missing.is_none(), "non-existent message should return None");
4205 }
4206
4207 #[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 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 let taken = client.take_recent_message(&pn_jid, "MISSING").await;
4246 assert!(taken.is_none(), "both primary and alternate miss → None");
4247 }
4248
4249 #[test]
4252 fn resolve_retry_chat_info_peer_device_with_recipient() {
4253 use wacore_binary::builder::NodeBuilder;
4254
4255 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 assert_eq!(info.chat.user, recipient.user);
4268 assert_eq!(info.chat.device(), 0, "chat should be bare");
4269 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 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 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 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 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 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 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 assert_eq!(info.original_from.device(), 33);
4356 assert_eq!(info.original_from.user, "5511999999999");
4357
4358 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 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}