Skip to content
Draft
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
1 change: 1 addition & 0 deletions queueber.capnp
Original file line number Diff line number Diff line change
Expand Up @@ -72,4 +72,5 @@ struct StoredItem {
struct LeaseEntry {
keys @0 :List(Data); # TODO: rename to ids which is what it is rn
expiryTsSecs @1 :UInt64;
leaseExpiryIndexKey @2 :Data;
}
4 changes: 3 additions & 1 deletion src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -211,7 +211,9 @@ impl crate::protocol::queue::Server for Server {
Promise::from_future(async move {
let removed = tokio::task::Builder::new()
.name("remove_in_progress_item")
.spawn_blocking(move || storage.remove_in_progress_item(id_owned.as_slice(), &lease))?
.spawn_blocking(move || {
storage.remove_in_progress_item(id_owned.as_slice(), &lease)
})?
.await
.map_err(Into::<Error>::into)??;

Expand Down
106 changes: 78 additions & 28 deletions src/storage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -90,10 +90,26 @@ impl Storage {

// Atomically insert the item and visibility index entry

let mut batch = WriteBatchWithTransaction::<true>::default();
batch.put(main_key.as_ref(), &stored_contents);
batch.put(visibility_index_key.as_ref(), main_key.as_ref());
self.db.write(batch)?;
// Retry transient RocksDB "Resource busy" errors seen under heavy concurrency in tests
const MAX_ATTEMPTS: usize = 50;
let mut attempt = 0usize;
loop {
let mut batch = WriteBatchWithTransaction::<true>::default();
batch.put(main_key.as_ref(), &stored_contents);
batch.put(visibility_index_key.as_ref(), main_key.as_ref());
match self.db.write(batch) {
Ok(_) => break,
Err(e) => {
let msg = e.as_ref().to_string();
if msg.contains("Resource busy") && attempt < MAX_ATTEMPTS {
attempt += 1;
std::thread::sleep(std::time::Duration::from_millis(1 << attempt));
continue;
}
return Err(Error::from(e));
}
}
}

tracing::debug!(
"inserted item (from parts): ({}: <contents len: {}>), (viz/{}: avail/{})",
Expand Down Expand Up @@ -205,7 +221,6 @@ impl Storage {
let stored_item = stored_item_message.get_root::<protocol::stored_item::Reader>()?;

debug_assert!(!stored_item.get_id()?.is_empty());
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()?,
Expand Down Expand Up @@ -239,7 +254,11 @@ impl Storage {
}

// Build the lease entry.
let lease_entry = build_lease_entry_message(lease_validity_secs, &polled_items)?;
let lease_entry = build_lease_entry_message(
lease_validity_secs,
&polled_items,
lease_expiry_index_key.as_bytes(),
)?;
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)?;

Expand Down Expand Up @@ -301,15 +320,20 @@ impl Storage {
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.
// delete the lease entry and its expiry index entry instead.
if keys.len() - 1 == 0 {
// Delete expiry index entry as well to avoid stale index
let existing_idx_key = lease_entry_reader.get_lease_expiry_index_key()?;
txn.delete(existing_idx_key)?;
txn.delete(lease_key.as_ref())?;
} else {
// Rebuild lease entry with remaining keys.
// Rebuild lease entry with remaining keys, preserving fields
// 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 builder = msg.init_root::<protocol::lease_entry::Builder>();
let mut out_keys = builder.init_keys(keys.len() as u32 - 1);
let mut builder = msg.init_root::<protocol::lease_entry::Builder>();
builder.set_expiry_ts_secs(lease_entry_reader.get_expiry_ts_secs());
builder.set_lease_expiry_index_key(lease_entry_reader.get_lease_expiry_index_key()?);
let mut out_keys = builder.reborrow().init_keys(keys.len() as u32 - 1);
let mut new_idx = 0;
#[allow(clippy::explicit_counter_loop)] // TODO: clean this up
for (_, k) in keys.iter().enumerate().filter(|(i, _)| *i != found_idx) {
Expand Down Expand Up @@ -430,36 +454,26 @@ impl Storage {
.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())?;

// Find current expiry index entry for this lease and delete it.
// TODO: add this index entry to the lease entry so we can do a point lookup.
let mut existing_idx_key: Option<Vec<u8>> = None;
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 {
existing_idx_key = Some(idx_key.to_vec());
break;
}
}
if let Some(k) = existing_idx_key {
txn.delete(&k)?;
}

// Update the lease entry's expiryTsSecs while preserving keys
// Update the lease entry, using the stored old index key to delete it
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::<protocol::lease_entry::Reader>()?;
let keys = lease_reader.get_keys()?;
let old_idx_key = lease_reader.get_lease_expiry_index_key()?;

// TODO: do this with set_root or some such / more efficiently.
// Write new index and delete old index
txn.put(new_idx_key.as_ref(), lease_key.as_ref())?;
txn.delete(old_idx_key)?;

// Rebuild lease entry with updated expiry and index key, preserving keys
let mut out = message::Builder::new_default();
let mut builder = out.init_root::<protocol::lease_entry::Builder>();
builder.set_expiry_ts_secs(expiry_ts_secs);
builder.set_lease_expiry_index_key(new_idx_key.as_bytes());
let mut out_keys = builder.reborrow().init_keys(keys.len());
for i in 0..keys.len() {
out_keys.set(i, keys.get(i)?);
Expand All @@ -481,12 +495,14 @@ type PolledItemOwnedReader =
fn build_lease_entry_message(
lease_validity_secs: u64,
polled_items: &[PolledItemOwnedReader],
lease_expiry_index_key_bytes: &[u8],
) -> Result<capnp::message::Builder<message::HeapAllocator>> {
// Build the lease entry. capnp lists aren't dynamically sized so we
// need to know how many to init before we start writing (?).
let mut lease_entry = message::Builder::new_default();
let mut lease_entry_builder = lease_entry.init_root::<protocol::lease_entry::Builder>();
lease_entry_builder.set_expiry_ts_secs(lease_validity_secs);
lease_entry_builder.set_lease_expiry_index_key(lease_expiry_index_key_bytes);
let mut lease_entry_keys = lease_entry_builder.init_keys(polled_items.len() as u32);
for (i, typed_item) in polled_items.iter().enumerate() {
let item_reader: protocol::polled_item::Reader = typed_item.get()?;
Expand Down Expand Up @@ -764,6 +780,40 @@ mod tests {
Ok(())
}

#[test]
fn extend_deletes_old_expiry_index_entry() -> Result<()> {
let _ = tracing_subscriber::fmt()
.with_max_level(tracing::Level::INFO)
.try_init();

let tmp = tempfile::tempdir().expect("tempdir");
let storage = Storage::new(tmp.path()).expect("storage");

storage.add_available_item_from_parts(b"x1", b"p", 0)?;
let (lease, items) = storage.get_next_available_entries_with_lease(1, 1)?;
assert_eq!(items.len(), 1);

// Extend twice and ensure only one expiry index entry exists for this lease
assert!(storage.extend_lease(&lease, 5)?);
assert!(storage.extend_lease(&lease, 5)?);

// Count expiry index entries for this lease
let mut count = 0usize;
for kv in storage.db.prefix_iterator(LeaseExpiryIndexKey::PREFIX) {
let (idx_key, _val) = kv?;
let (_ts, lbytes) = LeaseExpiryIndexKey::split_ts_and_lease(&idx_key)?;
if lbytes == lease {
count += 1;
}
}
assert_eq!(
count, 1,
"expected exactly one expiry index entry for lease"
);

Ok(())
}

#[test]
fn extend_unknown_lease_returns_false() -> Result<()> {
let _ = tracing_subscriber::fmt()
Expand Down
Loading
Loading