use crate::client::Client;
use crate::features::{MessageRetransmission, RetryRequestError};
use crate::message::RetryReason;
use crate::send::SendError;
use crate::types::events::Receipt;
use log::{debug, info, warn};
use wacore::types::message::MessageCategory;
use scopeguard;
use std::sync::Arc;
use wacore::iq::prekeys::{OneTimePreKeyNode, SignedPreKeyNode};
use wacore::libsignal::protocol::{PreKeyBundle, PublicKey};
use wacore::protocol::ProtocolNode;
use wacore::protocol::retry::{MAX_RETRY_COUNT, MIN_RETRY_FOR_BASE_KEY_CHECK};
use wacore::types::jid::JidExt;
use wacore_binary::JidExt as _;
#[cfg(test)]
use wacore_binary::NodeContent;
use wacore_binary::builder::NodeBuilder;
use wacore_binary::{Jid, Node, OwnedNodeRef};
use wacore_binary::{NodeContentRef, NodeRef};
use waproto::whatsapp as wa;
#[cfg(test)]
fn get_bytes_content(node: &Node) -> Option<&[u8]> {
match &node.content {
Some(NodeContent::Bytes(b)) => Some(b.as_slice()),
_ => None,
}
}
fn get_bytes_content_ref<'a>(node: &'a NodeRef<'_>) -> Option<&'a [u8]> {
match node.content.as_ref() {
Some(NodeContentRef::Bytes(b)) => Some(b.as_ref()),
_ => None,
}
}
const RECREATE_SESSION_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(3600);
#[derive(Clone, Copy)]
enum RetransmissionRoute {
Direct,
Group,
Status,
BroadcastList,
}
impl RetransmissionRoute {
const fn uses_sender_key(self) -> bool {
matches!(self, Self::Group | Self::Status)
}
}
#[inline]
fn is_own_account_jid(jid: &Jid, own_pn: Option<&Jid>, own_lid: Option<&Jid>) -> bool {
own_pn.is_some_and(|pn| jid.is_same_user_as(pn))
|| own_lid.is_some_and(|lid| jid.is_same_user_as(lid))
}
struct PreparedRetransmission {
route: RetransmissionRoute,
chat: Jid,
wire_requester: Jid,
encryption_jid: Jid,
message: wa::Message,
message_id: String,
retry_count: u8,
recipient: Option<Jid>,
group_info: Option<Arc<wacore::client::context::GroupInfo>>,
pre_encoded: Option<Arc<Vec<u8>>>,
}
fn validate_retransmission(
chat: &Jid,
requester: &Jid,
message_id: &str,
retry_count: u8,
recipient: Option<&Jid>,
) -> Result<RetransmissionRoute, SendError> {
if chat.is_empty() || requester.is_empty() {
return Err(SendError::InvalidRequest(
"retransmission JIDs must not be empty".into(),
));
}
if message_id.is_empty() {
return Err(SendError::InvalidRequest(
"retransmission message ID must not be empty".into(),
));
}
if !(1..MAX_RETRY_COUNT).contains(&retry_count) {
return Err(SendError::InvalidRequest(format!(
"retry count must be in 1..{MAX_RETRY_COUNT}"
)));
}
let requester_is_user = matches!(
requester.server,
wacore_binary::Server::Pn
| wacore_binary::Server::Lid
| wacore_binary::Server::Hosted
| wacore_binary::Server::HostedLid
| wacore_binary::Server::Bot
);
if !requester_is_user {
return Err(SendError::InvalidRequest(
"retransmission requester must be a user device JID".into(),
));
}
let route = if chat.is_group() {
RetransmissionRoute::Group
} else if chat.is_status_broadcast() {
if !matches!(
requester.server,
wacore_binary::Server::Pn | wacore_binary::Server::Lid
) {
return Err(SendError::InvalidRequest(
"status retransmission requester must be a PN or LID device".into(),
));
}
RetransmissionRoute::Status
} else if chat.is_broadcast_list() {
if !matches!(
requester.server,
wacore_binary::Server::Pn | wacore_binary::Server::Lid
) {
return Err(SendError::InvalidRequest(
"broadcast retransmission requester must be a PN or LID device".into(),
));
}
RetransmissionRoute::BroadcastList
} else if matches!(
chat.server,
wacore_binary::Server::Pn
| wacore_binary::Server::Lid
| wacore_binary::Server::Hosted
| wacore_binary::Server::HostedLid
| wacore_binary::Server::Bot
) {
RetransmissionRoute::Direct
} else {
return Err(SendError::InvalidRequest(
"unsupported retransmission chat class".into(),
));
};
if recipient.is_some() && !matches!(route, RetransmissionRoute::Direct) {
return Err(SendError::InvalidRequest(
"recipient is only valid for direct retransmissions".into(),
));
}
if recipient.is_some_and(|recipient| {
recipient.is_empty()
|| !matches!(
recipient.server,
wacore_binary::Server::Pn
| wacore_binary::Server::Lid
| wacore_binary::Server::Hosted
| wacore_binary::Server::HostedLid
| wacore_binary::Server::Bot
)
}) {
return Err(SendError::InvalidRequest(
"retransmission recipient must be a user JID".into(),
));
}
Ok(route)
}
pub(crate) enum RetryReceiptSendOutcome {
Sent { included_keys: bool },
Suppressed,
}
struct RetryChatInfo {
chat: Jid,
requester: Jid,
original_from: Jid,
recipient: Option<Jid>,
is_bot: bool,
is_fbid_bot_retry: bool,
}
fn is_fbid_bot_retry_jid(jid: &Jid) -> bool {
jid.server == wacore_binary::Server::Bot && jid.device() == 0
}
fn resolve_retry_chat_info(
receipt: &Receipt,
node: &NodeRef<'_>,
own_pn: Option<&Jid>,
own_lid: Option<&Jid>,
) -> Option<RetryChatInfo> {
let from = &receipt.source.chat;
if from.is_group() || from.is_status_broadcast() || from.is_broadcast_list() {
let participant = node.attrs().optional_jid("participant");
let is_fbid_bot_retry =
from.is_group() && participant.as_ref().is_some_and(is_fbid_bot_retry_jid);
let requester = participant.unwrap_or_else(|| receipt.source.sender.clone());
let is_bot = requester.is_bot();
Some(RetryChatInfo {
chat: from.clone(),
requester,
original_from: from.clone(),
recipient: node.attrs().optional_jid("recipient"),
is_bot,
is_fbid_bot_retry,
})
} else {
let recipient = node.attrs().optional_jid("recipient");
let is_bot = from.is_bot();
let is_peer = is_own_account_jid(from, own_pn, own_lid);
let chat = if is_bot && let Some(r) = recipient.as_ref() {
r.to_non_ad()
} else if is_peer {
match recipient.as_ref() {
Some(r) => r.to_non_ad(),
None => {
log::warn!("Ignoring peer device retry without recipient attr");
return None;
}
}
} else {
from.to_non_ad()
};
let requester = if from.device() == 0 && from.agent == 0 {
chat.clone()
} else {
from.clone()
};
Some(RetryChatInfo {
chat,
requester,
original_from: from.clone(),
recipient,
is_bot,
is_fbid_bot_retry: is_fbid_bot_retry_jid(from),
})
}
}
fn validate_retry_prekey_presence(
keys_node: &NodeRef<'_>,
is_fbid_bot_retry: bool,
) -> Result<(), anyhow::Error> {
if !is_fbid_bot_retry && keys_node.get_optional_child("key").is_none() {
anyhow::bail!("regular retry key bundle missing one-time prekey");
}
Ok(())
}
fn build_retry_processing_key(chat: &Jid, message_id: &str, participant_jid: &Jid) -> String {
let mut key = String::with_capacity(message_id.len() + 64);
chat.push_to(&mut key);
key.push(':');
key.push_str(message_id);
key.push(':');
participant_jid.push_to(&mut key);
key
}
impl Client {
async fn resolve_retransmission_encryption_jid(
&self,
route: RetransmissionRoute,
requester: &Jid,
) -> Result<Jid, anyhow::Error> {
if matches!(route, RetransmissionRoute::Status) && requester.is_pn() {
return match self.get_lid_pn_entry(requester).await? {
Some(mapping) => Ok(Jid {
user: wacore_binary::CompactString::new(&mapping.lid),
server: wacore_binary::Server::Lid,
device: requester.device,
agent: requester.agent,
integrator: requester.integrator,
}),
None => Ok(requester.clone()),
};
}
Ok(self.resolve_encryption_jid(requester).await)
}
pub async fn retransmit_message(
&self,
request: MessageRetransmission,
) -> Result<(), SendError> {
let route = validate_retransmission(
&request.chat,
&request.requester,
&request.message_id,
request.retry_count,
request.recipient.as_ref(),
)?;
if matches!(route, RetransmissionRoute::Direct) {
let snapshot = self.persistence_manager.get_device_snapshot();
let requester_is_local = is_own_account_jid(
&request.requester,
snapshot.pn.as_ref(),
snapshot.lid.as_ref(),
);
if request.recipient.is_some() {
if !requester_is_local && !request.requester.is_bot() {
return Err(SendError::InvalidRequest(
"a direct retransmission recipient is only valid for a local device or bot"
.into(),
));
}
} else if requester_is_local {
return Err(SendError::InvalidRequest(
"a direct retransmission to another local device requires a recipient".into(),
));
}
let routing_chat = request.recipient.as_ref().unwrap_or(&request.requester);
if !self
.jids_share_user_identity(&request.chat, routing_chat)
.await
.map_err(SendError::from_anyhow)?
{
return Err(SendError::InvalidRequest(
"direct retransmission chat does not match its routing identity".into(),
));
}
}
let group_info = if matches!(route, RetransmissionRoute::Group) {
Some(
self.groups()
.query_info_with_freshness(&request.chat, request.group_metadata_freshness)
.await?,
)
} else {
None
};
let encryption_jid = self
.resolve_retransmission_encryption_jid(route, &request.requester)
.await
.map_err(SendError::from_anyhow)?;
if route.uses_sender_key() {
let chat_key = request.chat.to_string();
self.mark_forget_sender_key(&chat_key, std::slice::from_ref(&encryption_jid))
.await
.map_err(SendError::from_anyhow)?;
}
let MessageRetransmission {
chat,
requester: wire_requester,
message,
message_id,
retry_count,
recipient,
group_metadata_freshness: _,
} = request;
let pre_encoded = Arc::new(waproto::codec::message_to_vec(&message));
self.add_recent_message(&chat, &message_id, &message, Some(Arc::clone(&pre_encoded)))
.await;
self.retransmit_message_prepared(PreparedRetransmission {
route,
wire_requester,
encryption_jid,
chat,
message,
message_id,
retry_count,
recipient,
group_info,
pre_encoded: Some(pre_encoded),
})
.await
.map_err(SendError::from_anyhow)
}
#[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)))]
pub(crate) async fn handle_retry_receipt(
self: &Arc<Self>,
receipt: &Receipt,
node: &Arc<OwnedNodeRef>,
) -> Result<(), anyhow::Error> {
let nr = node.get();
let retry_child = nr
.get_optional_child("retry")
.ok_or_else(|| anyhow::anyhow!("<retry> child missing from receipt"))?;
let message_id = retry_child
.get_attr("id")
.map(|v| v.as_str())
.ok_or_else(|| anyhow::anyhow!("<retry> missing 'id' attribute"))?
.into_owned();
let retry_count: u8 = retry_child
.get_attr("count")
.map(|v| v.as_str())
.and_then(|s| s.parse().ok())
.unwrap_or(1);
#[cfg(feature = "tracing")]
tracing::Span::current().record("count", retry_count);
if retry_count >= MAX_RETRY_COUNT {
debug!(
"Refusing retry #{} for message {} from {}: exceeds max attempts ({})",
retry_count,
message_id,
receipt.source.sender.observe(),
MAX_RETRY_COUNT
);
wacore::telemetry::retry_refused();
return Ok(());
}
let device_snapshot = self.persistence_manager.get_device_snapshot();
let Some(mut info) = resolve_retry_chat_info(
receipt,
nr,
device_snapshot.pn.as_ref(),
device_snapshot.lid.as_ref(),
) else {
return Ok(());
};
let route = match validate_retransmission(
&info.chat,
&info.requester,
&message_id,
retry_count,
info.recipient.as_ref(),
) {
Ok(route) => route,
Err(error) => {
debug!("Ignoring malformed retry request: {error}");
return Ok(());
}
};
let uses_sender_key = route.uses_sender_key();
let processing_key = build_retry_processing_key(&info.chat, &message_id, &info.requester);
if !self
.pending_retries
.lock()
.unwrap_or_else(|p| p.into_inner())
.insert(processing_key.clone())
{
log::debug!("Ignoring retry for {processing_key}: a retry is already in progress.");
return Ok(());
}
let pending = Arc::clone(&self.pending_retries);
let _guard = scopeguard::guard((), move |()| {
pending
.lock()
.unwrap_or_else(|p| p.into_inner())
.remove(&processing_key);
});
let sender_device_id = info.requester.device() as u32;
let device_known = self
.has_device(&info.requester.user, sender_device_id)
.await;
if !device_known {
wacore::telemetry::retry_unknown_device(if sender_device_id == 0 {
"primary"
} else {
"companion"
});
self.schedule_unknown_device_sync(info.requester.to_non_ad(), receipt.offline)
.await;
}
let (original_msg, alt_chat) = match self.peek_recent_message(&info.chat, &message_id).await
{
Some(result) => result,
None => match self.take_recent_message(&info.chat, &message_id).await {
Some(result) => {
self.add_recent_message(&info.chat, &message_id, &result.0, None)
.await;
result
}
None => {
log::debug!(
"Ignoring retry for message {message_id}: already handled or not found in cache."
);
return Ok(());
}
},
};
let resolved_jid = if let Some(alt_chat) = alt_chat
&& !uses_sender_key
&& !info.is_bot
{
let requester = &info.requester;
info.requester = Jid {
user: alt_chat.user,
server: alt_chat.server,
device: requester.device,
agent: requester.agent,
integrator: requester.integrator,
};
info.requester.clone()
} else {
self.resolve_retransmission_encryption_jid(route, &info.requester)
.await?
};
let keys_node_present = nr.get_optional_child("keys").is_some();
if wacore::protocol::retry::should_drop_unknown_device_retry(
keys_node_present,
device_known,
) {
warn!(
"handle_retry_receipt: device not found for device={}, user={}",
sender_device_id, info.requester.user
);
return Ok(());
}
let is_peer = is_own_account_jid(
&info.requester,
device_snapshot.pn.as_ref(),
device_snapshot.lid.as_ref(),
);
if uses_sender_key
&& !is_peer
&& let Some(policy) = self.retry_admission.get()
&& !policy.admit(&info.chat, &info.requester, retry_count)
{
debug!(
"Retry receipt from {} in {} dropped by RetryAdmission policy",
info.requester.observe(),
info.chat.observe()
);
return Ok(());
}
let cached_group_info = if info.chat.is_group() {
match self.groups().query_info(&info.chat).await {
Ok(gi) => Some(gi),
Err(e) => {
log::warn!(
"Failed to fetch group info for retry of msg {} in {}: {e}",
message_id,
info.chat.observe()
);
None
}
}
} else {
None
};
let mut rotated_sender_key = false;
if matches!(route, RetransmissionRoute::Group) && !info.requester.is_lid() {
let group_jid = info.chat.to_string();
let is_known_participant = cached_group_info
.as_ref()
.is_some_and(|g| g.participants.iter().any(|p| p.user == info.requester.user));
if !is_known_participant {
log::warn!(
"Unknown device {} in group {} — forcing full sender key rotation \
(matches WA Web's rotateKey behavior)",
info.requester.observe(),
group_jid
);
let _distribution_guard = self.group_distribution_lock(&info.chat).await;
let addressing_mode = cached_group_info.as_ref().map(|g| g.addressing_mode);
let jids_to_delete: Vec<_> = match addressing_mode {
Some(wacore::types::message::AddressingMode::Lid) => {
device_snapshot.lid.as_ref().into_iter().collect()
}
Some(wacore::types::message::AddressingMode::Pn) => {
device_snapshot.pn.as_ref().into_iter().collect()
}
None => device_snapshot
.lid
.as_ref()
.into_iter()
.chain(device_snapshot.pn.as_ref())
.collect(),
};
for own_jid in jids_to_delete {
use wacore::libsignal::store::sender_key_name::SenderKeyName;
let sk_name = SenderKeyName::from_parts(
&group_jid,
own_jid.to_protocol_address().as_str(),
);
self.signal_cache
.delete_sender_key(sk_name.cache_key())
.await;
}
if let Err(e) = self.reset_sender_key_device_tracking(&group_jid).await {
log::warn!("Failed to clear sender key devices for rotation: {}", e);
}
rotated_sender_key = true;
}
}
if rotated_sender_key {
self.flush_signal_cache_batch_safe_logged("unknown-participant rotation", None)
.await;
}
if !self
.update_local_signal_session(
&info,
&resolved_jid,
&message_id,
retry_count,
nr,
is_peer,
)
.await
{
return Ok(());
}
if nr.get_optional_child("keys").is_none() {
let signal_address = resolved_jid.to_protocol_address();
let lock = self.session_lock_for(signal_address.as_str()).await;
let guard = lock.lock().await;
if let Some(reason) = self
.should_recreate_session(retry_count, &resolved_jid)
.await
{
info!(
"Recreating session with {} for retry of {message_id}: {reason}",
resolved_jid.observe()
);
self.signal_cache.delete_session(&signal_address).await;
drop(guard);
self.flush_signal_cache_batch_safe_logged(
"should_recreate_session",
Some(&message_id),
)
.await;
}
}
if info.chat.is_group() && !self.resend_rate_limiter.try_acquire(&info.chat).await {
debug!(
"Throttling resend of {} to {}: per-chat resend rate cap reached",
message_id,
info.chat.observe()
);
return Ok(());
}
info!(
"Resending message {} to {} (retry #{})",
message_id,
info.chat.observe(),
retry_count
);
let wire_requester = if matches!(route, RetransmissionRoute::Direct) {
info.original_from
} else {
info.requester
};
self.retransmit_message_prepared(PreparedRetransmission {
route,
chat: info.chat,
wire_requester,
encryption_jid: resolved_jid,
message: original_msg,
message_id,
retry_count,
recipient: info.recipient,
group_info: cached_group_info,
pre_encoded: None,
})
.await?;
Ok(())
}
async fn send_retry_stanza(&self, stanza: Node) -> Result<(), anyhow::Error> {
self.persist_signal_state_pre_wire().await?;
self.send_node(stanza).await?;
Ok(())
}
async fn retransmit_message_prepared(
&self,
request: PreparedRetransmission,
) -> Result<(), anyhow::Error> {
let PreparedRetransmission {
route,
chat,
wire_requester,
encryption_jid,
message,
message_id,
retry_count,
recipient,
group_info,
pre_encoded,
} = request;
if matches!(route, RetransmissionRoute::Status) {
return self
.retransmit_status_message(
chat,
encryption_jid,
message,
message_id,
pre_encoded.as_deref().map(Vec::as_slice),
)
.await;
}
self.ensure_e2e_sessions_resolved(std::slice::from_ref(&encryption_jid))
.await?;
let signal_address = encryption_jid.to_protocol_address();
let session_mutex = self.session_lock_for(signal_address.as_str()).await;
let session_guard = session_mutex.lock().await;
let mut store_adapter = self.signal_adapter().await;
let device_snapshot = self.persistence_manager.get_device_snapshot();
let edit = wacore::types::message::EditAttribute::infer_from_message(&message);
let destination = match route {
RetransmissionRoute::Direct => wacore::send::PairwiseRetryDestination::Direct {
to: wire_requester,
recipient,
},
RetransmissionRoute::Group => {
let addressing_mode = group_info
.as_ref()
.map(|info| info.addressing_mode)
.unwrap_or_default();
wacore::send::PairwiseRetryDestination::Participant {
to: chat,
participant: wire_requester,
addressing_mode: Some(addressing_mode),
}
}
RetransmissionRoute::BroadcastList => {
wacore::send::PairwiseRetryDestination::Participant {
to: chat,
participant: wire_requester,
addressing_mode: None,
}
}
RetransmissionRoute::Status => unreachable!("status handled above"),
};
let stanza = wacore::send::prepare_pairwise_retry_stanza(
&mut store_adapter.session_store,
&mut store_adapter.identity_store,
wacore::send::PairwiseRetryRequest {
destination,
encryption_jid,
message: &message,
message_id,
retry_count,
account: device_snapshot.account.as_deref(),
edit,
pre_encoded: pre_encoded.as_deref().map(Vec::as_slice),
},
)
.await?;
drop(session_guard);
self.send_retry_stanza(stanza).await
}
async fn retransmit_status_message(
&self,
chat: Jid,
requester: Jid,
message: wa::Message,
message_id: String,
pre_encoded: Option<&[u8]>,
) -> Result<(), anyhow::Error> {
let snapshot = self.persistence_manager.get_device_snapshot();
let own_pn = snapshot
.pn
.as_ref()
.ok_or(crate::client::ClientError::NotLoggedIn)?;
let own_lid = snapshot
.lid
.as_ref()
.ok_or_else(|| anyhow::anyhow!("cannot retransmit status without a device LID"))?;
let is_sending_device = (requester.is_same_user_as(own_pn)
&& requester.device == own_pn.device)
|| (requester.is_same_user_as(own_lid) && requester.device == own_lid.device);
if is_sending_device {
anyhow::bail!("cannot retransmit a status to the sending device itself");
}
let chat_key = chat.to_string();
let distribution_guard = self.group_distribution_lock(&chat).await;
let group_info = wacore::client::context::GroupInfo::new(
Vec::new(),
wacore::types::message::AddressingMode::Lid,
);
let can_reuse_encoding = message.message_context_info.is_unset();
let encoded_fallback = (pre_encoded.is_none() && can_reuse_encoding)
.then(|| waproto::codec::message_to_vec(&message));
let encoded = pre_encoded
.filter(|_| can_reuse_encoding)
.or(encoded_fallback.as_deref());
let device_store = self.persistence_manager.get_device_arc().await;
let mut store_adapter = self.signal_adapter_from(device_store);
let mut stores = store_adapter.as_signal_stores();
let edit = wacore::types::message::EditAttribute::infer_from_message(&message);
let prepared = match wacore::send::prepare_group_stanza(
&*self.runtime,
&mut stores,
self,
wacore::send::GroupStanzaRequest {
group: &group_info,
own_jid: own_pn,
own_lid,
account: snapshot.account.as_deref(),
to: &chat,
message: &message,
message_id: &message_id,
force_distribution: false,
distribution_targets: Some(vec![requester]),
distribution_policy: wacore::send::SenderKeyDistributionPolicy::Required,
phash_devices: None,
edit: edit.as_ref(),
extra_nodes: &[],
pre_encoded: encoded,
},
)
.await
{
Ok(prepared) => prepared,
Err(error) => {
drop(distribution_guard);
if let Some(failure) =
error.downcast_ref::<wacore::send::RequiredSenderKeyDistributionError>()
{
for user in failure.stale_device_users() {
self.invalidate_device_cache(user).await;
}
}
return Err(error);
}
};
self.send_retry_stanza(prepared.node).await?;
self.update_sender_key_devices(&chat_key, &prepared.skdm_devices)
.await;
drop(distribution_guard);
for user in &prepared.stale_device_users {
self.invalidate_device_cache(user).await;
}
Ok(())
}
#[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)))]
async fn update_local_signal_session(
&self,
info: &RetryChatInfo,
resolved_jid: &Jid,
message_id: &str,
retry_count: u8,
node: &NodeRef<'_>,
is_peer: bool,
) -> bool {
if info.chat.is_group() || info.chat.is_status_broadcast() {
let group_jid = info.chat.to_string();
match self
.mark_forget_sender_key(&group_jid, std::slice::from_ref(resolved_jid))
.await
{
Ok(()) => {
let chat_type = if info.chat.is_status_broadcast() {
"status broadcast"
} else {
"group"
};
debug!(
"Marked {} for fresh SKDM in {} {} due to retry receipt",
resolved_jid.observe(),
chat_type,
group_jid
);
}
Err(e) => log::warn!(
"Failed to mark sender key forget for {} in {}: {}",
info.requester.observe(),
group_jid,
e
),
}
}
let keys_node_present = node.get_optional_child("keys").is_some();
let key_bundle_result = self
.process_retry_key_bundle(node, resolved_jid, is_peer, info.is_fbid_bot_retry)
.await;
let key_bundle_processed = key_bundle_result.is_ok();
if !key_bundle_processed && keys_node_present {
log::warn!(
"Key bundle present but rejected for {}: {:?} — aborting retry resend",
resolved_jid.observe(),
key_bundle_result.as_ref().err()
);
return false;
}
if !key_bundle_processed && !keys_node_present {
if let Err(ref e) = key_bundle_result {
log::debug!(
"No key bundle in retry receipt for {}: {}. Checking for reg ID mismatch.",
resolved_jid.observe(),
e
);
}
if let Some(received_reg_id) =
wacore::protocol::retry::extract_registration_id_from_node_ref(node)
{
let signal_address = resolved_jid.to_protocol_address();
let device_snapshot = self.persistence_manager.get_device_snapshot();
let session = self
.signal_cache
.peek_session(&signal_address, &*device_snapshot.backend)
.await
.ok()
.flatten();
if let Some(session) = session
&& let Ok(stored_reg_id) = session.remote_registration_id()
&& stored_reg_id != 0
&& stored_reg_id != received_reg_id
{
info!(
"Registration ID mismatch for {} (stored: {}, received: {}). \
Deleting session since no key bundle provided.",
wacore::types::jid::observe_protocol_address(&signal_address),
stored_reg_id,
received_reg_id
);
let lock = self.session_lock_for(signal_address.as_str()).await;
let _guard = lock.lock().await;
self.signal_cache.delete_session(&signal_address).await;
drop(_guard);
self.flush_signal_cache_batch_safe_logged(
"reg ID mismatch session deletion",
None,
)
.await;
}
}
}
let signal_address = resolved_jid.to_protocol_address();
let device_snapshot = self.persistence_manager.get_device_snapshot();
let session = self
.signal_cache
.peek_session(&signal_address, &*device_snapshot.backend)
.await
.ok()
.flatten();
let Some(session) = session else {
return true;
};
let Ok(current_base_key) = session.alice_base_key() else {
return true;
};
let addr_str = signal_address.as_str();
if retry_count == MIN_RETRY_FOR_BASE_KEY_CHECK {
match device_snapshot
.backend
.save_base_key(addr_str, message_id, current_base_key)
.await
{
Ok(()) => info!(
"Saved base key for {} at retry #{} for collision detection",
wacore::types::jid::observe_protocol_address(&signal_address),
retry_count
),
Err(e) => warn!(
"Failed to save base key for {}: {}",
wacore::types::jid::observe_protocol_address(&signal_address),
e
),
}
return true;
}
if retry_count > MIN_RETRY_FOR_BASE_KEY_CHECK {
match device_snapshot
.backend
.has_same_base_key(addr_str, message_id, current_base_key)
.await
{
Ok(true) => {
info!(
"Base key collision detected for {} (msg {}) at retry #{}. \
Session hasn't been regenerated. Forcing fresh session.",
wacore::types::jid::observe_protocol_address(&signal_address),
message_id,
retry_count
);
wacore::telemetry::base_key_collision();
let _ = device_snapshot
.backend
.delete_base_key(addr_str, message_id)
.await;
let lock = self.session_lock_for(signal_address.as_str()).await;
let _guard = lock.lock().await;
self.signal_cache.delete_session(&signal_address).await;
drop(_guard);
self.flush_signal_cache_batch_safe_logged(
"base key collision — forcing fresh session",
None,
)
.await;
}
Ok(false) => {
info!(
"Base key changed for {} (msg {}) at retry #{} - session regenerated",
wacore::types::jid::observe_protocol_address(&signal_address),
message_id,
retry_count
);
let _ = device_snapshot
.backend
.delete_base_key(addr_str, message_id)
.await;
}
Err(e) => {
warn!(
"Failed to check base key for {}: {}",
wacore::types::jid::observe_protocol_address(&signal_address),
e
);
}
}
}
true
}
async fn should_recreate_session(&self, retry_count: u8, jid: &Jid) -> Option<&'static str> {
self.should_recreate_session_at(retry_count, jid, wacore::time::Instant::now())
.await
}
async fn should_recreate_session_at(
&self,
retry_count: u8,
jid: &Jid,
now: wacore::time::Instant,
) -> Option<&'static str> {
let signal_address = jid.to_protocol_address();
let device_snapshot = self.persistence_manager.get_device_snapshot();
let has_session = match self
.signal_cache
.has_session(&signal_address, &*device_snapshot.backend)
.await
{
Ok(present) => present,
Err(e) => {
warn!(
"should_recreate_session: has_session failed for {}: {} — skipping recreate",
signal_address, e
);
return None;
}
};
let history = &self.session_recreate_history;
if !has_session {
history.insert(jid.clone(), now).await;
return Some("we don't have a Signal session with them");
}
if retry_count < MIN_RETRY_FOR_BASE_KEY_CHECK {
return None;
}
if let Some(prev) = history.get(jid).await
&& now.saturating_duration_since(prev) < RECREATE_SESSION_TIMEOUT
{
return None;
}
history.insert(jid.clone(), now).await;
Some("retry count > 1 and over an hour since last recreation")
}
#[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)))]
async fn process_retry_key_bundle(
&self,
node: &NodeRef<'_>,
requester_jid: &Jid,
is_peer: bool,
is_fbid_bot_retry: bool,
) -> Result<(), anyhow::Error> {
let keys_node = node
.get_optional_child("keys")
.ok_or_else(|| anyhow::anyhow!("<keys> child missing from retry receipt"))?;
validate_retry_prekey_presence(keys_node, is_fbid_bot_retry)?;
let registration_id =
wacore::protocol::retry::extract_registration_id_from_node_ref(node).unwrap_or(0);
if registration_id == 0 {
return Err(anyhow::anyhow!("Invalid registration ID in retry receipt"));
}
let signal_address = requester_jid.to_protocol_address();
{
let device_snapshot = self.persistence_manager.get_device_snapshot();
let session = self
.signal_cache
.peek_session(&signal_address, &*device_snapshot.backend)
.await
.ok()
.flatten();
if let Some(session) = session {
let existing_reg_id = session.remote_registration_id()?;
if existing_reg_id != 0 && existing_reg_id != registration_id {
if is_peer {
return Err(anyhow::anyhow!(
"Registration ID changed for peer device {} (was {}, now {}). \
This may indicate the device was reinstalled.",
signal_address,
existing_reg_id,
registration_id
));
}
info!(
"Registration ID changed for {} (was {}, now {}). Session will be replaced.",
signal_address, existing_reg_id, registration_id
);
}
}
}
let identity_bytes = keys_node
.get_optional_child("identity")
.and_then(get_bytes_content_ref)
.ok_or_else(|| anyhow::anyhow!("Missing identity key in retry receipt"))?;
let identity_key = PublicKey::from_djb_public_key_bytes(identity_bytes)?;
if requester_jid.device != 0
&& let Some(device_identity) = keys_node
.get_optional_child("device-identity")
.and_then(get_bytes_content_ref)
{
let fetched_identity: [u8; 32] = identity_bytes
.try_into()
.map_err(|_| anyhow::anyhow!("identity key in retry receipt is not 32 bytes"))?;
let account_identity = self.load_account_identity(requester_jid).await;
match wacore::adv::validate_adv_with_identity_key(
device_identity,
&fetched_identity,
account_identity.as_ref(),
) {
wacore::adv::AdvValidation::Valid => {}
wacore::adv::AdvValidation::Invalid => {
return Err(anyhow::anyhow!(
"device-identity ADV validation failed for companion {requester_jid}"
));
}
wacore::adv::AdvValidation::NoAccountKey => log::debug!(
"retry key bundle for companion {requester_jid} omits account_signature_key and no stored account identity; proceeding without ADV validation"
),
}
} else if requester_jid.device != 0 {
log::warn!(
"retry key bundle for companion {requester_jid} omits <device-identity>; proceeding without ADV validation"
);
}
let prekey_data = if let Some(key_ref) = keys_node.get_optional_child("key") {
let prekey_node = OneTimePreKeyNode::try_from_node_ref(key_ref)?;
let prekey_public = PublicKey::from_djb_public_key_bytes(&prekey_node.public_bytes)?;
Some((prekey_node.id.into(), prekey_public))
} else {
None
};
let skey_ref = keys_node
.get_optional_child("skey")
.ok_or_else(|| anyhow::anyhow!("Missing signed prekey in retry receipt"))?;
let signed_prekey = SignedPreKeyNode::try_from_node_ref(skey_ref)?;
let skey_public = PublicKey::from_djb_public_key_bytes(&signed_prekey.public_bytes)?;
let skey_signature: [u8; 64] = signed_prekey
.signature
.as_slice()
.try_into()
.map_err(|_| anyhow::anyhow!("Invalid signature length"))?;
let bundle = PreKeyBundle::new(
registration_id,
u32::from(requester_jid.device).into(),
prekey_data,
signed_prekey.id.into(),
skey_public,
skey_signature.into(),
identity_key.into(),
)?;
let mut adapter = self.signal_adapter().await;
let mut rng = rand::make_rng::<rand::rngs::StdRng>();
self.install_prekey_bundle_cached(requester_jid, &bundle, &mut adapter, &mut rng)
.await?;
self.flush_signal_cache_batch_safe().await?;
info!(
"Processed key bundle from retry receipt for {}",
signal_address
);
Ok(())
}
#[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)))]
pub(crate) async fn send_retry_receipt(
&self,
info: &crate::types::message::MessageInfo,
retry_count: u8,
reason: RetryReason,
force_include_keys: bool,
) -> Result<RetryReceiptSendOutcome, RetryRequestError> {
let device_snapshot = self.persistence_manager.get_device_snapshot();
if info.source.is_bot_authored_non_bot_chat() {
log::debug!(
"Skipping retry receipt for message {} from bot {} in non-bot chat {}",
info.id,
info.source.sender.observe(),
info.source.chat.observe()
);
return Ok(RetryReceiptSendOutcome::Suppressed);
}
debug!(
"Sending retry receipt #{} for message {} in chat {} from {} (reason: {:?})",
retry_count,
info.id,
info.source.chat.observe(),
info.source.sender.observe(),
reason
);
let mut retry_builder = NodeBuilder::new("retry")
.attr("v", "1")
.attr("id", info.id.clone())
.attr("t", info.timestamp.timestamp())
.attr("count", retry_count);
if reason != RetryReason::UnknownError {
retry_builder = retry_builder.attr("error", reason as u8);
}
let retry_node = retry_builder.build();
let registration_id_bytes = device_snapshot.registration_id.to_be_bytes().to_vec();
let registration_node = NodeBuilder::new("registration")
.bytes(registration_id_bytes)
.build();
let receipt_to = if info.source.is_group {
&info.source.chat
} else {
&info.source.sender
};
let include_keys = wacore::protocol::retry::should_include_keys_with_policy(
retry_count,
force_include_keys,
receipt_to.is_hosted(),
);
let keys_node = if include_keys {
let device_identity_bytes = waproto::codec::adv_signed_device_identity_to_vec(
device_snapshot.account.as_deref().ok_or_else(|| {
anyhow::anyhow!("Missing device account info for retry receipt")
})?,
);
let prekey_guard = self.prekey_upload_lock.lock().await;
let (new_prekey_id, new_prekey_public) = self.get_or_gen_single_pre_key().await?;
self.mark_single_prekey_uploaded(&prekey_guard, new_prekey_id)
.await?;
drop(prekey_guard);
Some(wacore::protocol::retry::build_retry_keys_node(
&device_snapshot.identity_key.public_key,
new_prekey_id,
&new_prekey_public,
device_snapshot.signed_pre_key_id,
&device_snapshot.signed_pre_key.public_key,
device_snapshot.signed_pre_key_signature.to_vec(),
device_identity_bytes,
))
} else {
None
};
let mut builder = NodeBuilder::new("receipt")
.attr("to", receipt_to)
.attr("id", info.id.clone())
.attr("type", "retry");
if info.source.is_group {
builder = builder.attr("participant", &info.source.sender);
}
if !info.source.is_group {
let is_from_own_account = device_snapshot
.pn
.as_ref()
.is_some_and(|pn| info.source.sender.is_same_user_as(pn))
|| device_snapshot
.lid
.as_ref()
.is_some_and(|lid| info.source.sender.is_same_user_as(lid));
if is_from_own_account {
if info.category == MessageCategory::Peer {
builder = builder.attr("category", MessageCategory::Peer.as_str());
} else {
let recipient = info.source.recipient.as_ref().unwrap_or(&info.source.chat);
builder = builder.attr("recipient", recipient);
}
}
}
let receipt_node = if let Some(keys) = keys_node {
builder
.children([retry_node, registration_node, keys])
.build()
} else {
builder.children([retry_node, registration_node]).build()
};
drop(device_snapshot);
self.send_node(receipt_node).await?;
Ok(RetryReceiptSendOutcome::Sent {
included_keys: include_keys,
})
}
#[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)))]
pub(crate) async fn send_enc_rekey_retry_receipt(
&self,
stanza_id: &str,
peer_jid: &Jid,
call_id: &str,
call_creator: &Jid,
retry_count: u8,
) -> Result<(), anyhow::Error> {
let device_snapshot = self.persistence_manager.get_device_snapshot();
let registration_id_bytes = device_snapshot.registration_id.to_be_bytes().to_vec();
let enc_rekey_node = NodeBuilder::new("enc_rekey")
.attr("call-creator", call_creator)
.attr("call-id", call_id)
.attr("count", retry_count)
.build();
let registration_node = NodeBuilder::new("registration")
.bytes(registration_id_bytes)
.build();
let receipt_node = NodeBuilder::new("receipt")
.attr("to", peer_jid)
.attr("id", stanza_id)
.attr("type", "enc_rekey_retry")
.children([enc_rekey_node, registration_node])
.build();
info!(
"Sending enc_rekey_retry receipt for call-id={} to {} (count={})",
call_id,
peer_jid.observe(),
retry_count
);
self.send_node(receipt_node).await?;
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::store::persistence_manager::PersistenceManager;
use crate::test_utils::MockHttpClient;
use std::borrow::Cow;
use std::sync::Arc;
use wacore::libsignal::protocol::{IdentityKeyPair, KeyPair};
use wacore::types::jid::JidExt as _;
use wacore_binary::{Jid, JidExt};
use waproto::whatsapp as wa;
fn resolve_retry_chat_info(
receipt: &Receipt,
node: &NodeRef<'_>,
own_pn: Option<&Jid>,
own_lid: Option<&Jid>,
) -> RetryChatInfo {
super::resolve_retry_chat_info(receipt, node, own_pn, own_lid)
.expect("retry should resolve a target chat")
}
fn maybe_resolve_retry_chat_info(
receipt: &Receipt,
node: &NodeRef<'_>,
own_pn: Option<&Jid>,
own_lid: Option<&Jid>,
) -> Option<RetryChatInfo> {
super::resolve_retry_chat_info(receipt, node, own_pn, own_lid)
}
async fn attach_mock_noise_socket(client: &Client) {
use crate::socket::NoiseSocket;
use crate::transport::mock::MockTransport;
use wacore::handshake::NoiseCipher;
let key = [0u8; 32];
let socket = NoiseSocket::new(
Arc::new(crate::runtime_impl::TokioRuntime),
Arc::new(MockTransport),
NoiseCipher::new(&key).expect("write cipher"),
NoiseCipher::new(&key).expect("read cipher"),
);
*client.noise_socket.lock().await = Some(Arc::new(socket));
}
async fn seed_retry_lease(
client: &Client,
address: &wacore::libsignal::protocol::ProtocolAddress,
durable: bool,
) {
use wacore::libsignal::protocol::SessionRecord;
let mut record = SessionRecord::new_fresh();
record.reserve_sender_chain_counters(0);
client.signal_cache.put_session(address, record).await;
if !durable {
return;
}
client.flush_signal_cache().await.expect("durable lease");
let snapshot = client.persistence_manager.get_device_snapshot();
let record = client
.signal_cache
.get_session(address, &*snapshot.backend)
.await
.expect("session read")
.expect("leased session");
assert!(record.reserved_sender_chain_index() > 0);
client.signal_cache.put_session(address, record).await;
}
#[tokio::test]
async fn retry_pre_wire_flush_failure_never_reaches_send_node() {
use std::sync::atomic::Ordering;
let client =
crate::test_utils::create_test_client_with_name("retry_pre_wire_failure").await;
attach_mock_noise_socket(&client).await;
let address = Jid::lid_device("100000000001035".to_string(), 7).to_protocol_address();
seed_retry_lease(&client, &address, false).await;
assert!(client.signal_cache.needs_pre_wire_flush().await);
client.inbound_commit_batch.reset();
client
.inbound_commit_batch
.fail_flushes
.store(true, Ordering::Release);
let id = "RETRY_PRE_WIRE_FAILURE";
let mut waiter =
client.wait_for_sent_node(crate::client::NodeFilter::tag("message").attr("id", id));
let result = client
.send_retry_stanza(NodeBuilder::new("message").attr("id", id).build())
.await;
client
.inbound_commit_batch
.fail_flushes
.store(false, Ordering::Release);
assert!(result.is_err(), "the failed durability gate must abort");
assert!(
waiter.try_recv().expect("waiter stays live").is_none(),
"send_node must not observe a stanza before durability"
);
assert!(
client.signal_cache.needs_pre_wire_flush().await,
"the failed reservation must remain gated"
);
}
#[tokio::test]
async fn retry_inside_durable_lease_skips_synchronous_full_flush() {
use std::sync::atomic::Ordering;
let client =
crate::test_utils::create_test_client_with_name("retry_covered_by_lease").await;
attach_mock_noise_socket(&client).await;
let address = Jid::lid_device("100000000001036".to_string(), 8).to_protocol_address();
seed_retry_lease(&client, &address, true).await;
assert!(!client.signal_cache.needs_pre_wire_flush().await);
client.inbound_commit_batch.reset();
client
.inbound_commit_batch
.fail_flushes
.store(true, Ordering::Release);
let id = "RETRY_COVERED_BY_LEASE";
let waiter =
client.wait_for_sent_node(crate::client::NodeFilter::tag("message").attr("id", id));
let result = client
.send_retry_stanza(NodeBuilder::new("message").attr("id", id).build())
.await;
client
.inbound_commit_batch
.fail_flushes
.store(false, Ordering::Release);
result.expect("an existing durable lease must not synchronously flush");
let sent = waiter.await.expect("retry stanza reached send_node");
assert_eq!(sent.attrs().required_string("id").unwrap(), id);
}
#[tokio::test]
async fn recent_message_cache_insert_and_take() {
let _ = env_logger::builder().is_test(true).try_init();
let backend = crate::test_utils::create_test_backend().await;
let pm = Arc::new(
PersistenceManager::new(backend)
.await
.expect("persistence manager should initialize"),
);
let mut config = crate::cache_config::CacheConfig::default();
config.recent_messages.capacity = 1_000;
let (client, _sync_rx) = Client::new_with_cache_config(
Arc::new(crate::runtime_impl::TokioRuntime),
pm.clone(),
Arc::new(crate::transport::mock::MockTransportFactory::new()),
Arc::new(MockHttpClient),
None,
config,
)
.await;
let chat: Jid = "120363021033254949@g.us"
.parse()
.expect("test JID should be valid");
let msg_id = "ABC123".to_string();
let msg = wa::Message {
conversation: Some("hello".into()),
..Default::default()
};
client.add_recent_message(&chat, &msg_id, &msg, None).await;
let taken = client.take_recent_message(&chat, &msg_id).await;
assert!(taken.is_some());
let (msg, alt_chat) = taken.unwrap();
assert!(alt_chat.is_none(), "primary key should match");
assert_eq!(msg.conversation.as_deref(), Some("hello"));
let taken_again = client.take_recent_message(&chat, &msg_id).await;
assert!(taken_again.is_none());
}
#[tokio::test]
async fn recent_message_db_only_round_trip() {
let _ = env_logger::builder().is_test(true).try_init();
let backend = crate::test_utils::create_test_backend().await;
let pm = Arc::new(
PersistenceManager::new(backend)
.await
.expect("persistence manager should initialize"),
);
let config = crate::cache_config::CacheConfig::default();
assert_eq!(
config.recent_messages.capacity, 0,
"this test asserts the DB-only (capacity 0) path"
);
let (client, _sync_rx) = Client::new_with_cache_config(
Arc::new(crate::runtime_impl::TokioRuntime),
pm.clone(),
Arc::new(crate::transport::mock::MockTransportFactory::new()),
Arc::new(MockHttpClient),
None,
config,
)
.await;
let chat: Jid = "120363021033254949@g.us"
.parse()
.expect("test JID should be valid");
let msg_id = "DBONLY1".to_string();
let msg = wa::Message {
conversation: Some("db-only".into()),
..Default::default()
};
client.add_recent_message(&chat, &msg_id, &msg, None).await;
let taken = client.take_recent_message(&chat, &msg_id).await;
assert!(
taken.is_some(),
"a DB-only stored message must be retrievable from the backend"
);
let (got, _alt) = taken.unwrap();
assert_eq!(got.conversation.as_deref(), Some("db-only"));
let again = client.take_recent_message(&chat, &msg_id).await;
assert!(again.is_none(), "take consumes the DB-only message");
}
#[tokio::test]
async fn peek_recent_message_does_not_consume() {
let _ = env_logger::builder().is_test(true).try_init();
let backend = crate::test_utils::create_test_backend().await;
let pm = Arc::new(
PersistenceManager::new(backend)
.await
.expect("persistence manager should initialize"),
);
let mut config = crate::cache_config::CacheConfig::default();
config.recent_messages.capacity = 1_000;
let (client, _sync_rx) = Client::new_with_cache_config(
Arc::new(crate::runtime_impl::TokioRuntime),
pm.clone(),
Arc::new(crate::transport::mock::MockTransportFactory::new()),
Arc::new(MockHttpClient),
None,
config,
)
.await;
let chat: Jid = "120363021033254949@g.us".parse().unwrap();
let msg_id = "PEEK1".to_string();
let msg = wa::Message {
conversation: Some("hi".into()),
..Default::default()
};
client.add_recent_message(&chat, &msg_id, &msg, None).await;
for _ in 0..2 {
let peeked = client.peek_recent_message(&chat, &msg_id).await;
let (m, alt) = peeked.expect("peek should find the cached message");
assert!(alt.is_none());
assert_eq!(m.conversation.as_deref(), Some("hi"));
}
assert!(client.take_recent_message(&chat, &msg_id).await.is_some());
}
#[test]
fn get_bytes_content_extracts_bytes() {
use wacore_binary::{Attrs, Node};
let node = Node {
tag: Cow::Borrowed("test"),
attrs: Attrs::new(),
content: Some(NodeContent::Bytes(vec![1, 2, 3, 4])),
};
assert_eq!(get_bytes_content(&node), Some(&[1, 2, 3, 4][..]));
let node_str = Node {
tag: Cow::Borrowed("test"),
attrs: Attrs::new(),
content: Some(NodeContent::String("hello".into())),
};
assert_eq!(get_bytes_content(&node_str), None);
let node_empty = Node {
tag: Cow::Borrowed("test"),
attrs: Attrs::new(),
content: None,
};
assert_eq!(get_bytes_content(&node_empty), None);
}
#[test]
fn peer_detection_logic() {
let our_jid = Jid::pn("559911112222");
let peer_jid = Jid::pn_device("559911112222", 1);
let other_jid = Jid::pn("559933334444");
assert_eq!(our_jid.user, peer_jid.user);
assert_ne!(our_jid.user, other_jid.user);
}
#[test]
fn retry_receipt_attributes_for_device_sync_vs_peer_vs_group() {
use wacore::types::message::{MessageCategory, MessageInfo, MessageSource};
use wacore_binary::builder::NodeBuilder;
let our_pn = Jid::pn("559999999999");
let our_lid = Jid::lid("100000000000001");
fn build_retry_receipt(info: &MessageInfo, our_pn: &Jid, our_lid: &Jid) -> Node {
let receipt_to = if info.source.is_group {
&info.source.chat
} else {
&info.source.sender
};
let mut builder = NodeBuilder::new("receipt")
.attr("to", receipt_to)
.attr("id", info.id.clone())
.attr("type", "retry");
if info.source.is_group {
builder = builder.attr("participant", &info.source.sender);
}
if !info.source.is_group {
let is_from_own_account = info.source.sender.is_same_user_as(our_pn)
|| info.source.sender.is_same_user_as(our_lid);
if is_from_own_account {
if info.category == MessageCategory::Peer {
builder = builder.attr("category", MessageCategory::Peer.as_str());
} else {
let recipient = info.source.recipient.as_ref().unwrap_or(&info.source.chat);
builder = builder.attr("recipient", recipient);
}
}
}
builder.build()
}
let recipient_lid = Jid::lid("200000000000002");
let device_sync_info = MessageInfo {
id: "DEVICE_SYNC_MSG_001".to_string(),
source: MessageSource {
chat: recipient_lid.clone(),
sender: our_lid.clone(),
is_from_me: true,
is_group: false,
recipient: Some(recipient_lid.clone()),
..Default::default()
},
category: MessageCategory::default(),
..Default::default()
};
let node = build_retry_receipt(&device_sync_info, &our_pn, &our_lid);
assert_eq!(
node.attrs
.get("recipient")
.map(|v| v == "200000000000002@lid"),
Some(true),
"Device sync DM should include recipient"
);
assert!(
node.attrs.get("category").is_none(),
"Device sync DM should NOT have category=peer"
);
assert!(
node.attrs.get("participant").is_none(),
"DM should NOT have participant"
);
let other_pn = Jid::pn("551188888888");
let peer_info = MessageInfo {
id: "PEER123".to_string(),
source: MessageSource {
chat: other_pn.clone(),
sender: our_pn.clone(),
is_from_me: true,
is_group: false,
recipient: None,
..Default::default()
},
category: MessageCategory::Peer,
..Default::default()
};
let node = build_retry_receipt(&peer_info, &our_pn, &our_lid);
assert_eq!(
node.attrs.get("category").map(|v| v == "peer"),
Some(true),
"Peer DM should have category=peer"
);
assert!(
node.attrs.get("recipient").is_none(),
"Peer DM should NOT have recipient"
);
let group_info = MessageInfo {
id: "GROUP123".to_string(),
source: MessageSource {
chat: "123456789@g.us".parse().unwrap(),
sender: our_lid.clone(),
is_from_me: true,
is_group: true,
recipient: None,
..Default::default()
},
category: MessageCategory::default(),
..Default::default()
};
let node = build_retry_receipt(&group_info, &our_pn, &our_lid);
assert!(
node.attrs.get("participant").is_some(),
"Group should have participant"
);
assert!(
node.attrs.get("category").is_none(),
"Group should NOT have category"
);
assert!(
node.attrs.get("recipient").is_none(),
"Group should NOT have recipient"
);
let other_dm_info = MessageInfo {
id: "OTHER123".to_string(),
source: MessageSource {
chat: other_pn.clone(),
sender: other_pn.clone(),
is_from_me: false,
is_group: false,
recipient: None,
..Default::default()
},
category: MessageCategory::default(),
..Default::default()
};
let node = build_retry_receipt(&other_dm_info, &our_pn, &our_lid);
assert!(
node.attrs.get("category").is_none(),
"DM from other should NOT have category"
);
assert!(
node.attrs.get("recipient").is_none(),
"DM from other should NOT have recipient"
);
}
#[test]
fn enc_rekey_retry_receipt_node_structure() {
use wacore_binary::builder::NodeBuilder;
let peer_jid: Jid = "5511999999999@s.whatsapp.net".parse().expect("peer JID");
let call_creator: Jid = "5511888888888@s.whatsapp.net".parse().expect("creator JID");
let call_id = "CALL-ABC-123";
let stanza_id = "3EB0AABBCCDD";
let retry_count: u8 = 2;
let registration_id: u32 = 12345;
let enc_rekey_node = NodeBuilder::new("enc_rekey")
.attr("call-creator", call_creator)
.attr("call-id", call_id)
.attr("count", retry_count)
.build();
let registration_node = NodeBuilder::new("registration")
.bytes(registration_id.to_be_bytes().to_vec())
.build();
let receipt_node = NodeBuilder::new("receipt")
.attr("to", peer_jid)
.attr("id", stanza_id)
.attr("type", "enc_rekey_retry")
.children([enc_rekey_node, registration_node])
.build();
assert_eq!(
receipt_node.attrs().optional_string("type").as_deref(),
Some("enc_rekey_retry"),
"receipt type must be enc_rekey_retry"
);
assert!(
receipt_node
.attrs
.get("to")
.is_some_and(|v| *v == "5511999999999@s.whatsapp.net"),
"receipt 'to' must be peer JID"
);
assert_eq!(
receipt_node.attrs().optional_string("id").as_deref(),
Some("3EB0AABBCCDD")
);
assert!(
receipt_node.get_optional_child("retry").is_none(),
"enc_rekey_retry must NOT contain <retry> child"
);
let enc_rekey = receipt_node
.get_optional_child("enc_rekey")
.expect("<enc_rekey> child must exist");
assert_eq!(
enc_rekey.attrs().optional_string("call-id").as_deref(),
Some("CALL-ABC-123")
);
assert!(
enc_rekey
.attrs
.get("call-creator")
.is_some_and(|v| *v == "5511888888888@s.whatsapp.net"),
"enc_rekey 'call-creator' must be creator JID"
);
assert_eq!(
enc_rekey.attrs().optional_string("count").as_deref(),
Some("2")
);
let registration = receipt_node
.get_optional_child("registration")
.expect("<registration> child must exist");
let reg_bytes = match ®istration.content {
Some(NodeContent::Bytes(b)) => b.clone(),
_ => panic!("registration must contain bytes"),
};
assert_eq!(
u32::from_be_bytes(reg_bytes.try_into().unwrap()),
12345,
"registration ID must be 4-byte big-endian"
);
}
#[test]
fn prekey_id_parsing() {
let id_bytes = [0x01, 0x02, 0x03];
let prekey_id = u32::from_be_bytes([0, id_bytes[0], id_bytes[1], id_bytes[2]]);
assert_eq!(prekey_id, 0x00010203);
let skey_id_bytes = [0xFF, 0xFE, 0xFD];
let skey_id = u32::from_be_bytes([0, skey_id_bytes[0], skey_id_bytes[1], skey_id_bytes[2]]);
assert_eq!(skey_id, 0x00FFFEFD);
}
#[tokio::test]
async fn base_key_store_operations() {
let _ = env_logger::builder().is_test(true).try_init();
let backend = crate::test_utils::create_test_backend().await;
let address = "12345.0:1";
let msg_id = "ABC123";
let base_key = vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
let result = backend.has_same_base_key(address, msg_id, &base_key).await;
assert!(result.is_ok());
assert!(!result.unwrap());
let save_result = backend.save_base_key(address, msg_id, &base_key).await;
assert!(save_result.is_ok());
let result = backend.has_same_base_key(address, msg_id, &base_key).await;
assert!(result.is_ok());
assert!(result.unwrap());
let different_key = vec![10, 9, 8, 7, 6, 5, 4, 3, 2, 1];
let result = backend
.has_same_base_key(address, msg_id, &different_key)
.await;
assert!(result.is_ok());
assert!(!result.unwrap());
let delete_result = backend.delete_base_key(address, msg_id).await;
assert!(delete_result.is_ok());
let result = backend.has_same_base_key(address, msg_id, &base_key).await;
assert!(result.is_ok());
assert!(!result.unwrap());
}
#[tokio::test]
async fn base_key_store_upsert() {
let _ = env_logger::builder().is_test(true).try_init();
let backend = crate::test_utils::create_test_backend().await;
let address = "12345.0:1";
let msg_id = "MSG001";
let first_key = vec![1, 2, 3];
let second_key = vec![4, 5, 6];
backend
.save_base_key(address, msg_id, &first_key)
.await
.unwrap();
assert!(
backend
.has_same_base_key(address, msg_id, &first_key)
.await
.unwrap()
);
assert!(
!backend
.has_same_base_key(address, msg_id, &second_key)
.await
.unwrap()
);
backend
.save_base_key(address, msg_id, &second_key)
.await
.unwrap();
assert!(
!backend
.has_same_base_key(address, msg_id, &first_key)
.await
.unwrap()
);
assert!(
backend
.has_same_base_key(address, msg_id, &second_key)
.await
.unwrap()
);
}
#[tokio::test]
async fn base_key_store_multiple_messages() {
let _ = env_logger::builder().is_test(true).try_init();
let backend = crate::test_utils::create_test_backend().await;
let address = "12345.0:1";
let msg_id_1 = "MSG001";
let msg_id_2 = "MSG002";
let key_1 = vec![1, 2, 3];
let key_2 = vec![4, 5, 6];
backend
.save_base_key(address, msg_id_1, &key_1)
.await
.unwrap();
backend
.save_base_key(address, msg_id_2, &key_2)
.await
.unwrap();
assert!(
backend
.has_same_base_key(address, msg_id_1, &key_1)
.await
.unwrap()
);
assert!(
!backend
.has_same_base_key(address, msg_id_1, &key_2)
.await
.unwrap()
);
assert!(
!backend
.has_same_base_key(address, msg_id_2, &key_1)
.await
.unwrap()
);
assert!(
backend
.has_same_base_key(address, msg_id_2, &key_2)
.await
.unwrap()
);
backend.delete_base_key(address, msg_id_1).await.unwrap();
assert!(
!backend
.has_same_base_key(address, msg_id_1, &key_1)
.await
.unwrap()
);
assert!(
backend
.has_same_base_key(address, msg_id_2, &key_2)
.await
.unwrap()
);
}
fn build_retry_receipt_without_keys() -> Node {
use wacore_binary::builder::NodeBuilder;
NodeBuilder::new("receipt").build()
}
fn build_retry_receipt_with_registration(reg_id: u32) -> Node {
use wacore_binary::builder::NodeBuilder;
NodeBuilder::new("receipt")
.children([NodeBuilder::new("registration")
.bytes(reg_id.to_be_bytes().to_vec())
.build()])
.build()
}
fn dm_retry_info(resolved_jid: &Jid) -> RetryChatInfo {
RetryChatInfo {
chat: resolved_jid.to_non_ad(),
requester: resolved_jid.clone(),
original_from: resolved_jid.clone(),
recipient: None,
is_bot: false,
is_fbid_bot_retry: false,
}
}
fn valid_serialized_session(remote_regid: u32, base_key: Vec<u8>) -> Vec<u8> {
use wacore::libsignal::protocol::{SessionRecord, SessionState};
use waproto::whatsapp::SessionStructure;
let state = SessionState::from_session_structure(SessionStructure {
session_version: Some(3),
local_identity_public: None,
remote_identity_public: None,
root_key: None,
previous_counter: Some(0),
sender_chain: buffa::MessageField::default(),
receiver_chains: vec![],
pending_pre_key: buffa::MessageField::default(),
remote_registration_id: Some(remote_regid),
local_registration_id: Some(0),
alice_base_key: Some(base_key),
needs_refresh: None,
pending_key_exchange: buffa::MessageField::default(),
});
SessionRecord::new(state)
.serialize()
.expect("serialize session record")
}
#[tokio::test]
async fn update_local_signal_session_preserves_dm_session_at_retry_1() {
let client =
crate::test_utils::create_test_client_with_failing_http("retry_preserve_retry_1").await;
let user = "100000000000088".to_string();
let resolved_jid = Jid::lid_device(user.clone(), 33);
let backend = client.persistence_manager.backend();
let device_0 = Jid::lid_device(user.clone(), 0).to_protocol_address();
let device_33 = Jid::lid_device(user, 33).to_protocol_address();
let session_bytes_33 = valid_serialized_session(4242, vec![0xAA; 32]);
let session_bytes_0 = valid_serialized_session(4243, vec![0xBB; 32]);
backend
.put_session(device_0.as_str(), &session_bytes_0)
.await
.unwrap();
backend
.put_session(device_33.as_str(), &session_bytes_33)
.await
.unwrap();
let node = build_retry_receipt_without_keys();
let node_ref = node.as_node_ref();
client
.update_local_signal_session(
&dm_retry_info(&resolved_jid),
&resolved_jid,
"MSG-RETRY-1",
1,
&node_ref,
false,
)
.await;
client.flush_signal_cache().await.unwrap();
assert!(
backend
.get_session(device_0.as_str())
.await
.unwrap()
.is_some(),
"non-requesting device session must be preserved"
);
assert!(
backend
.get_session(device_33.as_str())
.await
.unwrap()
.is_some(),
"requesting device session with valid record must be preserved at retry #1"
);
}
#[tokio::test]
async fn update_local_signal_session_deletes_on_regid_mismatch() {
let client =
crate::test_utils::create_test_client_with_failing_http("retry_regid_mismatch").await;
let resolved_jid = Jid::lid_device("100000000000099".to_string(), 17);
let signal_address = resolved_jid.to_protocol_address();
let backend = client.persistence_manager.backend();
let stored_regid = 4242u32;
let session_bytes = valid_serialized_session(stored_regid, vec![0xAA; 32]);
backend
.put_session(signal_address.as_str(), &session_bytes)
.await
.unwrap();
let received_regid = 0xDEAD_BEEFu32;
assert_ne!(stored_regid, received_regid);
let node = build_retry_receipt_with_registration(received_regid);
let node_ref = node.as_node_ref();
client
.update_local_signal_session(
&dm_retry_info(&resolved_jid),
&resolved_jid,
"MSG-REGID",
1,
&node_ref,
false,
)
.await;
client.flush_signal_cache().await.unwrap();
assert!(
backend
.get_session(signal_address.as_str())
.await
.unwrap()
.is_none(),
"session must be deleted when retry has no keys and reg IDs differ"
);
}
#[tokio::test]
async fn update_local_signal_session_handles_unparseable_session_gracefully() {
let client =
crate::test_utils::create_test_client_with_failing_http("retry_unparseable_session")
.await;
let resolved_jid = Jid::lid_device("100000000000099".to_string(), 17);
let signal_address = resolved_jid.to_protocol_address();
let backend = client.persistence_manager.backend();
backend
.put_session(signal_address.as_str(), b"invalid-session")
.await
.unwrap();
let node = build_retry_receipt_with_registration(0xDEAD_BEEF);
let node_ref = node.as_node_ref();
client
.update_local_signal_session(
&dm_retry_info(&resolved_jid),
&resolved_jid,
"MSG-REGID",
1,
&node_ref,
false,
)
.await;
client.flush_signal_cache().await.unwrap();
assert!(
backend
.get_session(signal_address.as_str())
.await
.unwrap()
.is_some(),
"unparseable bytes skip every branch; nothing should delete them"
);
}
#[tokio::test]
async fn update_local_signal_session_no_session_is_noop() {
let client =
crate::test_utils::create_test_client_with_failing_http("retry_no_session").await;
let resolved_jid = Jid::lid_device("100000000000199".to_string(), 42);
let node = build_retry_receipt_without_keys();
let node_ref = node.as_node_ref();
client
.update_local_signal_session(
&dm_retry_info(&resolved_jid),
&resolved_jid,
"MSG-NOSESS",
1,
&node_ref,
false,
)
.await;
}
#[tokio::test]
async fn update_local_signal_session_preserves_group_session_at_retry_1() {
let client =
crate::test_utils::create_test_client_with_failing_http("retry_group_preserve").await;
let resolved_jid = Jid::lid_device("100000000000088".to_string(), 33);
let signal_address = resolved_jid.to_protocol_address();
let backend = client.persistence_manager.backend();
let session_bytes = valid_serialized_session(9999, vec![0xCC; 32]);
backend
.put_session(signal_address.as_str(), &session_bytes)
.await
.unwrap();
let group_chat: Jid = "120363042537531116@g.us".parse().unwrap();
let info = RetryChatInfo {
chat: group_chat.clone(),
requester: resolved_jid.clone(),
original_from: group_chat,
recipient: None,
is_bot: false,
is_fbid_bot_retry: false,
};
let node = build_retry_receipt_without_keys();
let node_ref = node.as_node_ref();
client
.update_local_signal_session(&info, &resolved_jid, "MSG-GRP-1", 1, &node_ref, false)
.await;
client.flush_signal_cache().await.unwrap();
assert!(
backend
.get_session(signal_address.as_str())
.await
.unwrap()
.is_some(),
"group retry at #1 should not delete the session"
);
}
#[tokio::test]
async fn update_local_signal_session_cools_resolved_sender_key_namespace() {
let client = crate::test_utils::create_test_client_with_failing_http(
"retry_sender_key_resolved_namespace",
)
.await;
let group = "120363000000000006@g.us";
let requester_pn: Jid = "12025550108:33@s.whatsapp.net".parse().unwrap();
let resolved_lid: Jid = "100000000000088:33@lid".parse().unwrap();
client
.persistence_manager
.set_sender_key_status(
group,
&[
("12025550108:33@s.whatsapp.net", true),
("100000000000088:33@lid", true),
],
)
.await
.unwrap();
let rows = client
.persistence_manager
.get_sender_key_devices(group)
.await
.unwrap();
let cached = client
.sender_key_device_cache
.get_or_init(group, async {
Arc::new(crate::sender_key_device_cache::SenderKeyDeviceMap::from_db_rows(&rows))
})
.await;
let info = RetryChatInfo {
chat: group.parse().unwrap(),
requester: requester_pn,
original_from: group.parse().unwrap(),
recipient: None,
is_bot: false,
is_fbid_bot_retry: false,
};
let node = build_retry_receipt_without_keys();
assert!(
client
.update_local_signal_session(
&info,
&resolved_lid,
"MSG-GRP-NAMESPACE",
1,
&node.as_node_ref(),
false,
)
.await
);
assert_eq!(cached.device_has_key("100000000000088", 33), Some(false));
assert_eq!(cached.device_has_key("12025550108", 33), Some(true));
let persisted = crate::sender_key_device_cache::SenderKeyDeviceMap::from_db_rows(
&client
.persistence_manager
.get_sender_key_devices(group)
.await
.unwrap(),
);
assert_eq!(persisted.device_has_key("100000000000088", 33), Some(false));
assert_eq!(persisted.device_has_key("12025550108", 33), Some(true));
}
#[tokio::test]
async fn status_retransmission_resolution_is_cache_aside_with_pn_fallback() {
let client = crate::test_utils::create_test_client_with_failing_http(
"retry_status_requester_resolution",
)
.await;
client
.add_lid_pn_mapping(
"100000000000089",
"12025550109",
crate::lid_pn_cache::LearningSource::Usync,
)
.await
.unwrap();
client.lid_pn_cache.clear().await;
let mapped_pn: Jid = "12025550109:19@s.whatsapp.net".parse().unwrap();
let mapped = client
.resolve_retransmission_encryption_jid(RetransmissionRoute::Status, &mapped_pn)
.await
.unwrap();
assert_eq!(mapped, "100000000000089:19@lid".parse::<Jid>().unwrap());
let unmapped_pn: Jid = "12025550110:20@s.whatsapp.net".parse().unwrap();
let fallback = client
.resolve_retransmission_encryption_jid(RetransmissionRoute::Status, &unmapped_pn)
.await
.unwrap();
assert_eq!(fallback, unmapped_pn);
}
#[tokio::test]
async fn should_recreate_session_matrix() {
let client =
crate::test_utils::create_test_client_with_failing_http("should_recreate_session")
.await;
let jid_with = Jid::lid_device("999999999999991".to_string(), 3);
let jid_without = Jid::lid_device("999999999999992".to_string(), 3);
let session_bytes = valid_serialized_session(7777, vec![0xEE; 32]);
client
.persistence_manager
.backend()
.put_session(jid_with.to_protocol_address().as_str(), &session_bytes)
.await
.unwrap();
assert!(
client.should_recreate_session(1, &jid_with).await.is_none(),
"retry<2 with session present should not recreate"
);
assert!(
client
.session_recreate_history
.get(&jid_with)
.await
.is_none(),
"no-op path must not stamp the history"
);
assert!(
client
.should_recreate_session(2, &jid_with)
.await
.is_some_and(|r| r.contains("retry count > 1")),
"retry≥2 with cold history should recreate"
);
let after_first = client.session_recreate_history.get(&jid_with).await;
assert!(after_first.is_some(), "first recreate must stamp history");
assert!(
client.should_recreate_session(3, &jid_with).await.is_none(),
"retry≥2 within {}s should be throttled",
RECREATE_SESSION_TIMEOUT.as_secs()
);
let after_second = client.session_recreate_history.get(&jid_with).await;
assert_eq!(
after_first, after_second,
"throttled path must not re-stamp the history"
);
let stamp_then = after_first.expect("first recreate stamped history");
let well_past = stamp_then + RECREATE_SESSION_TIMEOUT + std::time::Duration::from_secs(1);
assert!(
client
.should_recreate_session_at(3, &jid_with, well_past)
.await
.is_some_and(|r| r.contains("over an hour")),
"entry past the throttle window must allow a fresh recreate"
);
assert!(
client
.should_recreate_session(0, &jid_without)
.await
.is_some_and(|r| r.contains("don't have a Signal session")),
"missing session should recreate"
);
}
#[tokio::test]
async fn session_recreate_history_is_capacity_bounded() {
let client =
crate::test_utils::create_test_client_with_failing_http("session_recreate_history_cap")
.await;
let now = wacore::time::Instant::now();
let cap: u64 = 256;
for i in 0..(cap * 2) {
let jid = Jid::lid_device(format!("{}", 900_000_000_000_000u64 + i), 3);
client.session_recreate_history.insert(jid, now).await;
}
client.session_recreate_history.run_pending_tasks().await;
let count = client.session_recreate_history.entry_count();
assert!(
count <= cap,
"capacity must bound the throttle history (got {count}, cap {cap}); \
a still-recent entry can be evicted under heavy peer load"
);
}
#[tokio::test]
async fn client_resend_rate_limiter_is_wired_and_tunable() {
let client =
crate::test_utils::create_test_client_with_failing_http("resend_rl_wired").await;
let chat: Jid = "120363021033254949@g.us".parse().unwrap();
client.set_resend_rate_limit(3, 0);
let mut allowed = 0;
for _ in 0..10 {
if client.resend_rate_limiter.try_acquire(&chat).await {
allowed += 1;
}
}
assert_eq!(allowed, 3, "client honors the configured per-chat burst");
assert_eq!(
client.stats().resends_throttled,
7,
"public counter tracks dropped resends"
);
client.set_resend_rate_limit(0, 0);
let other: Jid = "120363000000000001@g.us".parse().unwrap();
for _ in 0..50 {
assert!(client.resend_rate_limiter.try_acquire(&other).await);
}
}
#[tokio::test]
async fn handle_retry_receipt_drops_throttled_group_resend() {
use wacore_binary::builder::NodeBuilder;
let backend = crate::test_utils::create_test_backend().await;
let pm = Arc::new(PersistenceManager::new(backend).await.unwrap());
let mut config = crate::cache_config::CacheConfig::default();
config.recent_messages.capacity = 1_000;
let (client, _rx) = Client::new_with_cache_config(
Arc::new(crate::runtime_impl::TokioRuntime),
pm,
Arc::new(crate::transport::mock::MockTransportFactory::new()),
Arc::new(MockHttpClient),
None,
config,
)
.await;
let group: Jid = "120363021033254949@g.us".parse().unwrap();
let msg_id = "RLMSG001";
client
.add_recent_message(
&group,
msg_id,
&wa::Message {
conversation: Some("hi".into()),
..Default::default()
},
None,
)
.await;
client.set_resend_rate_limit(1, 0);
assert!(client.resend_rate_limiter.try_acquire(&group).await);
let node = NodeBuilder::new("receipt")
.attr("participant", "555000111@lid")
.children([NodeBuilder::new("retry")
.attr("id", msg_id)
.attr("count", "1")
.build()])
.build();
let node_ref = crate::test_utils::node_to_owned_ref(&node);
let receipt = Receipt::builder()
.source(crate::types::message::MessageSource {
chat: group.clone(),
sender: "555000111@lid".parse().unwrap(),
is_group: true,
..Default::default()
})
.message_ids(vec![msg_id.to_string()])
.timestamp(wacore::time::now_utc())
.r#type(crate::types::presence::ReceiptType::Retry)
.offline(false)
.build();
let result = client.handle_retry_receipt(&receipt, &node_ref).await;
assert!(
result.is_ok(),
"a throttled retry returns Ok(()), not an error"
);
assert_eq!(
client.stats().resends_throttled,
1,
"the resend was dropped by the limiter"
);
assert!(
client.peek_recent_message(&group, msg_id).await.is_some(),
"throttling keeps the message cached for the device's re-request"
);
assert_eq!(
client.pending_retries.lock().unwrap().len(),
0,
"the in-progress marker is cleared after the throttled return"
);
}
#[tokio::test]
async fn unknown_participant_rotation_is_durable_before_throttled_return() {
use wacore::libsignal::protocol::{SENDERKEY_MESSAGE_CURRENT_VERSION, SenderKeyRecord};
use wacore::libsignal::store::sender_key_name::SenderKeyName;
use wacore_binary::builder::NodeBuilder;
let backend = crate::test_utils::create_test_backend().await;
let pm = Arc::new(PersistenceManager::new(backend.clone()).await.unwrap());
let mut config = crate::cache_config::CacheConfig::default();
config.recent_messages.capacity = 1_000;
let (client, _rx) = Client::new_with_cache_config(
Arc::new(crate::runtime_impl::TokioRuntime),
pm,
Arc::new(crate::transport::mock::MockTransportFactory::new()),
Arc::new(MockHttpClient),
None,
config,
)
.await;
let own_lid: Jid = "100000000001040:13@lid".parse().unwrap();
client
.persistence_manager
.process_command(crate::store::commands::DeviceCommand::SetLid(Some(
own_lid.clone(),
)))
.await;
let group: Jid = "120363021033254950@g.us".parse().unwrap();
let group_id = group.to_string();
let sender_key_name =
SenderKeyName::from_parts(&group_id, own_lid.to_protocol_address().as_str());
let mut rng = rand::make_rng::<rand::rngs::StdRng>();
let key_pair = KeyPair::generate(&mut rng);
let mut record = SenderKeyRecord::new_empty();
record
.add_sender_key_state(
SENDERKEY_MESSAGE_CURRENT_VERSION,
9,
0,
&[7; 32],
key_pair.public_key,
Some(key_pair.private_key),
)
.unwrap();
client
.signal_cache
.put_sender_key(&sender_key_name, record)
.await;
client.flush_signal_cache().await.unwrap();
assert!(
backend
.get_sender_key(sender_key_name.cache_key())
.await
.unwrap()
.is_some()
);
let msg_id = "ROTATEFLUSH001";
client
.add_recent_message(
&group,
msg_id,
&wa::Message {
conversation: Some("hi".into()),
..Default::default()
},
None,
)
.await;
client.set_resend_rate_limit(1, 0);
assert!(client.resend_rate_limiter.try_acquire(&group).await);
let requester: Jid = "15551234002@s.whatsapp.net".parse().unwrap();
let node = NodeBuilder::new("receipt")
.attr("participant", &requester)
.children([NodeBuilder::new("retry")
.attr("id", msg_id)
.attr("count", "1")
.build()])
.build();
let node_ref = crate::test_utils::node_to_owned_ref(&node);
let receipt = Receipt::builder()
.source(crate::types::message::MessageSource {
chat: group.clone(),
sender: requester,
is_group: true,
..Default::default()
})
.message_ids(vec![msg_id.to_string()])
.timestamp(wacore::time::now_utc())
.r#type(crate::types::presence::ReceiptType::Retry)
.offline(false)
.build();
client
.handle_retry_receipt(&receipt, &node_ref)
.await
.unwrap();
assert!(
backend
.get_sender_key(sender_key_name.cache_key())
.await
.unwrap()
.is_none(),
"early retry return must not leave the retired key durable"
);
}
#[tokio::test]
async fn concurrent_same_peer_recreate_check_is_serialized() {
let client =
crate::test_utils::create_test_client_with_failing_http("concurrent_recreate").await;
let jid = Jid::lid_device("999999999999993".to_string(), 3);
let session_bytes = valid_serialized_session(8888, vec![0xCC; 32]);
client
.persistence_manager
.backend()
.put_session(jid.to_protocol_address().as_str(), &session_bytes)
.await
.unwrap();
let c1 = client.clone();
let j1 = jid.clone();
let task1 = async move {
let addr = j1.to_protocol_address();
let lock = c1.session_lock_for(addr.as_str()).await;
let _g = lock.lock().await;
c1.should_recreate_session(2, &j1).await.is_some()
};
let c2 = client.clone();
let j2 = jid.clone();
let task2 = async move {
let addr = j2.to_protocol_address();
let lock = c2.session_lock_for(addr.as_str()).await;
let _g = lock.lock().await;
c2.should_recreate_session(2, &j2).await.is_some()
};
let (a, b) = tokio::join!(task1, task2);
assert_eq!(
usize::from(a) + usize::from(b),
1,
"exactly one of two concurrent same-peer recreate checks may fire; \
the per-peer session lock serializes the non-atomic get+insert"
);
}
#[tokio::test]
async fn ensure_e2e_sessions_resolved_is_noop_when_session_exists() {
use std::sync::atomic::Ordering;
let client = crate::test_utils::create_test_client_with_failing_http(
"group_retry_ensure_sessions_noop",
)
.await;
client.offline_sync_completed.store(true, Ordering::Relaxed);
let resolved_jid = Jid::lid_device("100000000000199".to_string(), 17);
let signal_address = resolved_jid.to_protocol_address();
let session_bytes = valid_serialized_session(5555, vec![0xDD; 32]);
client
.persistence_manager
.backend()
.put_session(signal_address.as_str(), &session_bytes)
.await
.unwrap();
client
.ensure_e2e_sessions_resolved(std::slice::from_ref(&resolved_jid))
.await
.expect("no-op when session exists");
}
#[tokio::test]
async fn retry_key_bundle_requires_one_time_prekey_except_fbid_bot() {
let backend = crate::test_utils::create_test_backend().await;
let pm = Arc::new(
PersistenceManager::new(backend)
.await
.expect("persistence manager should initialize"),
);
let (client, _sync_rx) = Client::new(
Arc::new(crate::runtime_impl::TokioRuntime),
pm,
Arc::new(crate::transport::mock::MockTransportFactory::new()),
Arc::new(MockHttpClient),
None,
)
.await;
let mut rng = rand::make_rng::<rand::rngs::StdRng>();
let remote_identity = IdentityKeyPair::generate(&mut rng);
let signed_prekey = KeyPair::generate(&mut rng);
let signed_prekey_signature = remote_identity
.private_key()
.calculate_signature(&signed_prekey.public_key.serialize(), &mut rng)
.expect("signed prekey signature should be valid");
let regular_requester = Jid::pn_device("559922223333", 1);
let keys = NodeBuilder::new("keys")
.children([
NodeBuilder::new("type").bytes(vec![5]).build(),
NodeBuilder::new("identity")
.bytes(
remote_identity
.identity_key()
.public_key()
.public_key_bytes()
.to_vec(),
)
.build(),
SignedPreKeyNode::new(
100,
signed_prekey.public_key.public_key_bytes().to_vec(),
signed_prekey_signature.to_vec(),
)
.into_node(),
])
.build();
let receipt = NodeBuilder::new("receipt")
.children([
NodeBuilder::new("registration")
.bytes(12345u32.to_be_bytes().to_vec())
.build(),
keys,
])
.build();
let err = client
.process_retry_key_bundle(&receipt.as_node_ref(), ®ular_requester, false, false)
.await
.expect_err("regular retry without one-time prekey must be rejected");
assert!(
err.to_string()
.contains("regular retry key bundle missing one-time prekey")
);
let fbid_bot_requester = Jid::new("200000000000002", wacore_binary::Server::Bot);
client
.process_retry_key_bundle(&receipt.as_node_ref(), &fbid_bot_requester, false, true)
.await
.expect("fbid bot retry without one-time prekey should establish a session");
let snapshot = client.persistence_manager.get_device_snapshot();
let session = client
.signal_cache
.peek_session(
&fbid_bot_requester.to_protocol_address(),
&*snapshot.backend,
)
.await
.expect("session lookup should succeed");
assert!(session.is_some());
}
#[test]
fn bot_jid_detection() {
use wacore_binary::JidExt as _;
let regular_user: Jid = "1234567890@s.whatsapp.net".parse().unwrap();
assert!(!regular_user.is_bot());
let bot_server: Jid = "somebot@bot".parse().unwrap();
assert!(bot_server.is_bot());
let legacy_bot: Jid = "1313555123456@s.whatsapp.net".parse().unwrap();
assert!(legacy_bot.is_bot());
let legacy_bot2: Jid = "131655500123456@s.whatsapp.net".parse().unwrap();
assert!(legacy_bot2.is_bot());
let not_bot: Jid = "1313556123456@s.whatsapp.net".parse().unwrap();
assert!(!not_bot.is_bot());
}
#[test]
fn extract_registration_id_from_node_test() {
use wacore::protocol::retry::{
extract_registration_id_from_node, extract_registration_id_from_node_ref,
};
use wacore_binary::{Attrs, Node};
let reg_receipt = |bytes: Vec<u8>| Node {
tag: Cow::Borrowed("receipt"),
attrs: Attrs::new(),
content: Some(NodeContent::Nodes(vec![Node {
tag: Cow::Borrowed("registration"),
attrs: Attrs::new(),
content: Some(NodeContent::Bytes(bytes)),
}])),
};
let parent = reg_receipt(vec![0x00, 0x01, 0x02, 0x03]);
assert_eq!(extract_registration_id_from_node(&parent), Some(0x00010203));
assert_eq!(
extract_registration_id_from_node_ref(&parent.as_node_ref()),
Some(0x00010203)
);
let parent_short = reg_receipt(vec![0x01, 0x02, 0x03]);
assert_eq!(
extract_registration_id_from_node(&parent_short),
Some(0x00010203)
);
assert_eq!(
extract_registration_id_from_node_ref(&parent_short.as_node_ref()),
Some(0x00010203)
);
let parent_oversized = reg_receipt(vec![0x01, 0x02, 0x03, 0x04, 0x05]);
assert_eq!(extract_registration_id_from_node(&parent_oversized), None);
assert_eq!(
extract_registration_id_from_node_ref(&parent_oversized.as_node_ref()),
None
);
let parent_no_reg = Node {
tag: Cow::Borrowed("receipt"),
attrs: Attrs::new(),
content: Some(NodeContent::Nodes(vec![])),
};
assert_eq!(extract_registration_id_from_node(&parent_no_reg), None);
assert_eq!(
extract_registration_id_from_node_ref(&parent_no_reg.as_node_ref()),
None
);
let parent_empty = reg_receipt(vec![]);
assert_eq!(extract_registration_id_from_node(&parent_empty), None);
assert_eq!(
extract_registration_id_from_node_ref(&parent_empty.as_node_ref()),
None
);
}
#[test]
fn group_or_status_detection_for_sender_key_handling() {
use wacore_binary::JidExt as _;
let group: Jid = "120363021033254949@g.us".parse().unwrap();
let status: Jid = "status@broadcast".parse().unwrap();
let dm: Jid = "1234567890@s.whatsapp.net".parse().unwrap();
assert!(group.is_group() || group.is_status_broadcast());
assert!(status.is_group() || status.is_status_broadcast());
assert!(!(dm.is_group() || dm.is_status_broadcast()));
}
#[test]
fn retransmission_route_validation_is_strict_and_typed() {
let direct: Jid = "12025550100@s.whatsapp.net".parse().unwrap();
let requester: Jid = "12025550100:7@s.whatsapp.net".parse().unwrap();
let group: Jid = "120363000000000001@g.us".parse().unwrap();
let status = Jid::status_broadcast();
let broadcast: Jid = "1234567890@broadcast".parse().unwrap();
assert!(matches!(
validate_retransmission(&direct, &requester, "DM1", 1, Some(&direct)),
Ok(RetransmissionRoute::Direct)
));
assert!(matches!(
validate_retransmission(&group, &requester, "GROUP1", 1, None),
Ok(RetransmissionRoute::Group)
));
assert!(matches!(
validate_retransmission(&status, &requester, "STATUS1", 1, None),
Ok(RetransmissionRoute::Status)
));
assert!(matches!(
validate_retransmission(&broadcast, &requester, "BROADCAST1", 1, None),
Ok(RetransmissionRoute::BroadcastList)
));
for (id, count) in [("ZERO", 0), ("", 1), ("MAX", MAX_RETRY_COUNT)] {
assert!(
validate_retransmission(&direct, &requester, id, count, None).is_err(),
"invalid id/count pair must fail: {id:?}/{count}"
);
}
assert!(
validate_retransmission(&group, &requester, "GROUP2", 1, Some(&direct)).is_err(),
"recipient is only meaningful on a direct retry"
);
assert!(
validate_retransmission(&status, &group, "STATUS2", 1, None).is_err(),
"a group JID cannot be a requesting status device"
);
}
#[tokio::test]
async fn public_peer_retransmission_requires_a_recipient() {
let client = crate::test_utils::create_test_client().await;
let own_pn: Jid = "12025550100:13@s.whatsapp.net".parse().unwrap();
client
.persistence_manager
.process_command(crate::store::commands::DeviceCommand::SetId(Some(
own_pn.clone(),
)))
.await;
let chat: Jid = "12025550101@s.whatsapp.net".parse().unwrap();
let requester = own_pn.with_device(7);
let request = MessageRetransmission::new(
chat,
requester,
wa::Message::default(),
"PEER-RETRY-1".to_string(),
1,
);
let error = client
.retransmit_message(request)
.await
.expect_err("a peer route without its actual chat cannot be sent");
assert!(matches!(error, SendError::InvalidRequest(_)));
assert!(error.to_string().contains("requires a recipient"));
}
#[tokio::test]
async fn public_direct_retransmission_binds_chat_to_routing_identity() {
let client = crate::test_utils::create_test_client().await;
let chat = Jid::pn("12025550104");
let requester = Jid::pn_device("12025550105", 7);
let bot_requester: Jid = "200000000000002@bot".parse().unwrap();
for request in [
MessageRetransmission::new(
chat.clone(),
requester,
wa::Message::default(),
"DIRECT-CHAT-MISMATCH-1".to_string(),
1,
),
MessageRetransmission::new(
chat.clone(),
bot_requester,
wa::Message::default(),
"DIRECT-RECIPIENT-MISMATCH-1".to_string(),
1,
)
.with_recipient(Jid::pn("12025550106")),
] {
let error = client
.retransmit_message(request)
.await
.expect_err("an unrelated routing identity must be rejected");
assert!(matches!(error, SendError::InvalidRequest(_)));
assert!(error.to_string().contains("routing identity"));
}
}
#[tokio::test]
async fn public_direct_recipient_rejects_an_unrelated_requester() {
let client = crate::test_utils::create_test_client().await;
let chat = Jid::pn("12025550108");
let request = MessageRetransmission::new(
chat.clone(),
Jid::pn_device("12025550109", 7),
wa::Message::default(),
"DIRECT-RECIPIENT-SOURCE-1".to_string(),
1,
)
.with_recipient(chat);
let error = client
.retransmit_message(request)
.await
.expect_err("a normal remote user cannot declare a recipient route");
assert!(matches!(error, SendError::InvalidRequest(_)));
assert!(error.to_string().contains("local device or bot"));
}
#[tokio::test]
async fn direct_retransmission_chat_accepts_known_pn_lid_alias() {
let client = crate::test_utils::create_test_client().await;
let pn = Jid::pn("12025550107");
let lid = Jid::lid("100000000000107");
client
.lid_pn_cache
.add(&wacore::types::lid_pn::LidPnEntry {
lid: lid.user.as_str().into(),
phone_number: pn.user.as_str().into(),
created_at: 1,
learning_source: wacore::types::lid_pn::LearningSource::Usync,
})
.await;
assert!(client.jids_share_user_identity(&pn, &lid).await.unwrap());
assert!(client.jids_share_user_identity(&lid, &pn).await.unwrap());
}
#[tokio::test]
async fn public_retransmission_recaches_the_supplied_message() {
let mut config = crate::cache_config::CacheConfig::default();
config.recent_messages.capacity = 16;
let client = crate::test_utils::create_test_client_with_config(
"public_retransmission_cache",
Arc::new(MockHttpClient),
config,
)
.await;
let chat = Jid::pn("12025550103");
let requester = chat.with_device(7);
crate::test_utils::seed_peer_session(&client, &requester).await;
let message = wa::Message {
conversation: Some("retry me".into()),
..Default::default()
};
let message_id = "PUBLIC-RETRY-CACHE-1";
let result = client
.retransmit_message(MessageRetransmission::new(
chat.clone(),
requester,
message,
message_id.to_string(),
1,
))
.await;
assert!(result.is_err());
let (cached, alternate) = client
.peek_recent_message(&chat, message_id)
.await
.expect("a later retry count must find the retransmitted message");
assert!(alternate.is_none());
assert_eq!(cached.conversation.as_deref(), Some("retry me"));
}
#[test]
fn resolve_retry_chat_info_broadcast_uses_participant_device() {
let broadcast = "1234567890@broadcast";
let participant = "12025550101:9@s.whatsapp.net";
let node = NodeBuilder::new("receipt")
.attr("participant", participant)
.build();
let receipt = make_test_receipt(broadcast);
let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
assert!(info.chat.is_broadcast_list());
assert_eq!(info.requester, participant.parse::<Jid>().unwrap());
}
#[test]
fn retry_key_inclusion_matches_canonical_policy() {
use wacore::protocol::retry::{should_include_keys, should_include_keys_with_policy};
assert!(!should_include_keys(1, RetryReason::NoSession));
assert!(!should_include_keys(
1,
RetryReason::UnknownCompanionNoPrekey
));
assert!(should_include_keys_with_policy(1, true, false));
assert!(should_include_keys_with_policy(1, false, true));
assert!(should_include_keys(2, RetryReason::InvalidMessage));
assert!(should_include_keys(3, RetryReason::BadMac));
}
fn make_test_receipt(from: &str) -> Receipt {
Receipt::builder()
.source(crate::types::message::MessageSource {
chat: from.parse().unwrap(),
sender: from.parse().unwrap(),
..Default::default()
})
.message_ids(vec!["MSG001".to_string()])
.timestamp(wacore::time::now_utc())
.r#type(crate::types::presence::ReceiptType::Retry)
.offline(false)
.build()
}
#[test]
fn resolve_retry_chat_info_dm_with_device() {
use wacore_binary::builder::NodeBuilder;
let node = NodeBuilder::new("receipt").build();
let receipt = make_test_receipt("5511999999999:33@s.whatsapp.net");
let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
assert_eq!(info.chat.device(), 0);
assert_eq!(info.chat.user, "5511999999999");
assert!(info.chat.is_pn());
assert_eq!(info.requester.device(), 33);
assert_eq!(info.requester.user, "5511999999999");
}
#[test]
fn resolve_retry_chat_info_lid_dm_with_device() {
use wacore_binary::builder::NodeBuilder;
let node = NodeBuilder::new("receipt").build();
let receipt = make_test_receipt("236395184570386:5@lid");
let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
assert_eq!(info.chat.device(), 0);
assert_eq!(info.chat.user, "236395184570386");
assert!(info.chat.is_lid());
assert_eq!(info.requester.device(), 5);
assert_eq!(info.requester.user, "236395184570386");
assert!(info.requester.is_lid());
}
#[test]
fn resolve_retry_chat_info_forwards_recipient_attribute_verbatim() {
use wacore_binary::builder::NodeBuilder;
let node = NodeBuilder::new("receipt")
.attr("recipient", "5500000000123@s.whatsapp.net")
.build();
let receipt = make_test_receipt("100000000000456:5@lid");
let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
let recipient = info
.recipient
.as_ref()
.expect("recipient must be populated from the node attr");
assert_eq!(recipient.user, "5500000000123");
assert!(recipient.is_pn(), "recipient namespace must be PN");
assert_ne!(
recipient.user, info.chat.user,
"recipient must come from the node attr, not info.chat"
);
let node_no_recipient = NodeBuilder::new("receipt").build();
let info_no_recipient =
resolve_retry_chat_info(&receipt, &node_no_recipient.as_node_ref(), None, None);
assert!(
info_no_recipient.recipient.is_none(),
"missing `recipient` attr must propagate as None"
);
}
#[test]
fn resolve_retry_chat_info_dm_bare() {
use wacore_binary::builder::NodeBuilder;
let node = NodeBuilder::new("receipt").build();
let receipt = make_test_receipt("5511999999999@s.whatsapp.net");
let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
assert_eq!(info.chat.device(), 0);
assert_eq!(info.requester.device(), 0);
assert_eq!(info.chat, info.requester);
}
#[test]
fn resolve_retry_chat_info_group() {
use wacore_binary::builder::NodeBuilder;
let node = NodeBuilder::new("receipt")
.attr("from", "120363021033254949@g.us")
.attr("id", "MSG001")
.attr("participant", "236395184570386:33@lid")
.attr("type", "retry")
.build();
let receipt = Receipt::builder()
.source(crate::types::message::MessageSource {
chat: "120363021033254949@g.us".parse().unwrap(),
sender: "236395184570386:33@lid".parse().unwrap(),
..Default::default()
})
.message_ids(vec!["MSG001".to_string()])
.timestamp(wacore::time::now_utc())
.r#type(crate::types::presence::ReceiptType::Retry)
.offline(false)
.build();
let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
assert!(info.chat.is_group());
assert_eq!(info.chat.user, "120363021033254949");
assert!(info.requester.is_lid());
assert_eq!(info.requester.device(), 33);
}
#[test]
fn resolve_retry_chat_info_group_bot_device_marks_bot_namespace_only() {
use wacore_binary::builder::NodeBuilder;
let node = NodeBuilder::new("receipt")
.attr("participant", "somebot:4@bot")
.build();
let receipt = Receipt::builder()
.source(crate::types::message::MessageSource {
chat: "120363021033254949@g.us".parse().unwrap(),
sender: "somebot:4@bot".parse().unwrap(),
..Default::default()
})
.message_ids(vec!["MSG001".to_string()])
.timestamp(wacore::time::now_utc())
.r#type(crate::types::presence::ReceiptType::Retry)
.offline(false)
.build();
let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
assert!(info.chat.is_group());
assert!(info.is_bot);
assert!(!info.is_fbid_bot_retry);
}
#[test]
fn resolve_retry_chat_info_group_primary_fbid_bot_marks_bot_retry() {
use wacore_binary::builder::NodeBuilder;
let node = NodeBuilder::new("receipt")
.attr("participant", "somebot@bot")
.build();
let receipt = Receipt::builder()
.source(crate::types::message::MessageSource {
chat: "120363021033254949@g.us".parse().unwrap(),
sender: "somebot@bot".parse().unwrap(),
..Default::default()
})
.message_ids(vec!["MSG001".to_string()])
.timestamp(wacore::time::now_utc())
.r#type(crate::types::presence::ReceiptType::Retry)
.offline(false)
.build();
let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
assert!(info.chat.is_group());
assert!(info.is_bot);
assert!(info.is_fbid_bot_retry);
}
#[test]
fn resolve_retry_chat_info_status_broadcast() {
use wacore_binary::builder::NodeBuilder;
let node = NodeBuilder::new("receipt")
.attr("from", "status@broadcast")
.attr("id", "3EB06D00CAB92340790621")
.attr("participant", "236395184570386@lid")
.attr("type", "retry")
.build();
let receipt = make_test_receipt("status@broadcast");
let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
assert!(info.chat.is_status_broadcast());
assert!(info.requester.is_lid());
assert_eq!(info.requester.user, "236395184570386");
}
#[test]
fn resolve_retry_chat_info_status_broadcast_no_participant() {
use wacore_binary::builder::NodeBuilder;
let node = NodeBuilder::new("receipt")
.attr("from", "status@broadcast")
.attr("id", "MSG001")
.attr("type", "retry")
.build();
let receipt = make_test_receipt("status@broadcast");
let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
assert!(info.chat.is_status_broadcast());
assert!(info.requester.is_status_broadcast());
}
#[test]
fn retry_processing_key_per_participant() {
let msg_id = "3EB06D00CAB92340790621";
let status_chat = Jid::status_broadcast();
let status_participant_a: Jid = "236395184570386@lid".parse().unwrap();
let status_participant_b: Jid = "559985213786@s.whatsapp.net".parse().unwrap();
let status_key_a = build_retry_processing_key(&status_chat, msg_id, &status_participant_a);
let status_key_b = build_retry_processing_key(&status_chat, msg_id, &status_participant_b);
assert_ne!(
status_key_a, status_key_b,
"Different status participants must have different processing keys"
);
assert_eq!(
status_key_a,
build_retry_processing_key(&status_chat, msg_id, &status_participant_a),
"Same participant must produce the same key — any retry count for that \
participant serializes through pending_retries"
);
let dm_chat = Jid::pn("559911112222");
let dm_device_a = Jid::pn_device("559922223333", 1);
let dm_device_b = Jid::pn_device("559922223333", 2);
let dm_key_a = build_retry_processing_key(&dm_chat, msg_id, &dm_device_a);
let dm_key_b = build_retry_processing_key(&dm_chat, msg_id, &dm_device_b);
assert_ne!(
dm_key_a, dm_key_b,
"Different DM requester devices must have different processing keys"
);
assert_eq!(
dm_key_a,
build_retry_processing_key(&dm_chat, msg_id, &dm_device_a),
"Same DM requester device must produce the same processing key"
);
}
#[tokio::test]
async fn recent_message_cache_readd_after_take() {
let _ = env_logger::builder().is_test(true).try_init();
let backend = crate::test_utils::create_test_backend().await;
let pm = Arc::new(
PersistenceManager::new(backend)
.await
.expect("persistence manager should initialize"),
);
let mut config = crate::cache_config::CacheConfig::default();
config.recent_messages.capacity = 1_000;
let (client, _sync_rx) = Client::new_with_cache_config(
Arc::new(crate::runtime_impl::TokioRuntime),
pm.clone(),
Arc::new(crate::transport::mock::MockTransportFactory::new()),
Arc::new(MockHttpClient),
None,
config,
)
.await;
let msg = wa::Message {
extended_text_message: buffa::MessageField::some(wa::message::ExtendedTextMessage {
text: Some("status text".to_string()),
..Default::default()
}),
..Default::default()
};
for (chat, msg_id) in [
(Jid::status_broadcast(), "STATUS_MSG_001".to_string()),
(Jid::pn("559911112222"), "DM_MSG_001".to_string()),
] {
client.add_recent_message(&chat, &msg_id, &msg, None).await;
let taken = client.take_recent_message(&chat, &msg_id).await;
assert!(taken.is_some(), "First take should succeed for {chat}");
let (taken_msg, _) = taken.unwrap();
client
.add_recent_message(&chat, &msg_id, &taken_msg, None)
.await;
let taken2 = client.take_recent_message(&chat, &msg_id).await;
assert!(
taken2.is_some(),
"Second take should succeed after re-add for {chat}"
);
assert_eq!(
taken2
.unwrap()
.0
.extended_text_message
.as_option()
.unwrap()
.text
.as_deref(),
Some("status text")
);
}
}
#[tokio::test]
async fn dm_retry_message_lookup_uses_bare_jid() {
let _ = env_logger::builder().is_test(true).try_init();
let backend = crate::test_utils::create_test_backend().await;
let pm = Arc::new(
PersistenceManager::new(backend)
.await
.expect("persistence manager should initialize"),
);
let mut config = crate::cache_config::CacheConfig::default();
config.recent_messages.capacity = 1_000;
let (client, _sync_rx) = Client::new_with_cache_config(
Arc::new(crate::runtime_impl::TokioRuntime),
pm.clone(),
Arc::new(crate::transport::mock::MockTransportFactory::new()),
Arc::new(MockHttpClient),
None,
config,
)
.await;
let bare_jid: Jid = "5511999999999@s.whatsapp.net".parse().unwrap();
let msg_id = "RETRY_MSG_001";
let msg = wa::Message {
conversation: Some("test dm".into()),
..Default::default()
};
client
.add_recent_message(&bare_jid, msg_id, &msg, None)
.await;
let taken = client.take_recent_message(&bare_jid, msg_id).await;
assert!(taken.is_some(), "Lookup via bare JID should succeed");
let (msg_out, alt_chat) = taken.unwrap();
assert!(alt_chat.is_none(), "primary key should match for bare JID");
client
.add_recent_message(&bare_jid, msg_id, &msg_out, None)
.await;
let taken2 = client.take_recent_message(&bare_jid, msg_id).await;
assert!(
taken2.is_some(),
"Second lookup via bare JID should succeed after re-add"
);
}
#[tokio::test]
async fn alternate_key_lookup_pn_to_lid() {
let _ = env_logger::builder().is_test(true).try_init();
let backend = crate::test_utils::create_test_backend().await;
let pm = Arc::new(
PersistenceManager::new(backend)
.await
.expect("persistence manager should initialize"),
);
let mut config = crate::cache_config::CacheConfig::default();
config.recent_messages.capacity = 1_000;
let (client, _sync_rx) = Client::new_with_cache_config(
Arc::new(crate::runtime_impl::TokioRuntime),
pm.clone(),
Arc::new(crate::transport::mock::MockTransportFactory::new()),
Arc::new(MockHttpClient),
None,
config,
)
.await;
let pn_jid: Jid = "5511999999999@s.whatsapp.net".parse().unwrap();
let lid_jid: Jid = "236395184570386@lid".parse().unwrap();
let msg_id = "RETRY_ALT_001";
let msg = wa::Message {
conversation: Some("alternate key test".into()),
..Default::default()
};
client.add_recent_message(&pn_jid, msg_id, &msg, None).await;
client
.lid_pn_cache
.add(&wacore::types::lid_pn::LidPnEntry {
lid: lid_jid.user.as_str().into(),
phone_number: pn_jid.user.as_str().into(),
created_at: 0,
learning_source: wacore::types::lid_pn::LearningSource::Usync,
})
.await;
let taken = client.take_recent_message(&lid_jid, msg_id).await;
assert!(
taken.is_some(),
"Alternate PN key lookup should find message stored under PN"
);
let (msg_out, alt_chat) = taken.unwrap();
let alt_chat = alt_chat.expect("should be found via alternate key");
assert!(alt_chat.is_pn(), "alternate chat should be PN");
assert_eq!(alt_chat.user, pn_jid.user);
assert_eq!(msg_out.conversation.as_deref(), Some("alternate key test"));
}
#[tokio::test]
async fn swap_pn_lid_namespace_preserves_device() {
let _ = env_logger::builder().is_test(true).try_init();
let backend = crate::test_utils::create_test_backend().await;
let pm = Arc::new(
PersistenceManager::new(backend)
.await
.expect("persistence manager should initialize"),
);
let (client, _sync_rx) = Client::new(
Arc::new(crate::runtime_impl::TokioRuntime),
pm.clone(),
Arc::new(crate::transport::mock::MockTransportFactory::new()),
Arc::new(MockHttpClient),
None,
)
.await;
let pn_jid: Jid = "5511999999999@s.whatsapp.net".parse().unwrap();
let lid_jid: Jid = "236395184570386@lid".parse().unwrap();
client
.lid_pn_cache
.add(&wacore::types::lid_pn::LidPnEntry {
lid: lid_jid.user.as_str().into(),
phone_number: pn_jid.user.as_str().into(),
created_at: 0,
learning_source: wacore::types::lid_pn::LearningSource::Usync,
})
.await;
let lid_with_device: Jid = "236395184570386:5@lid".parse().unwrap();
let swapped = client.swap_pn_lid_namespace(&lid_with_device).await;
let swapped = swapped.expect("should resolve LID→PN");
assert!(swapped.is_pn());
assert_eq!(swapped.user, "5511999999999");
assert_eq!(swapped.device(), 5);
let pn_with_device: Jid = "5511999999999:3@s.whatsapp.net".parse().unwrap();
let swapped = client.swap_pn_lid_namespace(&pn_with_device).await;
let swapped = swapped.expect("should resolve PN→LID");
assert!(swapped.is_lid());
assert_eq!(swapped.user, "236395184570386");
assert_eq!(swapped.device(), 3);
let group: Jid = "120363021033254949@g.us".parse().unwrap();
assert!(client.swap_pn_lid_namespace(&group).await.is_none());
}
#[tokio::test]
async fn alternate_key_lookup_pn_input_server_changed() {
let _ = env_logger::builder().is_test(true).try_init();
let backend = crate::test_utils::create_test_backend().await;
let pm = Arc::new(
PersistenceManager::new(backend)
.await
.expect("persistence manager should initialize"),
);
let mut config = crate::cache_config::CacheConfig::default();
config.recent_messages.capacity = 1_000;
let (client, _sync_rx) = Client::new_with_cache_config(
Arc::new(crate::runtime_impl::TokioRuntime),
pm.clone(),
Arc::new(crate::transport::mock::MockTransportFactory::new()),
Arc::new(MockHttpClient),
None,
config,
)
.await;
let pn_jid: Jid = "5511999999999@s.whatsapp.net".parse().unwrap();
let lid_jid: Jid = "236395184570386@lid".parse().unwrap();
let msg_id = "RETRY_ALT_PN";
let msg = wa::Message {
conversation: Some("pn input alternate".into()),
..Default::default()
};
client.add_recent_message(&pn_jid, msg_id, &msg, None).await;
client
.lid_pn_cache
.add(&wacore::types::lid_pn::LidPnEntry {
lid: lid_jid.user.as_str().into(),
phone_number: pn_jid.user.as_str().into(),
created_at: 0,
learning_source: wacore::types::lid_pn::LearningSource::Usync,
})
.await;
let taken = client.take_recent_message(&pn_jid, msg_id).await;
assert!(
taken.is_some(),
"Should find message via server-changed path"
);
let (msg_out, alt_chat) = taken.unwrap();
let alt_chat = alt_chat.expect("should be alternate hit");
assert!(
alt_chat.is_pn(),
"alternate chat should be PN (the original input)"
);
assert_eq!(alt_chat.user, pn_jid.user);
assert_eq!(msg_out.conversation.as_deref(), Some("pn input alternate"));
}
#[tokio::test]
async fn no_alternate_without_mapping() {
let _ = env_logger::builder().is_test(true).try_init();
let backend = crate::test_utils::create_test_backend().await;
let pm = Arc::new(
PersistenceManager::new(backend)
.await
.expect("persistence manager should initialize"),
);
let mut config = crate::cache_config::CacheConfig::default();
config.recent_messages.capacity = 1_000;
let (client, _sync_rx) = Client::new_with_cache_config(
Arc::new(crate::runtime_impl::TokioRuntime),
pm.clone(),
Arc::new(crate::transport::mock::MockTransportFactory::new()),
Arc::new(MockHttpClient),
None,
config,
)
.await;
let lid_jid: Jid = "236395184570386@lid".parse().unwrap();
let msg_id = "RETRY_NO_ALT";
let msg = wa::Message {
conversation: Some("no alternate".into()),
..Default::default()
};
client
.add_recent_message(&lid_jid, msg_id, &msg, None)
.await;
let taken = client.take_recent_message(&lid_jid, msg_id).await;
assert!(taken.is_some());
let (_, alt_chat) = taken.unwrap();
assert!(alt_chat.is_none(), "primary hit should have no alt_chat");
let missing = client.take_recent_message(&lid_jid, "NONEXISTENT").await;
assert!(missing.is_none(), "non-existent message should return None");
}
#[tokio::test]
async fn alternate_key_both_miss() {
let _ = env_logger::builder().is_test(true).try_init();
let backend = crate::test_utils::create_test_backend().await;
let pm = Arc::new(
PersistenceManager::new(backend)
.await
.expect("persistence manager should initialize"),
);
let mut config = crate::cache_config::CacheConfig::default();
config.recent_messages.capacity = 1_000;
let (client, _sync_rx) = Client::new_with_cache_config(
Arc::new(crate::runtime_impl::TokioRuntime),
pm.clone(),
Arc::new(crate::transport::mock::MockTransportFactory::new()),
Arc::new(MockHttpClient),
None,
config,
)
.await;
let pn_jid: Jid = "5511999999999@s.whatsapp.net".parse().unwrap();
let lid_jid: Jid = "236395184570386@lid".parse().unwrap();
client
.lid_pn_cache
.add(&wacore::types::lid_pn::LidPnEntry {
lid: lid_jid.user.as_str().into(),
phone_number: pn_jid.user.as_str().into(),
created_at: 0,
learning_source: wacore::types::lid_pn::LearningSource::Usync,
})
.await;
let taken = client.take_recent_message(&pn_jid, "MISSING").await;
assert!(taken.is_none(), "both primary and alternate miss → None");
}
#[test]
fn resolve_retry_chat_info_peer_device_with_recipient() {
use wacore_binary::builder::NodeBuilder;
let our_pn: Jid = "5511999999999@s.whatsapp.net".parse().unwrap();
let recipient: Jid = "5522888888888@s.whatsapp.net".parse().unwrap();
let node = NodeBuilder::new("receipt")
.attr("recipient", "5522888888888@s.whatsapp.net")
.build();
let receipt = make_test_receipt("5511999999999:2@s.whatsapp.net");
let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), Some(&our_pn), None);
assert_eq!(info.chat.user, recipient.user);
assert_eq!(info.chat.device(), 0, "chat should be bare");
assert_eq!(info.requester.user, our_pn.user);
assert_eq!(info.requester.device(), 2);
}
#[test]
fn resolve_retry_chat_info_peer_device_without_recipient() {
use wacore_binary::builder::NodeBuilder;
let our_pn: Jid = "5511999999999@s.whatsapp.net".parse().unwrap();
let node = NodeBuilder::new("receipt").build();
let receipt = make_test_receipt("5511999999999:2@s.whatsapp.net");
let info =
maybe_resolve_retry_chat_info(&receipt, &node.as_node_ref(), Some(&our_pn), None);
assert!(info.is_none());
}
#[test]
fn resolve_retry_chat_info_bot_with_recipient() {
use wacore_binary::builder::NodeBuilder;
let node = NodeBuilder::new("receipt")
.attr("recipient", "5522888888888@s.whatsapp.net")
.build();
let receipt = make_test_receipt("131355500001@s.whatsapp.net");
let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
assert!(info.is_bot, "bot JID should be detected");
assert!(
!info.is_fbid_bot_retry,
"legacy PN bots use the regular retry parser"
);
assert_eq!(info.chat.user, "5522888888888");
assert_eq!(info.chat.device(), 0);
}
#[test]
fn resolve_retry_chat_info_bot_without_recipient() {
use wacore_binary::builder::NodeBuilder;
let node = NodeBuilder::new("receipt").build();
let receipt = make_test_receipt("131355500001@s.whatsapp.net");
let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
assert!(info.is_bot);
assert!(!info.is_fbid_bot_retry);
assert_eq!(info.chat.user, "131355500001");
}
#[test]
fn resolve_retry_chat_info_fbid_bot_dm_marks_bot_retry() {
use wacore_binary::builder::NodeBuilder;
let node = NodeBuilder::new("receipt")
.attr("recipient", "5522888888888@s.whatsapp.net")
.build();
let receipt = make_test_receipt("200000000000002@bot");
let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
assert!(info.is_bot);
assert!(info.is_fbid_bot_retry);
assert_eq!(info.chat.user, "5522888888888");
}
#[test]
fn resolve_retry_chat_info_preserves_original_from() {
use wacore_binary::builder::NodeBuilder;
let node = NodeBuilder::new("receipt").build();
let receipt = make_test_receipt("5511999999999:33@s.whatsapp.net");
let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, None);
assert_eq!(info.original_from.device(), 33);
assert_eq!(info.original_from.user, "5511999999999");
assert_eq!(info.chat.device(), 0);
assert_eq!(info.chat.user, "5511999999999");
}
#[test]
fn resolve_retry_chat_info_peer_via_lid() {
use wacore_binary::builder::NodeBuilder;
let our_lid: Jid = "236395184570386@lid".parse().unwrap();
let recipient: Jid = "5522888888888@s.whatsapp.net".parse().unwrap();
let node = NodeBuilder::new("receipt")
.attr("recipient", "5522888888888@s.whatsapp.net")
.build();
let receipt = make_test_receipt("236395184570386:5@lid");
let info = resolve_retry_chat_info(&receipt, &node.as_node_ref(), None, Some(&our_lid));
assert_eq!(info.chat.user, recipient.user);
assert_eq!(info.chat.device(), 0);
assert_eq!(info.requester.device(), 5);
}
}