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
12 changes: 12 additions & 0 deletions agent_docs/observability.md
Original file line number Diff line number Diff line change
Expand Up @@ -191,6 +191,18 @@ staying alive because the `Bytes` slices handed to `store_prekeys_batch` *are*
what the backend stores (one allocation instead of 812), and 41 KiB is that
map's `RawTable` at 1024 buckets. Nothing to optimise; do not re-derive it.

That 41 KiB deserves one clarification, because a heap profiler hands it to you
under a name that invites the wrong fix. dhat attributes the final table to
`hashbrown::RawTable::reserve_rehash`, the frame that happened to allocate it,
so a per-session diff reads "41.0 KiB in reserve_rehash" and looks like rehash
churn. It is not: 1024 buckets × (`size_of::<(u32, PreKeyEntry)>()` + 1 control
byte) = 41,984 B is the table that *stays*, and the intermediate tables are all
freed before the process peak. `store_prekeys_batch` does reserve for the batch
length (#1270), which cuts the call from 11 allocations / 84.1 KB to 3 / 42.1 KB
and its in-call transient high-water from 63.1 KB to 42.1 KB — but retained is
bit-identical at 42,072 B either way, because the final table is the same size.
Reserving is worth it for the allocator traffic; it will never move the 41 KiB.

**The rustls session cache is 5 KiB, not 44.** A whole retained
`default_tls_connector()` measures 14.0 KiB; disabling resumption entirely takes
it to 9.0 KiB, and sizing the store for the one host a factory dials takes it to
Expand Down
4 changes: 4 additions & 0 deletions wacore/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,10 @@ harness = false
name = "sender_key_derivation_benchmark"
harness = false

[[bench]]
name = "prekey_store_benchmark"
harness = false

[[bench]]
name = "voip_benchmark"
harness = false
Expand Down
74 changes: 74 additions & 0 deletions wacore/benches/prekey_store_benchmark.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,74 @@
//! Allocation accounting for the prekey batch write on the connect path.
//!
//! `upload_pre_keys_pass` generates `DEFAULT_WANTED_PRE_KEY_COUNT` (812) one-time
//! prekeys and hands them to `store_prekeys_batch` as one call. `divan::AllocProfiler`
//! is wired as the global allocator so each row reports allocation count and bytes
//! next to wall time -- the count is the signal here, since the map's growth is the
//! only thing this path allocates: the records themselves are `Bytes` slices of one
//! shared buffer the caller already owns, so they cost refcount bumps, not copies.
//!
//! The `populated` row is the second upload pass: the same batch size arriving at a
//! map that already holds a window. It exists so a reservation that over-grows an
//! already-populated table would show up as bytes here rather than silently.

use bytes::Bytes;
use divan::{Bencher, black_box};
use futures::executor::block_on;
use wacore::store::in_memory::InMemoryBackend;
use wacore::store::traits::SignalStore;

#[global_allocator]
static ALLOC: divan::AllocProfiler = divan::AllocProfiler::system();

fn main() {
divan::main();
}

/// `DEFAULT_WANTED_PRE_KEY_COUNT` in `src/prekeys.rs`, mirroring WA Web's
/// UPLOAD_KEYS_COUNT. The small row is there to show the growth is what scales.
const CONNECT_BATCH: usize = 812;

/// Upper bound on one encoded `PreKeyRecordStructure`, the same figure
/// `upload_pre_keys_pass` sizes its shared buffer with.
const RECORD_LEN: usize = 74;

/// The batch as the upload path actually hands it over: every record is a slice of
/// one contiguous buffer, so building it allocates once regardless of `count`.
fn batch(first_id: u32, count: usize) -> Vec<(u32, Bytes)> {
let shared = Bytes::from(vec![7u8; count * RECORD_LEN]);
(0..count)
.map(|i| {
(
first_id + i as u32,
shared.slice(i * RECORD_LEN..(i + 1) * RECORD_LEN),
)
})
.collect()
}

#[divan::bench(args = [64, CONNECT_BATCH])]
fn store_prekeys_batch(bencher: Bencher, count: usize) {
bencher
.with_inputs(|| (InMemoryBackend::new(), batch(1, count)))
.bench_refs(|(backend, keys)| {
block_on(backend.store_prekeys_batch(keys, false)).unwrap();

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Benchmark uses prohibited unwraps

The new benchmark calls .unwrap() here and again on lines 66 and 71, contrary to the repository-wide prohibition on .unwrap() outside tests. These calls also make benchmark failures panic without contextual error propagation.

Context Used: CLAUDE.md (source)

Prompt To Fix With AI
This is a comment left during a code review.
Path: wacore/benches/prekey_store_benchmark.rs
Line: 54

Comment:
**Benchmark uses prohibited unwraps**

The new benchmark calls `.unwrap()` here and again on lines 66 and 71, contrary to the repository-wide prohibition on `.unwrap()` outside tests. These calls also make benchmark failures panic without contextual error propagation.

**Context Used:** CLAUDE.md ([source](https://github.com/oxidezap/whatsapp-rust/blob/main/CLAUDE.md))

---

For each issue above, determine whether it is valid and should be fixed. If so, fix it directly.

Fix in Claude Code

black_box(&backend);
});
}

/// Prekey ids are minted from the monotonic NEXT_PK_ID counter, so a second pass
/// carries ids past the first window's -- every row is new, none overwrites.
#[divan::bench(name = "store_prekeys_batch/populated")]
fn store_prekeys_batch_populated(bencher: Bencher) {
bencher
.with_inputs(|| {
let backend = InMemoryBackend::new();
block_on(backend.store_prekeys_batch(&batch(1, CONNECT_BATCH), true)).unwrap();
let next = batch(CONNECT_BATCH as u32 + 1, CONNECT_BATCH);
(backend, next)
})
.bench_refs(|(backend, keys)| {
block_on(backend.store_prekeys_batch(keys, false)).unwrap();
black_box(&backend);
});
}
165 changes: 165 additions & 0 deletions wacore/src/store/in_memory.rs
Original file line number Diff line number Diff line change
Expand Up @@ -320,6 +320,38 @@ impl SignalStore for InMemoryBackend {

async fn store_prekeys_batch(&self, keys: &[(u32, Bytes)], _uploaded: bool) -> Result<()> {
let mut state = self.state.lock().await;
// Growing one insert at a time allocates and copies a whole table per
// rehash, and a connect-sized batch arriving at an empty map crosses
// the load factor eight times. The batch length is known, so the table
// can reach its final size in one allocation instead.
//
// Two things stop that reservation from over-growing a table, because a
// table grown for rows that were never added does not shrink back and
// this is meant to cost no retained bytes:
//
// 1. Subtract the rows already stored. A batch may legally overwrite
// ids, and no batch can overwrite more rows than exist, so
// `keys.len() - len()` is the floor on how many ids must be new.
// 2. Only reserve at all when the batch is strictly ascending, which
// proves its ids are distinct. Without that, 812 entries sharing one
// id would reserve a 1024-bucket table to hold a single row. Testing
// the order costs one pass of integer compares and no allocation;
// deduplicating properly would need a set, whose own allocation and
// 812 hashes cost more than the eight allocations being saved.
//
// Both are one-sided: they can only under-reserve and fall back to
// incremental growth, never inflate the resident table. The connect path
// satisfies both — the map is empty and `upload_pre_keys_pass` emits
// `gen_start + i`, so the whole batch is reserved and gets the full win.
//
// This does NOT shrink the table that stays resident: the final
// capacity is the same either way, so it buys allocator traffic and
// in-call headroom, not retained bytes.
let ascending = keys.windows(2).all(|pair| pair[0].0 < pair[1].0);
if ascending {
let at_least_new = keys.len().saturating_sub(state.prekeys.len());
state.prekeys.reserve(at_least_new);
}
for (id, record) in keys {
state.prekeys.insert(
*id,
Expand Down Expand Up @@ -1243,6 +1275,139 @@ mod tests {
);
}

/// A connect-sized batch: every id must be readable back, and a later batch
/// repeating an id must overwrite it rather than duplicate or drop it. This
/// is what pins `store_prekeys_batch` idempotent per id across the reserve.
#[tokio::test]
async fn store_prekeys_batch_stores_every_key() {
const COUNT: u32 = 812;
let backend = InMemoryBackend::new();

let batch: Vec<(u32, Bytes)> = (1..=COUNT)
.map(|id| (id, Bytes::from(format!("record-{id}"))))
.collect();
backend.store_prekeys_batch(&batch, false).await.unwrap();

for id in 1..=COUNT {
assert_eq!(
backend.load_prekey(id).await.unwrap(),
Some(Bytes::from(format!("record-{id}"))),
"prekey {id} must survive the batch write"
);
}
assert_eq!(backend.get_max_prekey_id().await.unwrap(), COUNT);

backend
.store_prekeys_batch(&[(7, Bytes::from_static(b"rewritten"))], true)
.await
.unwrap();
assert_eq!(
backend.load_prekey(7).await.unwrap(),
Some(Bytes::from_static(b"rewritten"))
);
assert_eq!(
backend.state.lock().await.prekeys.len(),
COUNT as usize,
"re-storing an existing id must not add a row"
);
}

/// Reserving for the batch length must leave the table exactly the size the
/// row count alone demands — the point of the reserve is to reach that size
/// in one allocation, not to reach a bigger one. The control map is grown
/// one insert at a time, which is the un-reserved shape; hashbrown sizes a
/// table from the element count alone, so the two must agree. The second
/// pass covers the reserve landing on an already-populated map, where
/// reserving the full batch length on top of the existing rows would be
/// visible as a doubled table.
#[tokio::test]
async fn store_prekeys_batch_reserve_does_not_over_grow_the_table() {
const COUNT: u32 = 812;
let backend = InMemoryBackend::new();
let mut control: HashMap<u32, ()> = HashMap::new();

for pass in 0..2u32 {
let first = pass * COUNT + 1;
let batch: Vec<(u32, Bytes)> = (first..first + COUNT)
.map(|id| (id, Bytes::from_static(b"record")))
.collect();
backend.store_prekeys_batch(&batch, false).await.unwrap();
for id in first..first + COUNT {
control.insert(id, ());
}

let state = backend.state.lock().await;
assert_eq!(state.prekeys.len(), control.len());
assert_eq!(
state.prekeys.capacity(),
control.capacity(),
"pass {pass}: the reserved table must match an incrementally grown one"
);
}
}

/// Replaying a stored window must not grow the table by one bucket. The
/// trait permits a batch to overwrite ids, and a table grown for rows that
/// were only overwritten never shrinks back — so a reservation taken on the
/// bare batch length would retain an extra table forever, which is exactly
/// the residency this change claims not to touch.
#[tokio::test]
async fn replaying_a_stored_batch_does_not_grow_the_table() {
const COUNT: u32 = 812;
let backend = InMemoryBackend::new();
let batch: Vec<(u32, Bytes)> = (1..=COUNT)
.map(|id| (id, Bytes::from_static(b"record")))
.collect();

backend.store_prekeys_batch(&batch, false).await.unwrap();
let settled = backend.state.lock().await.prekeys.capacity();

// Same ids twice more: every row is an overwrite, so nothing is added.
backend.store_prekeys_batch(&batch, true).await.unwrap();
backend.store_prekeys_batch(&batch, true).await.unwrap();

let state = backend.state.lock().await;
assert_eq!(state.prekeys.len(), COUNT as usize, "no rows were added");
assert_eq!(
state.prekeys.capacity(),
settled,
"an all-overwrite batch must not enlarge the table"
);
}

/// A batch whose ids repeat stores one row per distinct id, so sizing the
/// table from the batch length would leave it holding a table for rows that
/// never existed. The reservation is skipped unless the batch is strictly
/// ascending, which is what makes its ids provably distinct.
#[tokio::test]
async fn a_batch_of_repeated_ids_does_not_reserve_for_them() {
const COUNT: usize = 812;
let backend = InMemoryBackend::new();
let batch: Vec<(u32, Bytes)> = (0..COUNT)
.map(|i| (7, Bytes::from(format!("record-{i}"))))
.collect();

backend.store_prekeys_batch(&batch, false).await.unwrap();

// One row survives — the last write for id 7 — so the table must be
// sized for one row, not for the 812 entries that were handed over.
let mut control: HashMap<u32, ()> = HashMap::new();
control.insert(7, ());

let state = backend.state.lock().await;
assert_eq!(state.prekeys.len(), 1, "last write wins per id");
assert_eq!(
state.prekeys.capacity(),
control.capacity(),
"a repeated-id batch must not size the table by its length"
);
drop(state);
assert_eq!(
backend.load_prekey(7).await.unwrap(),
Some(Bytes::from(format!("record-{}", COUNT - 1)))
);
}

#[tokio::test]
async fn group_metadata_round_trip() {
use crate::store::traits::ProtocolStore;
Expand Down
Loading