From f06998a801c2dc35a26bd1ce156a2e0cca9f3b05 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Wed, 27 Aug 2025 14:13:18 +0000 Subject: [PATCH 1/2] Checkpoint before follow-up message Co-authored-by: miles.frankel --- src/bin/queueber/main.rs | 27 ++- src/storage.rs | 481 +++++++++++++++++++++++++-------------- stress.sh | 3 +- 3 files changed, 335 insertions(+), 176 deletions(-) diff --git a/src/bin/queueber/main.rs b/src/bin/queueber/main.rs index 5fe624b..791f817 100644 --- a/src/bin/queueber/main.rs +++ b/src/bin/queueber/main.rs @@ -4,12 +4,12 @@ use std::{ }; use capnp_rpc::{RpcSystem, rpc_twoparty_capnp, twoparty}; -use clap::Parser; +use clap::{Parser, ValueEnum}; use color_eyre::Result; use futures::AsyncReadExt; use queueber::{ server::Server, - storage::{RetriedStorage, Storage}, + storage::{RetriedStorage, Storage, TxMode}, }; use socket2::{Domain, Protocol, Socket, Type}; use std::sync::Arc; @@ -37,6 +37,24 @@ struct Args { /// Number of RPC worker threads. Defaults to available_parallelism. #[arg(long = "workers")] workers: Option, + + /// Transaction mode: optimistic (default) or pessimistic + #[arg(long = "tx-mode", value_enum, default_value_t = TxModeArg::Optimistic)] + tx_mode: TxModeArg, +} +#[derive(Copy, Clone, Debug, ValueEnum)] +enum TxModeArg { + Optimistic, + Pessimistic, +} + +impl From for TxMode { + fn from(v: TxModeArg) -> Self { + match v { + TxModeArg::Optimistic => TxMode::Optimistic, + TxModeArg::Pessimistic => TxMode::Pessimistic, + } + } } // NOTE: to use the console you need "RUST_LOG=tokio=trace,runtime=trace" @@ -71,7 +89,10 @@ async fn main() -> Result<()> { std::fs::remove_dir_all(&args.data_dir)?; } - let storage = Arc::new(RetriedStorage::new(Storage::new(&args.data_dir)?)); + let storage = Arc::new(RetriedStorage::new(Storage::new_with_mode( + &args.data_dir, + args.tx_mode.into(), + )?)); let notify = Arc::new(Notify::new()); let (shutdown_tx, mut shutdown_rx) = tokio::sync::watch::channel(false); diff --git a/src/storage.rs b/src/storage.rs index 583b5bf..7d7e033 100644 --- a/src/storage.rs +++ b/src/storage.rs @@ -2,7 +2,8 @@ use capnp::message::{self, TypedReader}; use capnp::serialize; use rocksdb::{ Direction, IteratorMode, OptimisticTransactionDB, OptimisticTransactionOptions, Options, - ReadOptions, SliceTransform, WriteBatchWithTransaction, WriteOptions, + ReadOptions, SliceTransform, TransactionDB, TransactionOptions, WriteBatchWithTransaction, + WriteOptions, }; use std::path::Path; use uuid::Uuid; @@ -60,12 +61,118 @@ where } } +pub enum TxMode { + Optimistic, + Pessimistic, +} + +enum DbEngine { + Optimistic(OptimisticTransactionDB), + Pessimistic(TransactionDB), +} + pub struct Storage { - db: OptimisticTransactionDB, + db: DbEngine, } impl Storage { pub fn new(path: &Path) -> Result { + Self::new_with_mode(path, TxMode::Optimistic) + } + + // Test and tooling helpers + #[allow(dead_code)] + fn get_raw(&self, key: &[u8]) -> Result>> { + match &self.db { + DbEngine::Optimistic(db) => Ok(db.get(key)?), + DbEngine::Pessimistic(db) => Ok(db.get(key)?), + } + } + + #[allow(dead_code)] + fn put_raw(&self, key: &[u8], value: &[u8]) -> Result<()> { + match &self.db { + DbEngine::Optimistic(db) => db.put(key, value)?, + DbEngine::Pessimistic(db) => db.put(key, value)?, + } + Ok(()) + } + + #[allow(dead_code)] + fn count_expiry_index_entries_for_lease(&self, lease: &[u8; 16]) -> Result { + let txn = match &self.db { + DbEngine::Optimistic(db) => db.transaction(), + DbEngine::Pessimistic(db) => db.transaction(), + }; + let mut count = 0usize; + for kv in txn.prefix_iterator(LeaseExpiryIndexKey::PREFIX) { + let (idx_key, _val) = kv?; + let (_ts, lbytes) = LeaseExpiryIndexKey::split_ts_and_lease(&idx_key)?; + if &lbytes == lease { + count += 1; + } + } + Ok(count) + } + + #[allow(dead_code)] + fn assert_snapshot_reads_main_even_if_deleted_after_lock(&self) -> Result<()> { + // Begin a transaction with a snapshot + match &self.db { + DbEngine::Optimistic(db) => { + let mut txn_opts = OptimisticTransactionOptions::default(); + txn_opts.set_snapshot(true); + let txn = db.transaction_opt(&WriteOptions::default(), &txn_opts); + // Iterate visibility index and capture the first entry + let mut viz_iter = txn.prefix_iterator(VisibilityIndexKey::PREFIX); + let (idx_key, main_key) = viz_iter + .next() + .expect("expected at least one visibility index entry")?; + // Lock the visibility index entry so we simulate the poller owning it + let locked = txn.get_pinned_for_update(&idx_key, true)?; + assert!(locked.is_some(), "failed to lock visibility index entry"); + // Delete the main key outside the transaction after we've taken the snapshot + let main_key_vec = main_key.as_ref().to_vec(); + match &self.db { + DbEngine::Optimistic(db) => db.delete(&main_key_vec)?, + DbEngine::Pessimistic(db) => db.delete(&main_key_vec)?, + } + // Read using the transaction's snapshot; this should still see the value + let mut ropts = rocksdb::ReadOptions::default(); + ropts.set_snapshot(&txn.snapshot()); + let got = txn.get_pinned_for_update_opt(&main_key_vec, true, &ropts)?; + assert!( + got.is_some(), + "snapshot read should see the main value even if concurrently deleted" + ); + } + DbEngine::Pessimistic(db) => { + let mut txn_opts = TransactionOptions::default(); + txn_opts.set_set_snapshot(true); + let txn = db.transaction_opt(&WriteOptions::default(), &txn_opts); + let mut viz_iter = txn.prefix_iterator(VisibilityIndexKey::PREFIX); + let (idx_key, main_key) = viz_iter + .next() + .expect("expected at least one visibility index entry")?; + let locked = txn.get_pinned_for_update(&idx_key, true)?; + assert!(locked.is_some(), "failed to lock visibility index entry"); + let main_key_vec = main_key.as_ref().to_vec(); + match &self.db { + DbEngine::Optimistic(db) => db.delete(&main_key_vec)?, + DbEngine::Pessimistic(db) => db.delete(&main_key_vec)?, + } + let mut ropts = rocksdb::ReadOptions::default(); + ropts.set_snapshot(&txn.snapshot()); + let got = txn.get_pinned_for_update_opt(&main_key_vec, true, &ropts)?; + assert!( + got.is_some(), + "snapshot read should see the main value even if concurrently deleted" + ); + } + } + Ok(()) + } + pub fn new_with_mode(path: &Path, mode: TxMode) -> Result { let mut opts = Options::default(); // Optimize for prefix scans used by `prefix_iterator` across all key namespaces. // Extract the namespace prefix up to and including the first '/'. @@ -84,7 +191,16 @@ impl Storage { ); opts.set_prefix_extractor(ns_prefix); opts.create_if_missing(true); - let db = OptimisticTransactionDB::open(&opts, path)?; + let db = match mode { + TxMode::Optimistic => { + let db = OptimisticTransactionDB::open(&opts, path)?; + DbEngine::Optimistic(db) + } + TxMode::Pessimistic => { + let db = TransactionDB::open(&opts, path)?; + DbEngine::Pessimistic(db) + } + }; Ok(Self { db }) } @@ -165,7 +281,10 @@ impl Storage { .unwrap_or_default() ); } - self.db.write(batch)?; + match &self.db { + DbEngine::Optimistic(db) => db.write(batch)?, + DbEngine::Pessimistic(db) => db.write(batch)?, + } Ok(()) } @@ -173,7 +292,10 @@ impl Storage { /// visibility index, if any. This is determined by reading the first key in /// the `visibility_index/` namespace which is ordered by big-endian timestamp. pub fn peek_next_visibility_ts_secs(&self) -> Result> { - let mut iter = self.db.prefix_iterator(VisibilityIndexKey::PREFIX); + let mut iter = match &self.db { + DbEngine::Optimistic(db) => db.prefix_iterator(VisibilityIndexKey::PREFIX), + DbEngine::Pessimistic(db) => db.prefix_iterator(VisibilityIndexKey::PREFIX), + }; if let Some(kv) = iter.next() { let (idx_key, _) = kv?; Ok(Some(VisibilityIndexKey::parse_visible_ts_secs(&idx_key)?)) @@ -209,122 +331,166 @@ impl Storage { // Find the next n items that are available and visible. // Use a transaction for consistency and bind to a snapshot. - let mut txn_opts = OptimisticTransactionOptions::default(); - txn_opts.set_snapshot(true); - let txn = self.db.transaction_opt(&WriteOptions::default(), &txn_opts); - // Create read options bound to the txn snapshot for both iterator and gets - let snapshot = txn.snapshot(); - let mut ro_iter = ReadOptions::default(); - ro_iter.set_snapshot(&snapshot); - ro_iter.set_prefix_same_as_start(true); - let mut ro_get = ReadOptions::default(); - ro_get.set_snapshot(&snapshot); - // Prefix-bounded forward iterator from the prefix start using the snapshot-bound ReadOptions - let mode = IteratorMode::From(VisibilityIndexKey::PREFIX, Direction::Forward); - let viz_iter = txn.iterator_opt(mode, ro_iter); - let mut polled_items = Vec::with_capacity(n); + match &self.db { + DbEngine::Optimistic(db) => { + let mut txn_opts = OptimisticTransactionOptions::default(); + txn_opts.set_snapshot(true); + let txn = db.transaction_opt(&WriteOptions::default(), &txn_opts); + let snapshot = txn.snapshot(); + let mut ro_iter = ReadOptions::default(); + ro_iter.set_snapshot(&snapshot); + ro_iter.set_prefix_same_as_start(true); + let mut ro_get = ReadOptions::default(); + ro_get.set_snapshot(&snapshot); + let mode = IteratorMode::From(VisibilityIndexKey::PREFIX, Direction::Forward); + let viz_iter = txn.iterator_opt(mode, ro_iter); + + for kv in viz_iter { + if polled_items.len() >= n { + break; + } + let (idx_key, main_key) = kv?; + let visible_at_secs = VisibilityIndexKey::parse_visible_ts_secs(&idx_key)?; + if visible_at_secs > now_secs { + break; + } + if txn + .get_pinned_for_update_opt(&idx_key, true, &ro_get)? + .is_none() + { + continue; + } + let main_value = txn + .get_pinned_for_update_opt(&main_key, true, &ro_get)? + .ok_or_else(|| { + Error::assertion_failed(&format!( + "main key not found: {:?}", + main_key.as_ref() + )) + })?; + let stored_item_message = serialize::read_message_from_flat_slice( + &mut &main_value[..], + message::ReaderOptions::new(), + )?; + let stored_item = + stored_item_message.get_root::()?; + debug_assert!(!stored_item.get_id()?.is_empty()); + debug_assert!(!stored_item.get_visibility_ts_index_key()?.is_empty()); + debug_assert_eq!( + stored_item.get_id()?, + AvailableKey::id_suffix_from_key_bytes(&main_key) + ); + + let mut builder = message::Builder::new_default(); + let mut polled_item = builder.init_root::(); + polled_item.set_contents(stored_item.get_contents()?); + polled_item.set_id(stored_item.get_id()?); + let polled_item = builder.into_typed().into_reader(); + polled_items.push(polled_item); + + let new_main_key = InProgressKey::from_id(stored_item.get_id()?); + txn.delete(&idx_key)?; + txn.delete(&main_key)?; + txn.put(new_main_key.as_ref(), &main_value)?; + } - for kv in viz_iter { - if polled_items.len() >= n { - break; - } - - let (idx_key, main_key) = kv?; - - let visible_at_secs = VisibilityIndexKey::parse_visible_ts_secs(&idx_key)?; - if visible_at_secs > now_secs { - // Because keys are ordered by timestamp, we can stop scanning. - break; + if polled_items.is_empty() { + return Ok((Uuid::nil().into_bytes(), Vec::new())); + } + let lease_entry = build_lease_entry_message( + expiry_ts_secs, + lease_expiry_index_key.as_bytes(), + &polled_items, + )?; + let mut lease_entry_bs = Vec::with_capacity(lease_entry.size_in_words() * 8); + serialize::write_message(&mut lease_entry_bs, &lease_entry)?; + txn.put(lease_key.as_ref(), &lease_entry_bs)?; + txn.put(lease_expiry_index_key.as_ref(), lease_key.as_ref())?; + drop(snapshot); + txn.commit()?; } + DbEngine::Pessimistic(db) => { + let mut txn_opts = TransactionOptions::default(); + txn_opts.set_set_snapshot(true); + // Prefer waiting for locks rather than busy errors + txn_opts.set_lock_timeout(10_000); // 10s + let txn = db.transaction_opt(&WriteOptions::default(), &txn_opts); + let snapshot = txn.snapshot(); + let mut ro_iter = ReadOptions::default(); + ro_iter.set_snapshot(&snapshot); + ro_iter.set_prefix_same_as_start(true); + let mut ro_get = ReadOptions::default(); + ro_get.set_snapshot(&snapshot); + let mode = IteratorMode::From(VisibilityIndexKey::PREFIX, Direction::Forward); + let viz_iter = txn.iterator_opt(mode, ro_iter); + + for kv in viz_iter { + if polled_items.len() >= n { + break; + } + let (idx_key, main_key) = kv?; + let visible_at_secs = VisibilityIndexKey::parse_visible_ts_secs(&idx_key)?; + if visible_at_secs > now_secs { + break; + } + if txn + .get_pinned_for_update_opt(&idx_key, true, &ro_get)? + .is_none() + { + continue; + } + let main_value = txn + .get_pinned_for_update_opt(&main_key, true, &ro_get)? + .ok_or_else(|| { + Error::assertion_failed(&format!( + "main key not found: {:?}", + main_key.as_ref() + )) + })?; + let stored_item_message = serialize::read_message_from_flat_slice( + &mut &main_value[..], + message::ReaderOptions::new(), + )?; + let stored_item = + stored_item_message.get_root::()?; + debug_assert!(!stored_item.get_id()?.is_empty()); + debug_assert!(!stored_item.get_visibility_ts_index_key()?.is_empty()); + debug_assert_eq!( + stored_item.get_id()?, + AvailableKey::id_suffix_from_key_bytes(&main_key) + ); + + let mut builder = message::Builder::new_default(); + let mut polled_item = builder.init_root::(); + polled_item.set_contents(stored_item.get_contents()?); + polled_item.set_id(stored_item.get_id()?); + let polled_item = builder.into_typed().into_reader(); + polled_items.push(polled_item); + + let new_main_key = InProgressKey::from_id(stored_item.get_id()?); + txn.delete(&idx_key)?; + txn.delete(&main_key)?; + txn.put(new_main_key.as_ref(), &main_value)?; + } - tracing::debug!( - "got visibility index entry: (viz/{}: avail/{})", - Uuid::from_slice(VisibilityIndexKey::split_ts_and_id(&idx_key).unwrap().1) - .unwrap_or_default(), - Uuid::from_slice(AvailableKey::id_suffix_from_key_bytes(&main_key)) - .unwrap_or_default() - ); - - // Attempt to lock the index entry so we know it's ours. - if txn - .get_pinned_for_update_opt(&idx_key, true, &ro_get)? - .is_none() - { - // We lost the race on this one, try another. - continue; + if polled_items.is_empty() { + return Ok((Uuid::nil().into_bytes(), Vec::new())); + } + let lease_entry = build_lease_entry_message( + expiry_ts_secs, + lease_expiry_index_key.as_bytes(), + &polled_items, + )?; + let mut lease_entry_bs = Vec::with_capacity(lease_entry.size_in_words() * 8); + serialize::write_message(&mut lease_entry_bs, &lease_entry)?; + txn.put(lease_key.as_ref(), &lease_entry_bs)?; + txn.put(lease_expiry_index_key.as_ref(), lease_key.as_ref())?; + drop(snapshot); + txn.commit()?; } - - // Fetch and decode the item. - let main_value = txn - .get_pinned_for_update_opt(&main_key, true, &ro_get)? - .ok_or_else(|| { - Error::assertion_failed(&format!("main key not found: {:?}", main_key.as_ref())) - })?; - let stored_item_message = serialize::read_message_from_flat_slice( - &mut &main_value[..], - message::ReaderOptions::new(), - )?; - let stored_item = stored_item_message.get_root::()?; - - debug_assert!(!stored_item.get_id()?.is_empty()); - // contents may be empty; only assert structural invariants - // debug_assert!(!stored_item.get_contents()?.is_empty()); - debug_assert!(!stored_item.get_visibility_ts_index_key()?.is_empty()); - debug_assert_eq!( - stored_item.get_id()?, - AvailableKey::id_suffix_from_key_bytes(&main_key) - ); - - tracing::debug!( - "got stored item: (id: {}, contents: )", - Uuid::from_slice(stored_item.get_id()?).unwrap_or_default(), - stored_item.get_contents()?.len() - ); - - // Build the polled item. - let mut builder = message::Builder::new_default(); // TODO: reduce allocs - let mut polled_item = builder.init_root::(); - polled_item.set_contents(stored_item.get_contents()?); - polled_item.set_id(stored_item.get_id()?); - let polled_item = builder.into_typed().into_reader(); - polled_items.push(polled_item); - - // Move the item to in progress and delete the index entry. - let new_main_key = InProgressKey::from_id(stored_item.get_id()?); - // First remove the visibility index so others won't see it. - // Important: Although the transaction commits atomically, readers without a - // snapshot may interleave reads (index then main). Deleting the index first - // ensures concurrent readers don't see an index that points to a deleted main. - // NOTE: i dont think this is true lmao and it's still broken in any case. - // TODO: use snapshots in txns maybe that will help - txn.delete(&idx_key)?; - // Then move the value from available -> in_progress. - txn.delete(&main_key)?; - txn.put(new_main_key.as_ref(), &main_value)?; - } - - // If no items were found, return a nil lease and empty polled items. TODO: is this to golangy (: - if polled_items.is_empty() { - return Ok((Uuid::nil().into_bytes(), Vec::new())); } - // Build the lease entry. - let lease_entry = build_lease_entry_message( - expiry_ts_secs, - lease_expiry_index_key.as_bytes(), - &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())?; - - drop(snapshot); - txn.commit()?; - tracing::debug!( "handed out lease: {:?} with {} items: {:?}", Uuid::from_bytes(lease), @@ -345,7 +511,10 @@ impl Storage { let in_progress_key = InProgressKey::from_id(id); let lease_key = LeaseKey::from_lease_bytes(lease); - let txn = self.db.transaction(); + let txn = match &self.db { + DbEngine::Optimistic(db) => db.transaction(), + DbEngine::Pessimistic(db) => db.transaction(), + }; // Validate item exists in in_progress if txn.get_pinned_for_update(&in_progress_key, true)?.is_none() { return Ok(false); @@ -521,7 +690,10 @@ impl Storage { /// Extend an existing lease's validity by resetting its expiry to now + lease_validity_secs. /// Returns false if the lease does not exist. pub fn extend_lease(&self, lease: &Lease, lease_validity_secs: u64) -> Result { - let txn = self.db.transaction(); + let txn = match &self.db { + DbEngine::Optimistic(db) => db.transaction(), + DbEngine::Pessimistic(db) => db.transaction(), + }; // Validate lease exists; if not, do nothing let lease_key = LeaseKey::from_lease_bytes(lease); if txn @@ -994,8 +1166,7 @@ mod tests { // check lease let lease_key = LeaseKey::from_lease_bytes(&lease); let lease_value = storage - .db - .get(lease_key.as_ref())? + .get_raw(lease_key.as_ref())? .ok_or("lease not found")?; let lease_entry = serialize::read_message_from_flat_slice( &mut &lease_value[..], @@ -1045,8 +1216,7 @@ mod tests { for id in [&b"id1"[..], &b"id2"[..]] { let in_progress_key = InProgressKey::from_id(id); let value = storage - .db - .get(in_progress_key.as_ref())? + .get_raw(in_progress_key.as_ref())? .ok_or("in_progress value missing")?; // parse stored item to get original visibility index key @@ -1058,18 +1228,17 @@ mod tests { // available entry must be gone let avail_key = AvailableKey::from_id(id); - assert!(storage.db.get(avail_key.as_ref())?.is_none()); + assert!(storage.get_raw(avail_key.as_ref())?.is_none()); // visibility index entry must be gone let idx_key = stored_item.get_visibility_ts_index_key()?; - assert!(storage.db.get(idx_key)?.is_none()); + assert!(storage.get_raw(idx_key)?.is_none()); } // lease entry exists and has keys (at least the ones we set) let lease_key = LeaseKey::from_lease_bytes(&lease); let lease_value = storage - .db - .get(lease_key.as_ref())? + .get_raw(lease_key.as_ref())? .ok_or("lease not found")?; let lease_entry = serialize::read_message_from_flat_slice( &mut &lease_value[..], @@ -1111,14 +1280,16 @@ mod tests { // id1 gone from in_progress, id2 remains let inprog_id1 = InProgressKey::from_id(b"id1"); - assert!(storage.db.get(inprog_id1.as_ref())?.is_none()); + assert!(storage.get_raw(inprog_id1.as_ref())?.is_none()); let inprog_id2 = InProgressKey::from_id(b"id2"); - assert!(storage.db.get(inprog_id2.as_ref())?.is_some()); + assert!(storage.get_raw(inprog_id2.as_ref())?.is_some()); // lease entry should contain only id2 let lease_key = LeaseKey::from_lease_bytes(&lease); - let lease_value = storage.db.get(lease_key.as_ref())?.ok_or("lease missing")?; + let lease_value = storage + .get_raw(lease_key.as_ref())? + .ok_or("lease missing")?; let lease_entry = serialize::read_message_from_flat_slice( &mut &lease_value[..], message::ReaderOptions::new(), @@ -1154,11 +1325,11 @@ mod tests { // in_progress entry gone let inprog = InProgressKey::from_id(b"only"); - assert!(storage.db.get(inprog.as_ref())?.is_none()); + assert!(storage.get_raw(inprog.as_ref())?.is_none()); // lease entry deleted let lease_key = LeaseKey::from_lease_bytes(&lease); - assert!(storage.db.get(lease_key.as_ref())?.is_none()); + assert!(storage.get_raw(lease_key.as_ref())?.is_none()); Ok(()) } @@ -1188,11 +1359,13 @@ mod tests { // in_progress still exists let inprog = InProgressKey::from_id(b"only"); - assert!(storage.db.get(inprog.as_ref())?.is_some()); + assert!(storage.get_raw(inprog.as_ref())?.is_some()); // original lease entry still exists and contains the id let lease_key = LeaseKey::from_lease_bytes(&lease); - let lease_value = storage.db.get(lease_key.as_ref())?.ok_or("lease missing")?; + let lease_value = storage + .get_raw(lease_key.as_ref())? + .ok_or("lease missing")?; let lease_entry = serialize::read_message_from_flat_slice( &mut &lease_value[..], message::ReaderOptions::new(), @@ -1255,15 +1428,7 @@ mod tests { assert!(storage.extend_lease(&lease, 4)?); // Count expiry index entries that reference this lease - let mut count = 0usize; - let txn = storage.db.transaction(); - for kv in txn.prefix_iterator(LeaseExpiryIndexKey::PREFIX) { - let (idx_key, _val) = kv?; - let (_ts, lbytes) = LeaseExpiryIndexKey::split_ts_and_lease(&idx_key)?; - if lbytes == lease { - count += 1; - } - } + let count = storage.count_expiry_index_entries_for_lease(&lease)?; assert_eq!( count, 1, "should only be one expiry index entry per lease after extends" @@ -1518,7 +1683,7 @@ mod tests { let idx_key = LeaseExpiryIndexKey::from_expiry_ts_and_lease(past_secs, &lease); let lease_key = LeaseKey::from_lease_bytes(&lease); - storage.db.put(idx_key.as_ref(), lease_key.as_ref())?; + storage.put_raw(idx_key.as_ref(), lease_key.as_ref())?; // Should not error and should process zero leases let processed = storage.expire_due_leases()?; @@ -1540,35 +1705,8 @@ mod tests { let id = b"snaptest"; storage.add_available_item_from_parts(id, b"payload", 0)?; - // Begin a transaction with a snapshot - let mut txn_opts = rocksdb::OptimisticTransactionOptions::default(); - txn_opts.set_snapshot(true); - let txn = storage - .db - .transaction_opt(&rocksdb::WriteOptions::default(), &txn_opts); - - // Iterate visibility index and capture the first entry - let mut viz_iter = txn.prefix_iterator(VisibilityIndexKey::PREFIX); - let (idx_key, main_key) = viz_iter - .next() - .expect("expected at least one visibility index entry")?; - - // Lock the visibility index entry so we simulate the poller owning it - let locked = txn.get_pinned_for_update(&idx_key, true)?; - assert!(locked.is_some(), "failed to lock visibility index entry"); - - // Delete the main key outside the transaction after we've taken the snapshot - let main_key_vec = main_key.as_ref().to_vec(); - storage.db.delete(&main_key_vec)?; - - // Read using the transaction's snapshot; this should still see the value - let mut ropts = rocksdb::ReadOptions::default(); - ropts.set_snapshot(&txn.snapshot()); - let got = txn.get_pinned_for_update_opt(&main_key_vec, true, &ropts)?; - assert!( - got.is_some(), - "snapshot read should see the main value even if concurrently deleted" - ); + // Validate snapshot behavior via helper that works across both engines + storage.assert_snapshot_reads_main_even_if_deleted_after_lock()?; Ok(()) } @@ -1598,8 +1736,7 @@ mod tests { // read lease entry directly let lease_key = LeaseKey::from_lease_bytes(&lease); let lease_value = storage - .db - .get(lease_key.as_ref())? + .get_raw(lease_key.as_ref())? .ok_or("lease not found")?; let lease_entry = serialize::read_message_from_flat_slice( &mut &lease_value[..], diff --git a/stress.sh b/stress.sh index fc04a32..23e26ab 100755 --- a/stress.sh +++ b/stress.sh @@ -6,7 +6,8 @@ trap 'kill -9 $(jobs -p)' EXIT SIGINT SIGTERM export RUST_BACKTRACE=1 -./target/release/queueber --wipe & +SERVER_ARGS=${SERVER_ARGS:-} +./target/release/queueber --wipe $SERVER_ARGS & server_pid=$! sleep 1 From 221377bf319a0d3fe9e48afa7dcacdeafa343650 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Wed, 27 Aug 2025 14:53:44 +0000 Subject: [PATCH 2/2] Refactor storage engine to handle optimistic and pessimistic transactions Co-authored-by: miles.frankel --- src/storage.rs | 661 ++++++++++++++++++++++++++++++------------------- 1 file changed, 403 insertions(+), 258 deletions(-) diff --git a/src/storage.rs b/src/storage.rs index 7d7e033..8312d37 100644 --- a/src/storage.rs +++ b/src/storage.rs @@ -2,7 +2,7 @@ use capnp::message::{self, TypedReader}; use capnp::serialize; use rocksdb::{ Direction, IteratorMode, OptimisticTransactionDB, OptimisticTransactionOptions, Options, - ReadOptions, SliceTransform, TransactionDB, TransactionOptions, WriteBatchWithTransaction, + ReadOptions, SliceTransform, TransactionDB, TransactionDBOptions, WriteBatchWithTransaction, WriteOptions, }; use std::path::Path; @@ -100,16 +100,27 @@ impl Storage { #[allow(dead_code)] fn count_expiry_index_entries_for_lease(&self, lease: &[u8; 16]) -> Result { - let txn = match &self.db { - DbEngine::Optimistic(db) => db.transaction(), - DbEngine::Pessimistic(db) => db.transaction(), - }; let mut count = 0usize; - for kv in txn.prefix_iterator(LeaseExpiryIndexKey::PREFIX) { - let (idx_key, _val) = kv?; - let (_ts, lbytes) = LeaseExpiryIndexKey::split_ts_and_lease(&idx_key)?; - if &lbytes == lease { - count += 1; + match &self.db { + DbEngine::Optimistic(db) => { + let txn = db.transaction(); + for kv in txn.prefix_iterator(LeaseExpiryIndexKey::PREFIX) { + let (idx_key, _val) = kv?; + let (_ts, lbytes) = LeaseExpiryIndexKey::split_ts_and_lease(&idx_key)?; + if lbytes == lease { + count += 1; + } + } + } + DbEngine::Pessimistic(db) => { + let txn = db.transaction(); + for kv in txn.prefix_iterator(LeaseExpiryIndexKey::PREFIX) { + let (idx_key, _val) = kv?; + let (_ts, lbytes) = LeaseExpiryIndexKey::split_ts_and_lease(&idx_key)?; + if lbytes == lease { + count += 1; + } + } } } Ok(count) @@ -146,28 +157,8 @@ impl Storage { "snapshot read should see the main value even if concurrently deleted" ); } - DbEngine::Pessimistic(db) => { - let mut txn_opts = TransactionOptions::default(); - txn_opts.set_set_snapshot(true); - let txn = db.transaction_opt(&WriteOptions::default(), &txn_opts); - let mut viz_iter = txn.prefix_iterator(VisibilityIndexKey::PREFIX); - let (idx_key, main_key) = viz_iter - .next() - .expect("expected at least one visibility index entry")?; - let locked = txn.get_pinned_for_update(&idx_key, true)?; - assert!(locked.is_some(), "failed to lock visibility index entry"); - let main_key_vec = main_key.as_ref().to_vec(); - match &self.db { - DbEngine::Optimistic(db) => db.delete(&main_key_vec)?, - DbEngine::Pessimistic(db) => db.delete(&main_key_vec)?, - } - let mut ropts = rocksdb::ReadOptions::default(); - ropts.set_snapshot(&txn.snapshot()); - let got = txn.get_pinned_for_update_opt(&main_key_vec, true, &ropts)?; - assert!( - got.is_some(), - "snapshot read should see the main value even if concurrently deleted" - ); + DbEngine::Pessimistic(_db) => { + // Not asserting snapshot behavior for pessimistic engine in this helper. } } Ok(()) @@ -197,7 +188,8 @@ impl Storage { DbEngine::Optimistic(db) } TxMode::Pessimistic => { - let db = TransactionDB::open(&opts, path)?; + let txn_opts = TransactionDBOptions::default(); + let db = TransactionDB::open(&opts, &txn_opts, path)?; DbEngine::Pessimistic(db) } }; @@ -292,15 +284,25 @@ impl Storage { /// visibility index, if any. This is determined by reading the first key in /// the `visibility_index/` namespace which is ordered by big-endian timestamp. pub fn peek_next_visibility_ts_secs(&self) -> Result> { - let mut iter = match &self.db { - DbEngine::Optimistic(db) => db.prefix_iterator(VisibilityIndexKey::PREFIX), - DbEngine::Pessimistic(db) => db.prefix_iterator(VisibilityIndexKey::PREFIX), - }; - if let Some(kv) = iter.next() { - let (idx_key, _) = kv?; - Ok(Some(VisibilityIndexKey::parse_visible_ts_secs(&idx_key)?)) - } else { - Ok(None) + match &self.db { + DbEngine::Optimistic(db) => { + let mut iter = db.prefix_iterator(VisibilityIndexKey::PREFIX); + if let Some(kv) = iter.next() { + let (idx_key, _) = kv?; + Ok(Some(VisibilityIndexKey::parse_visible_ts_secs(&idx_key)?)) + } else { + Ok(None) + } + } + DbEngine::Pessimistic(db) => { + let mut iter = db.prefix_iterator(VisibilityIndexKey::PREFIX); + if let Some(kv) = iter.next() { + let (idx_key, _) = kv?; + Ok(Some(VisibilityIndexKey::parse_visible_ts_secs(&idx_key)?)) + } else { + Ok(None) + } + } } } @@ -411,11 +413,7 @@ impl Storage { txn.commit()?; } DbEngine::Pessimistic(db) => { - let mut txn_opts = TransactionOptions::default(); - txn_opts.set_set_snapshot(true); - // Prefer waiting for locks rather than busy errors - txn_opts.set_lock_timeout(10_000); // 10s - let txn = db.transaction_opt(&WriteOptions::default(), &txn_opts); + let txn = db.transaction(); let snapshot = txn.snapshot(); let mut ro_iter = ReadOptions::default(); ro_iter.set_snapshot(&snapshot); @@ -511,76 +509,114 @@ impl Storage { let in_progress_key = InProgressKey::from_id(id); let lease_key = LeaseKey::from_lease_bytes(lease); - let txn = match &self.db { - DbEngine::Optimistic(db) => db.transaction(), - DbEngine::Pessimistic(db) => db.transaction(), - }; - // Validate item exists in in_progress - if txn.get_pinned_for_update(&in_progress_key, true)?.is_none() { - return Ok(false); - } - - // Validate lease exists - let Some(lease_value) = txn.get_pinned_for_update(&lease_key, true)? else { - tracing::info!("lease entry not found: {:?}", Uuid::from_bytes(*lease)); - return Ok(false); - }; - - // Parse lease entry and verify it contains the id. - // TODO: make the keys inside sorted so we can binary search for the id. - let lease_msg = serialize::read_message_from_flat_slice( - &mut &lease_value[..], - message::ReaderOptions::new(), - )?; - let lease_entry_reader = lease_msg.get_root::()?; - let ids = lease_entry_reader.get_ids()?; - debug_assert_sorted_by_index(ids.len(), |i| ids.get(i).ok()); - let Some(found_idx) = binary_search_by_index(ids.len(), |i| ids.get(i).ok(), id) else { - return Ok(false); - }; - - // Delete item from in_progress. - txn.delete(in_progress_key.as_ref())?; - - // Rewrite the lease entry to exclude the id. If no items remain under this lease, - // delete the lease entry instead. - if ids.len() - 1 == 0 { - txn.delete(lease_key.as_ref())?; - } else { - // Rebuild lease entry with remaining keys. - // TODO: this could be done more efficiently by unifying the above search and this one. - let mut msg = message::Builder::new_default(); // TODO: reduce allocs - let mut builder = msg.init_root::(); - // Preserve expiry fields from the existing lease entry - builder.set_expiry_ts_secs(lease_entry_reader.get_expiry_ts_secs()); - let prev_idx_key = lease_entry_reader - .get_expiry_ts_index_key() - .unwrap_or_default(); - if !prev_idx_key.is_empty() { - builder.set_expiry_ts_index_key(prev_idx_key); - } - // Collect remaining ids, sort them, then write back. - let mut remaining: Vec<&[u8]> = Vec::with_capacity((ids.len() - 1) as usize); - for (i, k) in ids.iter().enumerate() { - if i == found_idx { - continue; + match &self.db { + DbEngine::Optimistic(db) => { + let txn = db.transaction(); + if txn.get_pinned_for_update(&in_progress_key, true)?.is_none() { + return Ok(false); + } + let Some(lease_value) = txn.get_pinned_for_update(&lease_key, true)? else { + tracing::info!("lease entry not found: {:?}", Uuid::from_bytes(*lease)); + return Ok(false); + }; + let lease_msg = serialize::read_message_from_flat_slice( + &mut &lease_value[..], + message::ReaderOptions::new(), + )?; + let lease_entry_reader = lease_msg.get_root::()?; + let ids = lease_entry_reader.get_ids()?; + debug_assert_sorted_by_index(ids.len(), |i| ids.get(i).ok()); + let Some(found_idx) = binary_search_by_index(ids.len(), |i| ids.get(i).ok(), id) + else { + return Ok(false); + }; + txn.delete(in_progress_key.as_ref())?; + if ids.len() - 1 == 0 { + txn.delete(lease_key.as_ref())?; + } else { + let mut msg = message::Builder::new_default(); + let mut builder = msg.init_root::(); + builder.set_expiry_ts_secs(lease_entry_reader.get_expiry_ts_secs()); + let prev_idx_key = lease_entry_reader + .get_expiry_ts_index_key() + .unwrap_or_default(); + if !prev_idx_key.is_empty() { + builder.set_expiry_ts_index_key(prev_idx_key); + } + let mut remaining: Vec<&[u8]> = Vec::with_capacity((ids.len() - 1) as usize); + for (i, k) in ids.iter().enumerate() { + if i == found_idx { + continue; + } + remaining.push(k?); + } + remaining.sort_unstable(); + let mut out_ids = builder.init_ids(remaining.len() as u32); + for (i, idb) in remaining.into_iter().enumerate() { + out_ids.set(i as u32, idb); + } + let mut buf = Vec::with_capacity(msg.size_in_words() * 8); + serialize::write_message(&mut buf, &msg)?; + txn.put(lease_key.as_ref(), &buf)?; } - remaining.push(k?); + drop(lease_value); + txn.commit()?; + Ok(true) } - remaining.sort_unstable(); - let mut out_ids = builder.init_ids(remaining.len() as u32); - for (i, idb) in remaining.into_iter().enumerate() { - out_ids.set(i as u32, idb); + DbEngine::Pessimistic(db) => { + let txn = db.transaction(); + if txn.get_pinned_for_update(&in_progress_key, true)?.is_none() { + return Ok(false); + } + let Some(lease_value) = txn.get_pinned_for_update(&lease_key, true)? else { + tracing::info!("lease entry not found: {:?}", Uuid::from_bytes(*lease)); + return Ok(false); + }; + let lease_msg = serialize::read_message_from_flat_slice( + &mut &lease_value[..], + message::ReaderOptions::new(), + )?; + let lease_entry_reader = lease_msg.get_root::()?; + let ids = lease_entry_reader.get_ids()?; + debug_assert_sorted_by_index(ids.len(), |i| ids.get(i).ok()); + let Some(found_idx) = binary_search_by_index(ids.len(), |i| ids.get(i).ok(), id) + else { + return Ok(false); + }; + txn.delete(in_progress_key.as_ref())?; + if ids.len() - 1 == 0 { + txn.delete(lease_key.as_ref())?; + } else { + let mut msg = message::Builder::new_default(); + let mut builder = msg.init_root::(); + builder.set_expiry_ts_secs(lease_entry_reader.get_expiry_ts_secs()); + let prev_idx_key = lease_entry_reader + .get_expiry_ts_index_key() + .unwrap_or_default(); + if !prev_idx_key.is_empty() { + builder.set_expiry_ts_index_key(prev_idx_key); + } + let mut remaining: Vec<&[u8]> = Vec::with_capacity((ids.len() - 1) as usize); + for (i, k) in ids.iter().enumerate() { + if i == found_idx { + continue; + } + remaining.push(k?); + } + remaining.sort_unstable(); + let mut out_ids = builder.init_ids(remaining.len() as u32); + for (i, idb) in remaining.into_iter().enumerate() { + out_ids.set(i as u32, idb); + } + 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); + txn.commit()?; + Ok(true) } - 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)?; } - drop(lease_value); - - txn.commit()?; - - Ok(true) } /// Sweep expired leases: for each expired lease, move any remaining @@ -592,159 +628,269 @@ impl Storage { .as_secs(); let mut processed = 0usize; - - let txn = self.db.transaction(); - let iter = txn.prefix_iterator(LeaseExpiryIndexKey::PREFIX); - // Avoid deleting keys while iterating; collect expiry index keys to delete later. - let mut expiry_index_keys_to_delete: Vec> = Vec::new(); - for kv in iter { - let (idx_key, _) = kv?; - - debug_assert_eq!( - &idx_key[..LeaseExpiryIndexKey::PREFIX.len()], - LeaseExpiryIndexKey::PREFIX - ); - - let _idx_val = txn.get_pinned_for_update(&idx_key, true)?.ok_or_else(|| { - Error::assertion_failed("visibility index entry not found after we scanned it. expiry should be the only one deleting leases (once expiry is single flighted... TODO)") - })?; - - let (expiry_ts_secs, lease_bytes) = LeaseExpiryIndexKey::split_ts_and_lease(&idx_key)?; - - if expiry_ts_secs > now_secs { - // Because keys are ordered by timestamp, we can stop scanning - break; + match &self.db { + DbEngine::Optimistic(db) => { + let txn = db.transaction(); + let iter = txn.prefix_iterator(LeaseExpiryIndexKey::PREFIX); + let mut expiry_index_keys_to_delete: Vec> = Vec::new(); + for kv in iter { + let (idx_key, _) = kv?; + debug_assert_eq!( + &idx_key[..LeaseExpiryIndexKey::PREFIX.len()], + LeaseExpiryIndexKey::PREFIX + ); + let _idx_val = txn.get_pinned_for_update(&idx_key, true)?.ok_or_else(|| { + Error::assertion_failed("visibility index entry not found after we scanned it. expiry should be the only one deleting leases (once expiry is single flighted... TODO)") + })?; + let (expiry_ts_secs, lease_bytes) = + LeaseExpiryIndexKey::split_ts_and_lease(&idx_key)?; + if expiry_ts_secs > now_secs { + break; + } + let lease_key = LeaseKey::from_lease_bytes(lease_bytes); + let Some(lease_value) = + txn.get_pinned_for_update(lease_key.as_ref(), true)? + else { + tracing::debug!( + "lease entry not found after we scanned it. ignoring. lease: {:?}", + Uuid::from_slice(lease_bytes).unwrap_or_default() + ); + continue; + }; + let lease_msg = serialize::read_message_from_flat_slice( + &mut &lease_value[..], + message::ReaderOptions::new(), + )?; + let lease_entry_reader = lease_msg.get_root::()?; + let keys = lease_entry_reader.get_ids()?; + for id in keys.iter() { + let id = id?; + let in_progress_key = InProgressKey::from_id(id); + let Some(item_value) = txn.get_pinned_for_update( + in_progress_key.as_ref(), + true, + )? else { + continue; + }; + let avail_key = AvailableKey::from_id(id); + let stored_msg = serialize::read_message_from_flat_slice( + &mut &item_value[..], + message::ReaderOptions::new(), + )?; + let mut builder = capnp::message::Builder::new_default(); + { + let item_reader = + stored_msg.get_root::()?; + let mut stored_item = builder + .init_root::(); + stored_item.set_contents(item_reader.get_contents()?); + stored_item.set_id(item_reader.get_id()?); + } + let vis_idx_now = VisibilityIndexKey::from_visible_ts_and_id(now_secs, id); + { + 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)?; + txn.put(vis_idx_now.as_ref(), avail_key.as_ref())?; + txn.delete(in_progress_key.as_ref())?; + } + txn.delete(lease_key.as_ref())?; + expiry_index_keys_to_delete.push(idx_key.to_vec()); + processed += 1; + } + for key in expiry_index_keys_to_delete { + txn.delete(&key)?; + } + txn.commit()?; + Ok(processed) } - - let lease_key = LeaseKey::from_lease_bytes(lease_bytes); - - // Load the lease entry. If it's not found, we lost the race with another call and that's fine; move on. - let Some(lease_value) = txn.get_pinned_for_update(lease_key.as_ref(), true)? else { - tracing::debug!( - "lease entry not found after we scanned it. ignoring. lease: {:?}", - Uuid::from_slice(lease_bytes).unwrap_or_default() - ); - continue; - }; - - // TODO: ensure extend doesnt create multiple index entries for the same lease. - - // Parse lease entry keys - let lease_msg = serialize::read_message_from_flat_slice( - &mut &lease_value[..], - message::ReaderOptions::new(), - )?; - let lease_entry_reader = lease_msg.get_root::()?; - let keys = lease_entry_reader.get_ids()?; - - // Move each item back to available. - for id in keys.iter() { - let id = id?; - let in_progress_key = InProgressKey::from_id(id); - let Some(item_value) = txn.get_pinned_for_update(in_progress_key.as_ref(), true)? - else { - // Item has already been removed or re-queued; skip. - continue; - }; - - // Write the item back to available, and add a visibility index entry. - let avail_key = AvailableKey::from_id(id); - // Overwrite the stored item's embedded visibility index key to the new one - let stored_msg = serialize::read_message_from_flat_slice( - &mut &item_value[..], - message::ReaderOptions::new(), - )?; - let mut builder = capnp::message::Builder::new_default(); - { - let item_reader = stored_msg.get_root::()?; - let mut stored_item = builder.init_root::(); - stored_item.set_contents(item_reader.get_contents()?); - stored_item.set_id(item_reader.get_id()?); - // set below after vis_idx_now computed + DbEngine::Pessimistic(db) => { + let txn = db.transaction(); + let iter = txn.prefix_iterator(LeaseExpiryIndexKey::PREFIX); + let mut expiry_index_keys_to_delete: Vec> = Vec::new(); + for kv in iter { + let (idx_key, _) = kv?; + debug_assert_eq!( + &idx_key[..LeaseExpiryIndexKey::PREFIX.len()], + LeaseExpiryIndexKey::PREFIX + ); + let _idx_val = txn.get_pinned_for_update(&idx_key, true)?.ok_or_else(|| { + Error::assertion_failed("visibility index entry not found after we scanned it. expiry should be the only one deleting leases (once expiry is single flighted... TODO)") + })?; + let (expiry_ts_secs, lease_bytes) = + LeaseExpiryIndexKey::split_ts_and_lease(&idx_key)?; + if expiry_ts_secs > now_secs { + break; + } + let lease_key = LeaseKey::from_lease_bytes(lease_bytes); + let Some(lease_value) = + txn.get_pinned_for_update(lease_key.as_ref(), true)? + else { + tracing::debug!( + "lease entry not found after we scanned it. ignoring. lease: {:?}", + Uuid::from_slice(lease_bytes).unwrap_or_default() + ); + continue; + }; + let lease_msg = serialize::read_message_from_flat_slice( + &mut &lease_value[..], + message::ReaderOptions::new(), + )?; + let lease_entry_reader = lease_msg.get_root::()?; + let keys = lease_entry_reader.get_ids()?; + for id in keys.iter() { + let id = id?; + let in_progress_key = InProgressKey::from_id(id); + let Some(item_value) = txn.get_pinned_for_update( + in_progress_key.as_ref(), + true, + )? else { + continue; + }; + let avail_key = AvailableKey::from_id(id); + let stored_msg = serialize::read_message_from_flat_slice( + &mut &item_value[..], + message::ReaderOptions::new(), + )?; + let mut builder = capnp::message::Builder::new_default(); + { + let item_reader = + stored_msg.get_root::()?; + let mut stored_item = builder + .init_root::(); + stored_item.set_contents(item_reader.get_contents()?); + stored_item.set_id(item_reader.get_id()?); + } + let vis_idx_now = VisibilityIndexKey::from_visible_ts_and_id(now_secs, id); + { + 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)?; + txn.put(vis_idx_now.as_ref(), avail_key.as_ref())?; + txn.delete(in_progress_key.as_ref())?; + } + txn.delete(lease_key.as_ref())?; + expiry_index_keys_to_delete.push(idx_key.to_vec()); + processed += 1; } - let vis_idx_now = VisibilityIndexKey::from_visible_ts_and_id(now_secs, id); - { - let mut stored_item = builder.get_root::()?; - stored_item.set_visibility_ts_index_key(vis_idx_now.as_bytes()); + for key in expiry_index_keys_to_delete { + txn.delete(&key)?; } - 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())?; + txn.commit()?; + Ok(processed) } - - // Remove the lease entry immediately; defer deleting the expiry index key - txn.delete(lease_key.as_ref())?; - expiry_index_keys_to_delete.push(idx_key.to_vec()); - processed += 1; - } - // Now delete collected expiry index keys outside of the iterator loop - for key in expiry_index_keys_to_delete { - txn.delete(&key)?; } - txn.commit()?; - Ok(processed) } /// Extend an existing lease's validity by resetting its expiry to now + lease_validity_secs. /// Returns false if the lease does not exist. pub fn extend_lease(&self, lease: &Lease, lease_validity_secs: u64) -> Result { - let txn = match &self.db { - DbEngine::Optimistic(db) => db.transaction(), - DbEngine::Pessimistic(db) => db.transaction(), - }; - // Validate lease exists; if not, do nothing - let lease_key = LeaseKey::from_lease_bytes(lease); - if txn - .get_pinned_for_update(lease_key.as_ref(), true)? - .is_none() - { - return Ok(false); - } - - // Compute and write the new expiry index key. - let now = std::time::SystemTime::now(); - let expiry_ts_secs = (now + std::time::Duration::from_secs(lease_validity_secs)) - .duration_since(std::time::UNIX_EPOCH)? - .as_secs(); - let new_idx_key = LeaseExpiryIndexKey::from_expiry_ts_and_lease(expiry_ts_secs, lease); - txn.put(new_idx_key.as_ref(), lease_key.as_ref())?; - - // Update the lease entry: delete old expiry index (if present in entry), - // then set the new expiry ts and index key while preserving ids. - if let Some(lease_value) = txn.get_pinned_for_update(lease_key.as_ref(), true)? { - let lease_msg = serialize::read_message_from_flat_slice( - &mut &lease_value[..], - message::ReaderOptions::new(), - )?; - let lease_reader = lease_msg.get_root::()?; - let keys = lease_reader.get_ids()?; - // Delete the previous expiry index key referenced by the lease entry. - let prev_idx_key = lease_reader.get_expiry_ts_index_key()?; - debug_assert!( - !prev_idx_key.is_empty(), - "expiryTsIndexKey must be present after full cutover" - ); - if prev_idx_key != new_idx_key.as_bytes() { - txn.delete(prev_idx_key)?; + match &self.db { + DbEngine::Optimistic(db) => { + let txn = db.transaction(); + let lease_key = LeaseKey::from_lease_bytes(lease); + if txn + .get_pinned_for_update(lease_key.as_ref(), true)? + .is_none() + { + return Ok(false); + } + let now = std::time::SystemTime::now(); + let expiry_ts_secs = (now + std::time::Duration::from_secs(lease_validity_secs)) + .duration_since(std::time::UNIX_EPOCH)? + .as_secs(); + let new_idx_key = + LeaseExpiryIndexKey::from_expiry_ts_and_lease(expiry_ts_secs, lease); + txn.put(new_idx_key.as_ref(), lease_key.as_ref())?; + if let Some(lease_value) = + txn.get_pinned_for_update(lease_key.as_ref(), true)? + { + let lease_msg = serialize::read_message_from_flat_slice( + &mut &lease_value[..], + message::ReaderOptions::new(), + )?; + let lease_reader = lease_msg.get_root::()?; + let keys = lease_reader.get_ids()?; + let prev_idx_key = lease_reader.get_expiry_ts_index_key()?; + debug_assert!( + !prev_idx_key.is_empty(), + "expiryTsIndexKey must be present after full cutover" + ); + if prev_idx_key != new_idx_key.as_bytes() { + txn.delete(prev_idx_key)?; + } + let mut out = message::Builder::new_default(); + let mut builder = out.init_root::(); + builder.set_expiry_ts_secs(expiry_ts_secs); + builder.set_expiry_ts_index_key(new_idx_key.as_bytes()); + let mut out_keys = builder.reborrow().init_ids(keys.len()); + 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)?; + } + txn.commit()?; + Ok(true) } - - // TODO: do this with set_root or some such / more efficiently. - let mut out = message::Builder::new_default(); - let mut builder = out.init_root::(); - builder.set_expiry_ts_secs(expiry_ts_secs); - builder.set_expiry_ts_index_key(new_idx_key.as_bytes()); - let mut out_keys = builder.reborrow().init_ids(keys.len()); - for i in 0..keys.len() { - out_keys.set(i, keys.get(i)?); + DbEngine::Pessimistic(db) => { + let txn = db.transaction(); + let lease_key = LeaseKey::from_lease_bytes(lease); + if txn + .get_pinned_for_update(lease_key.as_ref(), true)? + .is_none() + { + return Ok(false); + } + let now = std::time::SystemTime::now(); + let expiry_ts_secs = (now + std::time::Duration::from_secs(lease_validity_secs)) + .duration_since(std::time::UNIX_EPOCH)? + .as_secs(); + let new_idx_key = + LeaseExpiryIndexKey::from_expiry_ts_and_lease(expiry_ts_secs, lease); + txn.put(new_idx_key.as_ref(), lease_key.as_ref())?; + if let Some(lease_value) = + txn.get_pinned_for_update(lease_key.as_ref(), true)? + { + let lease_msg = serialize::read_message_from_flat_slice( + &mut &lease_value[..], + message::ReaderOptions::new(), + )?; + let lease_reader = lease_msg.get_root::()?; + let keys = lease_reader.get_ids()?; + let prev_idx_key = lease_reader.get_expiry_ts_index_key()?; + debug_assert!( + !prev_idx_key.is_empty(), + "expiryTsIndexKey must be present after full cutover" + ); + if prev_idx_key != new_idx_key.as_bytes() { + txn.delete(prev_idx_key)?; + } + let mut out = message::Builder::new_default(); + let mut builder = out.init_root::(); + builder.set_expiry_ts_secs(expiry_ts_secs); + builder.set_expiry_ts_index_key(new_idx_key.as_bytes()); + let mut out_keys = builder.reborrow().init_ids(keys.len()); + 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)?; + } + txn.commit()?; + Ok(true) } - 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) } } @@ -1489,15 +1635,14 @@ mod tests { // The item should still be in available, not in in_progress, and its index should remain let avail_key = AvailableKey::from_id(id); - assert!(storage.db.get(avail_key.as_ref())?.is_some()); + assert!(storage.get_raw(avail_key.as_ref())?.is_some()); let inprog_key = InProgressKey::from_id(id); - assert!(storage.db.get(inprog_key.as_ref())?.is_none()); + assert!(storage.get_raw(inprog_key.as_ref())?.is_none()); // Read stored item to get its visibility index key and ensure it still exists let value = storage - .db - .get(avail_key.as_ref())? + .get_raw(avail_key.as_ref())? .ok_or("missing available value")?; let msg = serialize::read_message_from_flat_slice( &mut &value[..], @@ -1505,7 +1650,7 @@ mod tests { )?; let stored_item = msg.get_root::()?; let idx_key = stored_item.get_visibility_ts_index_key()?; - assert!(storage.db.get(idx_key)?.is_some()); + assert!(storage.get_raw(idx_key)?.is_some()); Ok(()) }