Skip to content
Merged
Show file tree
Hide file tree
Changes from 4 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
76 changes: 51 additions & 25 deletions src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1044,56 +1044,81 @@ fn build_pong(to: String, id: Option<&str>) -> wacore_binary::Node {
builder.build()
}

/// Compare decoded attribute values by their wire display without allocating.
#[inline]
fn value_refs_display_equal(
left: &wacore_binary::node::ValueRef<'_>,
right: &wacore_binary::node::ValueRef<'_>,
) -> bool {
use wacore_binary::node::ValueRef;

match (left, right) {
(ValueRef::String(left), ValueRef::String(right)) => left == right,
(ValueRef::Jid(left), ValueRef::Jid(right)) => left.display_eq_jid(right),
(ValueRef::String(left), ValueRef::Jid(right)) => right.display_eq(left),
(ValueRef::Jid(left), ValueRef::String(right)) => left.display_eq(right),
}
}

Comment thread
jlucaso1 marked this conversation as resolved.
/// Build an `<ack/>` for the given stanza, matching WA Web / whatsmeow behavior:
///
/// - `class` = original stanza tag
/// - `id`, `to` (flipped from `from`), `participant` copied from original
/// - `from` = own device PN, only for message acks
/// - `type` echoed for non-message stanzas (whatsmeow: `node.Tag != "message"`),
/// except `notification type="encrypt"` with `<identity/>` child (WA Web drops type there).
/// - `type` echoed when present, except `notification type="encrypt"` with
/// an `<identity/>` child
///
/// For receipt acks, WA Web uses `MAYBE_CUSTOM_STRING(ackString)` where
/// `ackString = maybeAttrString("type")` — so `type` is only included when
/// explicitly present on the incoming receipt (delivery receipts normally
/// have no type attribute, meaning the ack also has no type).
///
/// Encode an ack stanza directly to bytes, bypassing Node + marshal_auto.
/// Acks are the most frequent outbound stanza (~1 per inbound message).
fn encode_ack_bytes(
node: &wacore_binary::NodeRef<'_>,
own_device_pn: Option<&Jid>,
) -> Result<Option<Vec<u8>>, wacore_binary::error::BinaryError> {
) -> Result<Vec<u8>, crate::features::StanzaResponseError> {
use wacore_binary::encoder::{ByteWriter, EncodeNode, Encoder};

let Some(id_val) = node.get_attr("id") else {
return Ok(None);
};
let Some(from_val) = node.get_attr("from") else {
return Ok(None);
};
let id_val = node
.get_attr("id")
.ok_or(crate::features::StanzaResponseError::MissingAttribute("id"))?;
Comment thread
jlucaso1 marked this conversation as resolved.
Outdated
let from_val =
node.get_attr("from")
.ok_or(crate::features::StanzaResponseError::MissingAttribute(
"from",
))?;
let tag = node.tag.as_ref();
// WAWebReceiptAck: `participant: r && r !== e ? DEVICE_JID(r) : DROP_ATTR`.
// Drop the attribute when it would duplicate `to` (which is the flipped `from`).
let participant_val = node.get_attr("participant").filter(|p| {
let p_str = p.as_str();
let from_str = from_val.as_str();
p_str.as_ref() != from_str.as_ref()
// This is specific to the specialized receipt ACK. Generic ACKs and NACKs
// preserve participant even when it duplicates the destination.
let participant_val = node.get_attr("participant").filter(|participant| {
if tag != "receipt" {
return true;
}
!value_refs_display_equal(participant, from_val)
});
// Server expects `recipient` echoed back so it can route the ack to the
// origin companion/device (hosted-companion, peer, LID-routed stanzas).
// Dropping it makes the server close the stream with `<stream:error><ack/>`.
let recipient_val = node.get_attr("recipient");
let tag = node.tag.as_ref();

let typ_val = if tag != "message" && !is_encrypt_identity_notification(node) {
let typ_val = if !is_encrypt_identity_notification(node) {
node.get_attr("type")
Comment thread
jlucaso1 marked this conversation as resolved.
} else {
None
};

let include_from = tag == "message" && own_device_pn.is_some();
let own_device_pn = if tag == "message" {
Some(own_device_pn.ok_or(crate::features::StanzaResponseError::MissingLocalIdentity)?)
} else {
None
};

// Count attrs: class + id + to + optional(from, participant, recipient, type)
let attr_count = 3
+ usize::from(include_from)
+ usize::from(own_device_pn.is_some())
+ usize::from(participant_val.is_some())
+ usize::from(recipient_val.is_some())
+ usize::from(typ_val.is_some());
Expand Down Expand Up @@ -1161,15 +1186,15 @@ fn encode_ack_bytes(
participant: participant_val,
recipient: recipient_val,
typ: typ_val,
own_pn: if include_from { own_device_pn } else { None },
own_pn: own_device_pn,
tag_str: tag,
attr_count,
};

let mut buf = Vec::with_capacity(64);
let mut encoder = Encoder::new_vec(&mut buf)?;
encoder.write_node(&ack)?;
Ok(Some(buf))
Ok(buf)
}

/// Minimal `<message>` stanza carrying the attrs `encode_ack_bytes` needs,
Expand Down Expand Up @@ -1201,14 +1226,15 @@ fn build_ack_node(node: &wacore_binary::NodeRef<'_>, own_device_pn: Option<&Jid>
let id = node.get_attr("id")?.to_node_value();
let from_ref = node.get_attr("from")?;
let from = from_ref.to_node_value();
// Drop participant when it duplicates `to` (the flipped `from`).
let tag = node.tag.as_ref();
// Only the specialized receipt ACK drops a participant that duplicates
// `to` (the flipped `from`).
let participant = node
.get_attr("participant")
.filter(|p| p.as_str().as_ref() != from_ref.as_str().as_ref())
.filter(|participant| tag != "receipt" || !value_refs_display_equal(participant, from_ref))
.map(|v| v.to_node_value());
let recipient = node.get_attr("recipient").map(|v| v.to_node_value());
let tag = node.tag.as_ref();
let typ = if tag != "message" && !is_encrypt_identity_notification(node) {
let typ = if !is_encrypt_identity_notification(node) {
node.get_attr("type").map(|v| v.to_node_value())
} else {
None
Expand Down Expand Up @@ -1243,7 +1269,7 @@ fn is_encrypt_identity_notification(node: &wacore_binary::NodeRef<'_>) -> bool {
node.tag == "notification"
&& node
.get_attr("type")
.is_some_and(|v| v.as_str() == "encrypt")
.is_some_and(|value| value == "encrypt")
&& node.get_optional_child("identity").is_some()
}

Expand Down
5 changes: 2 additions & 3 deletions src/client/lifecycle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -990,9 +990,8 @@ impl Client {
self.is_connected.load(Ordering::Acquire)
}

/// Force the connected flag (tests only): the facade's connect path now gates on `is_connected`,
/// so a unit test driving `spawn_call`/`place_call` must mark the client connected first.
#[cfg(all(test, feature = "voip-runtime"))]
/// Force the connected flag for tests that exercise connected-only operations.
#[cfg(test)]
pub(crate) fn set_connected_for_test(&self, connected: bool) {
self.is_connected.store(connected, Ordering::Release);
}
Expand Down
48 changes: 36 additions & 12 deletions src/client/node_io.rs
Original file line number Diff line number Diff line change
Expand Up @@ -589,6 +589,17 @@ impl Client {
}
}

#[inline]
fn encode_ack_from_snapshot(
&self,
node: &wacore_binary::NodeRef<'_>,
) -> Result<Vec<u8>, crate::features::StanzaResponseError> {
let device = self.persistence_manager.get_device_snapshot();
let encoded = encode_ack_bytes(node, device.pn.as_ref());
drop(device);
encoded
}

/// Build and send an <ack/> node corresponding to the given stanza.
#[cfg_attr(
feature = "tracing",
Expand All @@ -604,10 +615,8 @@ impl Client {
if !self.is_connected() {
return Err(ClientError::NotConnected);
}
let own_pn = self.get_pn();
let buf = match encode_ack_bytes(node, own_pn.as_ref()) {
Ok(Some(buf)) => buf,
Ok(None) => return Ok(()),
let buf = match self.encode_ack_from_snapshot(node) {
Ok(buf) => buf,
Err(e) => {
log::warn!("Failed to encode ack: {e}");
return Ok(());
Expand All @@ -616,21 +625,38 @@ impl Client {
self.send_raw_bytes(buf).await
}

/// Confirm a received stanza using its original borrowed node.
///
/// Unlike the tolerant automatic receive path, malformed input is returned
/// to the caller and no successful outcome is reported unless the response
/// reaches the transport.
#[cfg_attr(
feature = "tracing",
tracing::instrument(name = "wa.conn.ack_explicit", level = "debug", skip_all, err(Debug))
)]
pub async fn acknowledge_stanza(
&self,
stanza: &wacore_binary::NodeRef<'_>,
) -> Result<(), crate::features::StanzaResponseError> {
let bytes = self.encode_ack_from_snapshot(stanza)?;
self.send_raw_bytes(bytes).await?;
Ok(())
}

/// Send a transport ack so the server stops replaying a stanza from the
/// offline queue. Awaitable so callers can order it after a retry receipt
/// in a single flushed task.
pub(crate) async fn send_transport_ack(&self, info: &crate::types::message::MessageInfo) {
let source = message_ack_source_node(info);
let own_pn = self.get_pn();
match encode_ack_bytes(&source.as_node_ref(), own_pn.as_ref()) {
Ok(Some(buf)) => {
let encoded = self.encode_ack_from_snapshot(&source.as_node_ref());
match encoded {
Ok(buf) => {
if let Err(e) = self.send_raw_bytes(buf).await
&& !e.is_transport_unavailable()
{
log::warn!("Failed to send transport ack for undecryptable message: {e:?}");
}
}
Ok(None) => {}
Err(e) => log::warn!("Failed to encode transport ack: {e}"),
}
}
Expand All @@ -655,10 +681,8 @@ impl Client {
self: &Arc<Self>,
node: &wacore_binary::NodeRef<'_>,
) {
let own_pn = self.get_pn();
let buf = match encode_ack_bytes(node, own_pn.as_ref()) {
Ok(Some(b)) => b,
Ok(None) => return,
let buf = match self.encode_ack_from_snapshot(node) {
Ok(buf) => buf,
Err(e) => {
log::warn!("Failed to encode node transport ack: {e}");
return;
Expand Down
Loading
Loading