Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
94 changes: 78 additions & 16 deletions src/client/sender_keys.rs
Original file line number Diff line number Diff line change
Expand Up @@ -63,38 +63,59 @@ impl Client {
}

/// Take a sent message for retry handling. Checks L1 cache first (if enabled),
/// then falls back to DB. Matches WA Web's getMessageTable().get() pattern.
/// then falls back to DB. On miss, tries an alternate PN/LID key to handle
/// mapping changes between send time and retry time (WAWebLidMigrationUtils
/// `getAlternateMsgKey`).
pub(crate) async fn take_recent_message(&self, to: &Jid, id: &str) -> Option<wa::Message> {
let primary_key = self.make_chat_message_id(to, id).await;
if let Some(msg) = self.try_take_by_key(&primary_key).await {
return Some(msg);
}

// Primary miss — try alternate PN<->LID key
if let Some(alt_key) = self.alternate_message_key(&primary_key).await {
log::debug!(
"Primary key miss for {}:{}, trying alternate {}",
primary_key.chat,
id,
alt_key.chat
);
if let Some(msg) = self.try_take_by_key(&alt_key).await {
return Some(msg);
}
}

None
}

/// Look up and consume a message by exact `ChatMessageId` (L1 cache then DB).
async fn try_take_by_key(
&self,
key: &wacore::types::message::ChatMessageId,
) -> Option<wa::Message> {
use prost::Message;
let key = self.make_chat_message_id(to, id).await;
let chat_str = key.chat.to_string();
let has_l1_cache = self.cache_config.recent_messages.capacity > 0;

// L1 cache check (if capacity > 0)
if has_l1_cache && let Some(bytes) = self.recent_messages.remove(&key).await {
if has_l1_cache && let Some(bytes) = self.recent_messages.remove(key).await {
if let Ok(msg) = wa::Message::decode(bytes.as_slice()) {
// Cache hit — consume the DB row in the background to avoid orphans.
// Note: if the background DB write from add_recent_message hasn't completed
// yet, this delete may run first and the write creates an orphan. This is
// harmless — periodic cleanup (sent_message_ttl_secs) purges it. The race
// window is negligible since retry receipts arrive seconds after send.
let backend = self.persistence_manager.backend();
let cs = chat_str.clone();
let mid = key.id.clone();
self.runtime
.spawn(Box::pin(async move {
if let Err(e) = backend.take_sent_message(&cs, &mid).await {
log::warn!("Failed to clean up sent message {cs}:{mid}: {e}");
if let Err(e) = backend.take_sent_message(&chat_str, &mid).await {
log::warn!("Failed to clean up sent message {chat_str}:{mid}: {e}");
}
}))
.detach();
return Some(msg);
}
// Cache decode failed — fall through to DB
log::warn!(
"Failed to decode cached message for {}:{}, trying DB",
to,
id
key.chat,
key.id
);
}

Expand All @@ -108,23 +129,64 @@ impl Client {
Ok(Some(bytes)) => match wa::Message::decode(bytes.as_slice()) {
Ok(msg) => Some(msg),
Err(e) => {
log::warn!("Failed to decode DB message for {}:{}: {}", to, id, e);
log::warn!(
"Failed to decode DB message for {}:{}: {}",
key.chat,
key.id,
e
);
None
}
},
Ok(None) => None,
Err(e) => {
log::warn!(
"Failed to read sent message from DB for {}:{}: {}",
to,
id,
key.chat,
key.id,
e
);
None
}
}
}

/// Compute the alternate PN<->LID message key for retry fallback.
/// Mirrors WAWebLidMigrationUtils `getAlternateMsgKey` for 1x1 chats:
/// swaps the chat JID between PN and LID namespaces.
async fn alternate_message_key(
&self,
key: &wacore::types::message::ChatMessageId,
) -> Option<wacore::types::message::ChatMessageId> {
use wacore::types::message::ChatMessageId;

// Only applies to user chats (not groups/newsletters)
if !key.chat.is_pn() && !key.chat.is_lid() {
return None;
}

let alt_chat = if key.chat.is_lid() {
let pn_user = self.lid_pn_cache.get_phone_number(&key.chat.user).await?;
Jid {
user: pn_user.into(),
server: wacore_binary::Server::Pn,
..Default::default()
}
} else {
let lid_user = self.lid_pn_cache.get_current_lid(&key.chat.user).await?;
Jid {
user: lid_user.into(),
server: wacore_binary::Server::Lid,
..Default::default()
}
};

Some(ChatMessageId {
chat: alt_chat,
id: key.id.clone(),
})
}

/// Store a sent message for retry handling. Always writes to DB; when L1 cache
/// is enabled (capacity > 0) also stores in-memory for fast retrieval.
/// In DB-only mode (capacity = 0), the DB write is awaited to guarantee persistence.
Expand Down
Loading
Loading