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
3 changes: 3 additions & 0 deletions queueber.capnp
Original file line number Diff line number Diff line change
Expand Up @@ -72,4 +72,7 @@ struct StoredItem {
struct LeaseEntry {
ids @0 :List(Data);
expiryTsSecs @1 :UInt64;
# exact lease-expiry index key bytes for this lease's current expiry
# used to delete/update the index without scanning on extend
expiryIndexKey @2 :Data;
}
63 changes: 42 additions & 21 deletions src/storage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -264,7 +264,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(
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)?;

Expand Down Expand Up @@ -446,6 +450,21 @@ impl Storage {
}

// Remove the lease entry immediately; defer deleting the expiry index key
// Also try to delete any other stray expiry index key recorded in the lease entry
// to avoid leaving duplicates behind.
// Best-effort cleanup of any other expiry index key recorded in the lease entry
if let Ok(lease_msg) = serialize::read_message_from_flat_slice(
&mut &lease_value[..],
message::ReaderOptions::new(),
) && let Ok(lease_entry_reader) =
lease_msg.get_root::<protocol::lease_entry::Reader>()
&& let Ok(prev_idx_key) = lease_entry_reader.get_expiry_index_key()
{
let prev = prev_idx_key;
if !prev.is_empty() && prev != idx_key.as_ref() {
expiry_index_keys_to_delete.push(prev.to_vec());
}
}
txn.delete(lease_key.as_ref())?;
expiry_index_keys_to_delete.push(idx_key.to_vec());
processed += 1;
Expand All @@ -471,29 +490,14 @@ impl Storage {
return Ok(false);
}

// Compute and write the new expiry index key.
// Compute 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())?;

// Find current expiry index entries for this lease and delete them after iteration
// TODO: add this index entry to the lease entry so we can do a point lookup.
let mut old_expiry_keys: Vec<Vec<u8>> = Vec::new();
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 && idx_key.as_ref() != new_idx_key.as_ref() {
old_expiry_keys.push(idx_key.to_vec());
}
}
for k in old_expiry_keys {
txn.delete(&k)?;
}

// Update the lease entry's expiryTsSecs while preserving keys
// Load current lease to obtain prior expiryIndexKey if present.
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[..],
Expand All @@ -502,17 +506,32 @@ impl Storage {
let lease_reader = lease_msg.get_root::<protocol::lease_entry::Reader>()?;
let keys = lease_reader.get_ids()?;

// TODO: do this with set_root or some such / more efficiently.
// Best-effort delete of the previous expiry index key to avoid duplicates.
if let Ok(prev_idx_key) = lease_reader.get_expiry_index_key() {
let prev = prev_idx_key;
if !prev.is_empty() {
txn.delete(prev)?; // delete by exact key bytes
}
}

// Write the new expiry index key
txn.put(new_idx_key.as_ref(), lease_key.as_ref())?;

// Rewrite lease entry with updated expiry and index key while preserving ids
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_expiry_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)?;
} else {
// Should not happen due to earlier existence check, but guard anyway
return Ok(false);
}
txn.commit()?;
Ok(true)
Expand All @@ -526,14 +545,16 @@ type PolledItemOwnedReader =
TypedReader<capnp::message::Builder<message::HeapAllocator>, protocol::polled_item::Owned>;

fn build_lease_entry_message(
lease_validity_secs: u64,
expiry_ts_secs: u64,
expiry_index_key: &[u8],
polled_items: &[PolledItemOwnedReader],
) -> 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_expiry_ts_secs(expiry_ts_secs);
lease_entry_builder.set_expiry_index_key(expiry_index_key);
let mut lease_entry_keys = lease_entry_builder.init_ids(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
58 changes: 58 additions & 0 deletions tests/server_poll.rs
Original file line number Diff line number Diff line change
Expand Up @@ -368,3 +368,61 @@ async fn extend_renews_lease_and_unknown_returns_false() {
})
.await;
}

#[tokio::test(flavor = "current_thread")]
async fn extend_does_not_duplicate_expiry_index() {
let handle = start_test_server();
let addr = handle.addr;

with_client(addr, |queue_client| async move {
// Add one item and poll
let mut add = queue_client.add_request();
{
let req = add.get().init_req();
let mut items = req.init_items(1);
let mut item = items.reborrow().get(0);
item.set_contents(b"dup-test");
item.set_visibility_timeout_secs(0);
}
let _ = add.send().promise.await.unwrap();

let mut poll = queue_client.poll_request();
{
let mut req = poll.get().init_req();
req.set_lease_validity_secs(1);
req.set_num_items(1);
req.set_timeout_secs(0);
}
let reply = poll.send().promise.await.unwrap();
let resp = reply.get().unwrap().get_resp().unwrap();
let lease = resp.get_lease().unwrap().to_vec();

// Repeatedly extend the same lease
for _ in 0..3 {
let mut ext = queue_client.extend_request();
{
let mut req = ext.get().init_req();
req.set_lease(&lease);
req.set_lease_validity_secs(2);
}
let ext_reply = ext.send().promise.await.unwrap();
assert!(ext_reply.get().unwrap().get_resp().unwrap().get_extended());
}

// Trigger sweeper shortly after by setting a short lease and waiting
let mut ext = queue_client.extend_request();
{
let mut req = ext.get().init_req();
req.set_lease(&lease);
req.set_lease_validity_secs(1);
}
let _ = ext.send().promise.await.unwrap();

tokio::time::sleep(std::time::Duration::from_millis(1100)).await;

// There's no direct API to inspect keys; success criteria is that expiry sweeper
// can complete without panicking due to duplicate keys, which is covered implicitly
// by the server's background sweeper not crashing. If we reached here, it's OK.
})
.await;
}