use super::{create_compaction_stream, pick_run_indexes};
use crate::{
AbstractTree, Config, KvSeparationOptions, SequenceNumberCounter, Table, TableId,
compaction::{Choice, CompactionStrategy, Input, state::CompactionState},
config::BlockSizePolicy,
version::Version,
};
use std::sync::Arc;
use test_log::test;
const TIGHT_SPACE_KEYS: u64 = 2_000;
fn tight_space_key(i: u64) -> String {
format!("key{i:08}")
}
struct FirstByteComparator;
impl crate::comparator::UserComparator for FirstByteComparator {
fn name(&self) -> &'static str {
"test-first-byte"
}
fn compare(&self, a: &[u8], b: &[u8]) -> core::cmp::Ordering {
a.first().cmp(&b.first())
}
}
#[test]
fn boundary_candidates_dedups_comparator_equal_keys() {
let cmp: crate::comparator::SharedComparator = Arc::new(FirstByteComparator);
let keys = vec![
crate::UserKey::from("a1"),
crate::UserKey::from("a2"),
crate::UserKey::from("b1"),
];
let out = super::boundary_candidates(keys, &cmp);
assert_eq!(
out.len(),
1,
"comparator-equal keys must collapse to a single boundary candidate",
);
assert_eq!(
out.first().and_then(|k| k.first()),
Some(&b'a'),
"the surviving boundary should be from the deduped a-group",
);
}
#[cfg(feature = "parallel")]
#[test]
fn failed_subcompaction_rolls_back_and_restores_inputs() -> crate::Result<()> {
use core::sync::atomic::Ordering;
const N: u64 = 4_000;
let key = |i: u64| format!("key_{i:08}");
let val = |i: u64, generation: u64| format!("g{generation}-{i}-{}", "x".repeat(40));
let dir = tempfile::tempdir()?;
let config = Config::new(
&dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.compaction_threads(4)
.subcompaction_min_bytes(0)
.with_kv_separation(Some(
KvSeparationOptions::default().separation_threshold(16),
));
let failpoint = config.fail_one_subcompaction.clone();
let tree = config.open()?;
for i in 0..N {
tree.insert(key(i), val(i, 0), i);
}
tree.flush_active_memtable(0)?;
tree.major_compact(4_096, 0)?;
for i in 0..N {
tree.insert(key(i), val(i, 1), N + i);
}
tree.flush_active_memtable(0)?;
let tables_before = tree.table_count();
failpoint.store(true, Ordering::SeqCst);
let result = tree.major_compact(u64::MAX, 0);
assert!(
result.is_err(),
"a failing sub-compaction range must abort the compaction",
);
assert!(
!failpoint.load(Ordering::SeqCst),
"the failpoint should have fired and disarmed itself",
);
assert_eq!(
tree.table_count(),
tables_before,
"rollback must leave nothing partially installed",
);
for i in 0..N {
assert_eq!(
tree.get(key(i), crate::MAX_SEQNO)?.as_deref(),
Some(val(i, 1).as_bytes()),
"value for {} must survive the rolled-back compaction",
key(i),
);
}
Ok(())
}
#[test]
fn tight_space_crash_after_first_slice_recovers_all_keys_on_reopen() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let mem = crate::fs::MemFs::with_capacity(u64::MAX);
let reopened = tight_space_crash_and_reopen(
dir.path(),
Arc::new(mem.clone()),
|used| mem.set_capacity(used + used / 4),
|| mem.punched_bytes(),
)?;
for i in 0..TIGHT_SPACE_KEYS {
assert!(
reopened
.get(tight_space_key(i).as_bytes(), crate::MAX_SEQNO)?
.is_some(),
"key {i} lost after a crash mid tight-space compaction + reopen",
);
}
Ok(())
}
#[test]
fn open_rewrites_a_sidecar_that_disagrees_with_the_manifest() -> crate::Result<()> {
use crate::fs::Fs;
let dir = tempfile::tempdir()?;
let mem = crate::fs::MemFs::with_capacity(u64::MAX);
let fs: Arc<dyn Fs> = Arc::new(mem.clone());
let reopened = tight_space_crash_and_reopen(
dir.path(),
Arc::clone(&fs),
|used| mem.set_capacity(used + used / 4),
|| mem.punched_bytes(),
)?;
let version = reopened.current_version();
let Some(restricted) = version
.iter_tables()
.find(|t| t.restrict_lower_bound().is_some())
else {
panic!("the crashed tight-space slice must leave a restricted table");
};
let path = restricted.path.clone();
let table_id = restricted.id();
let Some(bound) = restricted.restrict_lower_bound().cloned() else {
panic!("the table was selected by having a restriction bound");
};
let mut stale = bound.to_vec();
stale.truncate(bound.len().saturating_sub(1));
crate::restrict_bound::write(
&*fs,
&path,
None,
table_id,
&stale,
crate::fs::SyncMode::Normal,
)?;
drop(version);
drop(reopened);
let _tree = Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.with_shared_fs(Arc::new(mem))
.open()?;
let crate::restrict_bound::SidecarRead::Present(id, recorded) =
crate::restrict_bound::read(&*fs, &path, None)?
else {
panic!("the sidecar must be present after the open");
};
assert_eq!(id, table_id, "the sidecar names its own table");
assert_eq!(
recorded,
bound.to_vec(),
"an existing sidecar that disagrees with the manifest is republished \
from the manifest, which is the authority on the bound",
);
Ok(())
}
#[test]
fn open_rewrites_a_missing_restriction_sidecar() -> crate::Result<()> {
use crate::fs::Fs;
let dir = tempfile::tempdir()?;
let mem = crate::fs::MemFs::with_capacity(u64::MAX);
let fs: Arc<dyn Fs> = Arc::new(mem.clone());
let reopened = tight_space_crash_and_reopen(
dir.path(),
Arc::clone(&fs),
|used| mem.set_capacity(used + used / 4),
|| mem.punched_bytes(),
)?;
let version = reopened.current_version();
let Some(restricted) = version
.iter_tables()
.find(|t| t.restrict_lower_bound().is_some())
else {
panic!("the crashed tight-space slice must leave a restricted table");
};
let sidecar = crate::restrict_bound::sidecar_path(&restricted.path);
fs.remove_file(&sidecar)?;
assert!(
fs.metadata(&sidecar).is_err(),
"the fixture must start from a MISSING sidecar",
);
drop(version);
drop(reopened);
let _tree = Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.with_shared_fs(Arc::new(mem))
.open()?;
assert!(
fs.metadata(&sidecar).is_ok(),
"the open must republish the sidecar from the manifest'\u{2019}s own \
restriction, or a later manifest-loss repair republishes the whole \
input beside the slice output",
);
Ok(())
}
mod capfs {
use crate::fs::{Fs, FsCapabilities, FsDirEntry, FsFile, FsMetadata, FsOpenOptions, StdFs};
use crate::io;
use core::sync::atomic::{AtomicU64, Ordering};
use std::io::{Read as _, Seek as _, SeekFrom};
use std::path::Path;
use std::sync::Arc;
#[derive(Clone)]
pub(super) struct CapacityFs {
available: Arc<AtomicU64>,
punched: Arc<AtomicU64>,
link_count: Arc<AtomicU64>,
}
impl CapacityFs {
pub(super) fn new() -> Self {
Self {
available: Arc::new(AtomicU64::new(u64::MAX)),
punched: Arc::new(AtomicU64::new(0)),
link_count: Arc::new(AtomicU64::new(1)),
}
}
pub(super) fn set_link_count(&self, n: u64) {
self.link_count.store(n, Ordering::Relaxed);
}
pub(super) fn set_available_space(&self, bytes: u64) {
self.available.store(bytes, Ordering::Relaxed);
}
pub(super) fn punched_bytes(&self) -> u64 {
self.punched.load(Ordering::Relaxed)
}
}
impl Fs for CapacityFs {
fn open(&self, path: &Path, opts: &FsOpenOptions) -> io::Result<Box<dyn FsFile>> {
StdFs.open(path, opts)
}
fn create_dir_all(&self, path: &Path) -> io::Result<()> {
StdFs.create_dir_all(path)
}
fn read_dir(&self, path: &Path) -> io::Result<Vec<FsDirEntry>> {
StdFs.read_dir(path)
}
fn remove_file(&self, path: &Path) -> io::Result<()> {
StdFs.remove_file(path)
}
fn remove_dir_all(&self, path: &Path) -> io::Result<()> {
StdFs.remove_dir_all(path)
}
fn rename(&self, from: &Path, to: &Path) -> io::Result<()> {
StdFs.rename(from, to)
}
fn metadata(&self, path: &Path) -> io::Result<FsMetadata> {
StdFs.metadata(path)
}
fn sync_directory(&self, path: &Path) -> io::Result<()> {
StdFs.sync_directory(path)
}
fn exists(&self, path: &Path) -> io::Result<bool> {
StdFs.exists(path)
}
fn hard_link_count(&self, _path: &Path) -> io::Result<u64> {
Ok(self.link_count.load(Ordering::Relaxed))
}
fn backend_id(&self) -> Option<u64> {
StdFs.backend_id()
}
fn volume_id(&self, path: &Path) -> Option<u64> {
StdFs.volume_id(path)
}
fn available_space(&self, _path: &Path) -> io::Result<u64> {
Ok(self.available.load(Ordering::Relaxed))
}
fn capabilities(&self, path: &Path) -> FsCapabilities {
FsCapabilities {
punch_hole: true,
..StdFs.capabilities(path)
}
}
fn punch_hole(&self, path: &Path, offset: u64, len: u64) -> io::Result<()> {
if len == 0 {
return Ok(());
}
let mut f = StdFs.open(path, &FsOpenOptions::new().write(true))?;
let file_len = f.metadata()?.len;
let punch_len = len.min(file_len.saturating_sub(offset));
if punch_len == 0 {
return Ok(());
}
f.seek(SeekFrom::Start(offset))?;
std::io::copy(&mut std::io::repeat(0u8).take(punch_len), &mut f)?;
f.sync_all()?;
self.punched.fetch_add(punch_len, Ordering::Relaxed);
Ok(())
}
}
}
fn tight_space_crash_and_reopen(
dir: &std::path::Path,
shared_fs: Arc<dyn crate::fs::Fs>,
set_capacity: impl FnOnce(u64),
punched_bytes: impl Fn() -> u64,
) -> crate::Result<crate::AnyTree> {
use core::sync::atomic::Ordering;
let config = Config::new(
dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.with_shared_fs(Arc::clone(&shared_fs));
let failpoint = config.fail_tight_after_first_slice.clone();
let tree = match config.open()? {
crate::AnyTree::Standard(t) => t,
crate::AnyTree::Blob(_) => panic!("expected Standard tree"),
};
for i in 0..TIGHT_SPACE_KEYS {
tree.insert(tight_space_key(i).as_bytes(), vec![0xCDu8; 64], i);
}
tree.flush_active_memtable(0)?;
let used = tree.storage_stats()?.used_bytes;
set_capacity(used);
tree.update_runtime_config(|c| {
c.storage_admission_check = true;
c.tight_space_compaction = true;
})?;
failpoint.store(true, Ordering::SeqCst);
assert!(
tree.major_compact(64 * 1024 * 1024, 0).is_err(),
"the crash failpoint must abort the tight-space compaction",
);
assert!(
!failpoint.load(Ordering::SeqCst),
"the crash failpoint must have fired and disarmed",
);
assert!(
punched_bytes() > 0,
"the first slice must have punched before the crash",
);
drop(tree);
Config::new(
dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.with_shared_fs(shared_fs)
.open()
}
#[test]
fn manifest_loss_in_the_unwritten_sidecar_window_resolves_to_one_history() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule, Fs};
use crate::io::ErrorKind;
use core::sync::atomic::Ordering;
let dir = tempfile::tempdir()?;
let mem = crate::fs::MemFs::with_capacity(u64::MAX);
let fault = FaultFs::new(mem.clone());
fault.injector().arm(
FaultRule::new(FaultOp::Open, Fault::Error(ErrorKind::PermissionDenied))
.on_path(".restrict-bound"),
);
let config = Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.with_fs(fault);
let failpoint = config.fail_tight_after_first_slice.clone();
let tree = match config.open()? {
crate::AnyTree::Standard(t) => t,
crate::AnyTree::Blob(_) => panic!("expected Standard tree"),
};
for i in 0..TIGHT_SPACE_KEYS {
tree.insert(tight_space_key(i).as_bytes(), vec![0xCDu8; 64], i);
}
tree.flush_active_memtable(0)?;
let used = tree.storage_stats()?.used_bytes;
mem.set_capacity(used + used / 4);
tree.update_runtime_config(|c| {
c.storage_admission_check = true;
c.tight_space_compaction = true;
})?;
failpoint.store(true, Ordering::SeqCst);
assert!(
tree.major_compact(64 * 1024 * 1024, 0).is_err(),
"the crash failpoint must abort the tight-space compaction",
);
assert_eq!(
mem.punched_bytes(),
0,
"with the sidecar write refused, no punch may arm",
);
drop(tree);
for e in mem.read_dir(dir.path())? {
let is_version = e
.file_name
.strip_prefix('v')
.is_some_and(|rest| rest.parse::<u64>().is_ok());
if is_version || e.file_name == "current" {
mem.remove_file(&e.path)?;
}
}
let report = Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.with_shared_fs(Arc::new(mem.clone()))
.repair()?;
assert!(
report
.excluded_files
.iter()
.any(|(_, reason)| reason.contains("derived output")),
"the slice outputs must be excluded as derived — publishing both \
histories would double-apply merge operands: {report:?}",
);
assert_eq!(
report.unreadable, 0,
"a healthy redundancy exclusion is not an unreadable file: {report:?}",
);
let reopened = Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.with_shared_fs(Arc::new(mem))
.open()?;
for i in 0..TIGHT_SPACE_KEYS {
assert!(
reopened
.get(tight_space_key(i).as_bytes(), crate::MAX_SEQNO)?
.is_some(),
"key {i} lost across the sidecar window + manifest loss",
);
}
Ok(())
}
#[test]
fn restricted_view_passes_every_reconcile_gate() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let fs = capfs::CapacityFs::new();
let reopened = tight_space_crash_and_reopen(
dir.path(),
Arc::new(fs.clone()),
|used| fs.set_available_space(used / 4),
|| fs.punched_bytes(),
)?;
let integrity = crate::verify::verify_integrity(&reopened);
assert!(
integrity.is_ok(),
"verify_integrity must pass on a legitimately restricted table, got {:?}",
integrity.errors,
);
let version = reopened.current_version();
let Some(restricted) = version
.iter_tables()
.find(|t| t.restrict_lower_bound().is_some())
else {
panic!("the punched input must reopen as a restricted table");
};
restricted.verify_blob_links()?;
restricted.verify_tli_mirrors()?;
restricted.verify_block_layout()?;
if let Err((gate, e)) = restricted.verify_reconcile_gates(None, false) {
panic!("the healthy restricted suffix must pass every gate, {gate:?} refused it: {e}");
}
Ok(())
}
#[cfg(feature = "page_ecc")]
#[test]
fn unshare_for_heal_reproduces_the_source_faithfully() -> crate::Result<()> {
use crate::fs::{Fs, FsOpenOptions, SyncMode};
let dir = tempfile::tempdir()?;
let fs = capfs::CapacityFs::new();
let reopened = tight_space_crash_and_reopen(
dir.path(),
Arc::new(fs.clone()),
|used| fs.set_available_space(used / 4),
|| fs.punched_bytes(),
)?;
let version = reopened.current_version();
let Some(restricted) = version
.iter_tables()
.find(|t| t.restrict_lower_bound().is_some())
else {
panic!("the punched input must reopen as a restricted table");
};
assert!(
restricted.punch_offset()? > 0,
"the restricted table has a punched prefix",
);
let shared: Arc<dyn Fs> = Arc::new(fs);
let read_all = |path: &std::path::Path| -> crate::Result<alloc::vec::Vec<u8>> {
let f = shared.open(path, &FsOpenOptions::new().read(true))?;
let len = usize::try_from(f.metadata()?.len).unwrap_or(usize::MAX);
let mut buf = alloc::vec![0u8; len];
let mut off = 0usize;
while off < len {
let got = f.read_at(
buf.get_mut(off..).unwrap_or(&mut []),
u64::try_from(off).unwrap_or(u64::MAX),
)?;
if got == 0 {
break;
}
off += got;
}
Ok(buf)
};
let src = read_all(&restricted.path)?;
assert!(src.contains(&0), "the source has punched data blocks");
assert!(src.iter().any(|&b| b != 0), "the source has live blocks");
let source = shared.open(&restricted.path, &FsOpenOptions::new().read(true))?;
let _copy = match restricted.unshare_for_heal(source.as_ref(), SyncMode::Normal) {
Ok(copy) => copy,
Err(e) => panic!("unshare_for_heal must succeed on a restricted table: {e}"),
};
let copy = read_all(&restricted.path)?;
assert_eq!(
copy, src,
"the detached heal copy must reproduce the source byte-for-byte",
);
Ok(())
}
#[test]
fn tight_space_writes_a_restrict_bound_sidecar_matching_the_manifest_bound() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let fs = capfs::CapacityFs::new();
let reopened = tight_space_crash_and_reopen(
dir.path(),
Arc::new(fs.clone()),
|used| fs.set_available_space(used / 4),
|| fs.punched_bytes(),
)?;
let version = reopened.current_version();
let Some(restricted) = version
.iter_tables()
.find(|t| t.restrict_lower_bound().is_some())
else {
panic!("the punched input must reopen as a restricted table");
};
let Some(manifest_bound) = restricted.restrict_lower_bound().cloned() else {
panic!("restricted table has a manifest bound");
};
match crate::restrict_bound::read(&fs, &restricted.path, None)? {
crate::restrict_bound::SidecarRead::Present(_id, bound) => {
assert_eq!(
bound.as_slice(),
manifest_bound.as_ref(),
"the sidecar bound must equal the manifest restriction bound",
);
}
_ => panic!("a punched table must carry a valid .restrict-bound sidecar"),
}
Ok(())
}
#[test]
fn tight_space_sidecar_write_fault_is_nonfatal_and_recovers() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let capfs = capfs::CapacityFs::new();
let fault = FaultFs::new(capfs.clone());
let injector = fault.injector();
let shared: Arc<dyn crate::fs::Fs> = Arc::new(fault);
let config = Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.with_shared_fs(Arc::clone(&shared));
let tree = match config.open()? {
crate::AnyTree::Standard(t) => t,
crate::AnyTree::Blob(_) => panic!("expected Standard tree"),
};
for i in 0..TIGHT_SPACE_KEYS {
tree.insert(tight_space_key(i).as_bytes(), vec![0xCDu8; 64], i);
}
tree.flush_active_memtable(0)?;
let used = tree.storage_stats()?.used_bytes;
capfs.set_available_space(used / 4);
tree.update_runtime_config(|c| {
c.storage_admission_check = true;
c.tight_space_compaction = true;
})?;
injector.arm(
FaultRule::new(FaultOp::Open, Fault::Error(ErrorKind::Other)).on_path("restrict-bound"),
);
let result = match tree.major_compact(64 * 1024 * 1024, 0) {
Ok(r) => r,
Err(e) => panic!("a post-commit sidecar-write fault must not fail the compaction: {e:?}"),
};
assert_eq!(
result.action,
crate::compaction::CompactionAction::Merged,
"tight-space must have engaged and merged, got {result:?}",
);
assert_eq!(
capfs.punched_bytes(),
0,
"an input whose restrict-bound sidecar failed to land must stay unpunched",
);
drop(tree);
let reopened = match Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.with_shared_fs(Arc::clone(&shared))
.open()?
{
crate::AnyTree::Standard(t) => t,
crate::AnyTree::Blob(_) => panic!("expected Standard tree"),
};
for i in 0..TIGHT_SPACE_KEYS {
let got = reopened.get(tight_space_key(i).as_bytes(), crate::MAX_SEQNO)?;
assert_eq!(
got.as_deref(),
Some(&[0xCDu8; 64][..]),
"key {i} must read its latest value after the sidecar-fault reopen",
);
}
Ok(())
}
#[test]
fn tight_space_slice_defers_the_compaction_filter() -> crate::Result<()> {
use crate::compaction::filter::{
CompactionFilter, Context as FilterContext, Factory, ItemAccessor, Verdict,
};
let victim = tight_space_key(100);
struct DropVictim(Vec<u8>);
impl CompactionFilter for DropVictim {
fn filter_item(
&mut self,
item: ItemAccessor<'_>,
_ctx: &FilterContext,
) -> crate::Result<Verdict> {
if item.key().as_ref() == self.0.as_slice() {
Ok(Verdict::Destroy)
} else {
Ok(Verdict::Keep)
}
}
}
struct DropVictimFactory(Vec<u8>);
impl Factory for DropVictimFactory {
fn name(&self) -> &'static str {
"drop-victim"
}
fn make_filter(&self, _ctx: &FilterContext) -> Box<dyn CompactionFilter> {
Box::new(DropVictim(self.0.clone()))
}
}
let dir = tempfile::tempdir()?;
let capfs = capfs::CapacityFs::new();
let shared: Arc<dyn crate::fs::Fs> = Arc::new(capfs.clone());
let config = Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.with_compaction_filter_factory(Some(Arc::new(DropVictimFactory(
victim.clone().into_bytes(),
))))
.with_shared_fs(Arc::clone(&shared));
let tree = match config.open()? {
crate::AnyTree::Standard(t) => t,
crate::AnyTree::Blob(_) => panic!("expected Standard tree"),
};
for i in 0..TIGHT_SPACE_KEYS {
tree.insert(tight_space_key(i).as_bytes(), vec![0xCDu8; 64], i);
}
tree.flush_active_memtable(0)?;
let used = tree.storage_stats()?.used_bytes;
capfs.set_available_space(used / 4);
tree.update_runtime_config(|c| {
c.storage_admission_check = true;
c.tight_space_compaction = true;
})?;
let result = tree.major_compact(64 * 1024 * 1024, 0)?;
assert_eq!(
result.action,
crate::compaction::CompactionAction::Merged,
"tight-space must have engaged and merged, got {result:?}",
);
assert!(
tree.get(victim.as_bytes(), crate::MAX_SEQNO)?.is_some(),
"the tight-space slice must DEFER the compaction filter (its output \
must be a superset shadowing any surviving input prefix); the filter \
applies on the next normal compaction instead",
);
Ok(())
}
#[test]
fn tight_space_sidecar_fault_must_not_resurrect_evicted_deletes() -> crate::Result<()> {
use crate::AbstractTree;
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let capfs = capfs::CapacityFs::new();
let fault = FaultFs::new(capfs.clone());
let injector = fault.injector();
let shared: Arc<dyn crate::fs::Fs> = Arc::new(fault);
let config = Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.with_shared_fs(Arc::clone(&shared));
let tree = match config.open()? {
crate::AnyTree::Standard(t) => t,
crate::AnyTree::Blob(_) => panic!("expected Standard tree"),
};
let victim = tight_space_key(100);
for i in 0..TIGHT_SPACE_KEYS {
tree.insert(tight_space_key(i).as_bytes(), vec![0xCDu8; 64], i);
}
tree.flush_active_memtable(0)?;
tree.remove(victim.as_bytes(), TIGHT_SPACE_KEYS + 1);
tree.flush_active_memtable(0)?;
let used = tree.storage_stats()?.used_bytes;
capfs.set_available_space(used / 4);
tree.update_runtime_config(|c| {
c.storage_admission_check = true;
c.tight_space_compaction = true;
})?;
injector.arm(
FaultRule::new(FaultOp::Open, Fault::Error(ErrorKind::Other)).on_path("restrict-bound"),
);
let result = tree.major_compact(64 * 1024 * 1024, TIGHT_SPACE_KEYS + 2)?;
assert_eq!(
result.action,
crate::compaction::CompactionAction::Merged,
"tight-space must have engaged and merged, got {result:?}",
);
injector.clear();
assert_eq!(
tree.get(victim.as_bytes(), crate::MAX_SEQNO)?,
None,
"the tombstoned key is gone on the live tree",
);
drop(tree);
for e in shared.read_dir(dir.path())? {
let is_version = e
.file_name
.strip_prefix('v')
.is_some_and(|rest| rest.parse::<u64>().is_ok());
if is_version || e.file_name == "current" {
shared.remove_file(&e.path)?;
}
}
Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.with_shared_fs(Arc::clone(&shared))
.repair()?;
let reopened = match Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.with_shared_fs(Arc::clone(&shared))
.open()?
{
crate::AnyTree::Standard(t) => t,
crate::AnyTree::Blob(_) => panic!("expected Standard tree"),
};
assert_eq!(
reopened.get(victim.as_bytes(), crate::MAX_SEQNO)?,
None,
"a tombstone evicted by the last-level slice must never resurrect \
through an unpunched, sidecarless input after a manifest-loss repair",
);
Ok(())
}
#[test]
fn tight_space_install_failure_rolls_back_outputs_and_leaves_no_sidecar() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let capfs = capfs::CapacityFs::new();
let fault = FaultFs::new(capfs.clone());
let injector = fault.injector();
let shared: Arc<dyn crate::fs::Fs> = Arc::new(fault);
let config = Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.with_shared_fs(Arc::clone(&shared));
let tree = match config.open()? {
crate::AnyTree::Standard(t) => t,
crate::AnyTree::Blob(_) => panic!("expected Standard tree"),
};
for i in 0..TIGHT_SPACE_KEYS {
tree.insert(tight_space_key(i).as_bytes(), vec![0xCDu8; 64], i);
}
tree.flush_active_memtable(0)?;
let used = tree.storage_stats()?.used_bytes;
let tables_dir = dir.path().join(crate::file::TABLES_FOLDER);
let numeric_files = || -> Vec<(std::ffi::OsString, u64)> {
std::fs::read_dir(&tables_dir)
.into_iter()
.flatten()
.flatten()
.filter(|e| {
e.file_name()
.to_string_lossy()
.bytes()
.all(|b| b.is_ascii_digit())
})
.filter_map(|e| e.metadata().ok().map(|m| (e.file_name(), m.len())))
.collect()
};
let inputs: crate::HashSet<std::ffi::OsString> =
numeric_files().into_iter().map(|(name, _)| name).collect();
assert!(
!inputs.is_empty(),
"the flush produced at least one input SST"
);
let sidecar_count = || -> usize {
std::fs::read_dir(&tables_dir)
.into_iter()
.flatten()
.flatten()
.filter(|e| e.file_name().to_string_lossy().ends_with(".restrict-bound"))
.count()
};
capfs.set_available_space(used / 4);
tree.update_runtime_config(|c| {
c.storage_admission_check = true;
c.tight_space_compaction = true;
})?;
injector.arm(
FaultRule::new(FaultOp::Write, Fault::Error(ErrorKind::Other))
.on_path("edits")
.once(),
);
let result = tree.major_compact(64 * 1024 * 1024, 0);
assert!(
matches!(&result, Err(crate::Error::Io(e)) if e.kind() == ErrorKind::Other),
"the tight-space compaction must abort with the injected install fault, got {result:?}",
);
for (name, len) in numeric_files() {
assert!(
len == 0 || inputs.contains(&name),
"a finalized slice output ({name:?}, {len} bytes) was orphaned instead of \
rolled back after the install failed",
);
}
assert_eq!(
sidecar_count(),
0,
"an install failure must leave no `.restrict-bound` sidecar (it is written \
only after the install commits)",
);
Ok(())
}
#[test]
fn tight_space_slice_retains_a_tombstone_for_an_unrestricted_repair_survivor() -> crate::Result<()>
{
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
use core::sync::atomic::Ordering;
let dir = tempfile::tempdir()?;
let capfs = capfs::CapacityFs::new();
let fault = FaultFs::new(capfs.clone());
let injector = fault.injector();
let shared: Arc<dyn crate::fs::Fs> = Arc::new(fault);
let config = Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.with_shared_fs(Arc::clone(&shared));
let failpoint = config.fail_tight_after_first_slice.clone();
let tree = match config.open()? {
crate::AnyTree::Standard(t) => t,
crate::AnyTree::Blob(_) => panic!("expected Standard tree"),
};
for i in 0..TIGHT_SPACE_KEYS {
tree.insert(tight_space_key(i).as_bytes(), vec![0xCDu8; 64], i);
}
tree.flush_active_memtable(0)?;
let deleted = tight_space_key(0);
tree.remove(deleted.as_bytes(), TIGHT_SPACE_KEYS);
tree.flush_active_memtable(0)?;
let used = tree.storage_stats()?.used_bytes;
capfs.set_available_space(used / 4);
tree.update_runtime_config(|c| {
c.storage_admission_check = true;
c.tight_space_compaction = true;
})?;
injector.arm(
FaultRule::new(FaultOp::Open, Fault::Error(ErrorKind::Other)).on_path("restrict-bound"),
);
failpoint.store(true, Ordering::SeqCst);
assert!(
tree.major_compact(64 * 1024 * 1024, TIGHT_SPACE_KEYS + 1)
.is_err(),
"the crash failpoint must abort the tight-space compaction",
);
assert!(
!failpoint.load(Ordering::SeqCst),
"the failpoint must have fired after the first slice",
);
assert_eq!(
capfs.punched_bytes(),
0,
"the faulted-sidecar survivor must be left unpunched",
);
injector.clear();
drop(tree);
Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.with_shared_fs(Arc::clone(&shared))
.repair_with_salvage(true)?;
let reopened = match Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.with_shared_fs(shared)
.open()?
{
crate::AnyTree::Standard(t) => t,
crate::AnyTree::Blob(_) => panic!("expected Standard tree"),
};
assert!(
reopened
.get(deleted.as_bytes(), crate::MAX_SEQNO)?
.is_none(),
"a tombstone consumed by a tight-space slice must not be GC'd away, or manifest \
repair of the unrestricted survivor resurrects the deleted key",
);
Ok(())
}
const BLOB_RELOC_KEYS: u64 = 4_000;
const BLOB_RELOC_WATERMARK: u64 = 4 * BLOB_RELOC_KEYS;
fn blob_reloc_key(i: u64) -> String {
format!("key{i:08}")
}
fn blob_reloc_value(i: u64, generation: u8) -> Vec<u8> {
let mut s = (i + 1).wrapping_mul(0x9E37_79B9_7F4A_7C15) ^ (u64::from(generation) << 1);
(0..200u32)
.map(|_| {
s ^= s << 13;
s ^= s >> 7;
s ^= s << 17;
#[expect(
clippy::cast_possible_truncation,
reason = "xorshift byte extraction; the high bits are intentionally dropped"
)]
let byte = (s >> 24) as u8;
byte
})
.collect()
}
fn blob_reloc_kv_options() -> KvSeparationOptions {
KvSeparationOptions::default()
.separation_threshold(64)
.age_cutoff(1.0)
.staleness_threshold(0.1)
.file_target_size(48 * 1024)
}
fn blob_relocation_crash_and_reopen(
dir: &std::path::Path,
mem: &crate::fs::MemFs,
) -> crate::Result<crate::blob_tree::BlobTree> {
use core::sync::atomic::Ordering;
let config = Config::new(
dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.with_shared_fs(Arc::new(mem.clone()))
.with_kv_separation(Some(blob_reloc_kv_options()));
let failpoint = config.fail_tight_after_first_slice.clone();
let tree = match config.open()? {
crate::AnyTree::Blob(t) => t,
crate::AnyTree::Standard(_) => panic!("expected Blob tree"),
};
for i in 0..BLOB_RELOC_KEYS {
tree.insert(blob_reloc_key(i).as_bytes(), blob_reloc_value(i, 1), i);
}
tree.flush_active_memtable(0)?;
for i in (0..BLOB_RELOC_KEYS).step_by(2) {
tree.insert(
blob_reloc_key(i).as_bytes(),
blob_reloc_value(i, 2),
BLOB_RELOC_KEYS + i,
);
}
tree.flush_active_memtable(0)?;
tree.index.update_runtime_config(|c| {
c.storage_admission_check = true;
c.storage_limit_bytes = None;
})?;
tree.major_compact(64 * 1024 * 1024, BLOB_RELOC_WATERMARK)?;
let used = tree.storage_stats()?.used_bytes;
mem.set_capacity(used + used / 4);
tree.index.update_runtime_config(|c| {
c.tight_space_compaction = true;
})?;
failpoint.store(true, Ordering::SeqCst);
assert!(
tree.major_compact(64 * 1024 * 1024, BLOB_RELOC_WATERMARK)
.is_err(),
"the crash failpoint must abort the relocating tight-space compaction",
);
assert!(
!failpoint.load(Ordering::SeqCst),
"the failpoint should have fired and disarmed",
);
assert!(
mem.punched_bytes() > 0,
"the first relocated slice must have punched a stale blob prefix",
);
drop(tree);
match Config::new(
dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.with_kv_separation(Some(blob_reloc_kv_options()))
.with_shared_fs(Arc::new(mem.clone()))
.open()?
{
crate::AnyTree::Blob(t) => Ok(t),
crate::AnyTree::Standard(_) => panic!("expected Blob tree"),
}
}
#[test]
fn tight_space_blob_relocation_crash_after_first_slice_recovers_all_keys() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let mem = crate::fs::MemFs::with_capacity(u64::MAX);
let reopened = blob_relocation_crash_and_reopen(dir.path(), &mem)?;
for i in 0..BLOB_RELOC_KEYS {
let expected = blob_reloc_value(i, u8::from(i % 2 == 0) + 1);
assert_eq!(
reopened
.get(blob_reloc_key(i).as_bytes(), crate::MAX_SEQNO)?
.as_deref(),
Some(expected.as_slice()),
"key {i} wrong/lost after a crash mid blob-relocation + reopen",
);
}
Ok(())
}
#[test]
fn a_relocation_retry_resumes_at_the_committed_blob_frontier() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let mem = crate::fs::MemFs::with_capacity(u64::MAX);
let tree = blob_relocation_crash_and_reopen(dir.path(), &mem)?;
{
let version = tree.index.current_version();
assert!(
version.blob_files.iter().any(|bf| bf.live_data_start() > 0),
"the crashed relocation must leave a blob file with a committed frontier",
);
}
let punched_before = mem.punched_bytes();
tree.index.update_runtime_config(|c| {
c.storage_admission_check = true;
c.tight_space_compaction = true;
})?;
tree.major_compact(64 * 1024 * 1024, BLOB_RELOC_WATERMARK)?;
assert!(
mem.punched_bytes() > punched_before,
"the retry must relocate further and punch what it consumed",
);
for i in 0..BLOB_RELOC_KEYS {
let expected = blob_reloc_value(i, u8::from(i % 2 == 0) + 1);
assert_eq!(
tree.get(blob_reloc_key(i).as_bytes(), crate::MAX_SEQNO)?
.as_deref(),
Some(expected.as_slice()),
"key {i} wrong/lost after the relocation retry",
);
}
Ok(())
}
#[test]
fn tight_space_blob_reopen_failure_rolls_back_the_slice_outputs() -> crate::Result<()> {
use core::sync::atomic::Ordering;
const N: u64 = 4_000;
let k = |i: u64| format!("key{i:08}");
let val = |i: u64, generation: u8| -> Vec<u8> {
let mut s = (i + 1).wrapping_mul(0x9E37_79B9_7F4A_7C15) ^ (u64::from(generation) << 1);
(0..200u32)
.map(|_| {
s ^= s << 13;
s ^= s >> 7;
s ^= s << 17;
#[expect(
clippy::cast_possible_truncation,
reason = "xorshift byte extraction; the high bits are intentionally dropped"
)]
let byte = (s >> 24) as u8;
byte
})
.collect()
};
let dir = tempfile::tempdir()?;
let mem = crate::fs::MemFs::with_capacity(u64::MAX);
let shared: Arc<dyn crate::fs::Fs> = Arc::new(mem.clone());
let config = Config::new(
&dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.with_shared_fs(Arc::clone(&shared))
.with_kv_separation(Some(
KvSeparationOptions::default()
.separation_threshold(64)
.age_cutoff(1.0)
.staleness_threshold(0.1)
.file_target_size(48 * 1024),
));
let failpoint = config.fail_tight_blob_reopen.clone();
let tree = match config.open()? {
crate::AnyTree::Blob(t) => t,
crate::AnyTree::Standard(_) => panic!("expected Blob tree"),
};
for i in 0..N {
tree.insert(k(i).as_bytes(), val(i, 1), i);
}
tree.flush_active_memtable(0)?;
for i in (0..N).step_by(2) {
tree.insert(k(i).as_bytes(), val(i, 2), N + i);
}
tree.flush_active_memtable(0)?;
let gc_watermark = 4 * N;
tree.index.update_runtime_config(|c| {
c.storage_admission_check = true;
c.storage_limit_bytes = None;
})?;
tree.major_compact(64 * 1024 * 1024, gc_watermark)?;
let used = tree.storage_stats()?.used_bytes;
mem.set_capacity(used + used / 4);
tree.index.update_runtime_config(|c| {
c.tight_space_compaction = true;
})?;
let list_names = |folder: &std::path::Path| -> crate::Result<Vec<String>> {
let mut names: Vec<String> = if shared.exists(folder)? {
shared
.read_dir(folder)?
.into_iter()
.map(|e| e.file_name)
.collect()
} else {
Vec::new()
};
names.sort();
Ok(names)
};
let tables_before = list_names(&dir.path().join("tables"))?;
let blobs_before = list_names(&dir.path().join("blobs"))?;
failpoint.store(true, Ordering::SeqCst);
assert!(
tree.major_compact(64 * 1024 * 1024, gc_watermark).is_err(),
"the injected blob-reopen failure must abort the compaction",
);
assert!(
!failpoint.load(Ordering::SeqCst),
"the failpoint should have fired and disarmed",
);
let version = tree.current_version();
let referenced_tables: Vec<String> =
version.iter_tables().map(|t| t.id().to_string()).collect();
let referenced_blobs: Vec<String> = version
.blob_files
.iter()
.map(|bf| bf.id().to_string())
.collect();
for name in list_names(&dir.path().join("tables"))? {
assert!(
tables_before.contains(&name) || referenced_tables.contains(&name),
"leaked table file {name}: neither pre-existing nor referenced by \
the current version (rollback missed it)",
);
}
for name in list_names(&dir.path().join("blobs"))? {
assert!(
blobs_before.contains(&name) || referenced_blobs.contains(&name),
"leaked blob file {name}: neither pre-existing nor referenced by \
the current version (rollback missed it)",
);
}
Ok(())
}
#[test]
fn last_level_applies_and_gcs_below_watermark_range_tombstone() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let tree = Config::new(
&dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.compaction_threads(4)
.subcompaction_min_bytes(0)
.open()?;
let key = |i: u64| format!("k{i:04}");
let val = |i: u64| format!("v{i}-{}", "x".repeat(40));
for i in 0..200u64 {
tree.insert(key(i), val(i), i);
}
tree.flush_active_memtable(0)?;
tree.major_compact(4_096, 0)?;
tree.remove_range(
crate::UserKey::from("k0000"),
crate::UserKey::from("k0050"),
1000,
);
for i in 50..200u64 {
tree.insert(key(i), val(i), 1001 + i);
}
tree.flush_active_memtable(0)?;
tree.major_compact(u64::MAX, 5000)?;
for i in 0..50u64 {
assert_eq!(
tree.get(key(i), crate::MAX_SEQNO)?,
None,
"covered key {} must be physically gone after GC",
key(i),
);
}
for i in 50..200u64 {
assert!(
tree.get(key(i), crate::MAX_SEQNO)?.is_some(),
"uncovered key {} must survive",
key(i),
);
}
let remaining = super::collect_version_tombstones(&tree.current_version());
assert!(
remaining.is_empty(),
"a fully-applied below-watermark tombstone must be GC'd, found {remaining:?}",
);
Ok(())
}
#[test]
fn above_watermark_range_tombstone_is_retained() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let tree = Config::new(
&dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.open()?;
let key = |i: u64| format!("k{i:04}");
for i in 0..50u64 {
tree.insert(key(i), "v", i);
}
tree.flush_active_memtable(0)?;
tree.remove_range(
crate::UserKey::from("k0000"),
crate::UserKey::from("k0025"),
100,
);
tree.flush_active_memtable(0)?;
tree.major_compact(u64::MAX, 50)?;
let remaining = super::collect_version_tombstones(&tree.current_version());
assert!(
!remaining.is_empty(),
"an above-watermark tombstone must be retained, not GC'd",
);
Ok(())
}
#[test]
fn range_tombstone_at_exact_watermark_is_not_applied_or_gced() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let tree = Config::new(
&dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.compaction_threads(4)
.subcompaction_min_bytes(0)
.open()?;
let key = |i: u64| format!("k{i:04}");
let val = |i: u64| format!("v{i}-{}", "x".repeat(40));
for i in 0..200u64 {
tree.insert(key(i), val(i), i);
}
tree.flush_active_memtable(0)?;
tree.major_compact(4_096, 0)?;
tree.remove_range(
crate::UserKey::from("k0000"),
crate::UserKey::from("k0050"),
1000,
);
for i in 50..200u64 {
tree.insert(key(i), val(i), 1001 + i);
}
tree.flush_active_memtable(0)?;
tree.major_compact(u64::MAX, 1000)?;
for i in 0..50u64 {
assert_eq!(
tree.get(key(i), 1000)?.as_deref(),
Some(val(i).as_bytes()),
"covered key {} must survive: RT@watermark is invisible at read==watermark",
key(i),
);
}
let remaining = super::collect_version_tombstones(&tree.current_version());
assert!(
!remaining.is_empty(),
"a tombstone at the exact watermark must be retained, not GC'd",
);
Ok(())
}
#[test]
fn compaction_stream_run_not_found() -> crate::Result<()> {
let folder = tempfile::tempdir()?;
let tree = crate::Config::new(
folder,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.open()?;
tree.insert("a", "a", 0);
tree.flush_active_memtable(0)?;
assert!(
create_compaction_stream(
&tree.current_version(),
&[666],
0,
None,
crate::comparator::default_comparator()
)?
.is_none()
);
Ok(())
}
#[test]
#[expect(clippy::unwrap_used)]
fn compaction_stream_run() -> crate::Result<()> {
let folder = tempfile::tempdir()?;
let tree = crate::Config::new(
folder,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.open()?;
tree.insert("a", "a", 0);
tree.flush_active_memtable(0)?;
tree.insert("b", "b", 0);
tree.flush_active_memtable(0)?;
tree.insert("c", "c", 0);
tree.flush_active_memtable(0)?;
assert_eq!(
Some((0, 2)),
pick_run_indexes(
tree.current_version()
.level(0)
.unwrap()
.iter()
.next()
.unwrap(),
&[0, 1, 2],
)
);
Ok(())
}
#[test]
#[expect(clippy::unwrap_used)]
fn compaction_stream_run_2() -> crate::Result<()> {
let folder = tempfile::tempdir()?;
let tree = crate::Config::new(
folder,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.open()?;
tree.insert("a", "a", 0);
tree.flush_active_memtable(0)?;
tree.insert("b", "b", 0);
tree.flush_active_memtable(0)?;
tree.insert("c", "c", 0);
tree.flush_active_memtable(0)?;
assert_eq!(
Some((0, 0)),
pick_run_indexes(
tree.current_version()
.level(0)
.unwrap()
.iter()
.next()
.unwrap(),
&[0],
)
);
Ok(())
}
#[test]
#[expect(clippy::unwrap_used)]
fn compaction_stream_run_3() -> crate::Result<()> {
let folder = tempfile::tempdir()?;
let tree = crate::Config::new(
folder,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.open()?;
tree.insert("a", "a", 0);
tree.flush_active_memtable(0)?;
tree.insert("b", "b", 0);
tree.flush_active_memtable(0)?;
tree.insert("c", "c", 0);
tree.flush_active_memtable(0)?;
assert_eq!(
Some((2, 2)),
pick_run_indexes(
tree.current_version()
.level(0)
.unwrap()
.iter()
.next()
.unwrap(),
&[2],
)
);
Ok(())
}
#[test]
#[expect(clippy::unwrap_used)]
fn compaction_stream_run_4() -> crate::Result<()> {
let folder = tempfile::tempdir()?;
let tree = crate::Config::new(
folder,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.open()?;
tree.insert("a", "a", 0);
tree.flush_active_memtable(0)?;
tree.insert("b", "b", 0);
tree.flush_active_memtable(0)?;
tree.insert("c", "c", 0);
tree.flush_active_memtable(0)?;
assert_eq!(
None,
pick_run_indexes(
tree.current_version()
.level(0)
.unwrap()
.iter()
.next()
.unwrap(),
&[4],
)
);
Ok(())
}
#[test]
fn compaction_drop_tables() -> crate::Result<()> {
let folder = tempfile::tempdir()?;
let tree = crate::Config::new(
folder,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.open()?;
tree.insert("a", "a", 0);
tree.flush_active_memtable(0)?;
assert_eq!(1, tree.approximate_len());
assert_eq!(0, tree.sealed_memtable_count());
tree.insert("b", "a", 1);
tree.flush_active_memtable(0)?;
assert_eq!(2, tree.approximate_len());
assert_eq!(0, tree.sealed_memtable_count());
tree.insert("c", "a", 2);
tree.flush_active_memtable(0)?;
assert_eq!(3, tree.approximate_len());
assert_eq!(0, tree.sealed_memtable_count());
tree.compact(Arc::new(crate::compaction::Fifo::new(1, None)), 3)?;
assert_eq!(0, tree.table_count());
Ok(())
}
#[test]
fn blob_file_picking_simple() -> crate::Result<()> {
struct InPlaceStrategy(Vec<TableId>);
impl CompactionStrategy for InPlaceStrategy {
fn get_name(&self) -> &'static str {
"InPlaceCompaction"
}
fn choose(&self, _: &Version, _: &Config, _: &CompactionState) -> Choice {
Choice::Merge(Input {
table_ids: self.0.iter().copied().collect(),
dest_level: 6,
target_size: 64_000_000,
canonical_level: 6, })
}
}
let folder = tempfile::tempdir()?;
let tree = crate::Config::new(
folder,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(1))
.with_kv_separation(Some(
KvSeparationOptions::default()
.separation_threshold(1)
.age_cutoff(1.0)
.staleness_threshold(0.01)
.compression(crate::CompressionType::None),
))
.open()?;
tree.insert("a", "a", 0);
tree.insert("b", "b", 0);
tree.insert("c", "c", 0);
tree.flush_active_memtable(1_000)?;
assert_eq!(0, tree.sealed_memtable_count());
assert_eq!(1, tree.table_count());
assert_eq!(1, tree.blob_file_count());
tree.major_compact(1, 1_000)?;
assert_eq!(3, tree.table_count());
assert_eq!(1, tree.blob_file_count());
tree.drop_range("a"..="a")?;
assert_eq!(2, tree.table_count());
assert_eq!(1, tree.blob_file_count());
{
assert_eq!(
&{
let mut map = crate::HashMap::default();
map.insert(0, crate::blob_tree::FragmentationEntry::new(1, 1, 1));
map
},
&**tree.current_version().gc_stats(),
);
}
tree.compact(Arc::new(InPlaceStrategy(vec![2])), 1_000)?;
assert_eq!(2, tree.table_count());
assert_eq!(1, tree.blob_file_count());
{
assert_eq!(
&{
let mut map = crate::HashMap::default();
map.insert(0, crate::blob_tree::FragmentationEntry::new(1, 1, 1));
map
},
&**tree.current_version().gc_stats(),
);
}
tree.compact(Arc::new(InPlaceStrategy(vec![3, 4])), 1_000)?;
assert_eq!(1, tree.table_count());
assert_eq!(1, tree.blob_file_count());
{
assert_eq!(
crate::HashMap::default(),
**tree.current_version().gc_stats(),
);
}
Ok(())
}
#[expect(
clippy::expect_used,
clippy::indexing_slicing,
reason = "test asserts over known-good fixtures; failure surfaces via panic"
)]
#[test]
fn narrow_merge_candidates_for_full_run_are_adjacent_pairs_sorted_ascending() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let tree = Config::new(
&dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.open()?;
for i in 0..3_000u64 {
tree.insert(format!("k{i:08}"), "v".repeat(40), i);
}
tree.flush_active_memtable(0)?;
tree.major_compact(16 * 1024, 0)?;
let version = tree.current_version();
let run = version
.iter_levels()
.flat_map(|level| level.iter())
.find(|run| run.len() >= 3)
.expect("a bottom-level run with >= 3 tables");
let ordered: Vec<(TableId, u64)> = run
.iter()
.map(|t| Ok((t.id(), t.live_file_size()?)))
.collect::<crate::Result<_>>()?;
let payload = Input {
table_ids: ordered.iter().map(|(id, _)| *id).collect(),
dest_level: 6,
canonical_level: 6,
target_size: 64 * 1024 * 1024,
};
let candidates = super::narrow_merge_candidates(&version, &payload)?;
assert_eq!(
candidates.len(),
ordered.len() - 1,
"one candidate per run-adjacent pair"
);
for c in &candidates {
assert_eq!(c.table_ids.len(), 2, "each candidate is an adjacent pair");
assert_eq!(c.dest_level, 6, "destination preserved");
}
let combined = |c: &Input| -> crate::Result<u64> {
c.table_ids
.iter()
.filter_map(|id| version.get_table(*id))
.try_fold(0u64, |acc, t| t.live_file_size().map(|size| acc + size))
};
let sums: Vec<u64> = candidates
.iter()
.map(combined)
.collect::<crate::Result<_>>()?;
let mut sorted = sums.clone();
sorted.sort_unstable();
assert_eq!(sums, sorted, "candidates sorted ascending by SST size");
let smallest_pair = ordered
.windows(2)
.map(|w| w[0].1 + w[1].1)
.min()
.expect(">= 2 tables");
assert_eq!(sums[0], smallest_pair, "smallest-Σ pair is tried first");
Ok(())
}
#[test]
fn space_fits_two_layer_combines_shared_volume_outputs_and_separates_routed_ones()
-> crate::Result<()> {
use crate::fs::MemFs;
const MIB: u64 = 1024 * 1024;
let dir = tempfile::tempdir()?;
let cfg = Config::new(
&dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.with_shared_fs(Arc::new(MemFs::with_capacity(100 * MIB)));
assert!(
!super::space_fits_two_layer(&cfg, u64::MAX, 60 * MIB, 6, 60 * MIB),
"shared-volume outputs must be summed, not checked independently"
);
assert!(super::space_fits_two_layer(
&cfg,
u64::MAX,
60 * MIB,
6,
30 * MIB
));
assert!(!super::space_fits_two_layer(
&cfg,
80 * MIB,
50 * MIB,
6,
40 * MIB
));
let routed = Config::new(
&dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.with_shared_fs(Arc::new(MemFs::with_capacity(100 * MIB)))
.level_routes(vec![crate::config::LevelRoute {
levels: 6..7,
path: crate::path::PathBuf::from("/cold-tier"),
fs: Arc::new(MemFs::with_capacity(100 * MIB)),
}]);
assert!(
super::space_fits_two_layer(&routed, u64::MAX, 60 * MIB, 6, 60 * MIB),
"proven-independent volumes are checked independently"
);
assert!(!super::space_fits_two_layer(
&routed,
u64::MAX,
60 * MIB,
6,
130 * MIB
));
let shared: Arc<dyn crate::fs::Fs> = Arc::new(MemFs::with_capacity(100 * MIB));
let routed_same_mount = Config::new(
&dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.with_shared_fs(Arc::clone(&shared))
.level_routes(vec![crate::config::LevelRoute {
levels: 6..7,
path: crate::path::PathBuf::from("/same-mount-subdir"),
fs: Arc::clone(&shared),
}]);
assert!(
!super::space_fits_two_layer(&routed_same_mount, u64::MAX, 60 * MIB, 6, 60 * MIB),
"a route on the same volume must combine budgets, not admit each independently"
);
Ok(())
}
#[expect(
clippy::expect_used,
reason = "test asserts over known-good fixtures; failure surfaces via panic"
)]
#[test]
fn space_gate_for_merge_narrows_a_full_run_that_exceeds_free() -> crate::Result<()> {
use crate::fs::MemFs;
let dir = tempfile::tempdir()?;
let mem = MemFs::with_capacity(u64::MAX);
let any = Config::new(
&dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.with_shared_fs(Arc::new(mem.clone()))
.data_block_size_policy(BlockSizePolicy::all(512))
.open()?;
let crate::AnyTree::Standard(tree) = any else {
panic!("expected Standard tree");
};
for i in 0..3_000u64 {
tree.insert(format!("k{i:08}"), "v".repeat(40), i);
}
tree.flush_active_memtable(0)?;
tree.major_compact(16 * 1024, 0)?;
let version = tree.current_version();
let run = version
.iter_levels()
.flat_map(|level| level.iter())
.find(|run| run.len() >= 3)
.expect("a bottom-level run with >= 3 tables");
let run_sigma: u64 = run.iter().map(Table::file_size).sum();
let payload = Input {
table_ids: run.iter().map(Table::id).collect(),
dest_level: 6,
canonical_level: 6,
target_size: 64 * 1024 * 1024,
};
let probe_capacity = 1u64 << 40;
mem.set_capacity(probe_capacity);
let stored =
probe_capacity - crate::fs::Fs::available_space(&mem, dir.path()).unwrap_or(probe_capacity);
mem.set_capacity(stored + run_sigma - 1);
tree.update_runtime_config(|c| {
c.storage_admission_check = true;
c.storage_limit_bytes = None;
})?;
let opts = super::Options::from_tree(
&tree,
Arc::new(crate::compaction::major::Strategy::new(64 * 1024 * 1024)),
);
match super::space_gate_for_merge(&version, &opts, &payload)? {
super::SpaceGate::Narrowed(narrowed) => {
assert_eq!(narrowed.table_ids.len(), 2, "narrowed to an adjacent pair");
}
super::SpaceGate::Run => {
panic!("expected Narrowed, got Run (full run wrongly admitted)")
}
super::SpaceGate::Skip => panic!("expected Narrowed, got Skip (no pair admitted)"),
}
Ok(())
}
#[test]
fn a_retry_never_lowers_an_existing_restriction() -> crate::Result<()> {
use core::sync::atomic::Ordering;
let dir = tempfile::tempdir()?;
let mem = crate::fs::MemFs::with_capacity(u64::MAX);
let reopened = tight_space_crash_and_reopen(
dir.path(),
Arc::new(mem.clone()),
|used| mem.set_capacity(used + used / 4),
|| mem.punched_bytes(),
)?;
let (restricted_id, committed_bound) = {
let version = reopened.current_version();
let Some(table) = version
.iter_tables()
.find(|t| t.restrict_lower_bound().is_some())
else {
panic!("the crashed tight-space slice must leave a restricted table");
};
let Some(bound) = table.restrict_lower_bound().cloned() else {
panic!("the table was selected by having a restriction bound");
};
(table.id(), bound)
};
drop(reopened);
let config = Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.with_shared_fs(Arc::new(mem));
let failpoint = config.fail_tight_after_first_slice.clone();
let tree = match config.open()? {
crate::AnyTree::Standard(t) => t,
crate::AnyTree::Blob(_) => panic!("expected Standard tree"),
};
tree.update_runtime_config(|c| {
c.storage_admission_check = true;
c.tight_space_compaction = true;
})?;
failpoint.store(true, Ordering::SeqCst);
let _ = tree.major_compact(64 * 1024 * 1024, 0);
let version = tree.current_version();
let Some(after) = version.get_table(restricted_id) else {
panic!("the retry must still hold the restricted input, not consume it");
};
let Some(bound) = after.restrict_lower_bound() else {
panic!("a table that was restricted can never come back unrestricted");
};
assert!(
bound.as_ref() >= committed_bound.as_ref(),
"a retry must not serve a prefix an earlier slice already punched: \
bound went from {committed_bound:?} to {bound:?}",
);
let crate::restrict_bound::SidecarRead::Present(_, recorded) =
crate::restrict_bound::read(&*after.fs, &after.path, after.encryption.as_deref())?
else {
panic!("the committed restriction must have published its sidecar");
};
assert_eq!(
recorded,
bound.to_vec(),
"the sidecar must record the bound the manifest committed",
);
Ok(())
}
#[test]
fn tight_space_slices_when_only_the_quota_constrains() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
let dir = tempfile::tempdir()?;
let mem = crate::fs::MemFs::with_capacity(u64::MAX);
let faulty = FaultFs::new(mem.clone());
faulty.injector().arm(FaultRule::new(
FaultOp::AvailableSpace,
Fault::Error(crate::io::ErrorKind::Unsupported),
));
let tree = match Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.with_shared_fs(Arc::new(faulty))
.open()?
{
crate::AnyTree::Standard(t) => t,
crate::AnyTree::Blob(_) => panic!("expected Standard tree"),
};
for i in 0..TIGHT_SPACE_KEYS {
tree.insert(tight_space_key(i).as_bytes(), vec![0xCDu8; 64], i);
}
tree.flush_active_memtable(0)?;
let used = crate::storage_stats::compute_used_bytes(&tree.current_version())?;
tree.update_runtime_config(|c| {
c.storage_admission_check = true;
c.tight_space_compaction = true;
c.storage_limit_bytes = Some(used + used / 4);
})?;
tree.major_compact(64 * 1024 * 1024, 0)?;
assert!(
mem.punched_bytes() > 0,
"the slices must reclaim a prefix; a quota-only tree has no other way \
back from read-only",
);
Ok(())
}
#[test]
fn full_range_selectivity_is_one_over_a_restricted_view() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let mem = crate::fs::MemFs::with_capacity(u64::MAX);
let reopened = tight_space_crash_and_reopen(
dir.path(),
Arc::new(mem.clone()),
|used| mem.set_capacity(used + used / 4),
|| mem.punched_bytes(),
)?;
{
let version = reopened.current_version();
assert!(
version
.iter_tables()
.any(|t| t.restrict_lower_bound().is_some()),
"the crashed tight-space slice must leave a restricted table",
);
}
let card = reopened.approximate_range_cardinality::<&[u8], _>(.., crate::SeqNo::MAX)?;
assert!(
(card.selectivity - 1.0).abs() < 1e-9,
"a full-keyspace query selects everything the snapshot sees, got {}",
card.selectivity,
);
Ok(())
}
#[test]
fn space_gate_sizes_a_restricted_input_by_its_live_suffix() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let mem = crate::fs::MemFs::with_capacity(u64::MAX);
let reopened = tight_space_crash_and_reopen(
dir.path(),
Arc::new(mem.clone()),
|used| mem.set_capacity(used + used / 4),
|| mem.punched_bytes(),
)?;
let crate::AnyTree::Standard(tree) = reopened else {
panic!("expected Standard tree");
};
mem.set_capacity(u64::MAX);
let version = tree.current_version();
let Some(restricted) = version
.iter_tables()
.find(|t| t.restrict_lower_bound().is_some())
else {
panic!("the crashed tight-space slice must leave a restricted table");
};
let live = restricted.live_file_size()?;
assert!(
live < restricted.file_size(),
"the fixture's input must carry a punched, superseded prefix",
);
let payload = Input {
table_ids: core::iter::once(restricted.id()).collect(),
dest_level: 6,
canonical_level: 6,
target_size: 64 * 1024 * 1024,
};
let used = crate::storage_stats::compute_used_bytes(&version)?;
tree.update_runtime_config(|c| {
c.storage_admission_check = true;
c.tight_space_compaction = false;
c.storage_limit_bytes = Some(used + live);
})?;
let opts = super::Options::from_tree(
&tree,
Arc::new(crate::compaction::major::Strategy::new(64 * 1024 * 1024)),
);
match super::space_gate_for_merge(&version, &opts, &payload)? {
super::SpaceGate::Run => Ok(()),
super::SpaceGate::Narrowed(_) => {
panic!("expected Run, got Narrowed (a single-table merge cannot narrow)")
}
super::SpaceGate::Skip => {
panic!("expected Run, got Skip (the punched prefix was charged to the output)")
}
}
}
#[test]
fn compaction_outputs_inherit_the_deletion_pause() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let tree = match Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.open()?
{
crate::AnyTree::Standard(t) => t,
crate::AnyTree::Blob(_) => panic!("expected Standard tree"),
};
for i in 0..100u64 {
tree.insert(format!("k{i:05}").as_bytes(), vec![1u8; 16], i);
}
tree.flush_active_memtable(0)?;
for i in 100..200u64 {
tree.insert(format!("k{i:05}").as_bytes(), vec![1u8; 16], i);
}
tree.flush_active_memtable(0)?;
tree.major_compact(64 * 1024 * 1024, 0)?;
let version = tree.current_version();
let tables: Vec<_> = version.iter_tables().collect();
assert!(
!tables.is_empty(),
"the major compaction produced an output"
);
for t in &tables {
assert!(
t.deletion_pause.get().is_some(),
"compaction output table {} must inherit the deletion pause",
t.id(),
);
}
Ok(())
}
fn punch_blob_prefix_on_drop(
capfs: &capfs::CapacityFs,
link_count: u64,
pause_active: bool,
) -> crate::Result<(u64, u64)> {
use crate::{Config, KvSeparationOptions, SequenceNumberCounter};
let dir = tempfile::tempdir()?;
let tree = match Config::new(
&dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.with_shared_fs(Arc::new(capfs.clone()))
.with_kv_separation(Some(
KvSeparationOptions::default().separation_threshold(64),
))
.open()?
{
crate::AnyTree::Blob(t) => t,
crate::AnyTree::Standard(_) => panic!("expected Blob tree"),
};
for i in 0..64u64 {
tree.insert(format!("k{i:05}").as_bytes(), vec![b'v'; 256], i);
}
tree.flush_active_memtable(0)?;
let blob = {
let version = tree.index.current_version();
let Some(bf) = version.blob_files.iter().next().cloned() else {
panic!("the flush spilled at least one blob file");
};
bf
};
let physical = blob.physical_size()?;
assert!(physical > 0, "the blob file has bytes to reclaim");
capfs.set_link_count(link_count);
let pause_guard = if pause_active {
Some(tree.index.deletion_pause.acquire())
} else {
None
};
blob.mark_punch_on_drop(physical);
drop(blob);
drop(tree);
let during_pause = capfs.punched_bytes();
drop(pause_guard);
Ok((during_pause, capfs.punched_bytes()))
}
#[test]
fn reclaimed_blob_file_verifies_against_its_live_suffix() -> crate::Result<()> {
use crate::fs::Fs as _;
let dir = tempfile::tempdir()?;
let capfs = capfs::CapacityFs::new();
let tree = match Config::new(
&dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.with_shared_fs(Arc::new(capfs.clone()))
.with_kv_separation(Some(
KvSeparationOptions::default().separation_threshold(64),
))
.open()?
{
crate::AnyTree::Blob(t) => t,
crate::AnyTree::Standard(_) => panic!("expected Blob tree"),
};
for i in 0..64u64 {
tree.insert(format!("k{i:05}").as_bytes(), vec![b'v'; 256], i);
}
tree.flush_active_memtable(0)?;
let blob = {
let version = tree.index.current_version();
let Some(bf) = version.blob_files.iter().next().cloned() else {
panic!("the flush spilled at least one blob file");
};
bf
};
let frontier = blob.physical_size()? / 2;
let restricted = blob.reopen_restricted(frontier)?;
capfs.punch_hole(blob.path(), 0, frontier)?;
let got = crate::verify::stream_checksum_from(restricted.path(), restricted.live_data_start())?;
assert_eq!(
got,
restricted.checksum(),
"a reclaimed blob file must verify against its live suffix, not the \
punched prefix",
);
Ok(())
}
#[test]
fn blob_prefix_reclaim_skips_a_hard_linked_inode() -> crate::Result<()> {
let (owned, _) = punch_blob_prefix_on_drop(&capfs::CapacityFs::new(), 1, false)?;
assert!(
owned > 0,
"control: an exclusively-owned blob file must have its prefix reclaimed",
);
let (shared, _) = punch_blob_prefix_on_drop(&capfs::CapacityFs::new(), 2, false)?;
assert_eq!(
shared, 0,
"a blob file shared with a checkpoint must not be punched: the snapshot \
still references values in the reclaimed prefix",
);
Ok(())
}
#[test]
fn blob_prefix_reclaim_defers_across_a_deletion_pause() -> crate::Result<()> {
let (during, after) = punch_blob_prefix_on_drop(&capfs::CapacityFs::new(), 1, true)?;
assert_eq!(
during, 0,
"an active checkpoint pause must defer the reclaim, closing the \
probe-to-punch race against the checkpoint's link pass",
);
assert!(
after > 0,
"releasing the pause must carry out the deferred reclaim, not drop it",
);
Ok(())
}
#[test]
fn a_reopened_blob_view_needs_binding_to_carry_the_deletion_pause() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let tree = match Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.with_kv_separation(Some(
KvSeparationOptions::default().separation_threshold(16),
))
.open()?
{
crate::AnyTree::Blob(t) => t,
crate::AnyTree::Standard(_) => panic!("expected a blob tree"),
};
for i in 0..32u64 {
tree.insert(format!("key{i:06}").as_bytes(), vec![b'v'; 64], i);
}
tree.flush_active_memtable(0)?;
let original = {
let binding = tree.index.version_history.read().latest_version();
binding
.version
.blob_files
.iter()
.next()
.cloned()
.ok_or(crate::Error::Unrecoverable)?
};
assert!(
original.deletion_pause_for_test().is_some(),
"a flush-published blob file is bound when it is registered",
);
let reopened = original.reopen_restricted(0)?;
assert!(
reopened.deletion_pause_for_test().is_none(),
"the re-open starts from the file on disk, so it inherits nothing — \
the path that publishes it is what has to bind it",
);
reopened.bind_to_tree(&crate::table::TableSinks {
deletion_pause: &tree.index.deletion_pause,
heal_hints: &tree.index.heal_hints,
#[cfg(feature = "std")]
background_deleter: None,
});
assert!(
reopened.deletion_pause_for_test().is_some(),
"binding is what makes a re-opened view safe to publish",
);
Ok(())
}
#[test]
fn tight_space_declines_when_an_input_backend_cannot_punch() -> crate::Result<()> {
use crate::config::LevelRoute;
use crate::fs::{Fs, MemFs};
let dir = tempfile::tempdir()?;
let hot_dir = dir.path().join("hot");
let hot = Arc::new(MemFs::with_capacity(u64::MAX));
hot.set_punch_hole_supported(false);
hot.create_dir_all(&hot_dir)?;
let main = Arc::new(MemFs::with_capacity(u64::MAX));
let hot_fs: Arc<dyn Fs> = hot.clone();
let tree = match Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.with_shared_fs(main.clone())
.level_routes(vec![LevelRoute {
levels: 0..1,
path: hot_dir.clone(),
fs: hot_fs,
}])
.open()?
{
crate::AnyTree::Standard(t) => t,
crate::AnyTree::Blob(_) => panic!("expected Standard tree"),
};
for i in 0..TIGHT_SPACE_KEYS {
tree.insert(tight_space_key(i).as_bytes(), vec![0xCDu8; 64], i);
}
tree.flush_active_memtable(0)?;
main.set_capacity(64 * 1024);
tree.update_runtime_config(|c| {
c.storage_admission_check = true;
c.tight_space_compaction = true;
})?;
tree.major_compact(64 * 1024 * 1024, 0)?;
let sidecars = hot
.read_dir(&hot_dir.join("tables"))?
.into_iter()
.filter(|e| e.file_name.contains("restrict-bound"))
.count();
assert_eq!(
sidecars, 0,
"no slice may commit a restriction against an input whose backend \
cannot reclaim the prefix that restriction declares consumed",
);
for i in 0..TIGHT_SPACE_KEYS {
assert!(
tree.get(tight_space_key(i).as_bytes(), crate::MAX_SEQNO)?
.is_some(),
"key {i} lost",
);
}
Ok(())
}
#[test]
fn tight_space_budgets_blob_relocation_on_the_blob_volume() -> crate::Result<()> {
use crate::config::LevelRoute;
use crate::fs::{Fs, MemFs};
const N: u64 = 4_000;
let k = |i: u64| format!("key{i:08}");
let val = |i: u64, generation: u8| -> Vec<u8> {
let mut s = (i + 1).wrapping_mul(0x9E37_79B9_7F4A_7C15) ^ (u64::from(generation) << 1);
(0..200u32)
.map(|_| {
s ^= s << 13;
s ^= s >> 7;
s ^= s << 17;
#[expect(
clippy::cast_possible_truncation,
reason = "xorshift byte extraction; the high bits are intentionally dropped"
)]
let byte = (s >> 24) as u8;
byte
})
.collect()
};
let dir = tempfile::tempdir()?;
let cold_dir = dir.path().join("cold");
let cold = Arc::new(MemFs::with_capacity(u64::MAX));
cold.create_dir_all(&cold_dir)?;
let mem = Arc::new(MemFs::with_capacity(u64::MAX));
let cold_fs: Arc<dyn Fs> = cold;
let config = Config::new(
&dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.with_shared_fs(mem.clone())
.level_routes(vec![LevelRoute {
levels: 0..7,
path: cold_dir,
fs: cold_fs,
}])
.with_kv_separation(Some(
KvSeparationOptions::default()
.separation_threshold(64)
.age_cutoff(1.0)
.staleness_threshold(0.1)
.file_target_size(48 * 1024),
));
let tree = match config.open()? {
crate::AnyTree::Blob(t) => t,
crate::AnyTree::Standard(_) => panic!("expected Blob tree"),
};
for i in 0..N {
tree.insert(k(i).as_bytes(), val(i, 1), i);
}
tree.flush_active_memtable(0)?;
for i in (0..N).step_by(2) {
tree.insert(k(i).as_bytes(), val(i, 2), N + i);
}
tree.flush_active_memtable(0)?;
let gc_watermark = 4 * N;
tree.index.update_runtime_config(|c| {
c.storage_admission_check = true;
c.storage_limit_bytes = None;
})?;
tree.major_compact(64 * 1024 * 1024, gc_watermark)?;
const PROBE: u64 = 1 << 40;
mem.set_capacity(PROBE);
let blob_stored = PROBE - mem.available_space(dir.path())?;
mem.set_capacity(blob_stored + blob_stored / 4);
tree.index.update_runtime_config(|c| {
c.tight_space_compaction = true;
})?;
tree.major_compact(64 * 1024 * 1024, gc_watermark)?;
assert!(
mem.punched_bytes() > 0,
"the relocation must slice against the blob volume's free space and \
reclaim stale blob prefixes there",
);
for i in 0..N {
let expected = if i % 2 == 0 { val(i, 2) } else { val(i, 1) };
assert_eq!(
tree.get(k(i).as_bytes(), crate::MAX_SEQNO)?.as_deref(),
Some(expected.as_slice()),
"key {i} wrong/lost after the relocating tight-space merge",
);
}
Ok(())
}
#[test]
fn tight_space_declines_when_punching_inputs_frees_the_wrong_volume() -> crate::Result<()> {
use crate::config::LevelRoute;
use crate::fs::{Fs, MemFs};
let dir = tempfile::tempdir()?;
let hot_dir = dir.path().join("hot");
let hot = Arc::new(MemFs::with_capacity(u64::MAX));
hot.create_dir_all(&hot_dir)?;
let main = Arc::new(MemFs::with_capacity(u64::MAX));
let hot_fs: Arc<dyn Fs> = hot.clone();
let tree = match Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.with_shared_fs(main.clone())
.level_routes(vec![LevelRoute {
levels: 0..1,
path: hot_dir.clone(),
fs: hot_fs,
}])
.open()?
{
crate::AnyTree::Standard(t) => t,
crate::AnyTree::Blob(_) => panic!("expected Standard tree"),
};
for i in 0..TIGHT_SPACE_KEYS {
tree.insert(tight_space_key(i).as_bytes(), vec![0xCDu8; 64], i);
}
tree.flush_active_memtable(0)?;
main.set_capacity(64 * 1024);
tree.update_runtime_config(|c| {
c.storage_admission_check = true;
c.tight_space_compaction = true;
})?;
tree.major_compact(64 * 1024 * 1024, 0)?;
assert_eq!(
hot.punched_bytes(),
0,
"punching the routed inputs would free the wrong volume: the constrained \
destination gains every slice output and loses nothing",
);
let sidecars = hot
.read_dir(&hot_dir.join("tables"))?
.into_iter()
.filter(|e| e.file_name.contains("restrict-bound"))
.count();
assert_eq!(sidecars, 0, "no slice may commit against such an input");
for i in 0..TIGHT_SPACE_KEYS {
assert!(
tree.get(tight_space_key(i).as_bytes(), crate::MAX_SEQNO)?
.is_some(),
"key {i} lost",
);
}
Ok(())
}
#[test]
fn a_refused_sidecar_read_is_reported_not_walked_as_corruption() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule, Fs};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let capacity = capfs::CapacityFs::new();
let fault = FaultFs::new(capacity.clone());
let injector = fault.injector();
let fs: Arc<dyn Fs> = Arc::new(fault);
let reopened = tight_space_crash_and_reopen(
dir.path(),
Arc::clone(&fs),
|used| capacity.set_available_space(used / 4),
|| capacity.punched_bytes(),
)?;
let version = reopened.current_version();
let Some(restricted) = version
.iter_tables()
.find(|t| t.restrict_lower_bound().is_some())
else {
panic!("the punched input must reopen as a restricted table");
};
let path = (*restricted.path).clone();
let clean = crate::verify::verify_sst_file_with_fs(&fs, &path);
assert!(
clean.errors.is_empty(),
"a legitimately punched table verifies clean: {:?}",
clean.errors,
);
injector.arm(
FaultRule::new(FaultOp::Open, Fault::Error(ErrorKind::PermissionDenied))
.on_path(".restrict-bound"),
);
let report = crate::verify::verify_sst_file_with_fs(&fs, &path);
injector.clear();
assert!(
!report.errors.is_empty(),
"the refused read must be reported at all: {report:?}",
);
assert!(
report
.errors
.iter()
.all(|e| matches!(e, crate::verify::BlockVerifyError::SstFileUnreadable { .. })),
"a refused sidecar read must be reported, not turned into corruption \
findings over the punched prefix: {:?}",
report.errors,
);
Ok(())
}
#[test]
fn tight_space_engages_when_the_local_share_covers_the_output() -> crate::Result<()> {
use crate::config::LevelRoute;
use crate::fs::{Fs, MemFs};
let dir = tempfile::tempdir()?;
let hot_dir = dir.path().join("hot");
let hot = Arc::new(MemFs::with_capacity(u64::MAX));
hot.create_dir_all(&hot_dir)?;
let main = Arc::new(MemFs::with_capacity(u64::MAX));
let hot_fs: Arc<dyn Fs> = hot;
let tree = match Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.data_block_size_policy(BlockSizePolicy::all(512))
.with_shared_fs(main.clone())
.level_routes(vec![LevelRoute {
levels: 0..1,
path: hot_dir.clone(),
fs: hot_fs,
}])
.open()?
{
crate::AnyTree::Standard(t) => t,
crate::AnyTree::Blob(_) => panic!("expected Standard tree"),
};
for i in 0..TIGHT_SPACE_KEYS {
tree.insert(tight_space_key(i).as_bytes(), vec![0xCDu8; 64], i);
}
tree.flush_active_memtable(0)?;
tree.major_compact(64 * 1024 * 1024, 0)?;
let local_bytes: u64 = tree
.current_version()
.iter_tables()
.filter(|t| t.path.starts_with(dir.path().join("tables")))
.map(crate::table::Table::file_size)
.sum();
assert!(
local_bytes > 0,
"the first generation is on the destination"
);
for i in 0..64u64 {
tree.insert(
tight_space_key(i).as_bytes(),
vec![0xEFu8; 64],
TIGHT_SPACE_KEYS + i,
);
}
tree.flush_active_memtable(0)?;
let remote_bytes: u64 = tree
.current_version()
.iter_tables()
.filter(|t| t.path.starts_with(&hot_dir))
.map(crate::table::Table::file_size)
.sum();
assert!(
remote_bytes > 0 && remote_bytes < local_bytes,
"the routed share must be the smaller one ({remote_bytes} vs {local_bytes})",
);
let used = tree.storage_stats()?.used_bytes;
main.set_capacity(used + remote_bytes * 2);
tree.update_runtime_config(|c| {
c.storage_admission_check = true;
c.tight_space_compaction = true;
})?;
tree.major_compact(64 * 1024 * 1024, 0)?;
assert!(
main.punched_bytes() > 0,
"the destination-local inputs must be reclaimed rather than the pass \
declining over the routed remainder",
);
for i in 0..64u64 {
assert_eq!(
tree.get(tight_space_key(i).as_bytes(), crate::MAX_SEQNO)?
.as_deref(),
Some(&[0xEFu8; 64][..]),
"the newest generation survives the rewrite",
);
}
Ok(())
}