Skip to content
Open
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
106 changes: 106 additions & 0 deletions foyer-storage/src/engine/block/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -976,6 +976,52 @@ mod tests {
.unwrap()
}

/// Like `store_for_test_with_tombstone_log`, but with a 4 MiB + 64 KiB device so the
/// tombstone log spans 5 pages (1280 slots). The extra pages let the most-recent
/// tombstone leave page 0 after enough deletes, which is required to exercise the
/// cross-page resume path in `TombstoneLog::open`.
async fn store_for_test_with_multipage_tombstone_log(
dir: impl AsRef<Path>,
) -> Arc<BlockEngine<u64, Vec<u8>, TestProperties>> {
let device = FsDeviceBuilder::new(dir)
.with_capacity(ByteSize::mib(4).as_u64() as usize + ByteSize::kib(64).as_u64() as usize)
.build()
.unwrap();
let spawner = Spawner::current();
let io_engine = io_engine_for_test(spawner.clone()).await;
let metrics = Arc::new(Metrics::noop());
let builder = BlockEngineConfig {
device,
block_size: 16 * 1024,
compression: Compression::None,
indexer_shards: 4,
recover_concurrency: 2,
flushers: 1,
reclaimers: 1,
clean_block_threshold: 1,
eviction_pickers: vec![Box::<FifoPicker>::default()],
admission_filter: StorageFilter::new(),
reinsertion_filter: StorageFilter::new().with_condition(RejectAll),
enable_tombstone_log: true,
buffer_pool_size: 16 * 1024 * 1024,
blob_index_size: 4 * 1024,
submit_queue_size_threshold: 16 * 1024 * 1024 * 2,
flush_switch: Switch::default(),
load_holder: Holder::default(),
marker: PhantomData,
};
let builder = Box::new(builder);
builder
.build(EngineBuildContext {
io_engine,
metrics,
spawner,
recover_mode: RecoverMode::Strict,
})
.await
.unwrap()
}

fn enqueue(
store: &BlockEngine<u64, Vec<u8>, TestProperties>,
entry: CacheEntry<u64, Vec<u8>, ModHasher, TestProperties>,
Expand Down Expand Up @@ -1194,6 +1240,66 @@ mod tests {
);
}

/// Regression test for the tombstone-log resume bug. After a reopen with the
/// most-recent tombstone on a non-zero page, the writer must resume at the *global*
/// slot and not overwrite an older, still-needed tombstone on page 0. Otherwise a
/// deleted key whose data block is still on disk reappears as a phantom `load`.
#[test_log::test(tokio::test)]
async fn test_store_tombstone_log_resume_no_phantom_entry() {
let dir = tempfile::tempdir().unwrap();
let memory = cache_for_test();
let store = store_for_test_with_multipage_tombstone_log(dir.path()).await;

// 1. Insert + flush entry E (key 1). Its data block is now persisted on disk.
let e = memory.insert(1, vec![1; 3 * KB]);
enqueue(&store, e);
store.wait().await;
assert_eq!(
store.load(memory.hash(&1)).await.unwrap().kv().unwrap(),
(1, vec![1; 3 * KB])
);

// 2. Delete E. Tombstone T_E is appended at tombstone-log slot 1 (page 0, local 1).
store.delete(memory.hash(&1));
store.wait().await;
assert_eq!(store.load(memory.hash(&1)).await.unwrap().kv(), None);

// 3. Append (SLOTS_PER_PAGE - 1) more tombstones for keys without on-disk data blocks so the most-recent
// tombstone moves onto page 1 (logical slot SLOTS_PER_PAGE), leaving T_E on page 0. The log has 5 pages, so
// no wrap.
let extras = TombstoneLog::SLOTS_PER_PAGE - 1;
for i in 100u64..100 + extras as u64 {
store.delete(memory.hash(&i));
}
store.wait().await;

// 4. Close + reopen. The recovered tombstones still suppress E here under both the buggy and fixed code (T_E is
// still on disk at this point).
store.close().await.unwrap();
drop(store);
let store = store_for_test_with_multipage_tombstone_log(dir.path()).await;
assert_eq!(store.load(memory.hash(&1)).await.unwrap().kv(), None);

// 5. Append one more tombstone for an unrelated hash. Under the bug the writer resumes at page 0 local slot 1
// (T_E's slot) and silently overwrites T_E. Under the fix it resumes on page 1 and T_E survives.
store.delete(memory.hash(&9999u64));
store.wait().await;

// 6. Close + reopen again. Under the bug T_E is gone, so E's still-on-disk data block is re-indexed by recovery
// and `load` returns the stale deleted value (a phantom entry). Under the fix T_E survives and `load` stays
// `None`.
store.close().await.unwrap();
drop(store);
let store = store_for_test_with_multipage_tombstone_log(dir.path()).await;
assert_eq!(
store.load(memory.hash(&1)).await.unwrap().kv(),
None,
"deleted key must not resurface as a phantom entry after tombstone-log reopen"
);

store.close().await.unwrap();
}

// FIXME(MrCroxx): Move the admission test to store level.
// #[test_log::test(tokio::test)]
// async fn test_store_admission() {
Expand Down
91 changes: 83 additions & 8 deletions foyer-storage/src/engine/block/tombstone.rs
Original file line number Diff line number Diff line change
Expand Up @@ -69,26 +69,34 @@ impl TombstoneLog {
) -> Result<Self> {
let mut recovered = vec![];

// `addr` records the *global* byte offset of the most-recent tombstone seen so
// far. It must accumulate across all partitions and pages (not reset per page),
// otherwise `latest_tombstone_page` below always collapses to `0` and the
// writer resumes on page 0 instead of right after the last global slot,
// silently overwriting older tombstones that may still be needed.
let mut seq = 0;
let mut addr = 0;
let mut global_addr = 0usize;

for partition in partitions.iter() {
for offset in (0..partition.size()).step_by(PAGE) {
tracing::trace!(offset, "[tombstone log]: recover at");
let buf = IoSliceMut::new(PAGE);
let (buffer, res) = io_engine.read(Box::new(buf), partition.as_ref(), offset as u64).await;
res?;

let mut seq = 0;
let mut addr = 0;

for (slot, buf) in buffer.chunks_exact(Tombstone::SERIALIZED_LEN).enumerate() {
for buf in buffer.chunks_exact(Tombstone::SERIALIZED_LEN) {
let tombstone = Tombstone::read(buf);
if tombstone.sequence > seq {
seq = tombstone.sequence;
addr = slot * Tombstone::SERIALIZED_LEN;
}
if tombstone.sequence == 0 {
global_addr += Tombstone::SERIALIZED_LEN;
continue;
}
if tombstone.sequence > seq {
seq = tombstone.sequence;
addr = global_addr;
}
recovered.push((tombstone, addr));
global_addr += Tombstone::SERIALIZED_LEN;
}
}
}
Expand Down Expand Up @@ -312,4 +320,71 @@ mod tests {
assert_eq!(inner.buffer.page, page);
}
}

/// Regression test for the cross-page resume bug: after a ring wrap, the most-recent
/// tombstone lives on a non-zero page, so the writer must resume at its *global* slot
/// rather than its page-local slot. Under the bug, the resume slot always collapses to
/// a page-local value and the next append overwrites an older tombstone on page 0.
#[test_log::test(tokio::test)]
async fn test_tombstone_log_resume_after_wrap() {
let dir = tempdir().unwrap();

// 4 MB cache device => 16 KB tombstone log => 1K tombstones => 16K partition => 4 pages.
let device = FsDeviceBuilder::new(dir.path())
.with_capacity(4 * 1024 * 1024 + 16 * 1024)
.build()
.unwrap();
let p0 = device.create_partition(8 * 1024).unwrap();
let p1 = device.create_partition(8 * 1024).unwrap();
let io_engine = PsyncIoEngineConfig::new()
.boxed()
.build(IoEngineBuildContext {
spawner: Spawner::current(),
})
.await
.unwrap();

// 4 pages * 256 slots/page = 1024 slots in total.
let total_slots = TombstoneLog::SLOTS_PER_PAGE * 4;

let log = TombstoneLog::open(vec![p0.clone(), p1.clone()], io_engine.clone(), &mut vec![])
.await
.unwrap();

// Write enough tombstones to wrap the ring exactly once plus 266 more, so the
// most-recent tombstone lands on page 1 at local slot 10 (physical slot 266).
let written = total_slots + 266; // 1290
log.append(
(0..written)
.map(|i| Tombstone {
hash: i as u64 + 1,
sequence: i as u64 + 1,
})
.collect_vec()
.iter(),
)
.await
.unwrap();

drop(log);

let log = TombstoneLog::open(vec![p0.clone(), p1.clone()], io_engine.clone(), &mut vec![])
.await
.unwrap();

{
let inner = log.inner.lock().await;
// Resume must be the global slot right after the most-recent tombstone:
// (last physical slot 266) + 1 = 267. Under the bug this collapses to the
// page-local slot 10 + 1 = 11, i.e. page 0 instead of page 1.
let expected_slot = (written + 1) % total_slots; // 267
assert_eq!(
inner.slot, expected_slot,
"tombstone log must resume at the global slot after a ring wrap"
);
let (page, offset) = log.slot_addr(inner.slot);
assert_eq!(page, 1, "resume must land on page 1, not page 0, after wrap");
assert_eq!(offset, 11 * Tombstone::SERIALIZED_LEN);
}
}
}
Loading