From cd074e48410ffa89c9b4fe7172d22b5194a60328 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Tue, 26 Aug 2025 19:56:18 +0000 Subject: [PATCH 1/2] Reduce allocations using thread-local scratch buffer for serialization Co-authored-by: miles.frankel --- TODO.md | 2 +- src/server.rs | 13 +++------ src/storage.rs | 74 +++++++++++++++++++++++++++++++++++--------------- 3 files changed, 57 insertions(+), 32 deletions(-) diff --git a/TODO.md b/TODO.md index 7f2a819..1ce9031 100644 --- a/TODO.md +++ b/TODO.md @@ -35,7 +35,7 @@ - [ ] (minor) rename keys to ids in `LeaseEntry`, or actually put keys in there. either way - [ ] (minor) make `add_available_item_from_parts` wrap `add_available_items_from_parts`, not the other way around - [ ] (minor) use a mockable clock when generating uuidv7s -- [ ] (perf) reduce unnecessary allocs, such as when copying data or allocating buffers. some is called out in code comments +- [X] (perf) reduce unnecessary allocs, such as when copying data or allocating buffers. some is called out in code comments - [ ] (perf) sort the keys in `LeaseEntry` so we can do bsearch on them - [ ] (major) ensure `extend` doesnt create multiple index entries for the same lease. - [ ] (perf) add lease expiry index key to `LeaseEntry` so we don't have to do scans to find it when extending diff --git a/src/server.rs b/src/server.rs index f393022..25d0d40 100644 --- a/src/server.rs +++ b/src/server.rs @@ -116,9 +116,9 @@ impl crate::protocol::queue::Server for Server { // Generate ids upfront and copy request data into owned memory so we can move // it into a blocking task (capnp readers are not Send). - let ids: Vec> = items + let ids: Vec<[u8; 16]> = items .iter() - .map(|_| uuid::Uuid::now_v7().as_bytes().to_vec()) + .map(|_| uuid::Uuid::now_v7().into_bytes()) .collect(); let items_owned = items @@ -135,8 +135,6 @@ impl crate::protocol::queue::Server for Server { let storage = Arc::clone(&self.storage); let notify = Arc::clone(&self.notify); - let ids_for_resp = ids.clone(); - Promise::from_future(async move { // Offload RocksDB work to the blocking thread pool. let iter = ids @@ -150,11 +148,8 @@ impl crate::protocol::queue::Server for Server { } // Build the response on the RPC thread. - let mut ids_builder = results - .get() - .init_resp() - .init_ids(ids_for_resp.len() as u32); - for (i, id) in ids_for_resp.iter().enumerate() { + let mut ids_builder = results.get().init_resp().init_ids(ids.len() as u32); + for (i, id) in ids.iter().enumerate() { ids_builder.set(i as u32, id); } Ok(()) diff --git a/src/storage.rs b/src/storage.rs index 8eb37c7..375c239 100644 --- a/src/storage.rs +++ b/src/storage.rs @@ -14,6 +14,10 @@ use crate::errors::{Error, Result}; use crate::protocol; use rand::Rng; +thread_local! { + static SERIALIZE_SCRATCH: std::cell::RefCell> = const { std::cell::RefCell::new(Vec::new()) }; +} + pub struct Storage { db: OptimisticTransactionDB, } @@ -89,14 +93,20 @@ impl Storage { stored_item.set_contents(contents); stored_item.set_id(id); stored_item.set_visibility_ts_index_key(visibility_index_key.as_bytes()); - let mut stored_contents = Vec::with_capacity(simsg.size_in_words() * 8); - serialize::write_message(&mut stored_contents, &simsg)?; - // Atomically insert the item and visibility index entry - let mut batch = WriteBatchWithTransaction::::default(); - batch.put(main_key.as_ref(), &stored_contents); - batch.put(visibility_index_key.as_ref(), main_key.as_ref()); - self.db.write(batch)?; + SERIALIZE_SCRATCH.with(|cell| -> Result<()> { + let mut buf = cell.borrow_mut(); + buf.clear(); + buf.reserve(simsg.size_in_words() * 8); + serialize::write_message(&mut *buf, &simsg)?; + + // Atomically insert the item and visibility index entry + let mut batch = WriteBatchWithTransaction::::default(); + batch.put(main_key.as_ref(), &*buf); + batch.put(visibility_index_key.as_ref(), main_key.as_ref()); + self.db.write(batch)?; + Ok(()) + })?; tracing::debug!( "inserted item (from parts): ({}: ), (viz/{}: avail/{})", @@ -268,12 +278,17 @@ impl Storage { // Build the lease entry. let lease_entry = build_lease_entry_message(lease_validity_secs, &polled_items)?; - let mut lease_entry_bs = Vec::with_capacity(lease_entry.size_in_words() * 8); // TODO: avoid allocation - serialize::write_message(&mut lease_entry_bs, &lease_entry)?; - - // Write the lease entry and its expiry index - txn.put(lease_key.as_ref(), &lease_entry_bs)?; - txn.put(lease_expiry_index_key.as_ref(), lease_key.as_ref())?; + SERIALIZE_SCRATCH.with(|cell| -> Result<()> { + let mut buf = cell.borrow_mut(); + buf.clear(); + buf.reserve(lease_entry.size_in_words() * 8); + serialize::write_message(&mut *buf, &lease_entry)?; + + // Write the lease entry and its expiry index + txn.put(lease_key.as_ref(), &*buf)?; + txn.put(lease_expiry_index_key.as_ref(), lease_key.as_ref())?; + Ok(()) + })?; drop(snapshot); txn.commit()?; @@ -346,9 +361,14 @@ impl Storage { out_keys.set(new_idx as u32, k); new_idx += 1; } - let mut buf = Vec::with_capacity(msg.size_in_words() * 8); // TODO: reduce allocs - serialize::write_message(&mut buf, &msg)?; - txn.put(lease_key.as_ref(), &buf)?; + SERIALIZE_SCRATCH.with(|cell| -> Result<()> { + let mut buf = cell.borrow_mut(); + buf.clear(); + buf.reserve(msg.size_in_words() * 8); + serialize::write_message(&mut *buf, &msg)?; + txn.put(lease_key.as_ref(), &*buf)?; + Ok(()) + })?; } drop(lease_value); @@ -439,9 +459,14 @@ impl Storage { let mut stored_item = builder.get_root::()?; stored_item.set_visibility_ts_index_key(vis_idx_now.as_bytes()); } - let mut updated = Vec::with_capacity(builder.size_in_words() * 8); - serialize::write_message(&mut updated, &builder)?; - txn.put(avail_key.as_ref(), &updated)?; + SERIALIZE_SCRATCH.with(|cell| -> Result<()> { + let mut updated = cell.borrow_mut(); + updated.clear(); + updated.reserve(builder.size_in_words() * 8); + serialize::write_message(&mut *updated, &builder)?; + txn.put(avail_key.as_ref(), &*updated)?; + Ok(()) + })?; txn.put(vis_idx_now.as_ref(), avail_key.as_ref())?; txn.delete(in_progress_key.as_ref())?; } @@ -503,9 +528,14 @@ impl Storage { for i in 0..keys.len() { out_keys.set(i, keys.get(i)?); } - let mut buf = Vec::with_capacity(out.size_in_words() * 8); - serialize::write_message(&mut buf, &out)?; - txn.put(lease_key.as_ref(), &buf)?; + SERIALIZE_SCRATCH.with(|cell| -> Result<()> { + let mut buf = cell.borrow_mut(); + buf.clear(); + buf.reserve(out.size_in_words() * 8); + serialize::write_message(&mut *buf, &out)?; + txn.put(lease_key.as_ref(), &*buf)?; + Ok(()) + })?; } txn.commit()?; Ok(true) From 4abe6c45de50fe40be135393950de64d35645e30 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Tue, 26 Aug 2025 20:04:48 +0000 Subject: [PATCH 2/2] Remove thread_local serialization scratch buffer, use local Vec instead Co-authored-by: miles.frankel --- src/storage.rs | 72 +++++++++++++++----------------------------------- 1 file changed, 21 insertions(+), 51 deletions(-) diff --git a/src/storage.rs b/src/storage.rs index 375c239..16d88c9 100644 --- a/src/storage.rs +++ b/src/storage.rs @@ -14,10 +14,6 @@ use crate::errors::{Error, Result}; use crate::protocol; use rand::Rng; -thread_local! { - static SERIALIZE_SCRATCH: std::cell::RefCell> = const { std::cell::RefCell::new(Vec::new()) }; -} - pub struct Storage { db: OptimisticTransactionDB, } @@ -94,19 +90,14 @@ impl Storage { stored_item.set_id(id); stored_item.set_visibility_ts_index_key(visibility_index_key.as_bytes()); - SERIALIZE_SCRATCH.with(|cell| -> Result<()> { - let mut buf = cell.borrow_mut(); - buf.clear(); - buf.reserve(simsg.size_in_words() * 8); - serialize::write_message(&mut *buf, &simsg)?; + let mut buf = Vec::with_capacity(simsg.size_in_words() * 8); + serialize::write_message(&mut buf, &simsg)?; - // Atomically insert the item and visibility index entry - let mut batch = WriteBatchWithTransaction::::default(); - batch.put(main_key.as_ref(), &*buf); - batch.put(visibility_index_key.as_ref(), main_key.as_ref()); - self.db.write(batch)?; - Ok(()) - })?; + // Atomically insert the item and visibility index entry + let mut batch = WriteBatchWithTransaction::::default(); + batch.put(main_key.as_ref(), &buf); + batch.put(visibility_index_key.as_ref(), main_key.as_ref()); + self.db.write(batch)?; tracing::debug!( "inserted item (from parts): ({}: ), (viz/{}: avail/{})", @@ -278,17 +269,11 @@ impl Storage { // Build the lease entry. let lease_entry = build_lease_entry_message(lease_validity_secs, &polled_items)?; - SERIALIZE_SCRATCH.with(|cell| -> Result<()> { - let mut buf = cell.borrow_mut(); - buf.clear(); - buf.reserve(lease_entry.size_in_words() * 8); - serialize::write_message(&mut *buf, &lease_entry)?; - - // Write the lease entry and its expiry index - txn.put(lease_key.as_ref(), &*buf)?; - txn.put(lease_expiry_index_key.as_ref(), lease_key.as_ref())?; - Ok(()) - })?; + let mut lease_buf = Vec::with_capacity(lease_entry.size_in_words() * 8); + serialize::write_message(&mut lease_buf, &lease_entry)?; + // Write the lease entry and its expiry index + txn.put(lease_key.as_ref(), &lease_buf)?; + txn.put(lease_expiry_index_key.as_ref(), lease_key.as_ref())?; drop(snapshot); txn.commit()?; @@ -361,14 +346,9 @@ impl Storage { out_keys.set(new_idx as u32, k); new_idx += 1; } - SERIALIZE_SCRATCH.with(|cell| -> Result<()> { - let mut buf = cell.borrow_mut(); - buf.clear(); - buf.reserve(msg.size_in_words() * 8); - serialize::write_message(&mut *buf, &msg)?; - txn.put(lease_key.as_ref(), &*buf)?; - Ok(()) - })?; + let mut buf = Vec::with_capacity(msg.size_in_words() * 8); + serialize::write_message(&mut buf, &msg)?; + txn.put(lease_key.as_ref(), &buf)?; } drop(lease_value); @@ -459,14 +439,9 @@ impl Storage { let mut stored_item = builder.get_root::()?; stored_item.set_visibility_ts_index_key(vis_idx_now.as_bytes()); } - SERIALIZE_SCRATCH.with(|cell| -> Result<()> { - let mut updated = cell.borrow_mut(); - updated.clear(); - updated.reserve(builder.size_in_words() * 8); - serialize::write_message(&mut *updated, &builder)?; - txn.put(avail_key.as_ref(), &*updated)?; - Ok(()) - })?; + let mut updated = Vec::with_capacity(builder.size_in_words() * 8); + serialize::write_message(&mut updated, &builder)?; + txn.put(avail_key.as_ref(), &updated)?; txn.put(vis_idx_now.as_ref(), avail_key.as_ref())?; txn.delete(in_progress_key.as_ref())?; } @@ -528,14 +503,9 @@ impl Storage { for i in 0..keys.len() { out_keys.set(i, keys.get(i)?); } - SERIALIZE_SCRATCH.with(|cell| -> Result<()> { - let mut buf = cell.borrow_mut(); - buf.clear(); - buf.reserve(out.size_in_words() * 8); - serialize::write_message(&mut *buf, &out)?; - txn.put(lease_key.as_ref(), &*buf)?; - Ok(()) - })?; + let mut buf = Vec::with_capacity(out.size_in_words() * 8); + serialize::write_message(&mut buf, &out)?; + txn.put(lease_key.as_ref(), &buf)?; } txn.commit()?; Ok(true)