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
11 changes: 4 additions & 7 deletions src/client/sender_keys.rs
Original file line number Diff line number Diff line change
Expand Up @@ -186,13 +186,12 @@ impl Client {

/// Look up and consume a message by exact `ChatMessageId` (L1 cache then DB).
async fn try_take_by_key(&self, key: &ChatMessageId) -> Option<wa::Message> {
use prost::Message;
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 let Ok(msg) = wa::Message::decode(bytes.as_slice()) {
if let Ok(msg) = waproto::codec::message_decode(bytes.as_slice()) {
// Cache hit — consume the DB row in the background to avoid orphans.
let backend = self.persistence_manager.backend();
let mid = key.id.clone();
Expand All @@ -219,7 +218,7 @@ impl Client {
.take_sent_message(&chat_str, &key.id)
.await
{
Ok(Some(bytes)) => match wa::Message::decode(bytes.as_slice()) {
Ok(Some(bytes)) => match waproto::codec::message_decode(bytes.as_slice()) {
Ok(msg) => Some(msg),
Err(e) => {
log::warn!(
Expand Down Expand Up @@ -283,12 +282,11 @@ impl Client {
/// (capacity 0) or misses; the DB is intentionally not read here so the caller
/// can fall back to the consuming take + re-add path.
async fn peek_by_key(&self, key: &ChatMessageId) -> Option<wa::Message> {
use prost::Message;
if self.cache_config.recent_messages.capacity == 0 {
return None;
}
let bytes = self.recent_messages.get(key).await?;
match wa::Message::decode(bytes.as_slice()) {
match waproto::codec::message_decode(bytes.as_slice()) {
Ok(msg) => Some(msg),
Err(e) => {
log::warn!(
Expand All @@ -307,9 +305,8 @@ impl Client {
/// With L1 cache, the DB write is backgrounded since the cache serves reads immediately.
#[cfg_attr(feature = "tracing", tracing::instrument(name = "wa.session.add_recent_message", level = "debug", skip_all, fields(peer = %to.observe())))]
pub(crate) async fn add_recent_message(&self, to: &Jid, id: &str, msg: &wa::Message) {
use prost::Message;
let key = self.make_chat_message_id(to, id).await;
let bytes = msg.encode_to_vec();
let bytes = waproto::codec::message_to_vec(msg);
let has_l1_cache = self.cache_config.recent_messages.capacity > 0;

if has_l1_cache {
Expand Down
5 changes: 3 additions & 2 deletions src/features/newsletter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@ use wacore::WireEnum;

use crate::client::Client;
use crate::features::mex::{MexError, mex_request};
use prost::Message as ProtoMessage;
use wacore::iq::mex_operations::{
create_newsletter, fetch_all_newsletters_metadata, fetch_newsletter, join_newsletter,
leave_newsletter, update_newsletter, update_newsletter_user_setting,
Expand Down Expand Up @@ -665,7 +664,9 @@ fn parse_newsletter_messages_response(
msg_node
.get_optional_child("plaintext")
.and_then(|pt| match pt.content.as_deref() {
Some(NodeContentRef::Bytes(bytes)) => wa::Message::decode(bytes.as_ref()).ok(),
Some(NodeContentRef::Bytes(bytes)) => {
waproto::codec::message_decode(bytes.as_ref()).ok()
}
_ => None,
});

Expand Down
2 changes: 1 addition & 1 deletion src/message/msg_secret.rs
Original file line number Diff line number Diff line change
Expand Up @@ -689,7 +689,7 @@ impl Client {
},
};

let msg = match wa::Message::decode(plaintext.as_slice()) {
let msg = match waproto::codec::message_decode(plaintext.as_slice()) {
Ok(m) => m,
Err(e) => {
log::warn!(
Expand Down
2 changes: 1 addition & 1 deletion src/message/special.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ impl Client {
};

if let Some(bytes) = plaintext_node.content_bytes() {
match wa::Message::decode(bytes) {
match waproto::codec::message_decode(bytes) {
Ok(msg) => {
log::info!(
"[msg:{}] Received newsletter plaintext message from {}",
Expand Down
16 changes: 8 additions & 8 deletions src/pdo.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,6 @@
use crate::client::Client;
use crate::types::message::MessageInfo;
use log::{debug, info, warn};
use prost::Message;
use std::sync::Arc;
use wacore::types::message::{
ChatMessageId, EditAttribute, MessageCategory, MessageSource, MsgMetaInfo,
Expand Down Expand Up @@ -335,13 +334,14 @@ impl Client {
return;
};

let web_msg_info = match wa::WebMessageInfo::decode(web_message_info_bytes.as_slice()) {
Ok(info) => info,
Err(e) => {
warn!("Failed to decode WebMessageInfo from PDO response: {:?}", e);
return;
}
};
let web_msg_info =
match waproto::codec::web_message_info_decode(web_message_info_bytes.as_slice()) {
Ok(info) => info,
Err(e) => {
warn!("Failed to decode WebMessageInfo from PDO response: {:?}", e);
return;
}
};

let key = &web_msg_info.key;
let remote_jid_str = key.remote_jid.as_deref().unwrap_or("");
Expand Down
10 changes: 6 additions & 4 deletions src/send.rs
Original file line number Diff line number Diff line change
Expand Up @@ -351,7 +351,6 @@ pub(crate) fn build_newsletter_edit_node(
op: NewsletterEdit<'_>,
) -> Node {
use crate::types::message::EditAttribute;
use prost::Message as _;
let mut plaintext = NodeBuilder::new("plaintext");
let (edit, stanza_type, body) = match op {
NewsletterEdit::Edit(m) => {
Expand All @@ -361,7 +360,7 @@ pub(crate) fn build_newsletter_edit_node(
(
EditAttribute::AdminEdit,
wacore::send::stanza_type_from_message(m),
m.encode_to_vec(),
waproto::codec::message_to_vec(m),
)
}
NewsletterEdit::Revoke => (EditAttribute::AdminRevoke, "text", Vec::new()),
Expand Down Expand Up @@ -488,7 +487,6 @@ impl Client {
// Newsletters are not E2E encrypted — send as plaintext via SMAX stanza.
// Matches WA Web's OutMessagePublishNewsletterRequest + ContentType mixins.
if to.is_newsletter() {
use prost::Message as _;
let stanza_type = stanza_type_override
.map(StanzaType::as_wire)
.unwrap_or_else(|| wacore::send::stanza_type_from_message(&message));
Expand All @@ -497,7 +495,11 @@ impl Client {
if let Some(mt) = wacore::send::media_type_from_message(&message) {
plaintext_builder = plaintext_builder.attr("mediatype", mt);
}
let mut children = vec![plaintext_builder.bytes(message.encode_to_vec()).build()];
let mut children = vec![
plaintext_builder
.bytes(waproto::codec::message_to_vec(&message))
.build(),
];
children.extend(meta_node);
children.extend(options.extra_stanza_nodes);
let stanza = NodeBuilder::new("message")
Expand Down
5 changes: 2 additions & 3 deletions wacore/src/comment.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@
//! outer envelope.

use anyhow::{Result, ensure};
use prost::Message;
use waproto::whatsapp as wa;

use crate::secret_enc_addon::{AddonContext, ModificationType, decrypt_addon, encrypt_addon};
Expand Down Expand Up @@ -43,7 +42,7 @@ pub fn encrypt_comment_with_secret(
"message_secret must be {MESSAGE_SECRET_SIZE} bytes, got {}",
message_secret.len()
);
let plaintext = inner.encode_to_vec();
let plaintext = waproto::codec::message_to_vec(inner);
encrypt_addon(
&plaintext,
message_secret,
Expand Down Expand Up @@ -72,7 +71,7 @@ pub fn decrypt_comment_with_secret(
message_secret,
&comment_addon_ctx(parent_msg_id, parent_sender_jid, commenter_jid),
)?;
Ok(wa::Message::decode(&plaintext[..])?)
Ok(waproto::codec::message_decode(&plaintext[..])?)
}

#[cfg(test)]
Expand Down
7 changes: 3 additions & 4 deletions wacore/src/message_edit.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,6 @@
//! that already handle `protocolMessage.editedMessage` can reuse their code.

use anyhow::{Result, anyhow};
use prost::Message;

use crate::secret_enc_addon::{AddonContext, ModificationType, decrypt_addon, encrypt_addon};

Expand Down Expand Up @@ -66,7 +65,7 @@ pub fn encrypt_message_edit(
ctx: &MessageEditContext<'_>,
) -> Result<(Vec<u8>, [u8; IV_SIZE])> {
let mut plaintext = Vec::new();
inner_message.encode(&mut plaintext)?;
waproto::codec::message_encode_into(inner_message, &mut plaintext);
encrypt_addon(&plaintext, message_secret, &ctx.as_addon_ctx())
}

Expand Down Expand Up @@ -100,7 +99,7 @@ pub fn decrypt_secret_encrypted(
modification_type,
};
let plaintext = decrypt_addon(enc_payload, iv, message_secret, &addon)?;
waproto::whatsapp::Message::decode(&plaintext[..])
waproto::codec::message_decode(&plaintext[..])
.map_err(|e| anyhow!("Failed to decode inner secret-encrypted Message: {e}"))
}

Expand Down Expand Up @@ -180,6 +179,7 @@ pub fn decrypt_message_edit_with_fallback(
#[cfg(test)]
mod tests {
use super::*;
use prost::Message as _;
use waproto::whatsapp as wa;

fn make_inner_edit(new_text: &str) -> wa::Message {
Expand Down Expand Up @@ -289,7 +289,6 @@ mod tests {
#[test]
fn general_decrypt_roundtrips_non_edit_use_case() {
use crate::secret_enc_addon::{AddonContext, ModificationType, encrypt_addon};
use prost::Message as _;

// A POLL_EDIT envelope: same shape as MESSAGE_EDIT, different use-case.
let secret = [0x71u8; 32];
Expand Down
64 changes: 38 additions & 26 deletions wacore/src/messages.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,8 +24,8 @@ impl MessageUtils {
/// Encode + pad in a single pre-sized allocation.
pub fn encode_and_pad(msg: &wa::Message) -> Vec<u8> {
let pad = Self::random_pad_len();
let mut buf = Vec::with_capacity(msg.encoded_len() + pad as usize);
msg.encode(&mut buf).expect("encode into pre-sized Vec");
let mut buf = Vec::with_capacity(waproto::codec::message_encoded_len(msg) + pad as usize);
waproto::codec::message_encode_into(msg, &mut buf);
buf.resize(buf.len() + pad as usize, pad);
buf
}
Expand All @@ -47,10 +47,14 @@ impl MessageUtils {
) -> Vec<u8> {
let pad = Self::random_pad_len();
let extra_len = extra_context.map_or(0, |c| {
len_delimited_len(TAG_MESSAGE_CONTEXT_INFO, c.encoded_len())
len_delimited_len(
TAG_MESSAGE_CONTEXT_INFO,
waproto::codec::message_context_info_encoded_len(c),
)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
});
let mut buf = Vec::with_capacity(msg.encoded_len() + extra_len + pad as usize);
msg.encode(&mut buf).expect("encode into pre-sized Vec");
let mut buf =
Vec::with_capacity(waproto::codec::message_encoded_len(msg) + extra_len + pad as usize);
waproto::codec::message_encode_into(msg, &mut buf);
if let Some(c) = extra_context {
push_message_field(TAG_MESSAGE_CONTEXT_INFO, c, &mut buf);
}
Expand Down Expand Up @@ -94,7 +98,7 @@ impl MessageUtils {
let ctx = owned
.message_context_info
.get_or_insert_with(Default::default);
ctx.merge(extra.encode_to_vec().as_slice())
ctx.merge(waproto::codec::message_context_info_to_vec(extra).as_slice())
.expect("merge MessageContextInfo");
}
return Self::encode_dm_plaintexts_owned(owned, destination_jid);
Expand All @@ -107,18 +111,19 @@ impl MessageUtils {
const MAX_PAD: usize = 16;

let mci_field_len = extra_context.map_or(0, |m| {
len_delimited_len(TAG_MESSAGE_CONTEXT_INFO, m.encoded_len())
len_delimited_len(
TAG_MESSAGE_CONTEXT_INFO,
waproto::codec::message_context_info_encoded_len(m),
)
});
let content_len = message.encoded_len();
let content_len = waproto::codec::message_encoded_len(message);
let dest = destination_jid.as_bytes();

// recipient = content (encoded once) + the extra message_context_info field.
// Pre-size for content + the appended mci field + padding so it never
// reallocates; the content bytes are then spliced into the own-device buffer.
let mut recipient = Vec::with_capacity(content_len + mci_field_len + MAX_PAD);
message
.encode(&mut recipient)
.expect("encode into pre-sized Vec");
waproto::codec::message_encode_into(message, &mut recipient);

// own-device plaintext = Message { device_sent_message { destination_jid,
// message }, [message_context_info] }. The DeviceSentMessage length is
Expand Down Expand Up @@ -160,15 +165,16 @@ impl MessageUtils {
// mci struct (not a temp Vec): it is small and encoded straight into each buffer.
let mci = message.message_context_info.take();
let mci_field_len = mci.as_ref().map_or(0, |m| {
len_delimited_len(TAG_MESSAGE_CONTEXT_INFO, m.encoded_len())
len_delimited_len(
TAG_MESSAGE_CONTEXT_INFO,
waproto::codec::message_context_info_encoded_len(m),
)
});
let content_len = message.encoded_len();
let content_len = waproto::codec::message_encoded_len(&message);
let dest = destination_jid.as_bytes();

let mut recipient = Vec::with_capacity(content_len + mci_field_len + MAX_PAD);
message
.encode(&mut recipient)
.expect("encode into pre-sized Vec");
waproto::codec::message_encode_into(&message, &mut recipient);

let dsm_len = len_delimited_len(TAG_DSM_DESTINATION_JID, dest.len())
+ len_delimited_len(TAG_DSM_MESSAGE, content_len);
Expand Down Expand Up @@ -268,7 +274,7 @@ impl MessageUtils {
/// runtime-independent portion of `handle_decrypted_plaintext`.
pub fn decode_plaintext(padded_plaintext: &[u8], padding_version: u8) -> Result<wa::Message> {
let plaintext_slice = MessageUtils::unpad_message_ref(padded_plaintext, padding_version)?;
wa::Message::decode(plaintext_slice)
waproto::codec::message_decode(plaintext_slice)
.map_err(|e| anyhow::anyhow!("Failed to decode decrypted plaintext: {e}"))
}

Expand All @@ -282,12 +288,15 @@ pub struct DmPlaintexts {
pub own_devices: Vec<u8>,
}

// Protobuf field numbers spliced by `encode_dm_plaintexts`. A wrong tag changes
// the decoded result, so the `splice_*` differential tests pin them against prost.
const TAG_DEVICE_SENT_MESSAGE: u64 = 31; // Message.device_sent_message
const TAG_MESSAGE_CONTEXT_INFO: u64 = 35; // Message.message_context_info
const TAG_DSM_DESTINATION_JID: u64 = 1; // DeviceSentMessage.destination_jid
const TAG_DSM_MESSAGE: u64 = 2; // DeviceSentMessage.message
// Protobuf field numbers spliced by `encode_dm_plaintexts`, sourced from the
// generated schema tags so a .proto renumber breaks here at compile time
// instead of silently changing the wire payload. The `splice_*` differential
// tests still pin the hand-written framing itself against prost.
const TAG_DEVICE_SENT_MESSAGE: u64 = waproto::tags::message::DEVICE_SENT_MESSAGE as u64;
const TAG_MESSAGE_CONTEXT_INFO: u64 = waproto::tags::message::MESSAGE_CONTEXT_INFO as u64;
const TAG_DSM_DESTINATION_JID: u64 =
waproto::tags::message::device_sent_message::DESTINATION_JID as u64;
const TAG_DSM_MESSAGE: u64 = waproto::tags::message::device_sent_message::MESSAGE as u64;

/// Append a base-128 varint (protobuf wire format).
#[inline]
Expand Down Expand Up @@ -332,10 +341,13 @@ fn len_delimited_len(field: u64, payload_len: usize) -> usize {
/// straight into `out` (no intermediate `Vec`). Used for the small
/// `message_context_info` field on both plaintexts.
#[inline]
fn push_message_field<M: ProtoMessage>(field: u64, msg: &M, out: &mut Vec<u8>) {
fn push_message_field(field: u64, msg: &wa::MessageContextInfo, out: &mut Vec<u8>) {
push_varint((field << 3) | 2, out);
push_varint(msg.encoded_len() as u64, out);
msg.encode(out).expect("encode into Vec is infallible");
push_varint(
waproto::codec::message_context_info_encoded_len(msg) as u64,
out,
);
waproto::codec::message_context_info_encode_into(msg, out);
}

/// Wrap a message into a DeviceSentMessage for own-device sync, hoisting
Expand Down
4 changes: 2 additions & 2 deletions wacore/src/reporting_token.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,6 @@
use anyhow::{Result, anyhow};
use hkdf::Hkdf;
use hmac::{Hmac, KeyInit, Mac};
use prost::Message;
use sha2::Sha256;
use wacore_binary::Jid;
use wacore_binary::Node;
Expand Down Expand Up @@ -512,7 +511,7 @@ pub fn generate_reporting_token_content(message: &wa::Message) -> Option<Vec<u8>
if !should_include_reporting_token(message) {
return None;
}
let message_bytes = message.encode_to_vec();
let message_bytes = waproto::codec::message_to_vec(message);
extract_reporting_token_content(&message_bytes, REPORTING_FIELDS)
}

Expand Down Expand Up @@ -628,6 +627,7 @@ pub fn extract_message_secret(message: &wa::Message) -> Option<&[u8]> {
#[cfg(test)]
mod tests {
use super::*;
use prost::Message;

#[test]
fn test_generate_message_secret() {
Expand Down
5 changes: 3 additions & 2 deletions wacore/src/types/events.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,6 @@ use crate::types::message::MessageInfo;
use crate::types::presence::{ChatPresence, ChatPresenceMedia, ReceiptType};
use bytes::Bytes;
use chrono::{DateTime, Duration, Utc};
use prost::Message;
use serde::Serialize;
use std::fmt;
use std::sync::{Arc, Mutex, OnceLock, RwLock};
Expand Down Expand Up @@ -131,7 +130,9 @@ impl LazyHistorySync {
// Cheap refcount bump; the lock is released before decoding so a
// concurrent reader isn't blocked by the parse.
let raw = self.locked_raw().clone()?;
wa::HistorySync::decode(&raw[..]).ok().map(Box::new)
waproto::codec::history_sync_decode(&raw[..])
.ok()
.map(Box::new)
});
// Free the raw bytes only AFTER the owned proto is committed, so a
// concurrent clone never sees both gone (raw == None implies parsed set).
Expand Down
Loading
Loading