Skip to content
Merged
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
283 changes: 259 additions & 24 deletions wacore/benches/send_receive_benchmark.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,8 +16,8 @@ use wacore::client::context::{GroupInfo, SendContextResolver};
use wacore::messages::MessageUtils;
use wacore::runtime::{AbortHandle, Runtime};
use wacore::send::{
GroupStanzaRequest, SenderKeyDistributionPolicy, SignalStores, prepare_group_stanza,
prepare_peer_stanza,
DmStanzaRequest, GroupStanzaRequest, ResolvedDmDevices, SenderKeyDistributionPolicy,
SignalStores, prepare_dm_stanza, prepare_group_stanza, prepare_peer_stanza,
};
use wacore::types::jid::{JidExt, make_sender_key_name};
use wacore::types::message::AddressingMode;
Expand Down Expand Up @@ -344,6 +344,15 @@ impl User {
}
}

/// The same account on another device: same keys, distinct address. Real
/// DM fan-out targets several devices of one user, which `User::new` alone
/// cannot express since it derives the address from user+server.
fn with_device(mut self, device: u16) -> Self {
self.jid.device = device;
self.address = self.jid.to_protocol_address();
self
}

fn prekey_bundle(&self) -> PreKeyBundle {
PreKeyBundle::new(
self.identity.reg_id,
Expand Down Expand Up @@ -488,6 +497,30 @@ fn extract_skmsg_bytes(stanza: &Node) -> Vec<u8> {
}
}

/// Encrypt one padded text message from `alice` to `to`, returning the wire
/// bytes and the enc type the stanza would carry. Shared by both receive
/// fixtures so the ratchet and steady-state benches cannot drift apart on
/// padding or enc-type mapping.
fn encrypt_one(alice: &mut User, to: &ProtocolAddress) -> (Vec<u8>, String) {
futures::executor::block_on(async {
let ct = message_encrypt(
&MessageUtils::pad_message_v2(text_msg().encode_to_vec()),
to,
&mut alice.sessions,
&mut alice.identity,
)
.await
.unwrap();
match ct {
CiphertextMessage::SignalMessage(m) => (m.serialized().to_vec(), "msg".to_string()),
CiphertextMessage::PreKeySignalMessage(m) => {
(m.serialized().to_vec(), "pkmsg".to_string())
}
_ => panic!("unexpected ciphertext type in benchmark fixture"),
}
})
}

fn decrypt_dm(
ciphertext: &[u8],
enc_type: &str,
Expand Down Expand Up @@ -542,10 +575,37 @@ fn decrypt_group(
// DM setups
// ---------------------------------------------------------------------------

/// The `<device-identity>` blob a paired device carries.
///
/// Every send path takes `Option<&ADVSignedDeviceIdentity>`, and `None` changes
/// what they do — differently per path, so the distinction matters here:
///
/// - `prepare_peer_stanza` refuses, and pays a whole extra session checkout to
/// find out (`account.is_none() && pkmsg_would_be_emitted(..)`);
/// - `prepare_dm_stanza` and the `BestEffort` group policy are lenient: they
/// drop the `<device-identity>` child and send anyway;
/// - only `SenderKeyDistributionPolicy::Required` propagates the error.
///
/// A paired device always has one (`device_snapshot.account.as_deref()`), so a
/// bench passing `None` measures a shape production never sends — and on the
/// lenient paths it silently omits a stanza child rather than failing loudly.
///
/// Only the bytes matter here — it is serialised into the stanza, never
/// validated by the sender.
fn bench_account() -> wa::ADVSignedDeviceIdentity {
wa::ADVSignedDeviceIdentity {
details: Some(vec![0xAD; 32]),
account_signature_key: Some(vec![0xAC; 32]),
account_signature: Some(vec![0x51; 64]),
device_signature: Some(vec![0xD5; 64]),
}
}

struct DmSendData {
alice: User,
bob_jid: Jid,
msg: wa::Message,
account: wa::ADVSignedDeviceIdentity,
}

fn setup_dm_send() -> DmSendData {
Expand All @@ -556,6 +616,7 @@ fn setup_dm_send() -> DmSendData {
alice,
bob_jid: bob.jid,
msg: text_msg(),
account: bench_account(),
}
}

Expand All @@ -571,25 +632,13 @@ fn setup_dm_recv() -> DmRecvData {
let mut bob = User::new("5511999000002", "s.whatsapp.net");
establish_bidirectional(&mut alice, &mut bob);

// Encrypt a message (subsequent, not PreKey — realistic steady-state)
let (ciphertext, enc_type) = futures::executor::block_on(async {
let ct = message_encrypt(
&MessageUtils::pad_message_v2(text_msg().encode_to_vec()),
&bob.address,
&mut alice.sessions,
&mut alice.identity,
)
.await
.unwrap();

match ct {
CiphertextMessage::SignalMessage(msg) => (msg.serialized().to_vec(), "msg".to_string()),
CiphertextMessage::PreKeySignalMessage(msg) => {
(msg.serialized().to_vec(), "pkmsg".to_string())
}
_ => panic!("unexpected type"),
}
});
// A `msg` rather than a `pkmsg`, but NOT steady state: Bob has not decrypted
// anything on this sending chain yet, so this first one costs him a DH ratchet
// step (two X25519 agreements plus a key generation). That is a real shape —
// the first message after the peer replies — but it is ~40x a steady-state
// decrypt, so the two need separate benches or the common case is invisible.
let bob_addr = bob.address.clone();
let (ciphertext, enc_type) = encrypt_one(&mut alice, &bob_addr);

DmRecvData {
bob,
Expand All @@ -613,6 +662,11 @@ struct GrpSendData {
force_skdm: bool,
resolver: MockResolver,
msg: wa::Message,
/// Built in setup, not in the measured body: production borrows the cached
/// account rather than allocating one per send, so building it inside the
/// closure would charge every group baseline four `Vec` allocations that no
/// send performs.
account: wa::ADVSignedDeviceIdentity,
// Built once in setup so the measured body excludes thread-pool startup
// (building the pool inside the bench body would charge its syscalls to
// the encrypt path).
Expand Down Expand Up @@ -659,6 +713,7 @@ fn setup_group_send(n: usize) -> GrpSendData {
force_skdm: false,
resolver: MockResolver(devices),
msg: text_msg(),
account: bench_account(),
runtime: BenchRuntime,
}
}
Expand Down Expand Up @@ -737,6 +792,7 @@ fn setup_group_recv() -> GrpRecvData {

let runtime = BenchRuntime;
let message = text_msg();
let account = bench_account();
let result = futures::executor::block_on(prepare_group_stanza(
&runtime,
&mut stores,
Expand All @@ -745,7 +801,7 @@ fn setup_group_recv() -> GrpRecvData {
group: &group_info,
own_jid: &own_jid,
own_lid: &own_jid,
account: None,
account: Some(&account),
to: &group_jid,
message: &message,
message_id: "bench-grp-recv",
Expand Down Expand Up @@ -774,8 +830,149 @@ fn setup_group_recv() -> GrpRecvData {
// Benchmarks
// ===========================================================================

struct DmFanoutData {
alice: User,
to_jid: Jid,
msg: wa::Message,
account: wa::ADVSignedDeviceIdentity,
devices: ResolvedDmDevices,
resolver: MockResolver,
runtime: BenchRuntime,
}

/// The real 1:1 send: `prepare_dm_stanza`, which fans out to the recipient's
/// devices and our own companions and builds the `<participants>` tree.
///
/// `n_recipient_devices` counts the contact's devices; one own companion is
/// always added on top, because a paired account has at least the phone and
/// this bench should not measure the degenerate zero-companion shape. The
/// total fan-out is therefore `n_recipient_devices + 1` pairwise encryptions.
///
/// Sessions are established **bidirectionally**. `establish_session` alone
/// leaves `pending_pre_key` set, so every outbound stays a `pkmsg` and the
/// bench would measure first-message prekey wrapping plus `<device-identity>`
/// serialisation forever — the opposite of the repeat-send shape intended here.
fn setup_dm_fanout(n_recipient_devices: usize) -> DmFanoutData {
let mut alice = User::new("5511999000001", "s.whatsapp.net");
let own_jid = alice.jid.clone();
let to_base: Jid = "5511999000002@s.whatsapp.net".parse().unwrap();

let mut all_devices = Vec::with_capacity(n_recipient_devices + 1);
for i in 0..n_recipient_devices {
let mut device = to_base.clone();
device.device = i as u16;
let mut peer = User::new("5511999000002", "s.whatsapp.net").with_device(i as u16);
establish_bidirectional(&mut alice, &mut peer);
all_devices.push(device);
}
// One own companion, so the own-device half of the partition is exercised.
let mut companion = own_jid.clone();
companion.device = 42;
let mut own_peer = User::new("5511999000001", "s.whatsapp.net").with_device(42);
establish_bidirectional(&mut alice, &mut own_peer);
all_devices.push(companion);

// Assert the sessions really are acknowledged, so a fixture regression cannot
// quietly turn this back into a first-message benchmark: with `pending_pre_key`
// still set every outbound is a pkmsg, which measures prekey wrapping and
// `<device-identity>` serialisation instead of the repeat send.
for device in &all_devices {
let addr = device.to_protocol_address();
let emits_pkmsg = futures::executor::block_on(wacore::send::pkmsg_would_be_emitted(
&mut alice.sessions,
&addr,
))
.expect("session must be loadable in setup");
assert!(
!emits_pkmsg,
"fixture must use acknowledged sessions; {addr} still emits pkmsg"
);
}

let resolver = MockResolver(all_devices.clone());
let devices = ResolvedDmDevices::new(all_devices, &own_jid, None);
// Warm the phash memo in setup, as the per-recipient memo does for repeat
// sends in production; otherwise every iteration measures the cold path.
devices.phash();

DmFanoutData {
alice,
to_jid: to_base,
msg: text_msg(),
account: bench_account(),
devices,
resolver,
runtime: BenchRuntime,
}
}

/// One recipient device plus our companion: two pairwise encryptions.
fn setup_dm_fanout_1() -> DmFanoutData {
setup_dm_fanout(1)
}
/// Four recipient devices plus our companion: five pairwise encryptions.
fn setup_dm_fanout_4() -> DmFanoutData {
setup_dm_fanout(4)
Comment thread
jlucaso1 marked this conversation as resolved.
}

fn run_dm_fanout(d: &mut DmFanoutData) {
let own_jid = d.alice.jid.clone();
let mut stores = SignalStores {
sender_key_store: &mut d.alice.sender_keys,
session_store: &mut d.alice.sessions,
identity_store: &mut d.alice.identity,
prekey_store: &mut d.alice.prekeys,
signed_prekey_store: &d.alice.signed_prekeys,
};
let prepared = futures::executor::block_on(prepare_dm_stanza(
&d.runtime,
&mut stores,
&d.resolver,
DmStanzaRequest {
own_jid: &own_jid,
account: Some(&d.account),
to: &d.to_jid,
message: &d.msg,
message_id: "b-dm",
edit: None,
extra_nodes: &[],
devices: &d.devices,
pre_encoded: None,
},
))
.unwrap();
black_box(marshal(&prepared.node).unwrap());
}

/// The 1:1 path an actual message takes. Nothing covered this before: the only
/// DM-shaped bench measured `prepare_peer_stanza`, which is the own-device peer
/// path and skips the fan-out, the participants tree and the reporting token.
#[divan::bench]
fn bench_dm_send(bencher: divan::Bencher) {
bencher
.with_inputs(setup_dm_fanout_1)
.bench_refs(run_dm_fanout);
}

/// Multi-device recipient: the shape that actually spawns encrypt tasks.
/// Named for the total fan-out — four recipient devices plus one own companion,
/// so five pairwise encryptions.
#[divan::bench]
fn bench_dm_send_5_way_fanout(bencher: divan::Bencher) {
bencher
.with_inputs(setup_dm_fanout_4)
.bench_refs(run_dm_fanout);
}

/// A peer stanza: what we send to our OWN other devices, not a DM to a contact.
/// See `bench_dm_send` for the contact path.
///
/// `account` is `Some`, matching a paired device. With `None` the `&&` in
/// `prepare_peer_stanza` does not short-circuit and every send pays an extra
/// `SessionCheckout::load` + `commit()` just to pre-flight the pkmsg case —
/// work production never does, and enough of it to move the baseline.
#[divan::bench]
fn bench_peer_send(bencher: divan::Bencher) {
bencher.with_inputs(setup_dm_send).bench_refs(|d| {
let signal_addr = d.bob_jid.to_protocol_address();
let node = futures::executor::block_on(prepare_peer_stanza(
Expand All @@ -785,7 +982,7 @@ fn bench_dm_send(bencher: divan::Bencher) {
&signal_addr,
&d.msg,
"b-001",
None,
Some(&d.account),
))
.unwrap();
black_box(marshal(&node).unwrap());
Expand All @@ -804,6 +1001,44 @@ fn bench_dm_recv(bencher: divan::Bencher) {
});
}

/// Steady state: Bob has already ratcheted into Alice's current sending chain,
/// so this decrypt is symmetric-only — no curve operation at all. This is what
/// an active conversation actually pays per message, and it is roughly an order
/// of magnitude cheaper than the ratchet step `bench_dm_recv` measures.
fn setup_dm_recv_steady() -> DmRecvData {
let mut alice = User::new("5511999000001", "s.whatsapp.net");
let mut bob = User::new("5511999000002", "s.whatsapp.net");
establish_bidirectional(&mut alice, &mut bob);

let bob_addr = bob.address.clone();

// Burn the ratchet step on a throwaway message, so the one we hand back is
// decrypted with the chain already advanced.
let (first, first_type) = encrypt_one(&mut alice, &bob_addr);
decrypt_dm(&first, &first_type, &alice.address, &mut bob);

let (ciphertext, enc_type) = encrypt_one(&mut alice, &bob_addr);

DmRecvData {
bob,
alice_addr: alice.address,
ciphertext,
enc_type,
}
}

#[divan::bench]
fn bench_dm_recv_steady(bencher: divan::Bencher) {
bencher.with_inputs(setup_dm_recv_steady).bench_refs(|d| {
black_box(decrypt_dm(
&d.ciphertext,
&d.enc_type,
&d.alice_addr,
&mut d.bob,
));
});
}

fn run_group_send(d: &mut GrpSendData) {
let own_jid = d.alice.jid.clone();
// Warm sends (force_skdm=false) distribute no SKDM, so prepare_group_stanza
Expand Down Expand Up @@ -835,7 +1070,7 @@ fn run_group_send(d: &mut GrpSendData) {
group: &group_info,
own_jid: &own_jid,
own_lid: &own_jid,
account: None,
account: Some(&d.account),
to: &d.group_jid,
message: &d.msg,
message_id: "b-grp",
Expand Down
Loading