Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
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
2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,7 @@ portable-atomic = { version = "1", default-features = false, features = ["fallba
prost = { version = "0.14.3", default-features = false, features = ["std"] }
prost-build = { version = "0.14.3", default-features = false }
rand = "0.10"
serde = { version = "1.0", features = ["derive"] }
serde = { version = "1.0", features = ["derive", "rc"] }
serde-big-array = "0.5"
serde_json = { version = "1.0", default-features = false }
sha1 = { version = "0.11.0", default-features = false }
Expand Down
71 changes: 40 additions & 31 deletions src/message.rs
Original file line number Diff line number Diff line change
Expand Up @@ -35,17 +35,19 @@ pub(crate) use wacore::protocol::retry::RetryReason;

impl Client {
/// Dispatches a successfully parsed message to the event bus and sends a delivery receipt.
fn dispatch_parsed_message(self: &Arc<Self>, msg: wa::Message, info: &MessageInfo) {
fn dispatch_parsed_message(self: &Arc<Self>, msg: wa::Message, info: &Arc<MessageInfo>) {
use wacore::proto_helpers::MessageExt;

let mut info = info.clone();
if info.ephemeral_expiration.is_none() {
info.ephemeral_expiration = msg.get_base_message().get_ephemeral_expiration();
let mut info = Arc::clone(info);
if info.ephemeral_expiration.is_none()
&& msg.get_base_message().get_ephemeral_expiration().is_some()
{
Arc::make_mut(&mut info).ephemeral_expiration =
msg.get_base_message().get_ephemeral_expiration();
}

// Send delivery receipt immediately in the background.
let client_clone = self.clone();
let info_for_receipt = info.clone();
let info_for_receipt = Arc::clone(&info);
self.runtime
.spawn(Box::pin(async move {
client_clone.send_delivery_receipt(&info_for_receipt).await;
Expand All @@ -59,7 +61,11 @@ impl Client {

/// Handles a newsletter plaintext message.
/// Newsletters are not E2E encrypted and use the <plaintext> tag directly.
async fn handle_newsletter_message(self: &Arc<Self>, node: &NodeRef<'_>, info: &MessageInfo) {
async fn handle_newsletter_message(
self: &Arc<Self>,
node: &NodeRef<'_>,
info: &Arc<MessageInfo>,
) {
let Some(plaintext_node) = node.get_optional_child_by_tag(&["plaintext"]) else {
log::warn!(
"[msg:{}] Received newsletter message without <plaintext> child: {}",
Expand Down Expand Up @@ -97,12 +103,12 @@ impl Client {
/// * `decrypt_fail_mode` - Whether to show or hide the placeholder (matches WhatsApp Web's `hideFail`)
fn dispatch_undecryptable_event(
&self,
info: &MessageInfo,
info: Arc<MessageInfo>,
decrypt_fail_mode: crate::types::events::DecryptFailMode,
) {
self.core.event_bus.dispatch(Event::UndecryptableMessage(
crate::types::events::UndecryptableMessage {
info: info.clone(),
info,
is_unavailable: false,
unavailable_type: crate::types::events::UnavailableType::Unknown,
decrypt_fail_mode,
Expand All @@ -119,11 +125,11 @@ impl Client {
/// Returns `true` to be assigned to `dispatched_undecryptable` flag.
fn handle_decrypt_failure(
self: &Arc<Self>,
info: &MessageInfo,
info: &Arc<MessageInfo>,
reason: RetryReason,
decrypt_fail_mode: crate::types::events::DecryptFailMode,
) -> bool {
self.dispatch_undecryptable_event(info, decrypt_fail_mode);
self.dispatch_undecryptable_event(Arc::clone(info), decrypt_fail_mode);
self.spawn_retry_receipt(info, reason);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
true
}
Expand Down Expand Up @@ -190,9 +196,9 @@ impl Client {
/// # Arguments
/// * `info` - The message info for the failed message
/// * `reason` - The retry reason code (matches WhatsApp Web's RetryReason enum)
fn spawn_retry_receipt(self: &Arc<Self>, info: &MessageInfo, reason: RetryReason) {
fn spawn_retry_receipt(self: &Arc<Self>, info: &Arc<MessageInfo>, reason: RetryReason) {
let client = Arc::clone(self);
let info = info.clone();
let info = Arc::clone(info);

self.runtime.spawn(Box::pin(async move {
let cache_key = client
Expand Down Expand Up @@ -322,7 +328,7 @@ impl Client {
self.spawn_pdo_request_with_options(&info, true);
self.core.event_bus.dispatch(Event::UndecryptableMessage(
crate::types::events::UndecryptableMessage {
info: (*info).clone(),
info: Arc::clone(&info),
is_unavailable: true,
unavailable_type,
decrypt_fail_mode: crate::types::events::DecryptFailMode::Show,
Expand Down Expand Up @@ -566,7 +572,7 @@ impl Client {
info.id, info.source.sender
);
if !session_dispatched_undecryptable {
self.dispatch_undecryptable_event(&info, decrypt_fail_mode);
self.dispatch_undecryptable_event(Arc::clone(&info), decrypt_fail_mode);
}
}

Expand All @@ -591,7 +597,7 @@ impl Client {
// Dispatch UndecryptableMessage event for messages that failed to decrypt
// (This should not cause double-dispatching since process_session_enc_batch
// already returned dispatched_undecryptable=false for this case)
self.dispatch_undecryptable_event(&info, decrypt_fail_mode);
self.dispatch_undecryptable_event(Arc::clone(&info), decrypt_fail_mode);
// Do NOT send delivery receipt - transport ack is sufficient
}

Expand All @@ -603,7 +609,7 @@ impl Client {
async fn process_session_enc_batch<'n>(
self: Arc<Self>,
enc_nodes: &[&'n NodeRef<'n>],
info: &MessageInfo,
info: &Arc<MessageInfo>,
sender_encryption_jid: &Jid,
decrypt_fail_mode: crate::types::events::DecryptFailMode,
) -> (bool, bool, bool) {
Expand Down Expand Up @@ -993,7 +999,7 @@ impl Client {
async fn process_group_enc_batch<'n>(
self: Arc<Self>,
enc_nodes: &[&'n NodeRef<'n>],
info: &MessageInfo,
info: &Arc<MessageInfo>,
_sender_encryption_jid: &Jid,
decrypt_fail_mode: crate::types::events::DecryptFailMode,
) -> Result<(), DecryptionError> {
Expand Down Expand Up @@ -1094,7 +1100,7 @@ impl Client {
self.handle_unknown_device_sync(info).await;
}

self.dispatch_undecryptable_event(info, decrypt_fail_mode);
self.dispatch_undecryptable_event(Arc::clone(info), decrypt_fail_mode);
self.spawn_retry_receipt(info, retry_reason);
}
Err(e) => {
Expand Down Expand Up @@ -1151,7 +1157,7 @@ impl Client {
enc_type: &str,
padded_plaintext: &[u8],
padding_version: u8,
info: &MessageInfo,
info: &Arc<MessageInfo>,
) -> Result<(), anyhow::Error> {
let original_msg = wacore::messages::decode_plaintext(padded_plaintext, padding_version)?;
log::debug!(
Expand Down Expand Up @@ -1240,7 +1246,7 @@ impl Client {
rng: &mut rand::rngs::StdRng,
enc_type: &str,
padding_version: u8,
info: &MessageInfo,
info: &Arc<MessageInfo>,
) -> bool {
use wacore::libsignal::protocol::{UsePQRatchet, message_decrypt};

Expand Down Expand Up @@ -1666,14 +1672,14 @@ mod tests {
let sender_jid: Jid = "1234567890@s.whatsapp.net"
.parse()
.expect("test JID should be valid");
let info = MessageInfo {
let info = Arc::new(MessageInfo {
source: crate::types::message::MessageSource {
sender: sender_jid.clone(),
chat: sender_jid.clone(),
..Default::default()
},
..Default::default()
};
});

// Create a valid but undecryptable SignalMessage
let dummy_key = [0u8; 32];
Expand Down Expand Up @@ -1747,14 +1753,14 @@ mod tests {
let sender_jid: Jid = "0000000000000@s.whatsapp.net"
.parse()
.expect("test JID should be valid");
let info = MessageInfo {
let info = Arc::new(MessageInfo {
source: crate::types::message::MessageSource {
sender: sender_jid.clone(),
chat: sender_jid.clone(),
..Default::default()
},
..Default::default()
};
});

// Pre-store an empty (degenerate) session record in the signal cache.
// This simulates the bug scenario: record exists but has no usable ratchet state.
Expand Down Expand Up @@ -2721,14 +2727,14 @@ mod tests {
.parse()
.expect("test JID should be valid");

let info = MessageInfo {
let info = Arc::new(MessageInfo {
source: crate::types::message::MessageSource {
sender: sender_jid.clone(),
chat: sender_jid.clone(),
..Default::default()
},
..Default::default()
};
});

log::info!("Test: UntrustedIdentity scenario for {}", sender_jid);

Expand Down Expand Up @@ -2798,14 +2804,14 @@ mod tests {
.parse()
.expect("test JID should be valid");

let info = MessageInfo {
let info = Arc::new(MessageInfo {
source: crate::types::message::MessageSource {
sender: sender_jid.clone(),
chat: sender_jid.clone(),
..Default::default()
},
..Default::default()
};
});

log::info!("Test: Batch processing with multiple error messages");

Expand Down Expand Up @@ -2883,15 +2889,15 @@ mod tests {
.parse()
.expect("test JID should be valid");

let info = MessageInfo {
let info = Arc::new(MessageInfo {
source: crate::types::message::MessageSource {
sender: sender_phone.clone(),
chat: group_jid.clone(),
is_group: true,
..Default::default()
},
..Default::default()
};
});

log::info!("Test: Group context - error handling for {}", sender_phone);

Expand Down Expand Up @@ -4171,6 +4177,7 @@ mod tests {
);

// Call spawn_retry_receipt (this spawns a task, so we need to wait)
let info = Arc::new(info);
client.spawn_retry_receipt(&info, RetryReason::UnknownError);

// Give the spawned task time to execute
Expand Down Expand Up @@ -4205,6 +4212,7 @@ mod tests {
);

// Call spawn_retry_receipt - should NOT increment (already at max)
let info = Arc::new(info);
client.spawn_retry_receipt(&info, RetryReason::UnknownError);

// Give the spawned task time to execute
Expand Down Expand Up @@ -4758,6 +4766,7 @@ mod tests {

// WA Web retries revoked messages the same as any other — the revoke
// protocol message contains the target ID needed to process the deletion
let info = Arc::new(info);
client.spawn_retry_receipt(&info, RetryReason::NoSession);

// Wait for the spawned task to execute
Expand Down
35 changes: 17 additions & 18 deletions src/pdo.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@ use waproto::whatsapp as wa;

#[derive(Clone, Debug)]
pub struct PendingPdoRequest {
pub message_info: MessageInfo,
pub message_info: Arc<MessageInfo>,
pub requested_at: wacore::time::Instant,
}

Expand All @@ -56,7 +56,7 @@ impl Client {
/// * `Err` if we couldn't send the request (e.g., not logged in)
pub async fn send_pdo_placeholder_resend_request(
self: &Arc<Self>,
info: &MessageInfo,
info: &Arc<MessageInfo>,
) -> Result<(), anyhow::Error> {
let device_snapshot = self.persistence_manager.get_device_snapshot().await;

Expand Down Expand Up @@ -104,7 +104,7 @@ impl Client {
}

let pending = PendingPdoRequest {
message_info: info.clone(),
message_info: Arc::clone(info),
requested_at: wacore::time::Instant::now(),
};
self.pdo_pending_requests
Expand Down Expand Up @@ -335,7 +335,7 @@ impl Client {
pending.message_info
} else {
match self.message_info_from_web_message_info(&web_msg_info).await {
Ok(info) => info,
Ok(info) => Arc::new(info),
Err(e) => {
warn!(
"Failed to reconstruct MessageInfo from PDO response: {:?}",
Expand All @@ -346,24 +346,24 @@ impl Client {
}
};

// Extract the actual message content
let Some(message) = web_msg_info.message else {
warn!("PDO response WebMessageInfo missing message content");
return;
};

if message_info.ephemeral_expiration.is_none() {
{
use wacore::proto_helpers::MessageExt;
message_info.ephemeral_expiration =
message.get_base_message().get_ephemeral_expiration();
let mi = Arc::make_mut(&mut message_info);
if mi.ephemeral_expiration.is_none() {
mi.ephemeral_expiration = message.get_base_message().get_ephemeral_expiration();
}
mi.unavailable_request_id = if request_id.is_empty() {
None
} else {
Some(request_id.to_owned())
};
}

message_info.unavailable_request_id = if request_id.is_empty() {
None
} else {
Some(request_id.to_owned())
};

info!(
"Dispatching PDO-recovered message {} from {} via phone (request_id={})",
message_info.id, message_info.source.sender, request_id
Expand Down Expand Up @@ -456,10 +456,9 @@ impl Client {
/// This is used when we've exhausted retry attempts and need immediate PDO recovery.
pub(crate) fn spawn_pdo_request_with_options(
self: &Arc<Self>,
info: &MessageInfo,
info: &Arc<MessageInfo>,
immediate: bool,
) {
// Don't send PDO for our own messages or status broadcasts
if info.source.is_from_me {
return;
}
Expand All @@ -468,7 +467,7 @@ impl Client {
}

let client_clone = Arc::clone(self);
let info_clone = info.clone();
let info_clone = Arc::clone(info);

self.runtime
.spawn(Box::pin(async move {
Expand All @@ -493,7 +492,7 @@ impl Client {

/// Spawns a PDO request for a message that failed to decrypt.
/// This is called alongside the retry receipt to increase chances of recovery.
pub(crate) fn spawn_pdo_request(self: &Arc<Self>, info: &MessageInfo) {
pub(crate) fn spawn_pdo_request(self: &Arc<Self>, info: &Arc<MessageInfo>) {
self.spawn_pdo_request_with_options(info, false);
}
}
Expand Down
6 changes: 3 additions & 3 deletions wacore/src/types/events.rs
Original file line number Diff line number Diff line change
Expand Up @@ -372,7 +372,7 @@ pub enum Event {
QrScannedWithoutMultidevice(QrScannedWithoutMultidevice),
ClientOutdated(ClientOutdated),

Message(Box<wa::Message>, MessageInfo),
Message(Box<wa::Message>, Arc<MessageInfo>),
Receipt(Receipt),
UndecryptableMessage(UndecryptableMessage),
#[serde(skip)]
Expand Down Expand Up @@ -435,7 +435,7 @@ pub enum Event {
impl Event {
pub fn as_message(&self) -> Option<(&wa::Message, &MessageInfo)> {
if let Event::Message(msg, info) = self {
Some((msg, info))
Some((msg, &**info))
} else {
None
}
Expand Down Expand Up @@ -683,7 +683,7 @@ pub enum UnavailableType {

#[derive(Debug, Clone, Serialize)]
pub struct UndecryptableMessage {
pub info: MessageInfo,
pub info: Arc<MessageInfo>,
pub is_unavailable: bool,
pub unavailable_type: UnavailableType,
pub decrypt_fail_mode: DecryptFailMode,
Expand Down
Loading