diff --git a/queueber.capnp b/queueber.capnp index 4eab08b..b36c5ac 100644 --- a/queueber.capnp +++ b/queueber.capnp @@ -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; } diff --git a/src/storage.rs b/src/storage.rs index 8bcf6e7..4180e39 100644 --- a/src/storage.rs +++ b/src/storage.rs @@ -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)?; @@ -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::() + && 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; @@ -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::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[..], @@ -502,10 +506,22 @@ impl Storage { let lease_reader = lease_msg.get_root::()?; 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::(); 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)?); @@ -513,6 +529,9 @@ impl Storage { 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) @@ -526,14 +545,16 @@ type PolledItemOwnedReader = TypedReader, 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> { // 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::(); - 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()?; diff --git a/tests/server_poll.rs b/tests/server_poll.rs index 13b3cd2..2af4921 100644 --- a/tests/server_poll.rs +++ b/tests/server_poll.rs @@ -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; +}