-
-
Notifications
You must be signed in to change notification settings - Fork 127
perf(store): reserve the prekey map for a batch insert #1270
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from 3 commits
fa04c47
10d4057
cf70e5e
4dfd8b5
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| 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(); | ||
| 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); | ||
| }); | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -320,6 +320,27 @@ 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. | ||
| // | ||
| // Reserve the batch length MINUS the rows already stored, not the batch | ||
| // length: a batch may legally overwrite ids (the trait permits it), and | ||
| // a table grown for rows that were only overwritten never shrinks back. | ||
| // `keys.len() - len()` is the floor on how many ids must be new, since | ||
| // no batch can overwrite more rows than exist — so it can only | ||
| // under-reserve (falling back to incremental growth), never inflate the | ||
| // resident table. The connect path, where the map is empty and every id | ||
| // is freshly minted from the monotonic NEXT_PK_ID counter, reserves the | ||
| // whole batch and gets the full win; a replay of a stored window | ||
| // reserves nothing and leaves the table exactly as it found it. | ||
| // | ||
| // 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 at_least_new = keys.len().saturating_sub(state.prekeys.len()); | ||
| state.prekeys.reserve(at_least_new); | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
Fresh evidence in the revised code is that Useful? React with 👍 / 👎.
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Correct — fixed in I did not take the literal suggestion of reserving for distinct absent ids, because the dedup costs more than the optimization saves: a 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);
}One pass of integer compares, no allocation, no hashing. Both guards are one-sided in the same direction — they can only skip or shrink a reservation and fall back to incremental growth, never inflate the resident table — so the "costs no retained bytes" claim now holds for any caller rather than for well-behaved ones. The connect path satisfies both conditions ( Generated by Claude Code |
||
| for (id, record) in keys { | ||
| state.prekeys.insert( | ||
| *id, | ||
|
|
@@ -1243,6 +1264,106 @@ 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" | ||
| ); | ||
| } | ||
|
|
||
| #[tokio::test] | ||
| async fn group_metadata_round_trip() { | ||
| use crate::store::traits::ProtocolStore; | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
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