-
-
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 1 commit
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,16 @@ impl SignalStore for InMemoryBackend { | |
|
|
||
| async fn store_prekeys_batch(&self, keys: &[(u32, Bytes)], _uploaded: bool) -> Result<()> { | ||
| let mut state = self.state.lock().await; | ||
| // The batch length is known and every id in it is new: prekey ids are | ||
| // minted from the monotonic NEXT_PK_ID counter, so a batch never | ||
| // overwrites a stored row (unlike `put_msg_secrets`, where reserving a | ||
| // mostly-overwrite batch would grow the table for rows it never adds). | ||
| // Growing incrementally instead allocates and copies a whole table per | ||
| // rehash — a connect-sized batch crosses the load factor eight times. | ||
| // 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. | ||
| state.prekeys.reserve(keys.len()); | ||
|
greptile-apps[bot] marked this conversation as resolved.
Outdated
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.
When a caller replays or rewrites an existing batch, Useful? React with 👍 / 👎. |
||
| for (id, record) in keys { | ||
| state.prekeys.insert( | ||
| *id, | ||
|
|
@@ -1243,6 +1253,77 @@ 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" | ||
| ); | ||
| } | ||
| } | ||
|
|
||
| #[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