use crate::client::Client;
use crate::types::message::MessageInfo;
use log::{debug, info, warn};
use std::sync::Arc;
use wacore::types::message::{
ChatMessageId, EditAttribute, MessageCategory, MessageSource, MsgMetaInfo,
};
use wacore_binary::{Jid, JidExt};
use waproto::whatsapp as wa;
#[derive(Clone, Debug)]
pub struct PendingPdoRequest {
pub message_info: Arc<MessageInfo>,
pub requested_at: wacore::time::Instant,
}
fn self_peer_target(device: &wacore::store::Device) -> Result<Jid, crate::client::ClientError> {
if let Some(lid) = device.lid.as_ref() {
return Ok(Jid::lid(lid.user.clone()));
}
let pn = device
.pn
.as_ref()
.ok_or(crate::client::ClientError::NotLoggedIn)?;
Ok(Jid::pn(pn.user.clone()))
}
impl Client {
#[cfg_attr(feature = "tracing", tracing::instrument(name = "wa.pdo.placeholder_resend", level = "debug", skip_all, fields(chat = %info.source.chat.observe(), sender = %info.source.sender.observe(), msg_id = %info.id), err(Debug)))]
pub async fn send_pdo_placeholder_resend_request(
self: &Arc<Self>,
info: &Arc<MessageInfo>,
) -> Result<(), anyhow::Error> {
let device_snapshot = self.persistence_manager.get_device_snapshot();
let peer_target = self_peer_target(&device_snapshot)?;
let resolved_jid = self.resolve_encryption_jid(&info.source.chat).await;
let participant = if !info.source.is_from_me
&& (info.source.is_group || info.source.chat.server == wacore_binary::Server::Broadcast)
{
Some(self.resolve_encryption_jid(&info.source.sender).await)
} else {
None
};
let cache_chat = if !info.source.is_group && info.source.chat.is_lid() {
info.source
.sender_alt
.as_ref()
.map(|jid| jid.to_non_ad())
.unwrap_or_else(|| info.source.chat.clone())
} else {
info.source.chat.clone()
};
let cache_key = ChatMessageId::new(cache_chat, info.id.clone());
let claimed = Arc::new(std::sync::atomic::AtomicBool::new(false));
let claimed_clone = claimed.clone();
self.pdo_requested
.get_with(cache_key.clone(), async move {
claimed_clone.store(true, std::sync::atomic::Ordering::Release);
})
.await;
if !claimed.load(std::sync::atomic::Ordering::Acquire) {
debug!(
"PDO request already sent for message {} from {}; not re-requesting",
info.id,
info.source.sender.observe()
);
return Ok(());
}
if self.pdo_pending_requests.get(&cache_key).await.is_some() {
debug!(
"PDO request already pending for message {} from {}",
info.id,
info.source.sender.observe()
);
return Ok(());
}
let pending = PendingPdoRequest {
message_info: Arc::clone(info),
requested_at: wacore::time::Instant::now(),
};
self.pdo_pending_requests
.insert(cache_key.clone(), pending)
.await;
let message_key = wa::MessageKey {
remote_jid: Some(resolved_jid.to_string()),
from_me: Some(info.source.is_from_me),
id: Some(info.id.clone()),
participant: participant.map(|p| p.to_string()),
};
let pdo_request = wa::message::PeerDataOperationRequestMessage {
peer_data_operation_request_type: Some(
wa::message::PeerDataOperationRequestType::PLACEHOLDER_MESSAGE_RESEND,
),
placeholder_message_resend_request: vec![
wa::message::peer_data_operation_request_message::PlaceholderMessageResendRequest {
message_key: buffa::MessageField::some(message_key),
},
],
..Default::default()
};
let protocol_message = wa::message::ProtocolMessage {
r#type: Some(wa::message::protocol_message::Type::PEER_DATA_OPERATION_REQUEST_MESSAGE),
peer_data_operation_request_message: buffa::MessageField::some(pdo_request),
..Default::default()
};
let msg = wa::Message {
protocol_message: buffa::MessageField::some(protocol_message),
..Default::default()
};
info!(
"Sending PDO placeholder resend request for message {} from {} in {} to {}",
info.id,
info.source.sender.observe(),
info.source.chat.observe(),
peer_target.observe()
);
if let Err(e) = self
.ensure_e2e_sessions(std::slice::from_ref(&peer_target))
.await
{
self.pdo_pending_requests.remove(&cache_key).await;
self.pdo_requested.remove(&cache_key).await;
return Err(e);
}
if let Err(e) = self.send_peer_message(peer_target, &msg).await {
self.pdo_pending_requests.remove(&cache_key).await;
self.pdo_requested.remove(&cache_key).await;
warn!(
"Failed to send PDO request for message {}: {:?}",
info.id, e
);
return Err(e);
}
debug!("PDO request sent successfully for message {}", info.id);
Ok(())
}
#[cfg_attr(feature = "tracing", tracing::instrument(name = "wa.pdo.fetch_history", level = "debug", skip_all, fields(chat = %chat_jid.observe(), count), err(Debug)))]
pub async fn fetch_message_history(
self: &Arc<Self>,
chat_jid: &Jid,
oldest_msg_id: &str,
oldest_msg_from_me: bool,
oldest_msg_timestamp_ms: i64,
count: i32,
) -> Result<String, anyhow::Error> {
let device_snapshot = self.persistence_manager.get_device_snapshot();
let peer_target = self_peer_target(&device_snapshot)?;
let pdo_request = wa::message::PeerDataOperationRequestMessage {
peer_data_operation_request_type: Some(
wa::message::PeerDataOperationRequestType::HISTORY_SYNC_ON_DEMAND,
),
history_sync_on_demand_request: buffa::MessageField::some(
wa::message::peer_data_operation_request_message::HistorySyncOnDemandRequest {
chat_jid: Some(chat_jid.to_string()),
oldest_msg_id: Some(oldest_msg_id.to_string()),
oldest_msg_from_me: Some(oldest_msg_from_me),
oldest_msg_timestamp_ms: Some(oldest_msg_timestamp_ms),
on_demand_msg_count: Some(count),
..Default::default()
},
),
..Default::default()
};
let protocol_message = wa::message::ProtocolMessage {
r#type: Some(wa::message::protocol_message::Type::PEER_DATA_OPERATION_REQUEST_MESSAGE),
peer_data_operation_request_message: buffa::MessageField::some(pdo_request),
..Default::default()
};
let msg = wa::Message {
protocol_message: buffa::MessageField::some(protocol_message),
..Default::default()
};
info!(
"Sending PDO history sync on-demand request for chat {} (count={}) to {}",
chat_jid.observe(),
count,
peer_target.observe()
);
self.ensure_e2e_sessions(std::slice::from_ref(&peer_target))
.await?;
self.send_peer_message(peer_target, &msg).await
}
#[cfg_attr(feature = "tracing", tracing::instrument(name = "wa.pdo.send_peer_message", level = "debug", skip_all, fields(to = %to.observe()), err(Debug)))]
async fn send_peer_message(
self: &Arc<Self>,
to: Jid,
msg: &wa::Message,
) -> Result<String, anyhow::Error> {
let msg_id = self.generate_message_id();
self.send_message_impl(
to,
msg,
crate::send::SendPipelineOptions {
request_id: Some(&msg_id),
peer: true,
..Default::default()
},
)
.await?;
Ok(msg_id)
}
#[cfg_attr(feature = "tracing", tracing::instrument(name = "wa.pdo.handle_response", level = "debug", skip_all, fields(sender = %pdo_msg_info.source.sender.observe())))]
pub async fn handle_pdo_response(
self: &Arc<Self>,
response: &wa::message::PeerDataOperationRequestResponseMessage,
pdo_msg_info: &MessageInfo,
) {
if pdo_msg_info.source.sender.device != 0 {
debug!(
"Ignoring PDO response from non-primary device {}",
pdo_msg_info.source.sender.observe()
);
return;
}
let request_id = response.stanza_id.as_deref().unwrap_or("");
debug!(
"Received PDO response (request_id={}) with {} results",
request_id,
response.peer_data_operation_result.len()
);
for result in &response.peer_data_operation_result {
if let Some(placeholder_response) =
result.placeholder_message_resend_response.as_option()
{
self.handle_placeholder_resend_response(placeholder_response, request_id)
.await;
}
}
}
async fn handle_placeholder_resend_response(
self: &Arc<Self>,
response: &wa::message::peer_data_operation_request_response_message::peer_data_operation_result::PlaceholderMessageResendResponse,
request_id: &str,
) {
let Some(web_message_info_bytes) = &response.web_message_info_bytes else {
warn!("PDO placeholder response missing webMessageInfoBytes");
return;
};
let mut web_msg_info = match waproto::codec::web_message_info_decode(web_message_info_bytes)
{
Ok(info) => info,
Err(e) => {
warn!("Failed to decode WebMessageInfo from PDO response: {:?}", e);
return;
}
};
let Some(key) = web_msg_info.key.as_option() else {
warn!("PDO response WebMessageInfo missing key");
return;
};
let remote_jid_str = key.remote_jid.as_deref().unwrap_or("");
let msg_id = key.id.as_deref().unwrap_or("");
let cache_key = match remote_jid_str.parse::<Jid>() {
Ok(jid) => ChatMessageId::new(jid, msg_id.to_owned()),
Err(_) => {
warn!(
"PDO response has unparseable remote_jid: {}",
remote_jid_str
);
return;
}
};
let pending = self.pdo_pending_requests.remove(&cache_key).await;
let elapsed = pending
.as_ref()
.map(|p| p.requested_at.elapsed().as_millis())
.unwrap_or(0);
info!(
"Received PDO placeholder response for message {} (took {}ms)",
msg_id, elapsed
);
let mut message_info = if let Some(pending) = pending {
pending.message_info
} else {
match self.message_info_from_web_message_info(&web_msg_info).await {
Ok(info) => Arc::new(info),
Err(e) => {
warn!(
"Failed to reconstruct MessageInfo from PDO response: {:?}",
e
);
return;
}
}
};
let Some(message) = web_msg_info.message.take() else {
info!("PDO response WebMessageInfo missing message content");
return;
};
{
use wacore::proto_helpers::MessageExt;
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())
};
}
info!(
"Dispatching PDO-recovered message {} from {} via phone (request_id={})",
message_info.id,
message_info.source.sender.observe(),
request_id
);
self.core
.event_bus
.dispatch(wacore::types::events::Event::Messages(
wacore::types::events::MessageBatch::builder()
.messages(Arc::from([wacore::types::events::InboundMessage::builder(
)
.message(Arc::from(message))
.info(message_info)
.build()]))
.origin(wacore::types::events::BatchOrigin::Live)
.build(),
));
}
async fn message_info_from_web_message_info(
&self,
web_msg: &wa::WebMessageInfo,
) -> Result<MessageInfo, anyhow::Error> {
let Some(key) = web_msg.key.as_option() else {
anyhow::bail!("WebMessageInfo missing key");
};
self.message_info_from_web_message_parts(
key.remote_jid.as_deref(),
key.from_me,
key.id.as_deref(),
key.participant.as_deref(),
web_msg.message_timestamp,
web_msg.push_name.as_deref(),
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn message_info_from_web_message_parts(
&self,
remote_jid: Option<&str>,
from_me: Option<bool>,
id: Option<&str>,
participant: Option<&str>,
message_timestamp: Option<u64>,
push_name: Option<&str>,
) -> Result<MessageInfo, anyhow::Error> {
let remote_jid: Jid = remote_jid
.ok_or_else(|| anyhow::anyhow!("MessageKey missing remoteJid"))?
.parse()?;
let is_group = remote_jid.is_group();
let is_from_me = from_me.unwrap_or(false);
let sender = if let Some(p) = participant {
p.parse()?
} else if is_from_me {
self.persistence_manager
.get_device_snapshot()
.pn
.clone()
.unwrap_or_else(|| remote_jid.clone())
} else {
remote_jid.clone()
};
let timestamp = message_timestamp
.map(|ts| wacore::time::from_secs_or_now(ts as i64))
.unwrap_or_else(wacore::time::now_utc);
Ok(MessageInfo {
id: id.unwrap_or_default().to_owned(),
server_id: 0,
r#type: String::new(),
source: MessageSource {
chat: remote_jid,
sender,
sender_alt: None,
recipient_alt: None,
is_from_me,
is_group,
addressing_mode: None,
broadcast_list_owner: None,
recipient: None,
},
timestamp,
push_name: push_name.unwrap_or_default().to_owned(),
category: MessageCategory::default(),
multicast: false,
media_type: String::new(),
edit: EditAttribute::default(),
bot_info: None,
meta_info: MsgMetaInfo::default(),
verified_name: None,
device_sent_meta: None,
ephemeral_expiration: None,
is_offline: false,
unavailable_request_id: None,
server_timestamp_us: None,
verified_level: None,
verified_name_serial: None,
peer_recipient_pn: None,
comment_target: None,
bcl_participants: Vec::new(),
})
}
#[cfg_attr(feature = "tracing", tracing::instrument(name = "wa.pdo.run_request", level = "debug", skip_all, fields(chat = %info.source.chat.observe(), sender = %info.source.sender.observe(), msg_id = %info.id)))]
pub(crate) async fn run_pdo_request(self: &Arc<Self>, info: &Arc<MessageInfo>) -> bool {
const PDO_MAX_AGE_SECS: i64 = 14 * 24 * 60 * 60;
let age_secs = wacore::time::now_secs() - info.timestamp.timestamp();
if age_secs > PDO_MAX_AGE_SECS {
debug!(
"PDO request skipped for message {} (age {age_secs}s exceeds {PDO_MAX_AGE_SECS}s limit)",
info.id,
);
return true;
}
match self.send_pdo_placeholder_resend_request(info).await {
Ok(()) => true,
Err(e) => {
warn!(
"Failed to send PDO request for message {} from {}: {:?}",
info.id,
info.source.sender.observe(),
e
);
false
}
}
}
}
#[cfg(test)]
#[allow(clippy::disallowed_methods)]
mod tests {
use super::self_peer_target;
use wacore::store::Device;
use wacore_binary::{Jid, JidExt, Server};
fn empty_device() -> Device {
Device {
pn: None,
lid: None,
..Device::default()
}
}
#[test]
fn self_peer_target_prefers_lid_when_present() {
let mut device = empty_device();
device.pn = Some(Jid::pn_device("559999999999", 33));
device.lid = Some(Jid::lid_device("111111111111111", 33));
let target = self_peer_target(&device).expect("LID present");
assert_eq!(target.user, "111111111111111");
assert_eq!(target.server, Server::Lid);
assert_eq!(target.device, 0);
assert!(!target.is_ad());
}
#[test]
fn self_peer_target_falls_back_to_pn_without_lid() {
let mut device = empty_device();
device.pn = Some(Jid::pn_device("559999999999", 33));
let target = self_peer_target(&device).expect("PN present");
assert_eq!(target.user, "559999999999");
assert_eq!(target.server, Server::Pn);
assert_eq!(target.device, 0);
}
#[test]
fn self_peer_target_errors_when_no_identity_known() {
let device = empty_device();
assert!(
matches!(
self_peer_target(&device),
Err(crate::client::ClientError::NotLoggedIn)
),
"must require either PN or LID"
);
}
async fn setup_reconstruct_client() -> std::sync::Arc<crate::client::Client> {
use crate::test_utils::{MockHttpClient, create_test_backend};
use crate::{
client::Client, runtime_impl::TokioRuntime,
store::persistence_manager::PersistenceManager, transport::mock::MockTransportFactory,
};
use std::sync::Arc;
let backend = create_test_backend().await;
let pm = Arc::new(PersistenceManager::new(backend).await.unwrap());
let (client, _rx) = Client::new(
Arc::new(TokioRuntime),
pm,
Arc::new(MockTransportFactory::new()),
Arc::new(MockHttpClient),
None,
)
.await;
client
}
fn make_web_msg(
remote_jid: &str,
from_me: bool,
id: &str,
participant: Option<&str>,
) -> waproto::whatsapp::WebMessageInfo {
use waproto::whatsapp as wa;
wa::WebMessageInfo {
key: buffa::MessageField::some(wa::MessageKey {
remote_jid: Some(remote_jid.into()),
from_me: Some(from_me),
id: Some(id.into()),
participant: participant.map(|p| p.into()),
}),
..Default::default()
}
}
#[tokio::test]
async fn test_reconstruct_prefers_participant_for_status_broadcast() {
let client = setup_reconstruct_client().await;
let author_jid = "203040904720543@lid";
let web_msg = make_web_msg("status@broadcast", false, "STATUS_PDO_1", Some(author_jid));
let info = client
.message_info_from_web_message_info(&web_msg)
.await
.unwrap();
assert_eq!(info.source.chat.to_string(), "status@broadcast");
assert_eq!(info.source.sender.to_string(), author_jid);
}
#[tokio::test]
async fn test_reconstruct_dm_falls_back_to_remote_jid() {
let client = setup_reconstruct_client().await;
let peer = "5511999998888@s.whatsapp.net";
let web_msg = make_web_msg(peer, false, "DM_PDO_1", None);
let info = client
.message_info_from_web_message_info(&web_msg)
.await
.unwrap();
assert_eq!(info.source.chat.to_string(), peer);
assert_eq!(info.source.sender.to_string(), peer);
}
#[tokio::test]
async fn test_reconstruct_from_web_message_info_view() {
use buffa::Message as _;
use waproto::whatsapp as wa;
let client = setup_reconstruct_client().await;
let author_jid = "203040904720543@lid";
let mut web_msg = make_web_msg(
"status@broadcast",
false,
"STATUS_PDO_VIEW_1",
Some(author_jid),
);
web_msg.push_name = Some("Recovered Sender".to_string());
web_msg.message_timestamp = Some(1_700_000_000);
let encoded = web_msg.encode_to_vec();
let decoded = wa::WebMessageInfo::decode_from_slice(&encoded).expect("should decode");
let info = client
.message_info_from_web_message_info(&decoded)
.await
.unwrap();
assert_eq!(info.id, "STATUS_PDO_VIEW_1");
assert_eq!(info.source.chat.to_string(), "status@broadcast");
assert_eq!(info.source.sender.to_string(), author_jid);
assert_eq!(info.push_name, "Recovered Sender");
}
#[tokio::test]
async fn test_reconstruct_lid_migrated_dm_uses_lid_remote() {
let client = setup_reconstruct_client().await;
let peer_lid = "236395184570386@lid";
let web_msg = make_web_msg(peer_lid, false, "LID_DM_PDO_1", None);
let info = client
.message_info_from_web_message_info(&web_msg)
.await
.unwrap();
assert_eq!(info.source.chat.to_string(), peer_lid);
assert_eq!(info.source.sender.to_string(), peer_lid);
assert!(!info.source.is_group);
assert!(!info.source.is_from_me);
}
#[tokio::test]
async fn test_reconstruct_lid_migrated_dm_from_me_uses_own_pn() {
let client = setup_reconstruct_client().await;
let peer_lid = "236395184570386@lid";
let web_msg = make_web_msg(peer_lid, true, "LID_DM_FROM_ME_1", None);
let info = client
.message_info_from_web_message_info(&web_msg)
.await
.unwrap();
assert_eq!(info.source.chat.to_string(), peer_lid);
assert!(info.source.is_from_me);
}
fn make_group_message_info(
chat: &str,
sender: &str,
id: &str,
) -> std::sync::Arc<wacore::types::message::MessageInfo> {
use wacore::types::message::{MessageInfo, MessageSource};
std::sync::Arc::new(MessageInfo {
id: id.to_owned(),
source: MessageSource {
chat: chat.parse().expect("chat jid"),
sender: sender.parse().expect("sender jid"),
is_group: true,
..Default::default()
},
timestamp: wacore::time::now_utc(),
..Default::default()
})
}
async fn set_own_pn(client: &std::sync::Arc<crate::client::Client>) {
client
.persistence_manager
.process_command(crate::store::commands::DeviceCommand::SetId(Some(
"5511777776666:2@s.whatsapp.net".parse().expect("own jid"),
)))
.await;
}
#[tokio::test]
async fn pdo_request_skipped_when_already_requested() {
use wacore::types::message::ChatMessageId;
let client = setup_reconstruct_client().await;
set_own_pn(&client).await;
let info = make_group_message_info(
"120363000000000001@g.us",
"203040904720543@lid",
"PDO_ONCE_1",
);
let key = ChatMessageId::new(info.source.chat.clone(), info.id.clone());
client.pdo_requested.insert(key.clone(), ()).await;
let res = client.send_pdo_placeholder_resend_request(&info).await;
assert!(res.is_ok(), "gated path reports success: {res:?}");
assert!(
client.pdo_pending_requests.get(&key).await.is_none(),
"gated request must not register a pending entry"
);
}
#[tokio::test]
async fn pdo_request_failure_releases_once_per_message_slot() {
use wacore::types::message::ChatMessageId;
let client = setup_reconstruct_client().await;
set_own_pn(&client).await;
client
.offline_sync_completed
.store(true, std::sync::atomic::Ordering::Relaxed);
let info = make_group_message_info(
"120363000000000001@g.us",
"203040904720543@lid",
"PDO_ONCE_2",
);
let key = ChatMessageId::new(info.source.chat.clone(), info.id.clone());
let res = tokio::time::timeout(
std::time::Duration::from_secs(15),
client.send_pdo_placeholder_resend_request(&info),
)
.await
.expect("send attempt must resolve fast without a live transport");
assert!(res.is_err(), "no live transport, the send must fail");
assert!(
client.pdo_requested.get(&key).await.is_none(),
"failed send must release the once-per-message slot"
);
assert!(
client.pdo_pending_requests.get(&key).await.is_none(),
"failed send must clear the pending entry"
);
}
#[tokio::test]
async fn pdo_missing_content_response_clears_pending_but_keeps_memo() {
use buffa::Message as _;
use wacore::types::message::ChatMessageId;
let client = setup_reconstruct_client().await;
let chat = "5511999998888@s.whatsapp.net";
let msg_id = "PDO_ONCE_3";
let key = ChatMessageId::new(chat.parse().expect("chat jid"), msg_id.to_owned());
client.pdo_requested.insert(key.clone(), ()).await;
client
.pdo_pending_requests
.insert(
key.clone(),
super::PendingPdoRequest {
message_info: make_group_message_info(chat, chat, msg_id),
requested_at: wacore::time::Instant::now(),
},
)
.await;
let web_msg = waproto::whatsapp::WebMessageInfo {
key: buffa::MessageField::some(waproto::whatsapp::MessageKey {
remote_jid: Some(chat.to_owned()),
from_me: Some(false),
id: Some(msg_id.to_owned()),
participant: None,
}),
..Default::default()
};
let response = waproto::whatsapp::message::peer_data_operation_request_response_message::peer_data_operation_result::PlaceholderMessageResendResponse {
web_message_info_bytes: Some(web_msg.encode_to_vec()),
};
client
.handle_placeholder_resend_response(&response, "req-1")
.await;
assert!(
client.pdo_pending_requests.get(&key).await.is_none(),
"response consumes the pending slot"
);
assert!(
client.pdo_requested.get(&key).await.is_some(),
"memo must survive a content-less response"
);
}
}