use super::{BlobDropReason, DropReason, salvage_blob_file, salvage_sst};
use super::{SalvageOptions, salvage_sst_with_options};
use crate::comparator::default_comparator;
use crate::fs::{Fs, StdFs};
use crate::table::{Table, Writer};
use crate::{InternalValue, ValueType};
use alloc::sync::Arc;
use tempfile::tempdir;
use test_log::test;
#[cfg(feature = "std")]
fn reconcile_error(
table: &Table,
expected: crate::table::ReconcileGate,
prefix_extractor: Option<&Arc<dyn crate::prefix::PrefixExtractor>>,
) -> crate::Error {
match table.verify_reconcile_gates(prefix_extractor, false) {
Ok(()) => panic!("the forgery must be rejected, the reconcile pass accepted it"),
Err((gate, e)) => {
assert_eq!(gate, expected, "the wrong gate rejected the table: {e}");
e
}
}
}
#[cfg(feature = "std")]
fn reconcile_clean(
table: &Table,
prefix_extractor: Option<&Arc<dyn crate::prefix::PrefixExtractor>>,
) {
if let Err((gate, e)) = table.verify_reconcile_gates(prefix_extractor, false) {
panic!("an honest table must pass every gate, {gate:?} refused it: {e}");
}
}
#[test]
fn salvage_blob_file_drops_the_whole_tail_after_a_resync() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("blob_rot");
let dest = dir.path().join("blob_rot_salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
build_blob(
&source,
&fs,
&[
(b"aaaa", b"AAAAAAAA"),
(b"bbbb", b"BBBBBBBB"),
(b"cccc", b"CCCCCCCC"),
(b"dddd", b"DDDDDDDD"),
],
)?;
let frame_len = 42 + 4 + 8;
let kl_off = frame_len + 4 + 16 + 8;
let mut bytes = std::fs::read(&source)?;
let Some(slot) = bytes.get_mut(kl_off..kl_off + 2) else {
panic!("second frame's key_len within the file");
};
slot.copy_from_slice(&6u16.to_le_bytes());
std::fs::write(&source, &bytes)?;
let report = salvage_blob_file(
&source,
dest,
&fs,
0,
&default_comparator(),
0,
#[cfg(zstd_any)]
None,
)?;
assert_eq!(
report.records_salvaged, 1,
"only the frame before the rot is provable; the whole tail after the resync drops: {report:?}",
);
assert!(
report
.dropped
.iter()
.any(|d| matches!(&d.reason, BlobDropReason::Corrupt(msg) if msg.contains("HeaderCrcMismatch"))),
"the rotted frame drops as a header CRC mismatch: {report:?}",
);
assert_eq!(
report
.dropped
.iter()
.filter(
|d| matches!(&d.reason, BlobDropReason::Corrupt(msg) if msg.contains("surrendered"))
)
.count(),
1,
"the surrendered tail is recorded ONCE, not per tainted frame: {report:?}",
);
let Some(salvaged) = report.salvaged_path else {
panic!("one record was salvaged");
};
let keys: Vec<Vec<u8>> = BlobScanner::new(&salvaged, &*fs, 0)?
.map(|r| r.map(|e| e.key.to_vec()))
.collect::<crate::Result<_>>()?;
assert_eq!(keys, vec![b"aaaa".to_vec()]);
Ok(())
}
#[test]
fn salvage_recovers_all_blocks_under_a_forged_tli_tail() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(128);
for i in 0u64..64 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tli_tail_truncated(&source, 0, None)?;
let report = salvage_sst(&source, dest, &fs)?;
assert!(
report.dropped.is_empty(),
"every block is intact, nothing may drop: {report:?}",
);
assert_eq!(
report.entries_salvaged, 64,
"the block the forged tail hides must still be recovered: {report:?}",
);
Ok(())
}
#[test]
fn salvage_recovers_a_block_hidden_by_forged_tli_mirrors() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(128);
for i in 0u64..64 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tli_mirrors_truncated(&source, 0, None)?;
let report = salvage_sst(&source, dest, &fs)?;
assert!(
report.dropped.is_empty(),
"every block is intact, nothing may drop: {report:?}",
);
assert_eq!(
report.entries_salvaged, 64,
"the block both forged mirrors hide must still be recovered: {report:?}",
);
Ok(())
}
#[test]
fn salvage_recovers_all_blocks_under_a_reordered_tli() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(128);
for i in 0u64..64 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tli_mirrors_swap_first_two(&source, 0, None)?;
let report = salvage_sst(&source, dest, &fs)?;
assert!(
report.dropped.is_empty(),
"every block is intact and covered once, nothing may drop: {report:?}",
);
assert_eq!(
report.entries_salvaged, 64,
"a reordered index must not double-walk or lose blocks: {report:?}",
);
Ok(())
}
#[test]
fn salvage_recovers_physical_blocks_past_a_broken_index_partition() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_data_block_size(128)
.use_partitioned_index();
for i in 0u64..256 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let (index_pos, index_len) = {
let mut f = std::fs::File::open(&source)?;
let reader = crate::sfa::Reader::from_reader(&mut f)?;
let Some((pos, len)) = reader
.toc()
.iter()
.find(|e| e.name() == b"index")
.map(|e| (e.pos(), e.len()))
else {
panic!("a partitioned-index SST must carry an index section");
};
(pos, len)
};
let Ok(flip) = usize::try_from(index_pos + index_len / 2) else {
panic!("the index-section offset fits usize");
};
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(flip) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest, &fs)?;
assert_eq!(
report.entries_salvaged, 256,
"every data block must be recovered via the physical walk when the \
index enumeration breaks: {report:?}",
);
assert!(
report.is_complete(),
"a broken index the physical walk fully recovers around is not data loss: {report:?}",
);
Ok(())
}
#[test]
fn salvage_refuses_a_range_tombstone_hidden_as_a_recognized_section() -> crate::Result<()> {
use crate::UserKey;
use crate::config::BloomConstructionPolicy;
use crate::range_tombstone::RangeTombstone;
use crate::table::block::BlockType;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_bloom_policy(BloomConstructionPolicy::BitsPerKey(0.0));
for i in 0u64..8 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
writer.write_range_tombstone(RangeTombstone::new(
UserKey::from(b"key-002".as_slice()),
UserKey::from(b"key-005".as_slice()),
9,
));
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_duplicate_section_name(
&source,
b"range_tombstones",
b"filter",
BlockType::Filter,
)?;
let Err(err) = salvage_sst(&source, dest, &fs) else {
panic!("a hidden range tombstone must fail salvage");
};
let crate::Error::FeatureUnsupported(reason) = &err else {
panic!("the refusal must be FeatureUnsupported, got {err:?}");
};
assert!(
reason.contains("range tombstones"),
"the refusal must name range tombstones specifically, got {reason:?}",
);
Ok(())
}
#[test]
fn salvage_drops_indexed_blocks_after_a_break_when_the_index_is_untrusted() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_data_block_size(128)
.use_partitioned_index();
for i in 0u64..256 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let smash_offset = {
let table = open(source.clone(), &fs)?;
let offsets: Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|h| *h.as_ref().offset())
.collect();
let Some(&off) = offsets.get(offsets.len() / 2) else {
panic!("need several data blocks, got {}", offsets.len());
};
off
};
crate::test_forge::forge_flip_section_last_payload_byte(&source, b"tli_tail", None)?;
{
let mut bytes = std::fs::read(&source)?;
let Ok(at) = usize::try_from(smash_offset) else {
panic!("data block offset {smash_offset} fits usize");
};
let Some(b) = bytes.get_mut(at) else {
panic!("the block header at {at} lies within the file");
};
*b ^= 0xFF;
std::fs::write(&source, &bytes)?;
}
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(
reopen_get(dest.clone(), &fs, b"key-000")?.is_some(),
"the contiguous prefix before the broken chain must recover: {report:?}",
);
assert!(
reopen_get(dest, &fs, b"key-250")?.is_none(),
"an indexed block after a broken chain, under an untrusted index, must \
be dropped, not emitted: {report:?}",
);
assert!(
report.dropped.iter().any(|d| matches!(
&d.reason,
DropReason::HeaderCorrupted(msg) if msg.contains("broken")
)),
"the surrendered tail must be reported as a broken-chain drop: {report:?}",
);
Ok(())
}
#[test]
fn salvage_drops_a_persistently_unreadable_block_and_keeps_the_rest() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let clean: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer =
Writer::new(source.clone(), 0, 0, Arc::clone(&clean))?.use_data_block_size(128);
for i in 0u64..256 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let victim_offset = {
let table = open(source.clone(), &clean)?;
let offsets: Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|h| *h.as_ref().offset())
.collect();
let Some(&off) = offsets.get(offsets.len() / 2) else {
panic!("need several data blocks, got {}", offsets.len());
};
off
};
let fault = FaultFs::new(StdFs);
fault.injector().arm(
FaultRule::new(FaultOp::ReadAt, Fault::Error(ErrorKind::Other)).at_offset(victim_offset),
);
let faulting: Arc<dyn Fs> = Arc::new(fault);
let report = salvage_sst(&source, dest.clone(), &faulting)?;
assert!(
!report.dropped.is_empty(),
"a persistently unreadable block must be dropped, not abort the salvage: {report:?}",
);
assert!(
report.blocks_salvaged >= report.blocks_total - report.dropped.len(),
"every block except the dropped one is salvaged: {report:?}",
);
assert!(
reopen_get(dest, &clean, b"key-000")?.is_some(),
"an intact block must still be recovered past the dropped one: {report:?}",
);
Ok(())
}
#[test]
fn salvage_drops_the_gap_tail_after_an_unframeable_block() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_data_block_size(128)
.use_partitioned_index();
for i in 0u64..256 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let smash_offset = {
let table = open(source.clone(), &fs)?;
let offsets: Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|h| *h.as_ref().offset())
.collect();
let Some(&off) = offsets.get(offsets.len() / 3) else {
panic!("need several data blocks, got {}", offsets.len());
};
off
};
{
let mut bytes = std::fs::read(&source)?;
let Ok(at) = usize::try_from(smash_offset) else {
panic!("data block offset {smash_offset} fits usize");
};
let Some(b) = bytes.get_mut(at) else {
panic!("the block header at {at} lies within the file");
};
*b ^= 0xFF;
std::fs::write(&source, &bytes)?;
}
let (index_pos, index_len) = {
let mut f = std::fs::File::open(&source)?;
let reader = crate::sfa::Reader::from_reader(&mut f)?;
let Some((pos, len)) = reader
.toc()
.iter()
.find(|e| e.name() == b"index")
.map(|e| (e.pos(), e.len()))
else {
panic!("a partitioned-index SST must carry an index section");
};
(pos, len)
};
let Ok(flip) = usize::try_from(index_pos + index_len / 2) else {
panic!("the index-section offset fits usize");
};
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(flip) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(
reopen_get(dest.clone(), &fs, b"key-000")?.is_some(),
"the contiguous prefix before the broken boundary must recover: {report:?}",
);
assert!(
reopen_get(dest, &fs, b"key-250")?.is_none(),
"a block reached only by resync must be dropped, not emitted: {report:?}",
);
assert!(
report.dropped.iter().any(|d| matches!(
&d.reason,
DropReason::HeaderCorrupted(msg) if msg.contains("unanchored")
)),
"the resync tail must be reported as an unanchored drop: {report:?}",
);
Ok(())
}
#[test]
fn salvage_drops_the_tail_past_a_fake_oversized_block_header() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_data_block_size(128)
.use_partitioned_index();
for i in 0u64..256 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let (forge_off, forged_end) = {
let table = open(source.clone(), &fs)?;
let offsets: Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|h| *h.as_ref().offset())
.collect();
let n = offsets.len();
assert!(n >= 8, "need several data blocks, got {n}");
let (Some(&off), Some(&end)) = (offsets.get(n / 4), offsets.get(3 * n / 4)) else {
panic!("the quarter and three-quarter block boundaries exist");
};
(off, end)
};
crate::test_forge::forge_data_block_oversized_header(&source, forge_off, forged_end)?;
let (index_pos, index_len) = {
let mut f = std::fs::File::open(&source)?;
let reader = crate::sfa::Reader::from_reader(&mut f)?;
let Some((pos, len)) = reader
.toc()
.iter()
.find(|e| e.name() == b"index")
.map(|e| (e.pos(), e.len()))
else {
panic!("a partitioned-index SST must carry an index section");
};
(pos, len)
};
let Ok(flip) = usize::try_from(index_pos + index_len / 2) else {
panic!("the index-section offset fits usize");
};
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(flip) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(
reopen_get(dest.clone(), &fs, b"key-010")?.is_some(),
"the contiguous prefix before the fake header must recover: {report:?}",
);
assert!(
reopen_get(dest, &fs, b"key-128")?.is_none(),
"a block reached only by resync past the fake header must drop: {report:?}",
);
assert!(
report.entries_salvaged < 200,
"the resync tail must not be recovered, got {} of 256: {report:?}",
report.entries_salvaged,
);
assert!(
report.dropped.iter().any(|d| matches!(
&d.reason,
DropReason::HeaderCorrupted(msg) if msg.contains("unanchored")
)),
"the swallowed tail must be reported as an unanchored drop: {report:?}",
);
Ok(())
}
#[test]
fn salvage_recovers_an_interior_block_hidden_by_forged_tli_mirrors() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(128);
for i in 0u64..64 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tli_mirrors_drop_interior(&source, 0, None)?;
let report = salvage_sst(&source, dest, &fs)?;
assert!(
report.dropped.is_empty(),
"every block is intact, nothing may drop: {report:?}",
);
assert_eq!(
report.entries_salvaged, 64,
"the interior block both forged mirrors hide must still be recovered: {report:?}",
);
Ok(())
}
#[test]
fn verify_tli_mirrors_rejects_a_section_spanning_handle() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(128);
for i in 0u64..64 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
open(source.clone(), &fs)?.verify_tli_mirrors()?;
crate::test_forge::forge_tli_mirrors_span_single_handle(&source, 0)?;
let table = open(source, &fs)?;
let Err(err) = table.verify_tli_mirrors() else {
panic!("a handle spanning several physical blocks must fail the mirror gate");
};
assert!(
matches!(
err,
crate::Error::InvalidHeader(
"an index handle's size disagrees with its block's physical frame"
)
),
"the spanning handle must be rejected by the physical-frame check, got {err:?}",
);
Ok(())
}
#[test]
fn salvage_recovers_all_blocks_under_a_section_spanning_handle() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(128);
for i in 0u64..64 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tli_mirrors_span_single_handle(&source, 0)?;
let report = salvage_sst(&source, dest, &fs)?;
assert!(
report.dropped.is_empty(),
"every block is intact, nothing may drop: {report:?}",
);
assert_eq!(
report.entries_salvaged, 64,
"the blocks a spanning handle hides must still be recovered: {report:?}",
);
Ok(())
}
#[test]
fn salvage_ignores_an_index_handle_beyond_the_data_section() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(128);
let n = 64u32;
for i in 0..u64::from(n) {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tli_mirrors_offset_beyond_section(&source, 0)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(
report.entries_salvaged,
u64::from(n),
"every real block is recovered once the out-of-section handle is skipped: {report:?}",
);
assert!(
reopen_get(dest, &fs, b"key-060")?.is_some(),
"a late key must survive the bounded walk: {report:?}",
);
Ok(())
}
#[test]
fn salvage_drops_the_tail_after_an_unframeable_oversized_handle() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(128);
for i in 0u64..64 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let first_off = {
let table = open(source.clone(), &fs)?;
let Some(off) = table
.data_block_handles()
.filter_map(Result::ok)
.map(|h| *h.as_ref().offset())
.next()
else {
panic!("the source carries data blocks");
};
off
};
crate::test_forge::forge_tli_mirrors_span_single_handle(&source, 0)?;
{
let mut bytes = std::fs::read(&source)?;
let Ok(at) = usize::try_from(first_off) else {
panic!("data block offset {first_off} fits usize");
};
let Some(b) = bytes.get_mut(at) else {
panic!("the block header lies within the file");
};
*b ^= 0xFF;
std::fs::write(&source, &bytes)?;
}
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(
report.entries_salvaged, 0,
"nothing is anchored, so no entry is recovered: {report:?}",
);
assert!(
!dest.exists(),
"an all-dropped section produces no salvaged table: {report:?}",
);
assert!(
report.dropped.iter().any(|d| matches!(
&d.reason,
DropReason::HeaderCorrupted(msg) if msg.contains("broken")
)),
"the dropped section must be reported as a broken-chain drop: {report:?}",
);
Ok(())
}
#[test]
fn salvage_refuses_an_overflowing_data_section() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(128);
let n = 64u32;
for i in 0..u64::from(n) {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_section_len(&source, b"data", u64::MAX)?;
let Err(err) = salvage_sst(&source, dest.clone(), &fs) else {
panic!("an overflowing data section must fail salvage closed");
};
assert!(
matches!(err, crate::Error::FeatureUnsupported(msg) if msg.contains("TOC may hide a deletion section")),
"the refusal must name the TOC-coverage gate, not any other refusal, got {err:?}",
);
assert!(
!std::path::Path::new(&dest).exists(),
"no salvaged copy is produced when the catalogue may hide a deletion",
);
Ok(())
}
#[test]
fn salvage_refuses_a_corrupt_seqno_bounds_that_may_hide_a_deletion() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_seqno_in_index(true);
for i in 0u64..50 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
{
let pos = {
let mut f = std::fs::File::open(&source)?;
let reader = match crate::sfa::Reader::from_reader(&mut f) {
Ok(r) => r,
Err(e) => panic!("reading the SFA trailer failed: {e:?}"),
};
let Some(entry) = reader.toc().iter().find(|e| e.name() == b"seqno_bounds") else {
panic!("the source carries a seqno_bounds section");
};
let Ok(pos) = usize::try_from(entry.pos()) else {
panic!("pos fits usize");
};
pos
};
let mut bytes = std::fs::read(&source)?;
let at = pos + 40;
let Some(slot) = bytes.get_mut(at) else {
panic!("payload byte within the file");
};
*slot ^= 0xFF;
std::fs::write(&source, bytes)?;
}
let Err(err) = salvage_sst(&source, dest, &fs) else {
panic!("a corrupt seqno-bounds section with no visible deletion must fail salvage");
};
assert!(
matches!(err, crate::Error::FeatureUnsupported(_)),
"the refusal names the unsupported salvage, got {err:?}",
);
Ok(())
}
#[cfg(feature = "lz4")]
#[test]
fn salvage_arbitrates_divergent_meta_mirrors() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0u64..100 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tail_meta_value(&source, b"compression#data", &[1])?;
let report = salvage_sst(&source, dest, &fs)?;
assert!(report.blocks_total > 0, "the walk saw the data blocks");
assert_eq!(
report.blocks_salvaged, report.blocks_total,
"every block is recoverable through the intact MID mirror: {report:?}",
);
assert_eq!(
report.entries_salvaged, 100,
"all rows recovered: {report:?}"
);
assert!(
report.dropped.is_empty(),
"nothing should drop under the intact mirror: {report:?}",
);
assert!(report.salvaged_path.is_some(), "a copy was written");
for entry in std::fs::read_dir(dir.path())? {
let name = entry?.file_name().to_string_lossy().into_owned();
assert!(
!name.contains(".healtmp-"),
"the arbitration temp file must be renamed or removed, found {name:?}",
);
}
Ok(())
}
#[cfg(feature = "lz4")]
#[test]
fn salvage_removes_the_mid_copy_when_the_publish_dir_sync_fails() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fault = FaultFs::new(StdFs);
let injector = fault.injector();
let fs: Arc<dyn Fs> = Arc::new(fault);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0u64..100 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tail_meta_value(&source, b"compression#data", &[1])?;
injector.arm(FaultRule::new(FaultOp::SyncDirectory, Fault::Error(ErrorKind::Other)).skip(1));
let Err(err) = salvage_sst(&source, dest.clone(), &fs) else {
panic!("a failed publish directory sync must fail the salvage");
};
assert!(
err.to_string().contains("injected fault on SyncDirectory"),
"the salvage error must be the injected publish-sync fault, got {err:?}",
);
assert!(
!std::path::Path::new(&dest).exists(),
"the owned MID copy must be removed when the publish sync fails",
);
for entry in std::fs::read_dir(dir.path())? {
let name = entry?.file_name().to_string_lossy().into_owned();
assert!(
!name.contains(".healtmp-"),
"the MID temp copy must not survive either, found {name:?}",
);
}
Ok(())
}
#[cfg(feature = "lz4")]
#[test]
fn salvage_publishes_the_mid_copy_when_hard_link_is_unsupported() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fault = FaultFs::new(StdFs);
let injector = fault.injector();
let fs: Arc<dyn Fs> = Arc::new(fault);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0u64..100 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tail_meta_value(&source, b"compression#data", &[1])?;
injector.arm(FaultRule::new(
FaultOp::HardLink,
Fault::Error(ErrorKind::Unsupported),
));
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(
report.salvaged_path.as_deref(),
Some(dest.as_path()),
"the MID copy must be published via the rename fallback: {report:?}",
);
assert!(
report.entries_salvaged > 0,
"the recovered table must carry entries: {report:?}",
);
assert!(
std::path::Path::new(&dest).exists(),
"the recovered table must land at the destination",
);
Ok(())
}
#[cfg(feature = "lz4")]
#[test]
fn salvage_does_not_clobber_an_occupied_destination_when_the_mid_attempt_wins() -> crate::Result<()>
{
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(128);
for i in 0u64..64 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tail_meta_value(&source, b"compression#data", &[1])?;
std::fs::write(&dest, b"racing worker's file")?;
let result = salvage_sst(&source, dest.clone(), &fs);
assert!(
result.is_err(),
"publishing over an occupied destination must fail, got {result:?}",
);
assert_eq!(
std::fs::read(&dest)?,
b"racing worker's file",
"the occupant's bytes must survive the failed publish",
);
for entry in std::fs::read_dir(dir.path())? {
let name = entry?.file_name().to_string_lossy().into_owned();
assert!(
!name.contains(".healtmp-"),
"the MID temp copy must not leak, found {name:?}",
);
}
Ok(())
}
#[cfg(feature = "lz4")]
#[test]
fn salvage_arbitration_skips_a_stale_crash_artifact() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0u64..100 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tail_meta_value(&source, b"compression#data", &[1])?;
let stale = dest.with_extension("healtmp-0");
std::fs::write(&stale, b"crashed predecessor artifact")?;
let report = salvage_sst(&source, dest, &fs)?;
assert_eq!(
report.entries_salvaged, 100,
"the MID attempt must pick a fresh temp name past the stale artifact: {report:?}",
);
assert!(
std::path::Path::new(&stale).exists(),
"a foreign artifact is never reclaimed",
);
assert_eq!(
std::fs::read(&stale)?,
b"crashed predecessor artifact".to_vec(),
"the stale artifact's bytes stay untouched",
);
Ok(())
}
#[cfg(feature = "lz4")]
#[test]
fn salvage_arbitration_keeps_a_foreign_temp_it_did_not_create() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fault = FaultFs::new(StdFs);
let injector = fault.injector();
let fs: Arc<dyn Fs> = Arc::new(fault);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0u64..100 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tail_meta_value(&source, b"compression#data", &[1])?;
let foreign = dest.with_extension("healtmp-0");
let sentinel = b"concurrent salvage in-progress output".to_vec();
std::fs::write(&foreign, &sentinel)?;
injector.arm(
FaultRule::new(FaultOp::Metadata, Fault::Error(ErrorKind::Other))
.on_path("healtmp-0")
.once(),
);
let report = salvage_sst(&source, dest, &fs)?;
assert_eq!(
report.entries_salvaged, 100,
"the mid attempt recovers every entry: {report:?}",
);
assert!(
std::path::Path::new(&foreign).exists(),
"the foreign concurrent temp must not be discarded by this salvage",
);
assert_eq!(
std::fs::read(&foreign)?,
sentinel,
"the foreign temp's bytes stay untouched",
);
Ok(())
}
#[cfg(feature = "lz4")]
#[test]
fn salvage_propagates_a_transient_open_failure_from_the_mirror_probe() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fault = FaultFs::new(StdFs);
let injector = fault.injector();
let fs: Arc<dyn Fs> = Arc::new(fault);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0u64..100 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tail_meta_value(&source, b"compression#data", &[1])?;
injector.arm(
FaultRule::new(FaultOp::Open, Fault::Error(ErrorKind::Other))
.on_path("source")
.once(),
);
let result = salvage_sst(&source, dest, &fs);
assert!(
matches!(result, Err(crate::Error::Io(_))),
"a transient open failure in the mirror-divergence probe must propagate, not \
skip arbitration and salvage tail-only: {result:?}",
);
Ok(())
}
#[cfg(feature = "lz4")]
#[test]
fn salvage_propagates_an_environmental_failure_from_the_mirror_probe() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fault = FaultFs::new(StdFs);
let injector = fault.injector();
let fs: Arc<dyn Fs> = Arc::new(fault);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0u64..100 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tail_meta_value(&source, b"compression#data", &[1])?;
let mut reached = false;
for skip in 0..24u64 {
injector.clear();
injector.arm(
FaultRule::new(FaultOp::ReadAt, Fault::Error(ErrorKind::PermissionDenied))
.on_path("source")
.skip(skip)
.times(1),
);
let options = SalvageOptions {
encryption: None,
#[cfg(zstd_any)]
zstd_dictionary: None,
table_id: 0,
expected_stored_id: None,
output_id: None,
allow_delete_resurrection: false,
sync_mode: crate::fs::SyncMode::Normal,
prefix_extractor: None,
blob_rewrite: None,
progress: None,
};
match super::meta_mirrors_diverge(&source, &fs, &options) {
Ok(_) => {}
Err(crate::Error::Io(io)) if io.kind() == ErrorKind::PermissionDenied => {
reached = true;
}
Err(e) => panic!("unexpected probe failure at skip {skip}: {e:?}"),
}
}
assert!(
reached,
"no skip count made the probe read a mirror; the sweep proves nothing",
);
let _ = dest;
Ok(())
}
#[cfg(feature = "lz4")]
#[test]
fn salvage_refuses_to_overwrite_a_preexisting_destination() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0u64..100 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tail_meta_value(&source, b"compression#data", &[1])?;
let sentinel = b"unrelated pre-existing file".to_vec();
std::fs::write(&dest, &sentinel)?;
let result = salvage_sst(&source, dest.clone(), &fs);
assert!(
result.is_err(),
"an occupied destination must fail the salvage: {result:?}",
);
assert_eq!(
std::fs::read(&dest)?,
sentinel,
"the pre-existing destination must survive byte-for-byte",
);
for entry in std::fs::read_dir(dir.path())? {
let name = entry?.file_name().to_string_lossy().into_owned();
assert!(
!name.contains(".healtmp-"),
"the refused MID attempt must clean up its temp copy, found {name:?}",
);
}
Ok(())
}
#[test]
fn salvage_reencodes_all_blocks_when_meta_mirrors_diverge() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0u64..100 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tail_meta_value(&source, b"restart_interval#data", &[4])?;
let report = salvage_sst(&source, dest, &fs)?;
assert_eq!(
report.entries_salvaged, 100,
"all rows recovered: {report:?}"
);
assert!(report.salvaged_path.is_some(), "a copy was written");
assert_eq!(
report.blocks_copied_verbatim, 0,
"divergent mirrors must force the re-encode path: a byte-copied \
block would keep the original encoding under the chosen meta's \
forged layout: {report:?}",
);
Ok(())
}
#[test]
fn salvage_does_not_carry_a_forged_created_at_into_the_copy() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0u64..100 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let backdated: u128 = 1;
crate::test_forge::forge_tail_meta_value(&source, b"created_at", &backdated.to_le_bytes())?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(report.entries_salvaged, 100, "{report:?}");
let recovered = open(dest, &fs)?;
assert_ne!(
*recovered.metadata.created_at, backdated,
"the recovered copy must not carry the forged backdated created_at",
);
Ok(())
}
#[test]
fn verify_point_read_reachability_rejects_a_cross_block_seqno_inversion() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(1);
writer.write(InternalValue::from_components(
b"kkk",
b"v2",
2,
ValueType::Value,
))?;
writer.write(InternalValue::from_components(
b"kkk",
b"v1",
1,
ValueType::Value,
))?;
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let block1_off = {
let table = open(source.clone(), &fs)?;
let offsets: alloc::vec::Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|kh| *kh.as_ref().offset())
.collect();
let Some(&second) = offsets.get(1) else {
panic!("the key's versions must span two blocks, got {offsets:?}");
};
let Ok(off) = usize::try_from(second) else {
panic!("the block offset fits usize");
};
off
};
crate::test_forge::forge_raise_data_block_first_seqno(&source, block1_off, 3)?;
let table = open(source, &fs)?;
let err = reconcile_error(
&table,
crate::table::ReconcileGate::PointReadReachability,
None,
);
assert!(
matches!(&err, crate::Error::InvalidHeader(msg) if msg.contains("out of order")),
"a cross-block seqno inversion must be rejected, got {err:?}",
);
Ok(())
}
#[test]
fn salvage_drops_a_row_block_that_decodes_fewer_entries_than_declared() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0u64..3 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_inflated_item_count(&source)?;
let report = salvage_sst(&source, dest, &fs)?;
assert_eq!(
report.entries_salvaged, 0,
"an under-decoding block must not contribute recovered entries: {report:?}",
);
assert!(
report
.dropped
.iter()
.any(|d| matches!(d.reason, DropReason::DecodeError(_))),
"the count mismatch is a dropped DecodeError, not a clean recovery: {report:?}",
);
Ok(())
}
#[test]
fn salvage_recovers_a_block_with_multiple_versions_of_one_key() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
writer.write(InternalValue::from_components(
b"a".to_vec(),
b"a".to_vec(),
1,
ValueType::Value,
))?;
for seqno in [3u64, 2, 1] {
writer.write(InternalValue::from_components(
b"dup".to_vec(),
format!("v{seqno}").into_bytes(),
seqno,
ValueType::Value,
))?;
}
writer.write(InternalValue::from_components(
b"z".to_vec(),
b"z".to_vec(),
1,
ValueType::Value,
))?;
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(
report.is_complete(),
"a healthy SST with MVCC duplicates salvages cleanly: {report:?}",
);
assert_eq!(
report.entries_salvaged, 5,
"every version is recovered, including all 3 of `dup`: {report:?}",
);
let recovered = open(dest, &fs)?;
assert_eq!(
recovered.metadata.item_count, 5,
"all 5 entries (3 versions of `dup`) are recovered",
);
Ok(())
}
#[test]
fn salvage_recovers_a_reclaimable_weak_tombstone_pair() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
writer.write(InternalValue::from_components(
b"a".to_vec(),
b"a".to_vec(),
1,
ValueType::Value,
))?;
writer.write(InternalValue::from_components(
b"dup".to_vec(),
b"".to_vec(),
3,
ValueType::WeakTombstone,
))?;
writer.write(InternalValue::from_components(
b"dup".to_vec(),
b"v1".to_vec(),
1,
ValueType::Value,
))?;
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let report = salvage_sst(&source, dest, &fs)?;
assert!(
report.is_complete(),
"healthy SST salvages cleanly: {report:?}"
);
assert_eq!(
report.entries_salvaged, 3,
"the weak tombstone and both values are recovered: {report:?}",
);
assert!(
report.blocks_copied_verbatim >= 1,
"the clean block is copied verbatim: {report:?}",
);
Ok(())
}
#[test]
fn verify_locator_rejects_a_slot_hint_pointing_at_an_older_version() -> crate::Result<()> {
use crate::config::{LocatorPolicyEntry, LocatorPrecision};
use crate::runtime_config::{ChecksumAlgorithm, KvChecksumPolicy};
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_data_block_restart_interval(1)
.use_kv_checksums(KvChecksumPolicy::AllLevels, ChecksumAlgorithm::Xxh3_64)
.use_locator(LocatorPolicyEntry::Enabled {
precision: LocatorPrecision::Restart,
block_id_bits: None,
slot_bits: None,
});
writer.write(InternalValue::from_components(
b"key-a".to_vec(),
b"new".to_vec(),
2,
ValueType::Value,
))?;
writer.write(InternalValue::from_components(
b"key-a".to_vec(),
b"old".to_vec(),
1,
ValueType::Value,
))?;
writer.write(InternalValue::from_components(
b"key-b".to_vec(),
b"vb".to_vec(),
1,
ValueType::Value,
))?;
writer.write(InternalValue::from_components(
b"key-c".to_vec(),
b"vc".to_vec(),
1,
ValueType::Value,
))?;
assert!(writer.finish()?.is_some(), "source SST is non-empty");
reconcile_clean(&open(source.clone(), &fs)?, None);
let h = |k: &[u8]| crate::hash::hash64(k);
crate::test_forge::forge_locator_slots(
&source,
0,
&[
(h(b"key-a"), 0, 1),
(h(b"key-b"), 0, 2),
(h(b"key-c"), 0, 3),
],
)?;
let table = open(source, &fs)?;
let err = reconcile_error(&table, crate::table::ReconcileGate::Locator, None);
assert!(
matches!(
err,
crate::Error::InvalidHeader(
"locator slot hint does not resolve a key's newest version"
)
),
"the redirect must be rejected by the slot-hint check, got {err:?}",
);
Ok(())
}
#[test]
fn verify_locator_rejects_a_no_answer_locator() -> crate::Result<()> {
use crate::config::{LocatorPolicyEntry, LocatorPrecision};
use crate::runtime_config::{ChecksumAlgorithm, KvChecksumPolicy};
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_data_block_restart_interval(1)
.use_kv_checksums(KvChecksumPolicy::AllLevels, ChecksumAlgorithm::Xxh3_64)
.use_locator(LocatorPolicyEntry::Enabled {
precision: LocatorPrecision::Restart,
block_id_bits: None,
slot_bits: None,
});
for k in [b"key-a".as_slice(), b"key-b", b"key-c"] {
writer.write(InternalValue::from_components(
k.to_vec(),
b"v".to_vec(),
1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
reconcile_clean(&open(source.clone(), &fs)?, None);
let h = |k: &[u8]| crate::hash::hash64(k);
crate::test_forge::forge_locator_slots(
&source,
0,
&[
(h(b"key-a"), 0, 0),
(h(b"key-b"), 0, 1),
(h(b"key-c"), 99, 2),
],
)?;
let table = open(source, &fs)?;
let err = reconcile_error(&table, crate::table::ReconcileGate::Locator, None);
assert!(
matches!(
err,
crate::Error::InvalidHeader(
"locator gives no answer for a decoded key it should resolve"
)
),
"the miss must be rejected by the no-answer check, got {err:?}",
);
Ok(())
}
#[test]
fn verify_filter_rejects_missing_prefix_hashes() -> crate::Result<()> {
struct FixedLengthPrefix;
impl crate::prefix::PrefixExtractor for FixedLengthPrefix {
fn prefixes<'a>(&self, key: &'a [u8]) -> Box<dyn Iterator<Item = &'a [u8]> + 'a> {
key.get(..4).map_or_else(
|| Box::new(std::iter::empty()) as Box<dyn Iterator<Item = &'a [u8]>>,
|p| Box::new(std::iter::once(p)),
)
}
fn is_valid_scan_boundary(&self, prefix: &[u8]) -> bool {
prefix.len() == 4
}
}
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let extractor: Arc<dyn crate::prefix::PrefixExtractor> = Arc::new(FixedLengthPrefix);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0u32..100 {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let table = open(source, &fs)?;
reconcile_clean(&table, None);
let err = reconcile_error(
&table,
crate::table::ReconcileGate::Filter,
Some(&extractor),
);
assert!(
matches!(
err,
crate::Error::InvalidHeader(
"filter reports an existing key's prefix as definitely absent"
)
),
"the missing prefix hash must be the rejection reason, got {err:?}",
);
Ok(())
}
#[test]
fn verify_point_read_reachability_rejects_a_bucket_redirected_to_an_older_version()
-> crate::Result<()> {
use crate::runtime_config::{ChecksumAlgorithm, KvChecksumPolicy};
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_data_block_hash_ratio(2.0)
.use_data_block_restart_interval(1)
.use_kv_checksums(KvChecksumPolicy::AllLevels, ChecksumAlgorithm::Xxh3_64);
writer.write(InternalValue::from_components(
b"key-a".to_vec(),
b"new".to_vec(),
2,
ValueType::Value,
))?;
writer.write(InternalValue::from_components(
b"key-a".to_vec(),
b"old".to_vec(),
1,
ValueType::Value,
))?;
for k in [b"key-b", b"key-c"] {
writer.write(InternalValue::from_components(
k.to_vec(),
b"v".to_vec(),
1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
reconcile_clean(&open(source.clone(), &fs)?, None);
crate::test_forge::forge_hash_index_bucket(&source, b"key-a", 1, None)?;
let table = open(source, &fs)?;
let err = reconcile_error(
&table,
crate::table::ReconcileGate::PointReadReachability,
None,
);
assert!(
matches!(
err,
crate::Error::InvalidHeader(
"a decoded key's point_read does not return its newest version \
(an in-block index disagrees with the entries)"
)
),
"the redirect must be rejected by the newest-version check, got {err:?}",
);
Ok(())
}
#[test]
fn salvage_refuses_a_corrupt_filter_index_that_may_hide_a_deletion() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_partitioned_filter()
.use_meta_partition_size(3);
for i in 0u32..64 {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
{
let pos = {
let mut f = std::fs::File::open(&source)?;
let reader = match crate::sfa::Reader::from_reader(&mut f) {
Ok(r) => r,
Err(e) => panic!("reading the SFA trailer failed: {e:?}"),
};
let Some(entry) = reader.toc().iter().find(|e| e.name() == b"filter_tli") else {
panic!("the source carries a filter_tli section");
};
let Ok(pos) = usize::try_from(entry.pos()) else {
panic!("pos fits usize");
};
pos
};
let mut bytes = std::fs::read(&source)?;
let at =
pos + crate::table::block::Header::header_len(crate::table::block::BlockType::Index);
let Some(slot) = bytes.get_mut(at) else {
panic!("payload byte within the file");
};
*slot ^= 0xFF;
std::fs::write(&source, bytes)?;
}
assert!(
open(source.clone(), &fs).is_err(),
"the rotted filter index must fail a live open",
);
let Err(err) = salvage_sst(&source, dest, &fs) else {
panic!("a corrupt filter index with no visible deletion must fail salvage");
};
assert!(
matches!(err, crate::Error::FeatureUnsupported(_)),
"the refusal names the unsupported salvage, got {err:?}",
);
Ok(())
}
#[test]
fn salvage_preserves_prefix_filter_hashes() -> crate::Result<()> {
struct FixedLengthPrefix;
impl crate::prefix::PrefixExtractor for FixedLengthPrefix {
fn prefixes<'a>(&self, key: &'a [u8]) -> Box<dyn Iterator<Item = &'a [u8]> + 'a> {
key.get(..4).map_or_else(
|| Box::new(std::iter::empty()) as Box<dyn Iterator<Item = &'a [u8]>>,
|p| Box::new(std::iter::once(p)),
)
}
fn is_valid_scan_boundary(&self, prefix: &[u8]) -> bool {
prefix.len() == 4
}
}
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let extractor: Arc<dyn crate::prefix::PrefixExtractor> = Arc::new(FixedLengthPrefix);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_prefix_extractor(Some(Arc::clone(&extractor)));
for i in 0u32..100 {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let prefix_hash = crate::hash::hash64(b"key0");
{
let table = open(source.clone(), &fs)?;
assert!(
table.maybe_contains_prefix(prefix_hash)?,
"the source filter indexes the prefix",
);
}
let options = SalvageOptions {
prefix_extractor: Some(Arc::clone(&extractor)),
..SalvageOptions::default()
};
let report = salvage_sst_with_options(&source, dest.clone(), &fs, &options)?;
assert_eq!(
report.entries_salvaged, 100,
"all rows recovered: {report:?}"
);
let table = open(dest, &fs)?;
assert!(
table.maybe_contains_prefix(prefix_hash)?,
"the salvaged copy's filter must keep the source's prefix hashes",
);
Ok(())
}
#[test]
fn salvage_omits_the_filter_without_a_prefix_extractor() -> crate::Result<()> {
struct FixedLengthPrefix;
impl crate::prefix::PrefixExtractor for FixedLengthPrefix {
fn prefixes<'a>(&self, key: &'a [u8]) -> Box<dyn Iterator<Item = &'a [u8]> + 'a> {
key.get(..4).map_or_else(
|| Box::new(std::iter::empty()) as Box<dyn Iterator<Item = &'a [u8]>>,
|p| Box::new(std::iter::once(p)),
)
}
fn is_valid_scan_boundary(&self, prefix: &[u8]) -> bool {
prefix.len() == 4
}
}
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let extractor: Arc<dyn crate::prefix::PrefixExtractor> = Arc::new(FixedLengthPrefix);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_prefix_extractor(Some(Arc::clone(&extractor)));
for i in 0u32..100 {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(
report.entries_salvaged, 100,
"all rows recovered: {report:?}"
);
let prefix_hash = crate::hash::hash64(b"key0");
let table = open(dest, &fs)?;
assert!(
table.maybe_contains_prefix(prefix_hash)?,
"without the extractor the salvaged copy must not report a prefix as definitely absent",
);
Ok(())
}
fn iv(i: u32) -> InternalValue {
InternalValue::from_components(
format!("key{i:05}").into_bytes(),
format!("val{i:05}").into_bytes(),
1,
ValueType::Value,
)
}
fn section_pos(path: &std::path::Path, name: &[u8]) -> u64 {
let mut f = match std::fs::File::open(path) {
Ok(f) => f,
Err(e) => panic!("opening the source failed: {e:?}"),
};
let reader = match crate::sfa::Reader::from_reader(&mut f) {
Ok(r) => r,
Err(e) => panic!("reading the SFA trailer failed: {e:?}"),
};
let Some(entry) = reader.toc().iter().find(|e| e.name() == name) else {
panic!(
"source must carry a {} section",
String::from_utf8_lossy(name)
);
};
entry.pos()
}
fn open(path: std::path::PathBuf, fs: &Arc<dyn Fs>) -> crate::Result<Table> {
open_with_id(path, fs, 0)
}
fn open_with_id(
path: std::path::PathBuf,
fs: &Arc<dyn Fs>,
table_id: crate::TableId,
) -> crate::Result<Table> {
let checksum = crate::Checksum::from_raw(crate::repair::compute_table_checksum(&**fs, &path)?);
let mut params = crate::table::RecoverParams::new(
path,
checksum,
table_id,
Arc::clone(fs),
default_comparator(),
Arc::new(crate::cache::Cache::with_capacity_bytes(1 << 20)),
);
params.descriptor_table = Some(Arc::new(crate::descriptor_table::DescriptorTable::new(8)));
Table::recover(params)
}
#[cfg(feature = "encryption")]
fn open_encrypted(
path: std::path::PathBuf,
fs: &Arc<dyn Fs>,
encryption: Arc<dyn crate::encryption::EncryptionProvider>,
) -> crate::Result<Table> {
let checksum = crate::Checksum::from_raw(crate::repair::compute_table_checksum(&**fs, &path)?);
let mut params = crate::table::RecoverParams::new(
path,
checksum,
0,
Arc::clone(fs),
default_comparator(),
Arc::new(crate::cache::Cache::with_capacity_bytes(1 << 20)),
);
params.descriptor_table = Some(Arc::new(crate::descriptor_table::DescriptorTable::new(8)));
params.encryption = Some(encryption);
Table::recover(params)
}
fn reopen_item_count(path: std::path::PathBuf, fs: &Arc<dyn Fs>) -> crate::Result<u64> {
Ok(open(path, fs)?.metadata.item_count)
}
fn reopen_get(
path: std::path::PathBuf,
fs: &Arc<dyn Fs>,
key: &[u8],
) -> crate::Result<Option<crate::InternalValue>> {
open(path, fs)?.get(key, crate::MAX_SEQNO, crate::hash::hash64(key))
}
#[test]
fn salvage_preserves_the_source_linked_blob_files() -> crate::Result<()> {
use crate::AbstractTree;
let dir = tempdir()?;
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let crate::AnyTree::Blob(tree) = crate::Config::new(
dir.path(),
crate::SequenceNumberCounter::default(),
crate::SequenceNumberCounter::default(),
)
.with_kv_separation(Some(crate::KvSeparationOptions::default()))
.open()?
else {
unreachable!("kv separation configured");
};
let big = |i: u32| format!("{i:08}").repeat(512);
for i in 0u32..10 {
tree.insert(format!("key{i:05}"), big(i), u64::from(i) + 1);
}
tree.flush_active_memtable(10)?;
let source = {
let binding = tree.index.version_history.read().latest_version();
let Some(table) = binding.version.iter_tables().next() else {
panic!("flush produced one table");
};
(*table.path).clone()
};
drop(tree);
let dest = dir.path().join("salvaged");
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(report.is_complete(), "healthy SST: {report:?}");
let Some(source_links) = open(source, &fs)?.list_blob_file_references()? else {
panic!("the source carries a linked_blob_files section");
};
assert!(!source_links.is_empty(), "the source references blob files");
let Some(recovered_links) = open(dest, &fs)?.list_blob_file_references()? else {
panic!("the salvaged copy carries a linked_blob_files section");
};
assert_eq!(
recovered_links, source_links,
"the salvaged copy references the same blob files as the source",
);
Ok(())
}
#[test]
fn verify_blob_links_rejects_a_missing_section_with_live_indirections() -> crate::Result<()> {
use crate::AbstractTree;
let dir = tempdir()?;
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let crate::AnyTree::Blob(tree) = crate::Config::new(
dir.path(),
crate::SequenceNumberCounter::default(),
crate::SequenceNumberCounter::default(),
)
.with_kv_separation(Some(crate::KvSeparationOptions::default()))
.open()?
else {
unreachable!("kv separation configured");
};
let big = |i: u32| format!("{i:08}").repeat(512);
for i in 0u32..10 {
tree.insert(format!("key{i:05}"), big(i), u64::from(i) + 1);
}
tree.flush_active_memtable(10)?;
let source = {
let binding = tree.index.version_history.read().latest_version();
let Some(table) = binding.version.iter_tables().next() else {
panic!("flush produced one table");
};
(*table.path).clone()
};
drop(tree);
crate::test_forge::forge_section_omitted(&source, b"linked_blob_files")?;
let table = open(source, &fs)?;
assert!(
table.list_blob_file_references()?.is_none(),
"the forge must leave the table with no blob-link section",
);
let Err(err) = table.verify_blob_links() else {
panic!("a table with live indirections but no blob-link section must be rejected");
};
assert!(
matches!(
err,
crate::Error::InvalidHeader(
"table carries indirection entries but no linked_blob_files section"
)
),
"the rejection names the missing-blob-link-section reason, got {err:?}",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn verify_blob_links_rejects_a_present_empty_section() -> crate::Result<()> {
use crate::config::DeleteStrategy;
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 64u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 20, 40] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
crate::test_forge::forge_rename_and_replace_section(
&source,
b"delete_bitmap",
b"linked_blob_files",
&[0, 0, 0, 0],
)?;
let table = open(source, &fs)?;
let Err(err) = table.verify_blob_links() else {
panic!("a present zero-count linked_blob_files section must be rejected");
};
assert!(
matches!(
err,
crate::Error::InvalidHeader(
"linked_blob_files section is present but records no blob references"
)
),
"the rejection names the present-empty blob-link reason, got {err:?}",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_refuses_a_toc_that_omits_a_delete_bitmap() -> crate::Result<()> {
use crate::config::DeleteStrategy;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 64u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 20, 40] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
crate::test_forge::forge_section_omitted(&source, b"delete_bitmap")?;
let Err(err) = salvage_sst(&source, dest.clone(), &fs) else {
panic!("a TOC that hides a deletion section must be refused, not resurrected");
};
assert!(
matches!(err, crate::Error::FeatureUnsupported(msg) if msg.contains("TOC may hide a deletion section")),
"the refusal must name the hidden-deletion gate, not any other refusal, got {err:?}",
);
assert!(
!std::path::Path::new(&dest).exists(),
"no salvaged copy is produced when the deletion section may be hidden",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn verify_metadata_bounds_rejects_a_hidden_delete_bitmap() -> crate::Result<()> {
use crate::config::DeleteStrategy;
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 64u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 20, 40] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
crate::test_forge::forge_section_omitted(&source, b"delete_bitmap")?;
let table = open(source, &fs)?;
let err = reconcile_error(&table, crate::table::ReconcileGate::MetadataBounds, None);
assert!(
matches!(err, crate::Error::InvalidHeader(msg) if msg.contains("delete_bitmap count disagrees")),
"the rejection must name the delete_bitmap count mismatch, got {err:?}",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn verify_metadata_bounds_rejects_an_equal_cardinality_delete_bitmap_substitution()
-> crate::Result<()> {
use crate::config::DeleteStrategy;
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 64u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 20, 40] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
crate::test_forge::forge_delete_bitmap_substitute(&source, 0, &[6, 21, 41])?;
let table = open(source, &fs)?;
let err = reconcile_error(&table, crate::table::ReconcileGate::MetadataBounds, None);
assert!(
matches!(err, crate::Error::InvalidHeader(msg) if msg.contains("delete_bitmap contents disagree")),
"the rejection must name the delete_bitmap content-hash mismatch, got {err:?}",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_authenticates_the_delete_bitmap_before_masking() -> crate::Result<()> {
use crate::config::DeleteStrategy;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("dest");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 64u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 20, 40] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
crate::test_forge::forge_delete_bitmap_substitute(&source, 0, &[6, 21, 41])?;
let result = super::salvage_sst(&source, dest, &fs);
assert!(
matches!(&result, Err(crate::Error::InvalidHeader(msg)) if msg.contains("resurrect")),
"an unauthenticated delete-bitmap substitution must fail closed, got {result:?}",
);
Ok(())
}
#[test]
fn attempt_owns_temp_tracks_the_written_outcome_not_the_error_kind() {
let report = |salvaged_path| super::SalvageReport {
salvaged_path,
blocks_total: 1,
blocks_salvaged: 1,
blocks_copied_verbatim: 0,
entries_salvaged: 1,
entries_dropped_by_rewrite: 0,
dropped: alloc::vec::Vec::new(),
delete_rows_resurrected: false,
};
assert!(super::attempt_owns_temp(&Ok(report(Some(
std::path::PathBuf::from("/x")
)))));
assert!(!super::attempt_owns_temp(&Ok(report(None))));
assert!(!super::attempt_owns_temp(&Err(
crate::Error::FeatureUnsupported("range tombstones")
)));
}
#[test]
fn arbitrate_mirrors_propagates_a_transient_loss_to_an_incomplete_success() {
use super::{DropReason, DroppedBlock, MirrorArbitration, SalvageReport, arbitrate_mirrors};
let report = |dropped: usize| SalvageReport {
salvaged_path: Some(std::path::PathBuf::from("/x")),
blocks_total: 4,
blocks_salvaged: 4 - dropped,
blocks_copied_verbatim: 0,
entries_salvaged: 4,
entries_dropped_by_rewrite: 0,
dropped: (0..dropped)
.map(|i| DroppedBlock {
offset: i as u64,
section: b"data".to_vec(),
reason: DropReason::ChecksumMismatch,
key_range: None,
})
.collect(),
delete_rows_resurrected: false,
};
let complete = || Ok(report(0));
let incomplete = || Ok(report(1));
let transient = || {
Err(crate::Error::Io(crate::io::Error::from(
crate::io::ErrorKind::Interrupted,
)))
};
let persistent = || {
Err(crate::Error::Io(crate::io::Error::from(
crate::io::ErrorKind::Other,
)))
};
assert_eq!(
arbitrate_mirrors(&transient(), &incomplete()),
MirrorArbitration::Propagate,
);
assert_eq!(
arbitrate_mirrors(&incomplete(), &transient()),
MirrorArbitration::Propagate,
);
assert_eq!(
arbitrate_mirrors(&transient(), &complete()),
MirrorArbitration::PublishMid,
);
assert_eq!(
arbitrate_mirrors(&complete(), &transient()),
MirrorArbitration::PublishTail,
);
assert_eq!(
arbitrate_mirrors(&persistent(), &incomplete()),
MirrorArbitration::PublishMid,
);
}
#[test]
fn publish_from_temp_keeps_a_foreign_temp_on_an_erroring_attempt() -> crate::Result<()> {
let dir = tempdir()?;
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let temp = dir.path().join("dest.healtmp-0");
let dest = dir.path().join("dest");
std::fs::write(&temp, b"foreign process output")?;
let result = Err(crate::Error::FeatureUnsupported("deletion guard"));
let Err(err) = super::publish_from_temp(&fs, result, &temp, &dest, &SalvageOptions::default())
else {
panic!("an erroring attempt must propagate its error");
};
assert!(
matches!(err, crate::Error::FeatureUnsupported("deletion guard")),
"the original error propagates unchanged, got {err:?}",
);
assert!(
temp.exists(),
"the foreign temp must survive: this attempt never created it",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_refuses_a_present_but_empty_delete_bitmap() -> crate::Result<()> {
use crate::config::DeleteStrategy;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 64u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 20, 40] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
crate::test_forge::forge_delete_bitmap_empty(&source, 0)?;
let Err(err) = salvage_sst(&source, dest.clone(), &fs) else {
panic!("a present-but-empty delete bitmap must fail salvage closed");
};
assert!(
matches!(err, crate::Error::InvalidHeader(msg) if msg.contains("delete bitmap cannot be applied")),
"the refusal must name the unpositionable delete mask, got {err:?}",
);
assert!(
!std::path::Path::new(&dest).exists(),
"no salvaged copy is produced when the delete bitmap is corrupt to empty",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn verify_metadata_bounds_rejects_a_zero_count_with_a_live_bitmap() -> crate::Result<()> {
use crate::config::DeleteStrategy;
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 64u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 20, 40] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
crate::test_forge::forge_meta_value_both_mirrors(
&source,
b"descriptor#delete_bitmap_len",
&0u64.to_le_bytes(),
)?;
let table = open(source, &fs)?;
let err = reconcile_error(&table, crate::table::ReconcileGate::MetadataBounds, None);
assert!(
matches!(err, crate::Error::InvalidHeader(msg) if msg.contains("delete_bitmap count disagrees")),
"the rejection must name the delete_bitmap count mismatch, got {err:?}",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_refuses_a_reordered_columnar_index_with_deletes() -> crate::Result<()> {
use crate::config::DeleteStrategy;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 60, 120] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
crate::test_forge::forge_tli_mirrors_swap_first_two(&source, 0, None)?;
let Err(err) = salvage_sst(&source, dest.clone(), &fs) else {
panic!("a reordered columnar index with deletes must fail salvage closed");
};
assert!(
matches!(err, crate::Error::InvalidHeader(msg) if msg.contains("delete bitmap cannot be applied")),
"the refusal must name the unpositionable delete mask, got {err:?}",
);
assert!(
!std::path::Path::new(&dest).exists(),
"no salvaged copy is produced when the delete bitmap cannot be positioned",
);
Ok(())
}
#[test]
fn verify_seqno_bounds_rejects_a_present_empty_map_on_a_nonempty_table() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_seqno_in_index(true)
.use_data_block_size(128);
for i in 0u64..64 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_seqno_bounds_empty(&source, 0)?;
let table = open(source, &fs)?;
let err = reconcile_error(&table, crate::table::ReconcileGate::SeqnoBounds, None);
assert!(
matches!(
err,
crate::Error::InvalidHeader("seqno_bounds is missing a data block's entry")
),
"the rejection names the missing-block-entry reason, got {err:?}",
);
Ok(())
}
#[cfg(feature = "zstd")]
#[test]
fn verify_block_layout_rejects_a_present_empty_map_on_a_nonempty_table() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_data_block_size(256 * 1024)
.use_data_block_compression(crate::CompressionType::Zstd(19));
for i in 0u64..20_000 {
writer.write(InternalValue::from_components(
format!("key-{i:012}").into_bytes(),
format!("value-{i:08}-payload").into_bytes(),
1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_block_layout_empty(&source, 0, None)?;
let table = open(source, &fs)?;
let Err(err) = table.verify_block_layout() else {
panic!("an empty block_layout on a table with data blocks must be rejected");
};
assert!(
matches!(
err,
crate::Error::InvalidHeader(
"block_layout section is present but empty on a table with data blocks"
)
),
"the rejection names the present-empty block_layout reason, got {err:?}",
);
Ok(())
}
#[cfg(feature = "zstd")]
#[test]
fn reconcile_gates_with_a_shifted_block_layout_boundary_reject_it() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_data_block_size(256 * 1024)
.use_data_block_compression(crate::CompressionType::Zstd(19));
for i in 0u64..20_000 {
writer.write(InternalValue::from_components(
format!("key-{i:012}").into_bytes(),
format!("value-{i:08}-payload").into_bytes(),
1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let table = open(source.clone(), &fs)?;
assert!(
table.verify_reconcile_gates(None, false).is_ok(),
"the intact multi-inner-block table must reconcile clean",
);
drop(table);
crate::test_forge::forge_block_layout_shifted_end(&source, 0)?;
let table = open(source, &fs)?;
let Err((gate, err)) = table.verify_reconcile_gates(None, false) else {
panic!("a boundary that does not match the frame's inner blocks must be rejected");
};
assert!(
matches!(gate, crate::table::ReconcileGate::BlockLayout),
"the failure must be attributed to the block-layout gate, got {gate:?}",
);
assert!(
matches!(
err,
crate::Error::InvalidHeader("block_layout disagrees with the frames' inner blocks")
),
"the rejection names the frame disagreement, got {err:?}",
);
Ok(())
}
#[cfg(all(feature = "columnar", not(feature = "zstd")))]
#[test]
fn verify_block_layout_rejects_an_empty_map_without_zstd() -> crate::Result<()> {
use crate::config::DeleteStrategy;
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 64u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 20, 40] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
crate::test_forge::forge_delete_bitmap_as_empty_block_layout(&source, 0)?;
let table = open(source, &fs)?;
let Err(err) = table.verify_block_layout() else {
panic!("an empty block_layout on a non-zstd table with data blocks must be rejected");
};
assert!(
matches!(
err,
crate::Error::InvalidHeader(
"block_layout section is present but empty on a table with data blocks"
)
),
"the rejection names the present-empty block_layout reason, got {err:?}",
);
Ok(())
}
#[test]
fn verify_filter_rejects_a_present_empty_full_filter() -> crate::Result<()> {
use crate::config::BloomConstructionPolicy;
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_bloom_policy(BloomConstructionPolicy::BitsPerKey(10.0));
for i in 0u32..200 {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_filter_empty(&source, 0)?;
let table = open(source, &fs)?;
let err = reconcile_error(&table, crate::table::ReconcileGate::Filter, None);
assert!(
matches!(
err,
crate::Error::InvalidHeader(
"filter section is present but empty on a table with data blocks"
)
),
"the rejection names the present-empty filter reason, got {err:?}",
);
Ok(())
}
#[test]
fn verify_filter_rejects_an_empty_partition_for_an_existing_key() -> crate::Result<()> {
use crate::config::BloomConstructionPolicy;
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_partitioned_filter()
.use_meta_partition_size(3)
.use_bloom_policy(BloomConstructionPolicy::BitsPerKey(10.0));
for i in 0u32..64 {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_filter_first_partition_empty(&source)?;
let table = open(source, &fs)?;
let err = reconcile_error(&table, crate::table::ReconcileGate::Filter, None);
assert!(
matches!(
err,
crate::Error::InvalidHeader(
"filter partition is present but empty for an existing key"
)
),
"the rejection names the present-empty partition reason, got {err:?}",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn verify_zone_map_rejects_a_present_empty_map_on_a_nonempty_table() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true);
for i in 0u32..64 {
writer.write(iv(i))?;
}
assert!(
writer.finish()?.is_some(),
"source columnar+zonemap SST is non-empty"
);
crate::test_forge::forge_zone_map_empty(&source, 0)?;
let table = open(source, &fs)?;
let err = reconcile_error(&table, crate::table::ReconcileGate::ZoneMap, None);
assert!(
matches!(
err,
crate::Error::InvalidHeader(
"zone_map section is present but empty on a table with data blocks"
)
),
"the rejection names the present-empty zone_map reason, got {err:?}",
);
Ok(())
}
#[test]
fn verify_zone_map_rejects_a_forged_synthetic_column_id() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_zone_map(true);
for i in 0u32..64 {
writer.write(iv(i))?;
}
assert!(
writer.finish()?.is_some(),
"source row+zonemap SST is non-empty"
);
crate::test_forge::forge_zone_map_column_id(&source, 0)?;
let table = open(source, &fs)?;
let err = reconcile_error(&table, crate::table::ReconcileGate::ZoneMap, None);
assert!(
matches!(err, crate::Error::InvalidHeader(msg) if msg.contains("zone_map synthetic column")),
"the rejection must name the synthetic column identity, got {err:?}",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn verify_zone_map_rejects_a_forged_columnar_column_id() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true);
for i in 0u32..64 {
writer.write(iv(i))?;
}
assert!(
writer.finish()?.is_some(),
"source columnar+zonemap SST is non-empty"
);
crate::test_forge::forge_zone_map_column_id(&source, 0)?;
let table = open(source, &fs)?;
let err = reconcile_error(&table, crate::table::ReconcileGate::ZoneMap, None);
assert!(
matches!(err, crate::Error::InvalidHeader(msg) if msg.contains("per-column statistics")),
"the rejection must name the columnar per-column mismatch, got {err:?}",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvaged_columnar_table_keeps_per_column_zone_statistics() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(128);
for i in 0u32..64 {
writer.write(iv(i))?;
}
assert!(
writer.finish()?.is_some(),
"source columnar SST is non-empty"
);
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(
report.salvaged_path.is_some(),
"the clean SST salvages: {report:?}",
);
assert!(
report.blocks_copied_verbatim > 0,
"at least one clean columnar block is byte-copied verbatim: {report:?}",
);
let table = open(dest, &fs)?;
assert!(table.has_zone_map(), "the salvaged copy carries a zone map");
reconcile_clean(&table, None);
Ok(())
}
#[cfg(all(feature = "encryption", feature = "zstd"))]
#[test]
fn verify_block_layout_rejects_an_empty_map_on_an_encrypted_table() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let enc: Arc<dyn crate::encryption::EncryptionProvider> =
Arc::new(crate::encryption::Aes256GcmProvider::new(&[0x42; 32]));
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_data_block_size(256 * 1024)
.use_data_block_compression(crate::CompressionType::Zstd(19))
.use_encryption(Some(Arc::clone(&enc)));
for i in 0u64..20_000 {
writer.write(InternalValue::from_components(
format!("key-{i:012}").into_bytes(),
format!("value-{i:08}-payload").into_bytes(),
1,
ValueType::Value,
))?;
}
assert!(
writer.finish()?.is_some(),
"source encrypted SST is non-empty",
);
crate::test_forge::forge_block_layout_empty(&source, 0, Some(&*enc))?;
let table = open_encrypted(source, &fs, Arc::clone(&enc))?;
let Err(err) = table.verify_block_layout() else {
panic!("an empty block_layout on an encrypted table with data blocks must be rejected");
};
assert!(
matches!(
err,
crate::Error::InvalidHeader(
"block_layout section is present but empty on a table with data blocks"
)
),
"the rejection names the present-empty block_layout reason, got {err:?}",
);
Ok(())
}
#[test]
fn verify_metadata_bounds_rejects_a_raised_seqno_min_on_a_range_tombstone_table()
-> crate::Result<()> {
use crate::UserKey;
use crate::range_tombstone::RangeTombstone;
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0u64..8 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
writer.write_range_tombstone(RangeTombstone::new(
UserKey::from(b"key-002".as_slice()),
UserKey::from(b"key-005".as_slice()),
9,
));
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_meta_value_both_mirrors(&source, b"seqno#min", &5u64.to_le_bytes())?;
let table = open(source, &fs)?;
let err = reconcile_error(&table, crate::table::ReconcileGate::MetadataBounds, None);
assert!(
matches!(
err,
crate::Error::InvalidHeader("meta seqno#min is above the decoded minimum seqno")
),
"the rejection names the seqno#min branch specifically, got {err:?}",
);
Ok(())
}
#[test]
fn verify_metadata_bounds_keeps_a_real_weak_tombstone_matching_the_rt_sentinel() -> crate::Result<()>
{
use crate::UserKey;
use crate::range_tombstone::RangeTombstone;
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
writer.write(InternalValue::from_components(
b"key-000".to_vec(),
b"val-000".to_vec(),
4,
ValueType::Value,
))?;
writer.write(InternalValue::from_components(
b"key-001".to_vec(),
b"val-001".to_vec(),
5,
ValueType::Value,
))?;
writer.write(InternalValue::new_weak_tombstone(
UserKey::from(b"key-002".as_slice()),
3,
))?;
for i in 3u64..8 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 3,
ValueType::Value,
))?;
}
writer.write_range_tombstone(RangeTombstone::new(
UserKey::from(b"key-002".as_slice()),
UserKey::from(b"key-005".as_slice()),
3,
));
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_meta_value_both_mirrors(&source, b"seqno#min", &4u64.to_le_bytes())?;
let table = open(source, &fs)?;
let err = reconcile_error(&table, crate::table::ReconcileGate::MetadataBounds, None);
assert!(
matches!(
err,
crate::Error::InvalidHeader("meta seqno#min is above the decoded minimum seqno")
),
"the rejection names the seqno#min branch specifically, got {err:?}",
);
Ok(())
}
#[test]
fn verify_metadata_bounds_rejects_a_range_tombstone_count_mismatch() -> crate::Result<()> {
use crate::UserKey;
use crate::range_tombstone::RangeTombstone;
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0u64..8 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
writer.write_range_tombstone(RangeTombstone::new(
UserKey::from(b"key-002".as_slice()),
UserKey::from(b"key-005".as_slice()),
9,
));
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_section_omitted(&source, b"range_tombstones")?;
let table = open(source, &fs)?;
let err = reconcile_error(&table, crate::table::ReconcileGate::MetadataBounds, None);
assert!(
matches!(
err,
crate::Error::InvalidHeader(
"range_tombstones count disagrees with the recorded range_tombstone_count"
)
),
"the rejection names the RT count-mismatch reason, got {err:?}",
);
Ok(())
}
#[test]
fn verify_tli_mirrors_rejects_a_forged_partition_boundary() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_partitioned_index()
.use_data_block_size(128)
.use_meta_partition_size(2);
for i in 0u64..512 {
writer.write(InternalValue::from_components(
format!("key-{i:05}").into_bytes(),
format!("val-{i:05}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tli_mirrors_lower_first_separator(&source, 0, None)?;
let table = open(source, &fs)?;
let Err(err) = table.verify_tli_mirrors() else {
panic!("a forged partition boundary must be rejected");
};
assert!(
matches!(
err,
crate::Error::InvalidHeader(
"tli separator disagrees with its partition's last separator"
)
),
"the rejection names the partition-boundary mismatch, got {err:?}",
);
Ok(())
}
#[test]
fn salvage_rebuilds_blob_links_from_recovered_indirections() -> crate::Result<()> {
use crate::AbstractTree;
let dir = tempdir()?;
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let crate::AnyTree::Blob(tree) = crate::Config::new(
dir.path(),
crate::SequenceNumberCounter::default(),
crate::SequenceNumberCounter::default(),
)
.with_kv_separation(Some(crate::KvSeparationOptions::default()))
.open()?
else {
unreachable!("kv separation configured");
};
let big = |i: u32| format!("{i:08}").repeat(512);
for i in 0u32..10 {
tree.insert(format!("key{i:05}"), big(i), u64::from(i) + 1);
}
tree.flush_active_memtable(10)?;
let source = {
let binding = tree.index.version_history.read().latest_version();
let Some(table) = binding.version.iter_tables().next() else {
panic!("flush produced one table");
};
(*table.path).clone()
};
drop(tree);
let Some(true_links) = open(source.clone(), &fs)?.list_blob_file_references()? else {
panic!("the source carries a linked_blob_files section");
};
assert!(!true_links.is_empty(), "the source references blob files");
let pos = {
let mut f = std::fs::File::open(&source)?;
let reader = match crate::sfa::Reader::from_reader(&mut f) {
Ok(r) => r,
Err(e) => panic!("reading the SFA trailer failed: {e:?}"),
};
let Some(entry) = reader
.toc()
.iter()
.find(|e| e.name() == b"linked_blob_files")
else {
panic!("the source must carry a linked_blob_files section");
};
usize::try_from(entry.pos()).unwrap_or(usize::MAX)
};
let mut bytes = std::fs::read(&source)?;
let Some(count) = bytes.get_mut(pos..pos + 4) else {
panic!("linked_blob_files count header within the file");
};
count.copy_from_slice(&0u32.to_le_bytes());
std::fs::write(&source, &bytes)?;
let Some(forged) = open(source.clone(), &fs)?.list_blob_file_references()? else {
panic!("the forged section still parses");
};
assert!(forged.is_empty(), "the forged count hides every link");
let dest = dir.path().join("salvaged");
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(report.is_complete(), "data blocks are healthy: {report:?}");
let Some(recovered_links) = open(dest, &fs)?.list_blob_file_references()? else {
panic!("the salvaged copy must carry links derived from its indirections");
};
for link in &true_links {
assert!(
recovered_links
.iter()
.any(|l| l.blob_file_id == link.blob_file_id),
"blob file {} is referenced by recovered indirections but missing \
from the copy's links: {recovered_links:?}",
link.blob_file_id,
);
}
Ok(())
}
#[test]
fn salvage_drops_source_only_blob_links() -> crate::Result<()> {
use crate::AbstractTree;
let dir = tempdir()?;
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let crate::AnyTree::Blob(tree) = crate::Config::new(
dir.path(),
crate::SequenceNumberCounter::default(),
crate::SequenceNumberCounter::default(),
)
.with_kv_separation(Some(crate::KvSeparationOptions::default()))
.open()?
else {
unreachable!("kv separation configured");
};
let big = |i: u32| format!("{i:08}").repeat(512);
for i in 0u32..10 {
tree.insert(format!("key{i:05}"), big(i), u64::from(i) + 1);
}
tree.flush_active_memtable(10)?;
let source = {
let binding = tree.index.version_history.read().latest_version();
let Some(table) = binding.version.iter_tables().next() else {
panic!("flush produced one table");
};
(*table.path).clone()
};
drop(tree);
const FORGED_ID: u64 = 9_999_999;
let pos = {
let mut f = std::fs::File::open(&source)?;
let reader = match crate::sfa::Reader::from_reader(&mut f) {
Ok(r) => r,
Err(e) => panic!("reading the SFA trailer failed: {e:?}"),
};
let Some(entry) = reader
.toc()
.iter()
.find(|e| e.name() == b"linked_blob_files")
else {
panic!("the source must carry a linked_blob_files section");
};
usize::try_from(entry.pos()).unwrap_or(usize::MAX)
};
let mut bytes = std::fs::read(&source)?;
let Some(id_slot) = bytes.get_mut(pos + 4..pos + 12) else {
panic!("first link record within the file");
};
id_slot.copy_from_slice(&FORGED_ID.to_le_bytes());
std::fs::write(&source, &bytes)?;
let Some(forged) = open(source.clone(), &fs)?.list_blob_file_references()? else {
panic!("the forged section still parses");
};
assert!(
forged.iter().any(|l| l.blob_file_id == FORGED_ID),
"the source's section carries the forged id",
);
let dest = dir.path().join("salvaged");
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(report.is_complete(), "data blocks are healthy: {report:?}");
let Some(recovered_links) = open(dest, &fs)?.list_blob_file_references()? else {
panic!("the salvaged copy carries links derived from its indirections");
};
assert!(
!recovered_links.is_empty(),
"the recovered indirections reference the real blob file",
);
assert!(
!recovered_links.iter().any(|l| l.blob_file_id == FORGED_ID),
"a source-only id no recovered indirection references must not be \
copied into the salvaged table: {recovered_links:?}",
);
Ok(())
}
#[test]
fn salvage_recovers_with_derived_links_when_the_blob_link_section_is_unreadable()
-> crate::Result<()> {
use crate::AbstractTree;
let dir = tempdir()?;
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let crate::AnyTree::Blob(tree) = crate::Config::new(
dir.path(),
crate::SequenceNumberCounter::default(),
crate::SequenceNumberCounter::default(),
)
.with_kv_separation(Some(crate::KvSeparationOptions::default()))
.open()?
else {
unreachable!("kv separation configured");
};
let big = |i: u32| format!("{i:08}").repeat(512);
for i in 0u32..10 {
tree.insert(format!("key{i:05}"), big(i), u64::from(i) + 1);
}
tree.flush_active_memtable(10)?;
let source = {
let binding = tree.index.version_history.read().latest_version();
let Some(table) = binding.version.iter_tables().next() else {
panic!("flush produced one table");
};
(*table.path).clone()
};
drop(tree);
let pos = {
let mut f = std::fs::File::open(&source)?;
let reader = match crate::sfa::Reader::from_reader(&mut f) {
Ok(r) => r,
Err(e) => panic!("reading the SFA trailer failed: {e:?}"),
};
let Some(entry) = reader
.toc()
.iter()
.find(|e| e.name() == b"linked_blob_files")
else {
panic!("the source must carry a linked_blob_files section");
};
usize::try_from(entry.pos()).unwrap_or(usize::MAX)
};
let mut bytes = std::fs::read(&source)?;
let Some(count) = bytes.get_mut(pos..pos + 4) else {
panic!("linked_blob_files count header within the file");
};
count.copy_from_slice(&u32::MAX.to_le_bytes());
std::fs::write(&source, &bytes)?;
let dest = dir.path().join("salvaged");
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(report.is_complete(), "data blocks are healthy: {report:?}");
let Some(recovered_links) = open(dest, &fs)?.list_blob_file_references()? else {
panic!("the salvaged copy must carry links derived from its indirections");
};
assert!(
!recovered_links.is_empty(),
"the recovered indirections reference at least one blob file",
);
Ok(())
}
#[test]
fn salvage_preserves_a_nonzero_source_table_id() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
const TID: crate::TableId = 7;
let mut writer = Writer::new(source.clone(), TID, 0, Arc::clone(&fs))?.use_data_block_size(256);
let n = 200u32;
for i in 0..n {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(
report.is_complete(),
"a healthy non-zero-id SST salvages completely: {report:?}",
);
let recovered = open_with_id(dest, &fs, TID)?;
assert_eq!(
recovered.metadata.id, TID,
"the salvaged copy is stamped with the source's table id",
);
assert_eq!(recovered.metadata.item_count, u64::from(n));
Ok(())
}
#[test]
fn salvage_of_a_healthy_sst_recovers_every_block() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(256);
let n = 200u32;
for i in 0..n {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(
report.is_complete(),
"a healthy SST salvages with no dropped blocks: {report:?}",
);
assert!(
report.blocks_total >= 2,
"256-byte blocks over 200 entries should yield several data blocks, got {}",
report.blocks_total,
);
assert_eq!(
report.blocks_salvaged, report.blocks_total,
"every block of a healthy SST is salvaged",
);
assert_eq!(
report.entries_salvaged,
u64::from(n),
"every entry is recovered",
);
assert_eq!(
report.salvaged_path.as_deref(),
Some(dest.as_path()),
"a salvaged file is written when at least one block is recovered",
);
assert_eq!(
report.blocks_copied_verbatim, report.blocks_salvaged,
"a healthy SST's blocks are all copied verbatim",
);
assert_eq!(
reopen_item_count(dest, &fs)?,
u64::from(n),
"the salvaged SST reopens with the full item count",
);
Ok(())
}
#[test]
fn salvage_copies_a_clean_block_verbatim() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(256);
let n = 200u32;
for i in 0..n {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(
report.is_complete(),
"healthy SST salvages clean: {report:?}"
);
assert_eq!(
report.blocks_copied_verbatim, report.blocks_total,
"every clean block is copied verbatim, none re-encoded",
);
let first_block = |path: &std::path::Path| -> crate::Result<(usize, usize)> {
let table = open(path.to_path_buf(), &fs)?;
let Some(kh) = table.data_block_handles().find_map(Result::ok) else {
panic!("a non-empty SST has at least one data block");
};
let off = usize::try_from(*kh.as_ref().offset()).unwrap_or(usize::MAX);
Ok((off, kh.as_ref().size() as usize))
};
let (src_off, src_size) = first_block(&source)?;
let (dst_off, dst_size) = first_block(&dest)?;
assert_eq!(
src_size, dst_size,
"the verbatim copy preserves the block's on-disk size",
);
let src_bytes = std::fs::read(&source)?;
let dst_bytes = std::fs::read(&dest)?;
let src_block = src_bytes.get(src_off..src_off + src_size);
let dst_block = dst_bytes.get(dst_off..dst_off + dst_size);
assert!(
src_block.is_some() && src_block == dst_block,
"the clean block is copied byte-for-byte into the salvaged SST",
);
Ok(())
}
#[test]
fn salvage_drops_a_corrupted_block_and_keeps_the_rest() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(256);
let n = 200u32;
for i in 0..n {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let target = {
let table = open(source.clone(), &fs)?;
let offsets: alloc::vec::Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|kh| *kh.as_ref().offset())
.collect();
let Some(&second) = offsets.get(1) else {
panic!("source SST must have at least two data blocks, got {offsets:?}");
};
second
};
let header_len = crate::table::block::Header::header_len(crate::table::block::BlockType::Data);
let Ok(target_usize) = usize::try_from(target) else {
panic!("data block offset {target} does not fit usize on this target");
};
let flip = target_usize + header_len + 8;
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(flip) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(
!report.is_complete(),
"a corrupted block must be reported as dropped: {report:?}",
);
assert_eq!(
report.dropped.len(),
1,
"exactly the one corrupted block is dropped: {report:?}",
);
assert_eq!(
report.blocks_salvaged,
report.blocks_total - 1,
"every block but the corrupted one is recovered",
);
assert!(
report.entries_salvaged > 0 && report.entries_salvaged < u64::from(n),
"a partial key range is recovered, got {} of {n}",
report.entries_salvaged,
);
assert!(
report.dropped.first().is_some_and(|d| {
matches!(d.reason, DropReason::ChecksumMismatch) && d.key_range.is_some()
}),
"the dropped block reports a checksum mismatch and names the key range it lost: {report:?}",
);
assert_eq!(report.salvaged_path.as_deref(), Some(dest.as_path()));
assert_eq!(
reopen_item_count(dest, &fs)?,
report.entries_salvaged,
"the salvaged SST holds exactly the entries the report counted",
);
Ok(())
}
#[test]
fn salvage_counts_a_header_corrupt_block_in_blocks_total() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(256);
let n = 200u32;
for i in 0..n {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let target = {
let table = open(source.clone(), &fs)?;
let offsets: alloc::vec::Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|kh| *kh.as_ref().offset())
.collect();
let Some(&second) = offsets.get(1) else {
panic!("source SST must have at least two data blocks, got {offsets:?}");
};
second
};
let Ok(target_usize) = usize::try_from(target) else {
panic!("data block offset {target} does not fit usize on this target");
};
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(target_usize) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest, &fs)?;
assert!(
!report.dropped.is_empty(),
"the header-corrupt block must be reported dropped: {report:?}",
);
assert_eq!(
report.blocks_total,
report.blocks_salvaged + report.dropped.len(),
"every inspected block is either recovered or dropped, a header-corrupt \
block must count toward the total, not vanish from it: {report:?}",
);
Ok(())
}
#[cfg(feature = "page_ecc")]
#[test]
fn salvage_reencodes_an_ecc_recovered_block_rather_than_copying_it() -> crate::Result<()> {
use crate::table::block::{EccParams, Header};
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_data_block_size(256)
.use_ecc(Some(EccParams::RS_4_2));
let n = 200u32;
for i in 0..n {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let first_off = {
let table = open(source.clone(), &fs)?;
let Some(kh) = table.data_block_handles().find_map(Result::ok) else {
panic!("a non-empty SST has at least one data block");
};
*kh.as_ref().offset()
};
let pos = usize::try_from(first_off).unwrap_or(usize::MAX) + Header::MIN_LEN + 3;
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(pos) {
*b ^= 0x80;
}
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(
report.is_complete(),
"an ECC-recoverable block is healed, not dropped: {report:?}",
);
assert_eq!(
report.blocks_salvaged, report.blocks_total,
"every block is recovered",
);
assert_eq!(
report.entries_salvaged,
u64::from(n),
"every entry is recovered",
);
assert_eq!(
report.blocks_copied_verbatim,
report.blocks_salvaged - 1,
"the ECC-recovered block is re-encoded, not copied verbatim",
);
assert_eq!(reopen_item_count(dest, &fs)?, u64::from(n));
Ok(())
}
#[cfg(feature = "page_ecc")]
#[test]
fn salvage_regenerates_a_rotted_parity_trailer_rather_than_copying_it() -> crate::Result<()> {
use crate::coding::Decode;
use crate::table::block::{EccParams, Header, expected_parity_len};
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_data_block_size(256)
.use_ecc(Some(EccParams::RS_4_2));
let n = 200u32;
for i in 0..n {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let first_off = {
let table = open(source.clone(), &fs)?;
let Some(kh) = table.data_block_handles().find_map(Result::ok) else {
panic!("a non-empty SST has at least one data block");
};
usize::try_from(*kh.as_ref().offset()).unwrap_or(usize::MAX)
};
let mut bytes = std::fs::read(&source)?;
let Some(mut cursor) = bytes.get(first_off..) else {
panic!("first data block within the file");
};
let header = Header::decode_from(&mut cursor)?;
let trailer_pos =
first_off + Header::header_len(header.block_type) + header.data_length as usize;
let Some(slot) = bytes.get_mut(trailer_pos) else {
panic!("parity trailer within the file");
};
*slot ^= 0xFF;
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(
report.is_complete(),
"the payload is intact, every block is recovered: {report:?}",
);
assert_eq!(reopen_item_count(dest.clone(), &fs)?, u64::from(n));
let dest_bytes = std::fs::read(&dest)?;
let dest_table = open(dest, &fs)?;
let (ds, ps) = EccParams::RS_4_2.as_shards();
for kh in dest_table.data_block_handles() {
let kh = kh?;
let off = usize::try_from(*kh.as_ref().offset()).unwrap_or(usize::MAX);
let Some(mut cursor) = dest_bytes.get(off..) else {
panic!("block at offset {off} within the file");
};
let hdr = Header::decode_from(&mut cursor)?;
let hlen = Header::header_len(hdr.block_type);
let dl = hdr.data_length as usize;
let Some(payload) = dest_bytes.get(off + hlen..off + hlen + dl) else {
panic!("payload of block at offset {off} within the file");
};
let plen = expected_parity_len(hdr.data_length, EccParams::RS_4_2) as usize;
let Some(trailer) = dest_bytes.get(off + hlen + dl..off + hlen + dl + plen) else {
panic!("parity trailer of block at offset {off} within the file");
};
let fresh = crate::ecc::encode_parity(payload, ds, ps)?;
assert_eq!(
trailer,
fresh.as_slice(),
"block at offset {off}: the salvaged copy's parity matches its payload",
);
}
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_drops_a_corrupted_columnar_block_and_keeps_the_rest() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256);
let n = 200u32;
for i in 0..n {
writer.write(iv(i))?;
}
assert!(
writer.finish()?.is_some(),
"source columnar SST is non-empty"
);
let target = {
let table = open(source.clone(), &fs)?;
let offsets: alloc::vec::Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|kh| *kh.as_ref().offset())
.collect();
let Some(&second) = offsets.get(1) else {
panic!("source columnar SST must have at least two data blocks, got {offsets:?}");
};
second
};
let header_len = crate::table::block::Header::header_len(crate::table::block::BlockType::Data);
let Ok(target_usize) = usize::try_from(target) else {
panic!("data block offset {target} does not fit usize on this target");
};
let flip = target_usize + header_len + 8;
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(flip) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(
report.dropped.len(),
1,
"exactly the one corrupted columnar block is dropped: {report:?}",
);
assert_eq!(
report.blocks_salvaged,
report.blocks_total - 1,
"every columnar block but the corrupted one is recovered",
);
assert!(
report.entries_salvaged > 0 && report.entries_salvaged < u64::from(n),
"a partial key range is recovered, got {} of {n}",
report.entries_salvaged,
);
assert_eq!(report.salvaged_path.as_deref(), Some(dest.as_path()));
let recovered = open(dest, &fs)?;
assert_eq!(recovered.metadata.item_count, report.entries_salvaged);
assert!(
recovered.metadata.columnar,
"a columnar source salvages into a columnar copy, not a row-major one",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_drops_a_columnar_block_with_an_invalid_value_type() -> crate::Result<()> {
use crate::coding::{Decode, Encode};
use crate::table::block::Header;
use crate::table::columnar::{CodecId, ColumnBatch};
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256);
for i in 0..n {
writer.write(iv(i))?;
}
assert!(
writer.finish()?.is_some(),
"source columnar SST is non-empty"
);
let (block_off, block_size) = {
let table = open(source.clone(), &fs)?;
let Some(kh) = table.data_block_handles().filter_map(Result::ok).nth(1) else {
panic!("source must have at least two data blocks");
};
(
usize::try_from(*kh.as_ref().offset()).unwrap_or(usize::MAX),
kh.as_ref().size() as usize,
)
};
let mut bytes = std::fs::read(&source)?;
let Some(block) = bytes.get(block_off..block_off + block_size) else {
panic!("block range within the file");
};
let mut cursor = block;
let header = Header::decode_from(&mut cursor)?;
let header_len = Header::header_len(header.block_type);
let Some(payload) = block.get(header_len..header_len + header.data_length as usize) else {
panic!("payload range within the block");
};
let mut batch = ColumnBatch::decode(&payload.into())?;
let poisoned_rows = u64::from(batch.row_count);
let Some(col) = batch.columns.get_mut(2).filter(|c| !c.data.is_empty()) else {
panic!("value-type column present and non-empty");
};
let mut poisoned = col.data.to_vec();
let Some(first_byte) = poisoned.first_mut() else {
panic!("column is non-empty");
};
*first_byte = 0xFF;
col.data = poisoned.into();
let new_payload = batch.encode(CodecId::Plain)?;
assert_eq!(
new_payload.len(),
payload.len(),
"a one-byte in-place mutation re-encodes to the same length",
);
let new_header = Header {
checksum: crate::Checksum::from_raw(crate::hash::hash128(&new_payload)),
..header
};
let mut new_block = Vec::with_capacity(header_len + new_payload.len());
new_header.encode_into(&mut new_block)?;
assert_eq!(
new_block.len(),
header_len,
"header re-encodes to its length"
);
new_block.extend_from_slice(&new_payload);
let Some(target) = bytes.get_mut(block_off..block_off + new_block.len()) else {
panic!("block range within the file");
};
target.copy_from_slice(&new_block);
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(
report.dropped.len(),
1,
"exactly the poisoned block is dropped: {report:?}",
);
assert!(
matches!(
report.dropped.first().map(|d| &d.reason),
Some(DropReason::DecodeError(_))
),
"the invalid value-type tag classifies as a decode error: {report:?}",
);
assert_eq!(
report.entries_salvaged,
u64::from(n) - poisoned_rows,
"every row outside the poisoned block is recovered",
);
assert_eq!(reopen_item_count(dest, &fs)?, u64::from(n) - poisoned_rows);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_drops_a_columnar_block_with_out_of_order_keys() -> crate::Result<()> {
use crate::coding::{Decode, Encode};
use crate::table::block::Header;
use crate::table::columnar::{CodecId, ColumnBatch};
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256);
for i in 0..n {
writer.write(iv(i))?;
}
assert!(
writer.finish()?.is_some(),
"source columnar SST is non-empty"
);
let (block_off, block_size) = {
let table = open(source.clone(), &fs)?;
let Some(kh) = table.data_block_handles().filter_map(Result::ok).nth(1) else {
panic!("source must have at least two data blocks");
};
(
usize::try_from(*kh.as_ref().offset()).unwrap_or(usize::MAX),
kh.as_ref().size() as usize,
)
};
let mut bytes = std::fs::read(&source)?;
let Some(block) = bytes.get(block_off..block_off + block_size) else {
panic!("block range within the file");
};
let mut cursor = block;
let header = Header::decode_from(&mut cursor)?;
let header_len = Header::header_len(header.block_type);
let Some(payload) = block.get(header_len..header_len + header.data_length as usize) else {
panic!("payload range within the block");
};
let mut batch = ColumnBatch::decode(&payload.into())?;
let poisoned_rows = u64::from(batch.row_count);
assert!(batch.row_count >= 2, "block holds at least two rows");
{
let Some(key_col) = batch.columns.first_mut() else {
panic!("key column present");
};
let table_len = (batch.row_count as usize + 1) * 4;
let off = |data: &[u8], idx: usize| -> usize {
let Some(b) = data.get(idx * 4..idx * 4 + 4) else {
panic!("offset {idx} within the frame table");
};
u32::from_le_bytes(b.try_into().unwrap_or([0; 4])) as usize
};
let (o0, o1, o2) = (
off(&key_col.data, 0),
off(&key_col.data, 1),
off(&key_col.data, 2),
);
assert_eq!(o1 - o0, o2 - o1, "adjacent keys are equal-length");
let len = o1 - o0;
let Some(first) = key_col.data.get(table_len + o0..table_len + o0 + len) else {
panic!("first key within the column");
};
let first = first.to_vec();
let Some(second) = key_col.data.get(table_len + o1..table_len + o1 + len) else {
panic!("second key within the column");
};
let second = second.to_vec();
let mut swapped = key_col.data.to_vec();
let Some(dst0) = swapped.get_mut(table_len + o0..table_len + o0 + len) else {
panic!("first key range within the column");
};
dst0.copy_from_slice(&second);
let Some(dst1) = swapped.get_mut(table_len + o1..table_len + o1 + len) else {
panic!("second key range within the column");
};
dst1.copy_from_slice(&first);
key_col.data = swapped.into();
}
let new_payload = batch.encode(CodecId::Plain)?;
assert_eq!(
new_payload.len(),
payload.len(),
"an in-place key swap re-encodes to the same length",
);
let new_header = Header {
checksum: crate::Checksum::from_raw(crate::hash::hash128(&new_payload)),
..header
};
let mut new_block = Vec::with_capacity(header_len + new_payload.len());
new_header.encode_into(&mut new_block)?;
new_block.extend_from_slice(&new_payload);
let Some(target) = bytes.get_mut(block_off..block_off + new_block.len()) else {
panic!("block range within the file");
};
target.copy_from_slice(&new_block);
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(
report.dropped.len(),
1,
"exactly the out-of-order block is dropped: {report:?}",
);
assert!(
matches!(
report.dropped.first().map(|d| &d.reason),
Some(DropReason::DecodeError(_))
),
"the ordering violation classifies as a decode error: {report:?}",
);
assert_eq!(
report.entries_salvaged,
u64::from(n) - poisoned_rows,
"every row outside the poisoned block is recovered",
);
assert_eq!(reopen_item_count(dest, &fs)?, u64::from(n) - poisoned_rows);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn verify_point_read_reachability_rejects_a_reordered_columnar_block() -> crate::Result<()> {
use crate::coding::{Decode, Encode};
use crate::table::block::Header;
use crate::table::columnar::{CodecId, ColumnBatch};
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256);
for i in 0..n {
writer.write(iv(i))?;
}
assert!(
writer.finish()?.is_some(),
"source columnar SST is non-empty"
);
let (block_off, block_size) = {
let table = open(source.clone(), &fs)?;
let Some(kh) = table.data_block_handles().filter_map(Result::ok).nth(1) else {
panic!("source must have at least two data blocks");
};
(
usize::try_from(*kh.as_ref().offset()).unwrap_or(usize::MAX),
kh.as_ref().size() as usize,
)
};
let mut bytes = std::fs::read(&source)?;
let Some(block) = bytes.get(block_off..block_off + block_size) else {
panic!("block range within the file");
};
let mut cursor = block;
let header = Header::decode_from(&mut cursor)?;
let header_len = Header::header_len(header.block_type);
let Some(payload) = block.get(header_len..header_len + header.data_length as usize) else {
panic!("payload range within the block");
};
let mut batch = ColumnBatch::decode(&payload.into())?;
assert!(batch.row_count >= 2, "block holds at least two rows");
{
let Some(key_col) = batch.columns.first_mut() else {
panic!("key column present");
};
let table_len = (batch.row_count as usize + 1) * 4;
let off = |data: &[u8], idx: usize| -> usize {
let Some(b) = data.get(idx * 4..idx * 4 + 4) else {
panic!("offset {idx} within the frame table");
};
u32::from_le_bytes(b.try_into().unwrap_or([0; 4])) as usize
};
let (o0, o1, o2) = (
off(&key_col.data, 0),
off(&key_col.data, 1),
off(&key_col.data, 2),
);
assert_eq!(o1 - o0, o2 - o1, "adjacent keys are equal-length");
let len = o1 - o0;
let Some(first) = key_col.data.get(table_len + o0..table_len + o0 + len) else {
panic!("first key within the column");
};
let first = first.to_vec();
let Some(second) = key_col.data.get(table_len + o1..table_len + o1 + len) else {
panic!("second key within the column");
};
let second = second.to_vec();
let mut swapped = key_col.data.to_vec();
let Some(dst0) = swapped.get_mut(table_len + o0..table_len + o0 + len) else {
panic!("first key range within the column");
};
dst0.copy_from_slice(&second);
let Some(dst1) = swapped.get_mut(table_len + o1..table_len + o1 + len) else {
panic!("second key range within the column");
};
dst1.copy_from_slice(&first);
key_col.data = swapped.into();
}
let new_payload = batch.encode(CodecId::Plain)?;
let new_header = Header {
checksum: crate::Checksum::from_raw(crate::hash::hash128(&new_payload)),
..header
};
let mut new_block = Vec::with_capacity(header_len + new_payload.len());
new_header.encode_into(&mut new_block)?;
new_block.extend_from_slice(&new_payload);
let Some(target) = bytes.get_mut(block_off..block_off + new_block.len()) else {
panic!("block range within the file");
};
target.copy_from_slice(&new_block);
std::fs::write(&source, &bytes)?;
let table = open(source, &fs)?;
let err = reconcile_error(
&table,
crate::table::ReconcileGate::PointReadReachability,
None,
);
assert!(
matches!(
err,
crate::Error::InvalidHeader(
"columnar entries are out of order (a user key decreased, or an equal key's \
seqno did not strictly decrease) across the walk"
)
),
"the rejection names the columnar order violation, got {err:?}",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_drops_a_zero_row_columnar_block() -> crate::Result<()> {
use crate::coding::{Decode, Encode};
use crate::table::block::Header;
use crate::table::columnar::{
COL_SEQNO, COL_USER_KEY, COL_VALUE, COL_VALUE_TYPE, CodecId, Column, ColumnBatch, TypeTag,
};
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 9u32;
let build = |value_pad: usize| -> crate::Result<()> {
let _ = std::fs::remove_file(&source);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_columnar(true);
for i in 0..n {
writer.write(InternalValue::from_components(
format!("key{i:05}").into_bytes(),
format!("val{i:05}{}", "x".repeat(value_pad)).into_bytes(),
1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
Ok(())
};
build(0)?;
let payload_len = |src: &std::path::Path| -> crate::Result<usize> {
let bytes = std::fs::read(src)?;
let mut cursor = bytes.as_slice();
let header = Header::decode_from(&mut cursor)?;
Ok(header.data_length as usize)
};
let mut target_len = payload_len(&source)?;
if target_len % 2 != 0 {
build(1)?;
target_len = payload_len(&source)?;
}
assert_eq!(target_len % 2, 0, "an even payload length is reachable");
let empty_fixed = |id: u16, width: u8| Column {
column_id: id,
type_tag: TypeTag::Fixed(width),
validity: None,
data: Vec::new().into(),
};
let mut columns = vec![
Column {
column_id: COL_USER_KEY,
type_tag: TypeTag::Bytes,
validity: None,
data: vec![0u8; 4].into(),
},
empty_fixed(COL_SEQNO, 8),
empty_fixed(COL_VALUE_TYPE, 1),
empty_fixed(COL_VALUE, 1),
];
let Some(mut rem) = target_len.checked_sub(8 + 14 + 10 + 10 + 10) else {
panic!("source payload larger than the zero-row skeleton");
};
let mut next_id = COL_VALUE + 1;
while rem % 10 != 0 {
columns.push(Column {
column_id: next_id,
type_tag: TypeTag::Bytes,
validity: None,
data: vec![0u8; 4].into(),
});
next_id += 1;
let Some(next_rem) = rem.checked_sub(14) else {
panic!("remainder covers a Bytes column");
};
rem = next_rem;
}
while rem > 0 {
columns.push(empty_fixed(next_id, 1));
next_id += 1;
rem -= 10;
}
let batch = ColumnBatch {
row_count: 0,
columns,
};
let new_payload = batch.encode(CodecId::Plain)?;
assert_eq!(
new_payload.len(),
target_len,
"the zero-row batch pads to the original payload length",
);
let mut bytes = std::fs::read(&source)?;
let mut cursor = bytes.as_slice();
let header = Header::decode_from(&mut cursor)?;
let header_len = Header::header_len(header.block_type);
let new_header = Header {
checksum: crate::Checksum::from_raw(crate::hash::hash128(&new_payload)),
..header
};
let mut new_block: Vec<u8> = Vec::with_capacity(header_len + new_payload.len());
new_header.encode_into(&mut new_block)?;
new_block.extend_from_slice(&new_payload);
let Some(target) = bytes.get_mut(..new_block.len()) else {
panic!("block range within the file");
};
target.copy_from_slice(&new_block);
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(
report.dropped.len(),
1,
"the zero-row block is dropped as malformed: {report:?}",
);
assert_eq!(report.blocks_salvaged, 0, "{report:?}");
assert_eq!(
report.salvaged_path, None,
"an SST whose only block is empty reports nothing salvaged",
);
assert!(
fs.metadata(&dest).is_err(),
"no destination file is left behind",
);
Ok(())
}
#[test]
fn salvage_drops_a_row_block_with_out_of_order_keys() -> crate::Result<()> {
use crate::coding::Encode;
use crate::comparator::default_comparator;
use crate::table::block::Header;
use crate::table::block::decoder::ParsedItem as _;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_data_block_size(256)
.use_data_block_restart_interval(1)
.use_data_block_hash_ratio(0.0);
for i in 0..n {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let (block_off, poisoned_rows, new_block) = {
let table = open(source.clone(), &fs)?;
let Some(kh) = table.data_block_handles().filter_map(Result::ok).nth(1) else {
panic!("source must have at least two data blocks");
};
let block_off = usize::try_from(*kh.as_ref().offset()).unwrap_or(usize::MAX);
let sb = table.salvage_load_block(kh.as_ref(), crate::table::block::BlockType::Data)?;
let header = sb.block.header;
let db = crate::table::DataBlock::from_loaded(sb.block, false)?;
let iter = db.try_iter(default_comparator())?;
let mut entries: alloc::vec::Vec<crate::InternalValue> =
iter.map(|p| p.materialize(db.as_slice())).collect();
assert!(entries.len() >= 2, "block holds at least two rows");
entries.swap(0, 1);
let mut new_payload: alloc::vec::Vec<u8> = alloc::vec::Vec::new();
crate::table::DataBlock::encode_into(&mut new_payload, &entries, 1, 0.0)?;
assert_eq!(
new_payload.len(),
header.data_length as usize,
"an adjacent equal-length key swap re-encodes to the same length",
);
let header_len = Header::header_len(header.block_type);
let new_header = Header {
checksum: crate::Checksum::from_raw(crate::hash::hash128(&new_payload)),
..header
};
let mut new_block: alloc::vec::Vec<u8> =
alloc::vec::Vec::with_capacity(header_len + new_payload.len());
new_header.encode_into(&mut new_block)?;
new_block.extend_from_slice(&new_payload);
(block_off, entries.len() as u64, new_block)
};
let mut bytes = std::fs::read(&source)?;
let Some(target) = bytes.get_mut(block_off..block_off + new_block.len()) else {
panic!("block range within the file");
};
target.copy_from_slice(&new_block);
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(
report.dropped.len(),
1,
"exactly the out-of-order block is dropped: {report:?}",
);
assert!(
matches!(
report.dropped.first().map(|d| &d.reason),
Some(DropReason::DecodeError(_))
),
"the ordering violation classifies as a decode error: {report:?}",
);
assert_eq!(
report.entries_salvaged,
u64::from(n) - poisoned_rows,
"every row outside the poisoned block is recovered",
);
assert_eq!(reopen_item_count(dest, &fs)?, u64::from(n) - poisoned_rows);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_does_not_resurrect_deletes_in_an_index_omitted_block() -> crate::Result<()> {
use crate::config::DeleteStrategy;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
let deleted = [n - 1, n - 2, n - 3];
for pos in deleted {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
let hidden_rows = {
let table = open(source.clone(), &fs)?;
let Some(last) = table.data_block_handles().filter_map(Result::ok).last() else {
panic!("source must have data blocks");
};
let Some(batch) = table.load_columnar_block_masked(last.as_ref())? else {
panic!("the last block holds live rows");
};
u64::from(batch.row_count) + deleted.len() as u64
};
crate::test_forge::forge_zone_map_drop_last_entry(&source, 0)?;
crate::test_forge::forge_tli_mirrors_truncated(&source, 0, None)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(
report.dropped.len(),
1,
"the index-omitted block has no verified delete position and must \
drop, never re-emit live: {report:?}",
);
assert_eq!(
report.entries_salvaged,
u64::from(n) - deleted.len() as u64 - (hidden_rows - deleted.len() as u64),
"every indexed block's surviving rows are recovered: {report:?}",
);
for pos in deleted {
let key = format!("key{pos:05}");
assert!(
reopen_get(dest.clone(), &fs, key.as_bytes())?.is_none(),
"a positionally deleted row must not resurrect via the gap walk",
);
}
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_drops_an_out_of_order_columnar_block_in_a_delete_bearing_sst() -> crate::Result<()> {
use crate::coding::{Decode, Encode};
use crate::config::DeleteStrategy;
use crate::table::block::Header;
use crate::table::columnar::{CodecId, ColumnBatch};
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
let deletes = [5u32, 50, 150];
for pos in deletes {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
let (block_off, block_size) = {
let table = open(source.clone(), &fs)?;
let Some(kh) = table.data_block_handles().filter_map(Result::ok).nth(1) else {
panic!("source must have at least two data blocks");
};
(
usize::try_from(*kh.as_ref().offset()).unwrap_or(usize::MAX),
kh.as_ref().size() as usize,
)
};
let mut bytes = std::fs::read(&source)?;
let Some(block) = bytes.get(block_off..block_off + block_size) else {
panic!("block range within the file");
};
let mut cursor = block;
let header = Header::decode_from(&mut cursor)?;
let header_len = Header::header_len(header.block_type);
let Some(payload) = block.get(header_len..header_len + header.data_length as usize) else {
panic!("payload range within the block");
};
let mut batch = ColumnBatch::decode(&payload.into())?;
let poisoned_rows = u64::from(batch.row_count);
assert!(batch.row_count >= 2, "block holds at least two rows");
{
let Some(key_col) = batch.columns.first_mut() else {
panic!("key column present");
};
let table_len = (batch.row_count as usize + 1) * 4;
let off = |data: &[u8], idx: usize| -> usize {
let Some(b) = data.get(idx * 4..idx * 4 + 4) else {
panic!("offset {idx} within the frame table");
};
u32::from_le_bytes(b.try_into().unwrap_or([0; 4])) as usize
};
let (o0, o1, o2) = (
off(&key_col.data, 0),
off(&key_col.data, 1),
off(&key_col.data, 2),
);
assert_eq!(o1 - o0, o2 - o1, "adjacent keys are equal-length");
let len = o1 - o0;
let Some(first) = key_col.data.get(table_len + o0..table_len + o0 + len) else {
panic!("first key within the column");
};
let first = first.to_vec();
let Some(second) = key_col.data.get(table_len + o1..table_len + o1 + len) else {
panic!("second key within the column");
};
let second = second.to_vec();
let mut swapped = key_col.data.to_vec();
let Some(dst0) = swapped.get_mut(table_len + o0..table_len + o0 + len) else {
panic!("first key range within the column");
};
dst0.copy_from_slice(&second);
let Some(dst1) = swapped.get_mut(table_len + o1..table_len + o1 + len) else {
panic!("second key range within the column");
};
dst1.copy_from_slice(&first);
key_col.data = swapped.into();
}
let new_payload = batch.encode(CodecId::Plain)?;
assert_eq!(
new_payload.len(),
payload.len(),
"an in-place key swap re-encodes to the same length",
);
let new_header = Header {
checksum: crate::Checksum::from_raw(crate::hash::hash128(&new_payload)),
..header
};
let mut new_block = Vec::with_capacity(header_len + new_payload.len());
new_header.encode_into(&mut new_block)?;
new_block.extend_from_slice(&new_payload);
let Some(target) = bytes.get_mut(block_off..block_off + new_block.len()) else {
panic!("block range within the file");
};
target.copy_from_slice(&new_block);
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(
report.dropped.len(),
1,
"exactly the out-of-order block is dropped: {report:?}",
);
assert!(
matches!(
report.dropped.first().map(|d| &d.reason),
Some(DropReason::DecodeError(_))
),
"the ordering violation classifies as a decode error: {report:?}",
);
assert_eq!(
report.entries_salvaged,
u64::from(n) - poisoned_rows - deletes.len() as u64,
"rows outside the poisoned block are recovered with deletes applied",
);
for pos in deletes {
assert!(
reopen_get(dest.clone(), &fs, format!("key{pos:05}").as_bytes())?.is_none(),
"the deleted key at position {pos} stays masked in the recovered copy",
);
}
assert!(
reopen_get(dest, &fs, b"key00051")?.is_some(),
"a neighbouring live key reads back from the recovered copy",
);
Ok(())
}
#[test]
fn salvage_sst_errors_and_discards_the_dest_on_a_write_failure() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fault = FaultFs::new(StdFs);
let injector = fault.injector();
let fs: Arc<dyn Fs> = Arc::new(fault);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0..200u32 {
writer.write(InternalValue::from_components(
format!("key{i:05}").into_bytes(),
vec![0xAB; 1_024],
1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
injector.arm(
FaultRule::new(FaultOp::Write, Fault::Error(ErrorKind::Other))
.on_path("salvaged")
.once(),
);
let result = salvage_sst(&source, dest.clone(), &fs);
injector.clear();
assert!(
result.is_err(),
"a destination write failure errors the whole salvage: {result:?}",
);
assert!(
fs.metadata(&dest).is_err(),
"the partial destination is removed on a write failure",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_refuses_a_delete_bitmap_relabeled_to_a_full_filter() -> crate::Result<()> {
use crate::config::{BloomConstructionPolicy, DeleteStrategy};
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256)
.use_bloom_policy(BloomConstructionPolicy::BitsPerKey(0.0))
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 50, 150] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
crate::test_forge::forge_duplicate_section_name(
&source,
b"delete_bitmap",
b"filter",
crate::table::block::BlockType::Filter,
)?;
let Err(err) = salvage_sst(&source, dest.clone(), &fs) else {
panic!("a delete_bitmap relabeled to a filter must fail salvage");
};
let crate::Error::FeatureUnsupported(reason) = &err else {
panic!("the refusal must be FeatureUnsupported, got {err:?}");
};
assert!(
reason.contains("rebuildable"),
"the refusal must name the degraded rebuildable section, got {reason:?}",
);
assert!(
fs.metadata(&dest).is_err(),
"no destination file is left behind by the refused salvage",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_refuses_a_delete_bitmap_renamed_to_a_filter_without_reroling() -> crate::Result<()> {
use crate::config::{BloomConstructionPolicy, DeleteStrategy};
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256)
.use_bloom_policy(BloomConstructionPolicy::BitsPerKey(0.0))
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 50, 150] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
crate::test_forge::forge_duplicate_section_name(
&source,
b"delete_bitmap",
b"filter",
crate::table::block::BlockType::DeleteBitmap,
)?;
let Err(err) = salvage_sst(&source, dest.clone(), &fs) else {
panic!("a delete_bitmap renamed to a filter without re-roling must fail salvage");
};
let crate::Error::FeatureUnsupported(reason) = &err else {
panic!("the refusal must be FeatureUnsupported, got {err:?}");
};
assert!(
reason.contains("rebuildable"),
"the refusal must name the degraded rebuildable section, got {reason:?}",
);
assert!(
fs.metadata(&dest).is_err(),
"no destination file is left behind by the refused salvage",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_fails_closed_on_a_corrupt_delete_bitmap_by_default() -> crate::Result<()> {
use crate::config::DeleteStrategy;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 50, 150] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
let (db_pos, db_len) = {
let mut f = std::fs::File::open(&source)?;
let reader = match crate::sfa::Reader::from_reader(&mut f) {
Ok(r) => r,
Err(e) => panic!("reading the SFA trailer failed: {e:?}"),
};
let Some(entry) = reader.toc().iter().find(|e| e.name() == b"delete_bitmap") else {
panic!("source must carry a delete_bitmap section");
};
(entry.pos(), entry.len())
};
let flip = usize::try_from(db_pos + db_len / 2).unwrap_or(0);
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(flip) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let result = salvage_sst(&source, dest.clone(), &fs);
assert!(
result.is_err(),
"default salvage refuses to resurrect deleted rows: {result:?}",
);
assert!(
fs.metadata(&dest).is_err(),
"no destination file is left behind by the refused salvage",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_fails_closed_on_an_unpositionable_delete_bitmap_by_default() -> crate::Result<()> {
use crate::config::DeleteStrategy;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 50, 150] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
let (zm_pos, zm_len) = {
let mut f = std::fs::File::open(&source)?;
let reader = match crate::sfa::Reader::from_reader(&mut f) {
Ok(r) => r,
Err(e) => panic!("reading the SFA trailer failed: {e:?}"),
};
let Some(entry) = reader.toc().iter().find(|e| e.name() == b"zone_map") else {
panic!("source must carry a zone_map section");
};
(entry.pos(), entry.len())
};
let flip = usize::try_from(zm_pos + zm_len / 2).unwrap_or(0);
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(flip) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let result = salvage_sst(&source, dest.clone(), &fs);
assert!(
result.is_err(),
"default salvage refuses to resurrect deleted rows: {result:?}",
);
assert!(
fs.metadata(&dest).is_err(),
"no destination file is left behind by the refused salvage",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_tolerates_a_corrupt_delete_bitmap_as_all_live() -> crate::Result<()> {
use crate::config::DeleteStrategy;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 50, 150] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
let (db_pos, db_len) = {
let mut f = std::fs::File::open(&source)?;
let reader = match crate::sfa::Reader::from_reader(&mut f) {
Ok(r) => r,
Err(e) => panic!("reading the SFA trailer failed: {e:?}"),
};
let Some(entry) = reader.toc().iter().find(|e| e.name() == b"delete_bitmap") else {
panic!("source must carry a delete_bitmap section");
};
(entry.pos(), entry.len())
};
let flip = usize::try_from(db_pos + db_len / 2).unwrap_or(0);
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(flip) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
assert!(
open(source.clone(), &fs).is_err(),
"normal recovery must fail closed on a corrupt delete-bitmap",
);
let options = SalvageOptions {
allow_delete_resurrection: true,
..SalvageOptions::default()
};
let report = salvage_sst_with_options(&source, dest.clone(), &fs, &options)?;
assert!(
report.is_complete(),
"the data blocks are intact; only the sidecar was corrupt: {report:?}",
);
assert_eq!(
report.entries_salvaged,
u64::from(n),
"every row is recovered live, the corrupt bitmap is ignored",
);
assert_eq!(
report.blocks_copied_verbatim, 0,
"a delete-bearing SST is never copied verbatim, even with a degraded bitmap: {report:?}",
);
assert_eq!(reopen_item_count(dest, &fs)?, u64::from(n));
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_tolerates_a_persistently_unreadable_delete_bitmap_as_all_live() -> crate::Result<()> {
use crate::config::DeleteStrategy;
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let clean: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&clean))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 50, 150] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
let db_pos = {
let mut f = std::fs::File::open(&source)?;
let reader = match crate::sfa::Reader::from_reader(&mut f) {
Ok(r) => r,
Err(e) => panic!("reading the SFA trailer failed: {e:?}"),
};
let Some(entry) = reader.toc().iter().find(|e| e.name() == b"delete_bitmap") else {
panic!("source must carry a delete_bitmap section");
};
entry.pos()
};
let fault = FaultFs::new(StdFs);
fault
.injector()
.arm(FaultRule::new(FaultOp::ReadAt, Fault::Error(ErrorKind::Other)).at_offset(db_pos));
let fs: Arc<dyn Fs> = Arc::new(fault);
assert!(
salvage_sst(&source, dest.clone(), &fs).is_err(),
"default salvage must fail closed on a persistently unreadable delete-bitmap",
);
let options = SalvageOptions {
allow_delete_resurrection: true,
..SalvageOptions::default()
};
let report = salvage_sst_with_options(&source, dest.clone(), &fs, &options)?;
assert_eq!(
report.entries_salvaged,
u64::from(n),
"every row is recovered live once the unreadable bitmap degrades: {report:?}",
);
assert_eq!(reopen_item_count(dest, &fs)?, u64::from(n));
Ok(())
}
#[test]
fn recover_degrades_a_persistently_unreadable_zone_map_section() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("source");
let clean: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&clean))?.use_zone_map(true);
for i in 0..64u32 {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let zm_pos = section_pos(&source, b"zone_map");
let fault = FaultFs::new(StdFs);
fault
.injector()
.arm(FaultRule::new(FaultOp::ReadAt, Fault::Error(ErrorKind::Other)).at_offset(zm_pos));
let fs: Arc<dyn Fs> = Arc::new(fault);
let table = open(source, &fs)?;
assert!(
table.zone_map.is_empty(),
"the persistently unreadable zone map degrades to an empty map",
);
Ok(())
}
#[test]
fn recover_degrades_a_persistently_unreadable_seqno_bounds_section() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("source");
let clean: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer =
Writer::new(source.clone(), 0, 0, Arc::clone(&clean))?.use_seqno_in_index(true);
for i in 0..64u32 {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let sb_pos = section_pos(&source, b"seqno_bounds");
let fault = FaultFs::new(StdFs);
fault
.injector()
.arm(FaultRule::new(FaultOp::ReadAt, Fault::Error(ErrorKind::Other)).at_offset(sb_pos));
let fs: Arc<dyn Fs> = Arc::new(fault);
let table = open(source, &fs)?;
assert!(
table.seqno_bounds.is_empty(),
"the persistently unreadable seqno-bounds section degrades to an empty map",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_fails_closed_on_an_unreadable_block_in_a_delete_bearing_sst() -> crate::Result<()> {
use crate::config::DeleteStrategy;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 50, 150] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
let target = {
let table = open(source.clone(), &fs)?;
let offsets: alloc::vec::Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|kh| *kh.as_ref().offset())
.collect();
let Some(&second) = offsets.get(1) else {
panic!("source must have at least two data blocks, got {offsets:?}");
};
second
};
let flip = usize::try_from(target).unwrap_or(0) + 16;
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(flip) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let result = salvage_sst(&source, dest.clone(), &fs);
assert!(
result.is_err(),
"unverifiable delete positions fail the default salvage: {result:?}",
);
assert!(
fs.metadata(&dest).is_err(),
"no destination file is left behind by the refused salvage",
);
let options = SalvageOptions {
allow_delete_resurrection: true,
..SalvageOptions::default()
};
let report = salvage_sst_with_options(&source, dest.clone(), &fs, &options)?;
assert_eq!(
report.dropped.len(),
1,
"exactly the corrupt block is dropped: {report:?}",
);
assert!(
report.entries_salvaged < u64::from(n),
"the dropped block's rows are lost: {report:?}",
);
assert_eq!(
reopen_item_count(dest.clone(), &fs)?,
report.entries_salvaged,
"the recovered copy reopens with every salvaged row live",
);
for pos in [5u32, 50, 150] {
assert!(
reopen_get(dest.clone(), &fs, format!("key{pos:05}").as_bytes())?.is_some(),
"the opt-in resurrects the deleted key at position {pos}",
);
}
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_fails_closed_on_a_zero_row_block_in_a_delete_bearing_sst() -> crate::Result<()> {
use crate::coding::{Decode, Encode};
use crate::config::DeleteStrategy;
use crate::table::block::Header;
use crate::table::columnar::{
COL_SEQNO, COL_USER_KEY, COL_VALUE, COL_VALUE_TYPE, CodecId, Column, ColumnBatch, TypeTag,
};
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 50, 150] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
let (block_off, block_size) = {
let table = open(source.clone(), &fs)?;
let Some(kh) = table.data_block_handles().filter_map(Result::ok).nth(1) else {
panic!("source must have at least two data blocks");
};
(
usize::try_from(*kh.as_ref().offset()).unwrap_or(usize::MAX),
kh.as_ref().size() as usize,
)
};
let mut bytes = std::fs::read(&source)?;
let Some(block) = bytes.get(block_off..block_off + block_size) else {
panic!("block range within the file");
};
let mut cursor = block;
let header = Header::decode_from(&mut cursor)?;
let header_len = Header::header_len(header.block_type);
let target_len = header.data_length as usize;
assert_eq!(
target_len % 2,
0,
"the padded skeleton needs an even length"
);
let empty_fixed = |id: u16, width: u8| Column {
column_id: id,
type_tag: TypeTag::Fixed(width),
validity: None,
data: Vec::new().into(),
};
let mut columns = vec![
Column {
column_id: COL_USER_KEY,
type_tag: TypeTag::Bytes,
validity: None,
data: vec![0u8; 4].into(),
},
empty_fixed(COL_SEQNO, 8),
empty_fixed(COL_VALUE_TYPE, 1),
empty_fixed(COL_VALUE, 1),
];
let Some(mut rem) = target_len.checked_sub(8 + 14 + 10 + 10 + 10) else {
panic!("source payload larger than the zero-row skeleton");
};
let mut next_id = COL_VALUE + 1;
while rem % 10 != 0 {
columns.push(Column {
column_id: next_id,
type_tag: TypeTag::Bytes,
validity: None,
data: vec![0u8; 4].into(),
});
next_id += 1;
let Some(next_rem) = rem.checked_sub(14) else {
panic!("remainder covers a Bytes column");
};
rem = next_rem;
}
while rem > 0 {
columns.push(empty_fixed(next_id, 1));
next_id += 1;
rem -= 10;
}
let new_payload = ColumnBatch {
row_count: 0,
columns,
}
.encode(CodecId::Plain)?;
assert_eq!(new_payload.len(), target_len, "length-preserving forgery");
let new_header = Header {
checksum: crate::Checksum::from_raw(crate::hash::hash128(&new_payload)),
..header
};
let mut new_block: Vec<u8> = Vec::with_capacity(header_len + new_payload.len());
new_header.encode_into(&mut new_block)?;
new_block.extend_from_slice(&new_payload);
let Some(target) = bytes.get_mut(block_off..block_off + new_block.len()) else {
panic!("block range within the file");
};
target.copy_from_slice(&new_block);
let zm_pos = {
let mut f = std::io::Cursor::new(&bytes);
let reader = match crate::sfa::Reader::from_reader(&mut f) {
Ok(r) => r,
Err(e) => panic!("reading the SFA trailer failed: {e:?}"),
};
let Some(entry) = reader.toc().iter().find(|e| e.name() == b"zone_map") else {
panic!("source must carry a zone_map section");
};
usize::try_from(entry.pos()).unwrap_or(usize::MAX)
};
let Some(mut zm_cursor) = bytes.get(zm_pos..) else {
panic!("zone_map section within the file");
};
let zm_header = Header::decode_from(&mut zm_cursor)?;
let zm_header_len = Header::header_len(zm_header.block_type);
let zm_payload_range =
zm_pos + zm_header_len..zm_pos + zm_header_len + zm_header.data_length as usize;
{
let Some(payload) = bytes.get_mut(zm_payload_range.clone()) else {
panic!("zone_map payload within the file");
};
let read_u32 = |data: &[u8], at: usize| -> u32 {
let Some(b) = data.get(at..at + 4) else {
panic!("u32 at {at} within the zone map payload");
};
u32::from_le_bytes(b.try_into().unwrap_or([0; 4]))
};
let read_u16 = |data: &[u8], at: usize| -> u16 {
let Some(b) = data.get(at..at + 2) else {
panic!("u16 at {at} within the zone map payload");
};
u16::from_le_bytes(b.try_into().unwrap_or([0; 2]))
};
let mut at = 4; at += 8;
let block1_cols = read_u16(payload, at);
at += 2;
for _ in 0..block1_cols {
at += 10; at += 4; let min_len = read_u32(payload, at) as usize;
at += 4 + min_len;
let max_len = read_u32(payload, at) as usize;
at += 4 + max_len;
}
at += 8; at += 2; at += 10; let claimed = read_u32(payload, at);
assert!(claimed > 0, "the second block originally holds rows");
let Some(rc) = payload.get_mut(at..at + 4) else {
panic!("second block's first row_count within the zone map payload");
};
rc.copy_from_slice(&0u32.to_le_bytes());
}
let new_zm_checksum = crate::Checksum::from_raw(crate::hash::hash128(
bytes.get(zm_payload_range).unwrap_or(&[]),
));
let new_zm_header = Header {
checksum: new_zm_checksum,
..zm_header
};
let mut zm_hdr_bytes: Vec<u8> = Vec::with_capacity(zm_header_len);
new_zm_header.encode_into(&mut zm_hdr_bytes)?;
let Some(zm_dst) = bytes.get_mut(zm_pos..zm_pos + zm_header_len) else {
panic!("zone_map header within the file");
};
zm_dst.copy_from_slice(&zm_hdr_bytes);
std::fs::write(&source, &bytes)?;
let result = salvage_sst(&source, dest.clone(), &fs);
assert!(
result.is_err(),
"a zero-row block in a delete-bearing SST fails the default salvage: {result:?}",
);
assert!(
fs.metadata(&dest).is_err(),
"no destination file is left behind by the refused salvage",
);
Ok(())
}
#[test]
fn salvage_drops_a_row_block_with_a_stale_kv_digest() -> crate::Result<()> {
use crate::coding::{Decode, Encode};
use crate::runtime_config::{ChecksumAlgorithm, KvChecksumPolicy};
use crate::table::block::Header;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_data_block_size(256)
.use_kv_checksums(KvChecksumPolicy::AllLevels, ChecksumAlgorithm::Xxh3_64);
for i in 0..n {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let (block_off, block_size) = {
let table = open(source.clone(), &fs)?;
let Some(kh) = table.data_block_handles().filter_map(Result::ok).nth(1) else {
panic!("source must have at least two data blocks");
};
(
usize::try_from(*kh.as_ref().offset()).unwrap_or(usize::MAX),
kh.as_ref().size() as usize,
)
};
let mut bytes = std::fs::read(&source)?;
let Some(block) = bytes.get(block_off..block_off + block_size) else {
panic!("block range within the file");
};
let mut cursor = block;
let header = Header::decode_from(&mut cursor)?;
let header_len = Header::header_len(header.block_type);
let payload_range =
block_off + header_len..block_off + header_len + header.data_length as usize;
{
let Some(payload) = bytes.get_mut(payload_range.clone()) else {
panic!("payload range within the block");
};
let Some(b) = payload.get_mut(12) else {
panic!("entry byte within the payload");
};
*b ^= 0xFF;
}
let new_checksum = crate::Checksum::from_raw(crate::hash::hash128(
bytes.get(payload_range).unwrap_or(&[]),
));
let new_header = Header {
checksum: new_checksum,
..header
};
let mut hdr_bytes: Vec<u8> = Vec::with_capacity(header_len);
new_header.encode_into(&mut hdr_bytes)?;
let Some(hdr_dst) = bytes.get_mut(block_off..block_off + header_len) else {
panic!("block header within the file");
};
hdr_dst.copy_from_slice(&hdr_bytes);
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(
report.dropped.len(),
1,
"the block with a stale per-KV digest is dropped: {report:?}",
);
assert!(
report.entries_salvaged < u64::from(n),
"the tampered block's rows are not laundered into the copy: {report:?}",
);
assert_eq!(reopen_item_count(dest, &fs)?, report.entries_salvaged);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_fails_closed_on_an_undecodable_checksum_clean_block_with_deletes() -> crate::Result<()> {
use crate::coding::{Decode, Encode};
use crate::config::DeleteStrategy;
use crate::table::block::Header;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 50, 150] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
let (block_off, block_size) = {
let table = open(source.clone(), &fs)?;
let Some(kh) = table.data_block_handles().filter_map(Result::ok).nth(1) else {
panic!("source must have at least two data blocks");
};
(
usize::try_from(*kh.as_ref().offset()).unwrap_or(usize::MAX),
kh.as_ref().size() as usize,
)
};
let mut bytes = std::fs::read(&source)?;
let Some(block) = bytes.get(block_off..block_off + block_size) else {
panic!("block range within the file");
};
let mut cursor = block;
let header = Header::decode_from(&mut cursor)?;
let header_len = Header::header_len(header.block_type);
let payload_range =
block_off + header_len..block_off + header_len + header.data_length as usize;
{
let Some(payload) = bytes.get_mut(payload_range.clone()) else {
panic!("payload range within the block");
};
let Some(tag) = payload.get_mut(10) else {
panic!("first column's type byte within the payload");
};
*tag = 0xEE;
}
let new_checksum = crate::Checksum::from_raw(crate::hash::hash128(
bytes.get(payload_range).unwrap_or(&[]),
));
let new_header = Header {
checksum: new_checksum,
..header
};
let mut hdr_bytes = Vec::with_capacity(header_len);
new_header.encode_into(&mut hdr_bytes)?;
let Some(hdr_dst) = bytes.get_mut(block_off..block_off + header_len) else {
panic!("block header within the file");
};
hdr_dst.copy_from_slice(&hdr_bytes);
std::fs::write(&source, &bytes)?;
let result = salvage_sst(&source, dest.clone(), &fs);
assert!(
result.is_err(),
"unverifiable delete positions fail the default salvage: {result:?}",
);
assert!(
fs.metadata(&dest).is_err(),
"no destination file is left behind by the refused salvage",
);
let options = SalvageOptions {
allow_delete_resurrection: true,
..SalvageOptions::default()
};
let report = salvage_sst_with_options(&source, dest.clone(), &fs, &options)?;
assert_eq!(
report.dropped.len(),
1,
"exactly the poisoned block is dropped: {report:?}",
);
assert!(
report.entries_salvaged < u64::from(n),
"the poisoned block's rows are lost: {report:?}",
);
assert_eq!(
reopen_item_count(dest.clone(), &fs)?,
report.entries_salvaged,
"the recovered copy reopens with every salvaged row live",
);
for pos in [5u32, 50, 150] {
assert!(
reopen_get(dest.clone(), &fs, format!("key{pos:05}").as_bytes())?.is_some(),
"the opt-in resurrects the deleted key at position {pos}",
);
}
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_fails_closed_on_a_zone_map_with_wrong_row_counts() -> crate::Result<()> {
use crate::coding::{Decode, Encode};
use crate::config::DeleteStrategy;
use crate::table::block::Header;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 50, 150] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
let zm_pos = {
let mut f = std::fs::File::open(&source)?;
let reader = match crate::sfa::Reader::from_reader(&mut f) {
Ok(r) => r,
Err(e) => panic!("reading the SFA trailer failed: {e:?}"),
};
let Some(entry) = reader.toc().iter().find(|e| e.name() == b"zone_map") else {
panic!("source must carry a zone_map section");
};
usize::try_from(entry.pos()).unwrap_or(usize::MAX)
};
let mut bytes = std::fs::read(&source)?;
let Some(mut cursor) = bytes.get(zm_pos..) else {
panic!("zone_map section within the file");
};
let header = Header::decode_from(&mut cursor)?;
let header_len = Header::header_len(header.block_type);
let payload_range = zm_pos + header_len..zm_pos + header_len + header.data_length as usize;
{
let Some(payload) = bytes.get_mut(payload_range.clone()) else {
panic!("zone_map payload within the file");
};
let Some(rc) = payload.get_mut(24..28) else {
panic!("first row_count within the zone map payload");
};
let claimed = u32::from_le_bytes(rc.try_into().unwrap_or([0; 4]));
assert!(claimed >= 2, "the first block holds at least two rows");
rc.copy_from_slice(&(claimed - 1).to_le_bytes());
}
let new_checksum = crate::Checksum::from_raw(crate::hash::hash128(
bytes.get(payload_range).unwrap_or(&[]),
));
let new_header = Header {
checksum: new_checksum,
..header
};
let mut hdr_bytes = Vec::with_capacity(header_len);
new_header.encode_into(&mut hdr_bytes)?;
let Some(hdr_dst) = bytes.get_mut(zm_pos..zm_pos + header_len) else {
panic!("zone_map header within the file");
};
hdr_dst.copy_from_slice(&hdr_bytes);
std::fs::write(&source, &bytes)?;
let result = salvage_sst(&source, dest.clone(), &fs);
assert!(
result.is_err(),
"a mispositioning zone map fails the default salvage: {result:?}",
);
assert!(
fs.metadata(&dest).is_err(),
"no destination file is left behind by the refused salvage",
);
let options = SalvageOptions {
allow_delete_resurrection: true,
..SalvageOptions::default()
};
let report = salvage_sst_with_options(&source, dest.clone(), &fs, &options)?;
assert!(
report.is_complete(),
"the data blocks are intact; only the zone map lies: {report:?}",
);
assert_eq!(
report.entries_salvaged,
u64::from(n),
"every row is recovered live under the opt-in",
);
assert_eq!(reopen_item_count(dest.clone(), &fs)?, u64::from(n));
for pos in [5u32, 50, 150] {
assert!(
reopen_get(dest.clone(), &fs, format!("key{pos:05}").as_bytes())?.is_some(),
"the opt-in resurrects the deleted key at position {pos}",
);
}
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_ignores_a_delete_bitmap_without_a_readable_zone_map() -> crate::Result<()> {
use crate::config::DeleteStrategy;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 50, 150] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
let (zm_pos, zm_len) = {
let mut f = std::fs::File::open(&source)?;
let reader = match crate::sfa::Reader::from_reader(&mut f) {
Ok(r) => r,
Err(e) => panic!("reading the SFA trailer failed: {e:?}"),
};
let Some(entry) = reader.toc().iter().find(|e| e.name() == b"zone_map") else {
panic!("source must carry a zone_map section");
};
(entry.pos(), entry.len())
};
let flip = usize::try_from(zm_pos + zm_len / 2).unwrap_or(0);
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(flip) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
assert!(
open(source.clone(), &fs).is_err(),
"normal recovery must reject a bitmap with no readable zone map",
);
let options = SalvageOptions {
allow_delete_resurrection: true,
..SalvageOptions::default()
};
let report = salvage_sst_with_options(&source, dest.clone(), &fs, &options)?;
assert!(
report.is_complete(),
"the data blocks are intact; only the zone map was corrupt: {report:?}",
);
assert_eq!(
report.entries_salvaged,
u64::from(n),
"every row is recovered live once the unpositionable bitmap is ignored",
);
assert_eq!(reopen_item_count(dest, &fs)?, u64::from(n));
Ok(())
}
#[test]
fn salvage_sst_errors_when_the_source_cannot_be_opened() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0..50 {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let mut bytes = std::fs::read(&source)?;
bytes.truncate(bytes.len() / 2);
std::fs::write(&source, &bytes)?;
assert!(
salvage_sst(&source, dest.clone(), &fs).is_err(),
"an unparseable container must fail salvage, not write a partial file",
);
assert!(
!dest.exists(),
"no destination is written on an open failure"
);
Ok(())
}
#[test]
fn salvage_sst_recovers_nothing_when_the_only_block_is_corrupt() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0..8 {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let target = {
let table = open(source.clone(), &fs)?;
let offsets: alloc::vec::Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|kh| *kh.as_ref().offset())
.collect();
let Some(&only) = offsets.first() else {
panic!("expected a single data block, got {offsets:?}");
};
assert_eq!(
offsets.len(),
1,
"expected a single data block, got {offsets:?}"
);
only
};
let flip = usize::try_from(target).unwrap_or(0) + 16;
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(flip) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(report.blocks_salvaged, 0, "the only block was corrupt");
assert_eq!(report.entries_salvaged, 0, "no entries recovered");
assert_eq!(report.dropped.len(), 1, "the dropped block is reported");
assert!(
report.salvaged_path.is_none(),
"nothing recoverable means no file is written",
);
assert!(!dest.exists(), "no destination file on an empty salvage");
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_skips_a_wholly_deleted_block() -> crate::Result<()> {
use crate::config::DeleteStrategy;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let n = 200u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_data_block_size(256)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
let deleted = 60u32;
for pos in 0..deleted {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(
report.is_complete(),
"wholly-deleted blocks are skipped, not dropped: {report:?}",
);
assert!(
report.blocks_salvaged < report.blocks_total,
"at least one leading block was wholly deleted and skipped: {report:?}",
);
assert_eq!(
report.entries_salvaged,
u64::from(n - deleted),
"every live row is recovered, the deleted prefix is skipped",
);
assert_eq!(reopen_item_count(dest, &fs)?, u64::from(n - deleted));
Ok(())
}
#[test]
fn salvage_rejects_an_sst_with_range_tombstones() -> crate::Result<()> {
use crate::UserKey;
use crate::range_tombstone::RangeTombstone;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0..20 {
writer.write(iv(i))?;
}
writer.write_range_tombstone(RangeTombstone::new(
UserKey::from(b"key00005".as_slice()),
UserKey::from(b"key00010".as_slice()),
2,
));
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let result = salvage_sst(&source, dest.clone(), &fs);
assert!(
matches!(result, Err(crate::Error::FeatureUnsupported(_))),
"an SST with range tombstones must fail closed, got {result:?}",
);
assert!(
!dest.exists(),
"no salvaged file is written when salvage fails closed",
);
Ok(())
}
#[test]
fn salvage_sst_reads_and_writes_through_the_injected_fs() -> crate::Result<()> {
use crate::fs::MemFs;
let fs: Arc<dyn Fs> = Arc::new(MemFs::new());
let dir = std::path::absolute("/memfs")?;
fs.create_dir_all(&dir)?;
let source = dir.join("source");
let dest = dir.join("salvaged");
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(256);
let n = 200u32;
for i in 0..n {
writer.write(iv(i))?;
}
assert!(
writer.finish()?.is_some(),
"in-memory source SST is non-empty"
);
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(
report.is_complete(),
"a healthy in-memory SST salvages with no dropped blocks: {report:?}",
);
assert_eq!(
report.entries_salvaged,
u64::from(n),
"every entry is recovered through the in-memory backend",
);
assert_eq!(report.salvaged_path.as_deref(), Some(dest.as_path()));
assert_eq!(
reopen_item_count(dest, &fs)?,
u64::from(n),
"the salvaged SST reopens through the same in-memory backend",
);
Ok(())
}
#[cfg(any(feature = "encryption", zstd_any))]
fn corrupt_second_data_block(
source: &std::path::Path,
fs: &Arc<dyn Fs>,
table_id: crate::table::TableId,
encryption: Option<Arc<dyn crate::encryption::EncryptionProvider>>,
#[cfg(zstd_any)] zstd_dictionary: Option<Arc<crate::compression::ZstdDictionary>>,
) -> crate::Result<()> {
let checksum = crate::Checksum::from_raw(crate::repair::compute_table_checksum(&**fs, source)?);
let table = {
let mut params = crate::table::RecoverParams::new(
source.to_path_buf(),
checksum,
table_id,
Arc::clone(fs),
default_comparator(),
Arc::new(crate::cache::Cache::with_capacity_bytes(1 << 20)),
);
params.descriptor_table = Some(Arc::new(crate::descriptor_table::DescriptorTable::new(8)));
params.encryption = encryption;
#[cfg(zstd_any)]
{
params.zstd_dictionary = zstd_dictionary;
}
Table::recover(params)?
};
let offsets: alloc::vec::Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|kh| *kh.as_ref().offset())
.collect();
let Some(&second) = offsets.get(1) else {
panic!("source SST must have at least two data blocks, got {offsets:?}");
};
let flip = usize::try_from(second).unwrap_or(0) + 16;
let mut bytes = std::fs::read(source)?;
if let Some(b) = bytes.get_mut(flip) {
*b ^= 0xFF;
}
std::fs::write(source, &bytes)?;
Ok(())
}
#[cfg(feature = "encryption")]
#[test]
fn salvage_recovers_an_encrypted_sst_with_the_provider() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let enc: Arc<dyn crate::encryption::EncryptionProvider> =
Arc::new(crate::encryption::Aes256GcmProvider::new(&[0x42; 32]));
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_data_block_size(256)
.use_encryption(Some(Arc::clone(&enc)));
let n = 200u32;
for i in 0..n {
writer.write(iv(i))?;
}
assert!(
writer.finish()?.is_some(),
"source encrypted SST is non-empty"
);
corrupt_second_data_block(
&source,
&fs,
0,
Some(Arc::clone(&enc)),
#[cfg(zstd_any)]
None,
)?;
assert!(
salvage_sst(&source, dest.clone(), &fs).is_err(),
"an encrypted SST must not salvage without the provider",
);
let options = SalvageOptions {
encryption: Some(Arc::clone(&enc)),
#[cfg(zstd_any)]
zstd_dictionary: None,
table_id: 0,
expected_stored_id: None,
output_id: None,
allow_delete_resurrection: false,
sync_mode: crate::fs::SyncMode::Normal,
prefix_extractor: None,
blob_rewrite: None,
progress: None,
};
let report = salvage_sst_with_options(&source, dest.clone(), &fs, &options)?;
assert_eq!(
report.dropped.len(),
1,
"exactly the corrupt block drops: {report:?}"
);
assert!(
report.entries_salvaged > 0 && report.entries_salvaged < u64::from(n),
"a partial key range is recovered, got {} of {n}",
report.entries_salvaged,
);
let checksum = crate::Checksum::from_raw(crate::repair::compute_table_checksum(&*fs, &dest)?);
let reopened = {
let mut params = crate::table::RecoverParams::new(
dest,
checksum,
0,
Arc::clone(&fs),
default_comparator(),
Arc::new(crate::cache::Cache::with_capacity_bytes(1 << 20)),
);
params.descriptor_table = Some(Arc::new(crate::descriptor_table::DescriptorTable::new(8)));
params.encryption = Some(Arc::clone(&enc));
Table::recover(params)?
};
assert_eq!(
reopened.metadata.item_count, report.entries_salvaged,
"the encrypted salvaged copy reopens with exactly the recovered entries",
);
Ok(())
}
#[cfg(zstd_any)]
#[test]
fn salvage_recovers_a_dictionary_sst_with_the_dictionary() -> crate::Result<()> {
use crate::CompressionType;
use crate::compression::ZstdDictionary;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let samples: alloc::vec::Vec<u8> = (0..4000u32).map(|i| (i % 251) as u8).collect();
let dict = Arc::new(ZstdDictionary::new(&samples));
let compression = CompressionType::ZstdDict {
level: 3,
dict_id: dict.id(),
};
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_data_block_size(256)
.use_data_block_compression(compression)
.use_zstd_dictionary(Some(Arc::clone(&dict)));
let n = 200u32;
for i in 0..n {
writer.write(iv(i))?;
}
assert!(
writer.finish()?.is_some(),
"source dictionary SST is non-empty"
);
corrupt_second_data_block(&source, &fs, 0, None, Some(Arc::clone(&dict)))?;
assert!(
salvage_sst(&source, dest.clone(), &fs).is_err(),
"a dictionary SST must not salvage without the dictionary",
);
let options = SalvageOptions {
encryption: None,
zstd_dictionary: Some(Arc::clone(&dict)),
table_id: 0,
expected_stored_id: None,
output_id: None,
allow_delete_resurrection: false,
sync_mode: crate::fs::SyncMode::Normal,
prefix_extractor: None,
blob_rewrite: None,
progress: None,
};
let report = salvage_sst_with_options(&source, dest.clone(), &fs, &options)?;
assert_eq!(
report.dropped.len(),
1,
"exactly the corrupt block drops: {report:?}"
);
assert!(
report.entries_salvaged > 0 && report.entries_salvaged < u64::from(n),
"a partial key range is recovered, got {} of {n}",
report.entries_salvaged,
);
let checksum = crate::Checksum::from_raw(crate::repair::compute_table_checksum(&*fs, &dest)?);
let reopened = {
let mut params = crate::table::RecoverParams::new(
dest,
checksum,
0,
Arc::clone(&fs),
default_comparator(),
Arc::new(crate::cache::Cache::with_capacity_bytes(1 << 20)),
);
params.descriptor_table = Some(Arc::new(crate::descriptor_table::DescriptorTable::new(8)));
params.zstd_dictionary = Some(Arc::clone(&dict));
Table::recover(params)?
};
assert_eq!(
reopened.metadata.item_count, report.entries_salvaged,
"the dictionary salvaged copy reopens with exactly the recovered entries",
);
Ok(())
}
#[cfg(zstd_any)]
#[test]
fn a_blob_salvage_under_the_wrong_dictionary_fails_closed() -> crate::Result<()> {
use crate::CompressionType;
use crate::compression::ZstdDictionary;
let dir = tempdir()?;
let source = dir.path().join("dict_blob");
let dest = dir.path().join("dict_blob_salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let samples_a: alloc::vec::Vec<u8> = (0..4000u32).map(|i| (i % 251) as u8).collect();
let dict_a = Arc::new(ZstdDictionary::new(&samples_a));
let compression = CompressionType::ZstdDict {
level: 3,
dict_id: dict_a.id(),
};
let mut writer = BlobWriter::new(&source, 0, 0, &*fs)?
.use_compression(compression)
.use_zstd_dictionary(Some(Arc::clone(&dict_a)));
for i in 0..8u32 {
let key = format!("key{i:04}");
let value: alloc::vec::Vec<u8> = (0..512u32).map(|j| ((i + j) % 251) as u8).collect();
writer.write(key.as_bytes(), 0, &value)?;
}
writer.finish()?;
let samples_b: alloc::vec::Vec<u8> = (0..4000u32).map(|i| (i % 13) as u8).collect();
let dict_b = Arc::new(ZstdDictionary::new(&samples_b));
assert_ne!(dict_a.id(), dict_b.id(), "the fixture needs distinct ids");
let Err(err) = salvage_blob_file(
&source,
dest,
&fs,
0,
&default_comparator(),
0,
Some(&dict_b),
) else {
panic!("a mismatched dictionary must fail the salvage, not empty it");
};
assert!(
matches!(
&err,
crate::Error::ZstdDictMismatch { expected, got }
if *expected == dict_a.id() && *got == Some(dict_b.id())
),
"the failure names both ids: {err:?}",
);
Ok(())
}
#[cfg(feature = "encryption")]
#[test]
fn salvage_recovers_an_encrypted_sst_with_a_nonzero_table_id() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let enc: Arc<dyn crate::encryption::EncryptionProvider> =
Arc::new(crate::encryption::Aes256GcmProvider::new(&[0x37; 32]));
const TID: crate::table::TableId = 7;
let mut writer = Writer::new(source.clone(), TID, 0, Arc::clone(&fs))?
.use_data_block_size(256)
.use_encryption(Some(Arc::clone(&enc)));
let n = 200u32;
for i in 0..n {
writer.write(iv(i))?;
}
assert!(
writer.finish()?.is_some(),
"source encrypted SST is non-empty"
);
corrupt_second_data_block(
&source,
&fs,
TID,
Some(Arc::clone(&enc)),
#[cfg(zstd_any)]
None,
)?;
let wrong = SalvageOptions {
encryption: Some(Arc::clone(&enc)),
#[cfg(zstd_any)]
zstd_dictionary: None,
table_id: 0,
expected_stored_id: None,
output_id: None,
allow_delete_resurrection: false,
sync_mode: crate::fs::SyncMode::Normal,
prefix_extractor: None,
blob_rewrite: None,
progress: None,
};
let recovered_wrong = salvage_sst_with_options(&source, dest.clone(), &fs, &wrong)
.map_or(0, |r| r.entries_salvaged);
assert_eq!(
recovered_wrong, 0,
"the wrong table id cannot decrypt the AAD-bound encrypted source",
);
let options = SalvageOptions {
encryption: Some(Arc::clone(&enc)),
#[cfg(zstd_any)]
zstd_dictionary: None,
table_id: TID,
expected_stored_id: None,
output_id: None,
allow_delete_resurrection: false,
sync_mode: crate::fs::SyncMode::Normal,
prefix_extractor: None,
blob_rewrite: None,
progress: None,
};
let report = salvage_sst_with_options(&source, dest.clone(), &fs, &options)?;
assert_eq!(
report.dropped.len(),
1,
"exactly the corrupt block drops: {report:?}"
);
assert!(
report.entries_salvaged > 0 && report.entries_salvaged < u64::from(n),
"a partial key range is recovered, got {} of {n}",
report.entries_salvaged,
);
let checksum = crate::Checksum::from_raw(crate::repair::compute_table_checksum(&*fs, &dest)?);
let reopened = {
let mut params = crate::table::RecoverParams::new(
dest,
checksum,
TID,
Arc::clone(&fs),
default_comparator(),
Arc::new(crate::cache::Cache::with_capacity_bytes(1 << 20)),
);
params.descriptor_table = Some(Arc::new(crate::descriptor_table::DescriptorTable::new(8)));
params.encryption = Some(Arc::clone(&enc));
Table::recover(params)?
};
assert_eq!(
reopened.metadata.item_count, report.entries_salvaged,
"the recovered copy reopens under the same table id with the recovered entries",
);
Ok(())
}
use crate::vlog::blob_file::scanner::Scanner as BlobScanner;
use crate::vlog::blob_file::writer::Writer as BlobWriter;
fn build_blob(
path: &std::path::Path,
fs: &Arc<dyn Fs>,
records: &[(&[u8], &[u8])],
) -> crate::Result<()> {
let mut writer = BlobWriter::new(path, 0, 0, &**fs)?;
for (k, v) in records {
writer.write(k, 0, v)?;
}
writer.finish()?;
Ok(())
}
#[test]
fn divergent_non_derivable_metadata_is_never_arbitrated() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer =
Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_bulk_ingested(Some(true));
for i in 0..32u32 {
writer.write(crate::InternalValue::from_components(
format!("k{i:05}").into_bytes(),
b"v".to_vec(),
0,
crate::ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "the source SST is non-empty");
crate::test_forge::forge_tail_meta_value(&source, b"descriptor#bulk_ingested", &[0])?;
let result = salvage_sst(&source, dest, &fs);
assert!(
result.is_err(),
"a disagreement in provenance the walk cannot re-derive must refuse \
arbitration, not publish the mirror that happens to decode: {result:?}",
);
Ok(())
}
#[test]
fn an_environmental_read_never_breaks_the_physical_tiling_silently() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let mut fault_reached_the_walk = false;
for skip in 0..48u64 {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let plain: Arc<dyn Fs> = Arc::new(StdFs);
{
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&plain))?
.use_data_block_size(128)
.use_partitioned_index();
for i in 0..256u32 {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "the source SST is non-empty");
}
{
let (index_pos, index_len) = {
let mut f = std::fs::File::open(&source)?;
let reader = crate::sfa::Reader::from_reader(&mut f)?;
let Some((pos, len)) = reader
.toc()
.iter()
.find(|e| e.name() == b"index")
.map(|e| (e.pos(), e.len()))
else {
panic!("the SST carries an index section");
};
(pos, len)
};
let Ok(flip) = usize::try_from(index_pos + index_len / 2) else {
panic!("the index-section offset fits usize");
};
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(flip) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
}
let fault = FaultFs::new(StdFs);
fault.injector().arm(
FaultRule::new(FaultOp::ReadAt, Fault::Error(ErrorKind::PermissionDenied))
.skip(skip)
.times(1),
);
let fs: Arc<dyn Fs> = Arc::new(fault);
match salvage_sst(&source, dest, &fs) {
Ok(report) => {
assert!(
!report
.dropped
.iter()
.any(|d| format!("{:?}", d.reason).contains("PermissionDenied")),
"an environmental failure must never be recorded as a \
dropped block (skip {skip}): {report:?}",
);
}
Err(e) => {
assert!(
matches!(&e, crate::Error::Io(io) if io.kind() == ErrorKind::PermissionDenied),
"the only acceptable failure is the propagated \
environmental error (skip {skip}): {e:?}",
);
fault_reached_the_walk = true;
}
}
}
assert!(
fault_reached_the_walk,
"no skip count reached the salvage read path; the sweep proves nothing",
);
Ok(())
}
#[test]
fn an_environmental_read_failure_never_becomes_a_lossy_sst_salvage() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let mut fault_reached_the_walk = false;
for skip in 0..64u64 {
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let plain: Arc<dyn Fs> = Arc::new(StdFs);
{
let mut writer =
Writer::new(source.clone(), 0, 0, Arc::clone(&plain))?.use_data_block_size(256);
for i in 0..64u32 {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "the source SST is non-empty");
}
let fault = FaultFs::new(StdFs);
fault.injector().arm(
FaultRule::new(FaultOp::ReadAt, Fault::Error(ErrorKind::PermissionDenied))
.skip(skip)
.times(1),
);
let fs: Arc<dyn Fs> = Arc::new(fault);
match salvage_sst(&source, dest, &fs) {
Ok(report) => {
assert!(
!report
.dropped
.iter()
.any(|d| format!("{:?}", d.reason).contains("PermissionDenied")),
"an environmental failure must never be recorded as a \
dropped block (skip {skip}): {report:?}",
);
}
Err(e) => {
assert!(
matches!(&e, crate::Error::Io(io) if io.kind() == ErrorKind::PermissionDenied),
"the only acceptable failure is the propagated \
environmental error (skip {skip}): {e:?}",
);
fault_reached_the_walk = true;
}
}
}
assert!(
fault_reached_the_walk,
"no skip count reached the salvage read path; the sweep proves nothing",
);
Ok(())
}
#[test]
fn an_environmental_read_failure_never_becomes_a_lossy_blob_salvage() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let mut fault_reached_the_walk = false;
for skip in 0..64u64 {
let dir = tempdir()?;
let source = dir.path().join("blob_env");
let dest = dir.path().join("blob_env_salvaged");
let plain: Arc<dyn Fs> = Arc::new(StdFs);
build_blob(
&source,
&plain,
&[
(b"aaaa", b"AAAAAAAA"),
(b"bbbb", b"BBBBBBBB"),
(b"cccc", b"CCCCCCCC"),
],
)?;
let fault = FaultFs::new(StdFs);
fault.injector().arm(
FaultRule::new(FaultOp::Read, Fault::Error(ErrorKind::PermissionDenied))
.skip(skip)
.times(1),
);
let fs: Arc<dyn Fs> = Arc::new(fault);
let result = salvage_blob_file(
&source,
dest,
&fs,
0,
&default_comparator(),
0,
#[cfg(zstd_any)]
None,
);
match result {
Ok(report) => {
assert!(
!report
.dropped
.iter()
.any(|d| matches!(&d.reason, BlobDropReason::Corrupt(msg) if msg.contains("PermissionDenied"))),
"an environmental failure must never be recorded as a \
corrupt drop (skip {skip}): {report:?}",
);
}
Err(e) => {
assert!(
matches!(&e, crate::Error::Io(io) if io.kind() == ErrorKind::PermissionDenied),
"the only acceptable failure is the propagated \
environmental error (skip {skip}): {e:?}",
);
fault_reached_the_walk = true;
}
}
}
assert!(
fault_reached_the_walk,
"no skip count reached the salvage read path; the sweep proves nothing",
);
Ok(())
}
fn scan_blob(path: &std::path::Path, fs: &Arc<dyn Fs>) -> crate::Result<Vec<(Vec<u8>, Vec<u8>)>> {
Ok(BlobScanner::new(path, &**fs, 0)?
.filter_map(Result::ok)
.map(|e| (e.key.to_vec(), e.value.to_vec()))
.collect())
}
#[test]
fn salvage_blob_file_recovers_every_record_of_a_healthy_file() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("blob_source");
let dest = dir.path().join("blob_salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let records: Vec<(&[u8], &[u8])> = vec![
(b"k0", b"v0"),
(b"k1", b"v1"),
(b"k2", b"v2"),
(b"k3", b"v3"),
];
build_blob(&source, &fs, &records)?;
let report = salvage_blob_file(
&source,
dest.clone(),
&fs,
0,
&default_comparator(),
0,
#[cfg(zstd_any)]
None,
)?;
assert!(
report.is_complete(),
"a healthy blob file drops nothing: {report:?}"
);
assert_eq!(report.records_salvaged, 4);
assert_eq!(report.salvaged_path.as_deref(), Some(dest.as_path()));
let recovered = scan_blob(&dest, &fs)?;
let expected: Vec<(Vec<u8>, Vec<u8>)> = records
.iter()
.map(|(k, v)| (k.to_vec(), v.to_vec()))
.collect();
assert_eq!(
recovered, expected,
"every record round-trips through salvage"
);
Ok(())
}
#[cfg(feature = "lz4")]
#[test]
fn publish_fails_when_the_temp_name_cannot_be_dropped() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fault = FaultFs::new(StdFs);
fault.injector().arm(
FaultRule::new(
FaultOp::RemoveFile,
Fault::Error(ErrorKind::PermissionDenied),
)
.on_path("healtmp"),
);
let fs: Arc<dyn Fs> = Arc::new(fault);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0u64..10 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tail_meta_value(&source, b"compression#data", &[1])?;
let result = salvage_sst(&source, dest.clone(), &fs);
assert!(
result.as_ref().is_err_and(
|e| matches!(e, crate::Error::Io(e) if e.kind() == ErrorKind::PermissionDenied)
),
"a stuck temp name must fail the salvage, not be shrugged off: {result:?}",
);
assert!(
!dest.exists(),
"the unpublished destination must not survive the failure",
);
Ok(())
}
#[cfg(feature = "lz4")]
#[test]
fn publish_refuses_to_replace_a_concurrently_created_destination_without_hard_links()
-> crate::Result<()> {
use crate::fs::{Fs, FsDirEntry, FsFile, FsMetadata, FsOpenOptions};
use crate::io;
use std::path::Path;
#[derive(Debug)]
struct RacingProbeFs(StdFs);
impl Fs for RacingProbeFs {
fn open(&self, path: &Path, opts: &FsOpenOptions) -> io::Result<Box<dyn FsFile>> {
self.0.open(path, opts)
}
fn create_dir_all(&self, path: &Path) -> io::Result<()> {
self.0.create_dir_all(path)
}
fn read_dir(&self, path: &Path) -> io::Result<Vec<FsDirEntry>> {
self.0.read_dir(path)
}
fn remove_file(&self, path: &Path) -> io::Result<()> {
self.0.remove_file(path)
}
fn remove_dir_all(&self, path: &Path) -> io::Result<()> {
self.0.remove_dir_all(path)
}
fn rename(&self, from: &Path, to: &Path) -> io::Result<()> {
self.0.rename(from, to)
}
fn metadata(&self, path: &Path) -> io::Result<FsMetadata> {
self.0.metadata(path)
}
fn sync_directory(&self, path: &Path) -> io::Result<()> {
self.0.sync_directory(path)
}
fn exists(&self, _path: &Path) -> io::Result<bool> {
Ok(false)
}
}
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("dest");
let concurrent = b"concurrently published content".to_vec();
std::fs::write(&dest, &concurrent)?;
let fs: Arc<dyn Fs> = Arc::new(RacingProbeFs(StdFs));
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0u64..10 {
writer.write(InternalValue::from_components(
format!("key-{i:03}").into_bytes(),
format!("val-{i:03}").into_bytes(),
i + 1,
ValueType::Value,
))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_tail_meta_value(&source, b"compression#data", &[1])?;
let result = salvage_sst(&source, dest.clone(), &fs);
assert!(
result.as_ref().is_err_and(
|e| matches!(e, crate::Error::Io(e) if e.kind() == crate::io::ErrorKind::AlreadyExists)
),
"the publish must refuse the taken destination, not replace it: {result:?}",
);
assert_eq!(
std::fs::read(&dest)?,
concurrent,
"the concurrently published file must survive byte-for-byte",
);
Ok(())
}
#[test]
fn salvage_blob_file_accepts_a_bare_relative_destination_on_memfs() -> crate::Result<()> {
let fs: Arc<dyn Fs> = Arc::new(crate::fs::MemFs::new());
let source = std::path::Path::new("blob_source");
let dest = std::path::PathBuf::from("recovered");
let records: Vec<(&[u8], &[u8])> = vec![(b"k0", b"v0"), (b"k1", b"v1")];
build_blob(source, &fs, &records)?;
let report = salvage_blob_file(
source,
dest.clone(),
&fs,
0,
&default_comparator(),
0,
#[cfg(zstd_any)]
None,
)?;
assert!(report.is_complete(), "nothing to drop: {report:?}");
assert_eq!(report.records_salvaged, 2);
assert!(
fs.exists(&dest)?,
"the salvaged destination survives its entry sync",
);
Ok(())
}
#[test]
fn salvage_blob_file_drops_a_frame_whose_key_regresses() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("blob_source");
let dest = dir.path().join("blob_salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let records: Vec<(&[u8], &[u8])> = vec![(b"k0", b"v0"), (b"k2", b"v2"), (b"k1", b"v1")];
build_blob(&source, &fs, &records)?;
let report = salvage_blob_file(
&source,
dest.clone(),
&fs,
0,
&default_comparator(),
0,
#[cfg(zstd_any)]
None,
)?;
assert_eq!(
report.records_salvaged, 2,
"the order-regressing frame is dropped, not re-emitted: {report:?}",
);
let [dropped] = report.dropped.as_slice() else {
panic!(
"exactly the regressing frame drops, got {:?}",
report.dropped
);
};
assert!(
matches!(dropped.reason, BlobDropReason::Corrupt(_)),
"the drop is recorded as corruption: {:?}",
dropped.reason,
);
let recovered = scan_blob(&dest, &fs)?;
assert_eq!(
recovered,
vec![
(b"k0".to_vec(), b"v0".to_vec()),
(b"k2".to_vec(), b"v2".to_vec())
],
"only the in-order prefix survives, still sorted",
);
Ok(())
}
#[test]
fn salvage_blob_file_keeps_reverse_ordered_records_under_a_reverse_comparator() -> crate::Result<()>
{
struct ReverseComparator;
impl crate::comparator::UserComparator for ReverseComparator {
fn name(&self) -> &'static str {
"reverse-lexicographic-test"
}
fn compare(&self, a: &[u8], b: &[u8]) -> core::cmp::Ordering {
b.cmp(a)
}
}
let dir = tempdir()?;
let source = dir.path().join("blob_source");
let dest = dir.path().join("blob_salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let records: Vec<(&[u8], &[u8])> = vec![(b"k2", b"v2"), (b"k1", b"v1"), (b"k0", b"v0")];
build_blob(&source, &fs, &records)?;
let comparator: crate::comparator::SharedComparator = Arc::new(ReverseComparator);
let report = salvage_blob_file(
&source,
dest,
&fs,
0,
&comparator,
0,
#[cfg(zstd_any)]
None,
)?;
assert_eq!(
report.records_salvaged, 3,
"descending order is NOT regressing under a reverse comparator: {report:?}",
);
assert!(
report.dropped.is_empty(),
"no record is dropped as out-of-order: {report:?}",
);
Ok(())
}
#[test]
fn salvage_blob_file_syncs_the_destination_directory() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultInjector, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("blob_source");
let out = dir.path().join("blobdest");
std::fs::create_dir_all(&out)?;
let dest = out.join("blob_salvaged");
let injector = Arc::new(FaultInjector::new());
let fs: Arc<dyn Fs> = Arc::new(FaultFs::with_injector(StdFs, Arc::clone(&injector)));
let records: Vec<(&[u8], &[u8])> = vec![(b"k0", b"v0"), (b"k1", b"v1")];
build_blob(&source, &fs, &records)?;
injector.arm(
FaultRule::new(FaultOp::SyncDirectory, Fault::Error(ErrorKind::Other)).on_path("blobdest"),
);
let Err(err) = salvage_blob_file(
&source,
dest,
&fs,
0,
&default_comparator(),
0,
#[cfg(zstd_any)]
None,
) else {
panic!("the destination-directory fsync fault must surface");
};
assert!(
err.to_string().contains("injected fault on SyncDirectory"),
"the salvage error must be the injected destination-directory sync fault, got {err:?}",
);
Ok(())
}
#[test]
fn salvage_blob_file_keeps_a_racing_dest_created_after_the_existence_probe() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("blob_source");
let dest = dir.path().join("blob_dest");
let plain: Arc<dyn Fs> = Arc::new(StdFs);
build_blob(&source, &plain, &[(b"k0", b"v0")])?;
std::fs::write(&dest, b"racing worker's blob")?;
let fault = FaultFs::new(StdFs);
fault.injector().arm(
FaultRule::new(FaultOp::Metadata, Fault::Error(ErrorKind::NotFound)).on_path("blob_dest"),
);
let fs: Arc<dyn Fs> = Arc::new(fault);
let result = salvage_blob_file(
&source,
dest.clone(),
&fs,
0,
&default_comparator(),
0,
#[cfg(zstd_any)]
None,
);
assert!(
result.is_err(),
"the destination is taken, the salvage fails: {result:?}",
);
assert_eq!(
std::fs::read(&dest)?,
b"racing worker's blob",
"the racing worker's file survives the failed salvage",
);
Ok(())
}
#[test]
fn salvage_load_block_reencodes_when_the_verbatim_reread_fails() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
use crate::table::block::BlockType;
let dir = tempdir()?;
let source = dir.path().join("source");
let fault = FaultFs::new(StdFs);
let injector = fault.injector();
let fs: Arc<dyn Fs> = Arc::new(fault);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0..50u32 {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let table = open(source, &fs)?;
let Some(kh) = table.data_block_handles().find_map(Result::ok) else {
panic!("source has at least one data block");
};
let handle = *kh.as_ref();
injector.arm(
FaultRule::new(FaultOp::ReadAt, Fault::Error(ErrorKind::Other))
.on_path("source")
.skip(1)
.once(),
);
let result = table.salvage_load_block(&handle, BlockType::Data);
injector.clear();
let sb = match result {
Ok(sb) => sb,
Err(e) => panic!("a failed verbatim re-read falls back to re-encode, got Err({e:?})"),
};
assert!(
sb.verbatim.is_none(),
"the re-read was never verified, so the block must not be byte-copied",
);
assert!(
!sb.block.data.is_empty(),
"the verified first read's decoded payload is preserved for re-encoding",
);
Ok(())
}
#[test]
fn tli_structure_authenticated_propagates_a_transient_read_failure() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("source");
let fault = FaultFs::new(StdFs);
let injector = fault.injector();
let fs: Arc<dyn Fs> = Arc::new(fault);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0..50u32 {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let table = open(source, &fs)?;
injector.arm(
FaultRule::new(FaultOp::Open, Fault::Error(ErrorKind::Interrupted))
.on_path("source")
.once(),
);
let result = table.tli_structure_authenticated();
injector.clear();
assert!(
matches!(result, Err(crate::Error::Io(_))),
"a transient failure authenticating the index structure must propagate, not be \
read as an untrusted index: {result:?}",
);
Ok(())
}
#[test]
fn tli_structure_authenticated_degrades_a_persistent_read_failure() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("source");
let fault = FaultFs::new(StdFs);
let injector = fault.injector();
let fs: Arc<dyn Fs> = Arc::new(fault);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
for i in 0..50u32 {
writer.write(iv(i))?;
}
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let table = open(source, &fs)?;
injector.arm(FaultRule::new(FaultOp::Open, Fault::Error(ErrorKind::Other)).on_path("source"));
let result = table.tli_structure_authenticated();
injector.clear();
assert!(
matches!(result, Ok(false)),
"a persistent authentication failure must degrade to an untrusted index so the \
salvage walk can still recover the readable data: {result:?}",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn delete_positions_verified_propagates_a_transient_block_read_failure() -> crate::Result<()> {
use crate::config::DeleteStrategy;
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("source");
let fault = FaultFs::new(StdFs);
let injector = fault.injector();
let fs: Arc<dyn Fs> = Arc::new(fault);
let n = 64u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 20, 40] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
let table = open(source, &fs)?;
injector.arm(
FaultRule::new(FaultOp::ReadAt, Fault::Error(ErrorKind::Interrupted))
.on_path("source")
.once(),
);
let result = table.delete_positions_verified();
injector.clear();
assert!(
matches!(result, Err(crate::Error::Io(_))),
"a transient read during delete-position validation must propagate, not be read \
as a persistent unpositionable mask: {result:?}",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn delete_positions_verified_propagates_an_environmental_block_read_failure() -> crate::Result<()> {
use crate::config::DeleteStrategy;
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("source");
let fault = FaultFs::new(StdFs);
let injector = fault.injector();
let fs: Arc<dyn Fs> = Arc::new(fault);
let n = 64u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 20, 40] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
let table = open(source, &fs)?;
injector.arm(
FaultRule::new(FaultOp::ReadAt, Fault::Error(ErrorKind::PermissionDenied))
.on_path("source")
.once(),
);
let result = table.delete_positions_verified();
injector.clear();
assert!(
matches!(result, Err(crate::Error::Io(ref e)) if e.kind() == ErrorKind::PermissionDenied),
"a refused read during delete-position validation must propagate, not be read \
as a persistent unpositionable mask: {result:?}",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn delete_positions_verified_degrades_a_persistent_block_read_failure() -> crate::Result<()> {
use crate::config::DeleteStrategy;
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("source");
let fault = FaultFs::new(StdFs);
let injector = fault.injector();
let fs: Arc<dyn Fs> = Arc::new(fault);
let n = 64u32;
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.delete_strategy(DeleteStrategy::MergeOnRead);
for i in 0..n {
writer.write(iv(i))?;
}
for pos in [5u32, 20, 40] {
writer.delete_bitmap_mut().insert(pos);
}
assert!(
writer.finish()?.is_some(),
"source columnar+deletes SST is non-empty",
);
let table = open(source, &fs)?;
injector.arm(
FaultRule::new(FaultOp::ReadAt, Fault::Error(ErrorKind::Other))
.on_path("source")
.once(),
);
let result = table.delete_positions_verified();
injector.clear();
assert!(
matches!(result, Ok(false)),
"a persistent read during delete-position validation must degrade to an \
unpositionable mask so the resurrection opt-in can take effect: {result:?}",
);
Ok(())
}
#[test]
fn salvage_blob_file_keeps_a_preexisting_dest_on_open_failure() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("blob_source");
let dest = dir.path().join("blob_dest");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
build_blob(&source, &fs, &[(b"k0", b"v0")])?;
std::fs::write(&dest, b"pre-existing destination bytes")?;
let result = salvage_blob_file(
&source,
dest.clone(),
&fs,
0,
&default_comparator(),
0,
#[cfg(zstd_any)]
None,
);
assert!(
result.is_err(),
"an already-existing destination fails the salvage: {result:?}",
);
assert_eq!(
std::fs::read(&dest)?,
b"pre-existing destination bytes",
"a pre-existing destination file survives the failed salvage",
);
Ok(())
}
#[test]
fn salvage_blob_file_reports_an_offset_remap_for_every_salvaged_record() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("blob_source");
let dest = dir.path().join("blob_salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let records: Vec<(&[u8], &[u8])> = vec![
(b"k0", b"v0-payload"),
(b"k1", b"v1-payload"),
(b"k2", b"v2-payload"),
(b"k3", b"v3-payload"),
];
build_blob(&source, &fs, &records)?;
let source_offsets: Vec<u64> = BlobScanner::new(&source, &*fs, 0)?
.filter_map(Result::ok)
.map(|e| e.offset)
.collect();
assert_eq!(source_offsets.len(), 4, "four source records");
{
let Some(&second) = source_offsets.get(1) else {
panic!("second record offset");
};
let flip = usize::try_from(second).unwrap_or(0) + 45;
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(flip) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
}
let report = salvage_blob_file(
&source,
dest.clone(),
&fs,
0,
&default_comparator(),
0,
#[cfg(zstd_any)]
None,
)?;
assert_eq!(report.records_salvaged, 1, "{report:?}");
assert_eq!(
report.dropped.len(),
2,
"the corrupt record + the surrendered tail (recorded once) drop: {report:?}"
);
let dest_records: Vec<(u64, u32)> = BlobScanner::new(&dest, &*fs, 0)?
.filter_map(Result::ok)
.map(|e| {
let value_start = e.offset
+ u64::try_from(crate::vlog::blob_file::writer::BLOB_HEADER_LEN).unwrap_or(0)
+ u64::try_from(e.key.len()).unwrap_or(0);
(
e.offset,
u32::try_from(e.frame_end - value_start).unwrap_or(0),
)
})
.collect();
let expected: Vec<(u64, super::BlobRecordRelocation)> = std::iter::once(&0usize)
.zip(&dest_records)
.map(|(&src_idx, &(offset, on_disk_size))| {
(
source_offsets.get(src_idx).copied().unwrap_or(u64::MAX),
super::BlobRecordRelocation {
offset,
on_disk_size,
},
)
})
.collect();
assert_eq!(
report.offset_remap, expected,
"the remap maps each surviving source frame to its compacted target",
);
let Some(&dropped_src) = source_offsets.get(1) else {
panic!("second record offset");
};
assert!(
report
.offset_remap
.iter()
.all(|(src, _)| *src != dropped_src),
"a dropped record has no remap target: {report:?}",
);
Ok(())
}
#[test]
fn salvage_blob_file_removes_the_partial_dest_when_a_write_fails() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempdir()?;
let source = dir.path().join("blob_source");
let dest = dir.path().join("blob_salvaged");
let fault = FaultFs::new(StdFs);
let injector = fault.injector();
let fs: Arc<dyn Fs> = Arc::new(fault);
let records: Vec<(&[u8], &[u8])> = vec![(b"k0", b"v0"), (b"k1", b"v1")];
build_blob(&source, &fs, &records)?;
injector.arm(
FaultRule::new(FaultOp::Write, Fault::Error(ErrorKind::Other)).on_path("blob_salvaged"),
);
let result = salvage_blob_file(
&source,
dest.clone(),
&fs,
0,
&default_comparator(),
0,
#[cfg(zstd_any)]
None,
);
assert!(
result.is_err(),
"a failed destination write must error the salvage",
);
assert!(
!std::path::Path::new(&dest).exists(),
"the partial destination is removed on a write failure",
);
Ok(())
}
#[test]
fn salvage_blob_file_drops_a_corrupt_record_and_keeps_the_rest() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("blob_source");
let dest = dir.path().join("blob_salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let records: Vec<(&[u8], &[u8])> = vec![
(b"k0", b"value-zero"),
(b"k1", b"value-one"),
(b"k2", b"value-two"),
(b"k3", b"value-three"),
];
build_blob(&source, &fs, &records)?;
let Some(second_frame_end) = BlobScanner::new(&source, &*fs, 0)?
.filter_map(Result::ok)
.nth(1)
.map(|e| e.frame_end)
else {
panic!("source blob must have at least two records");
};
let flip = usize::try_from(second_frame_end - 1).unwrap_or(0);
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(flip) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let report = salvage_blob_file(
&source,
dest.clone(),
&fs,
0,
&default_comparator(),
0,
#[cfg(zstd_any)]
None,
)?;
assert_eq!(
report.dropped.len(),
2,
"the corrupt record + the surrendered tail (recorded once) drop: {report:?}"
);
assert!(
matches!(
report.dropped.first().map(|d| &d.reason),
Some(BlobDropReason::ChecksumMismatch)
),
"the corrupt record reports a checksum mismatch: {report:?}",
);
assert_eq!(
report
.dropped
.iter()
.filter(
|d| matches!(&d.reason, BlobDropReason::Corrupt(m) if m.contains("surrendered"))
)
.count(),
1,
"the surrendered tail after the resync is recorded ONCE: {report:?}",
);
assert_eq!(
report.records_salvaged, 1,
"only k0 is recovered; k1 (corrupt) and the surrendered tail (k2, k3) are not"
);
let recovered = scan_blob(&dest, &fs)?;
let keys: Vec<Vec<u8>> = recovered.iter().map(|(k, _)| k.clone()).collect();
assert_eq!(
keys,
vec![b"k0".to_vec()],
"only the provable frame before the corruption survives",
);
Ok(())
}
#[test]
fn salvage_blob_file_removes_the_empty_dest_when_nothing_is_recoverable() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("blob_source");
let dest = dir.path().join("blob_salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let records: Vec<(&[u8], &[u8])> = vec![(b"k0", b"value-zero"), (b"k1", b"value-one")];
build_blob(&source, &fs, &records)?;
let frame_ends: Vec<u64> = BlobScanner::new(&source, &*fs, 0)?
.filter_map(Result::ok)
.map(|e| e.frame_end)
.collect();
assert_eq!(frame_ends.len(), 2, "source blob holds two records");
let mut bytes = std::fs::read(&source)?;
for end in frame_ends {
let flip = usize::try_from(end - 1).unwrap_or(0);
if let Some(b) = bytes.get_mut(flip) {
*b ^= 0xFF;
}
}
std::fs::write(&source, &bytes)?;
let report = salvage_blob_file(
&source,
dest.clone(),
&fs,
0,
&default_comparator(),
0,
#[cfg(zstd_any)]
None,
)?;
assert_eq!(report.records_salvaged, 0, "{report:?}");
assert_eq!(report.dropped.len(), 2, "both records drop: {report:?}");
assert_eq!(
report.salvaged_path, None,
"nothing recoverable yields no salvaged path",
);
assert!(
fs.metadata(&dest).is_err(),
"the empty destination placeholder is removed",
);
Ok(())
}
#[test]
fn salvage_blob_file_stops_at_a_smashed_frame_and_keeps_the_prefix() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("blob_source");
let dest = dir.path().join("blob_salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let records: Vec<(&[u8], &[u8])> = vec![
(b"k0", b"value-zero"),
(b"k1", b"value-one"),
(b"k2", b"value-two"),
];
build_blob(&source, &fs, &records)?;
let Some(last_start) = BlobScanner::new(&source, &*fs, 0)?
.filter_map(Result::ok)
.nth(1)
.map(|e| e.frame_end)
else {
panic!("source blob must have at least two records");
};
let mut bytes = std::fs::read(&source)?;
let at = usize::try_from(last_start).unwrap_or(0);
let Some(magic) = bytes.get_mut(at..at + 4) else {
panic!("last record's frame magic within the file");
};
magic.copy_from_slice(b"????");
std::fs::write(&source, &bytes)?;
let report = salvage_blob_file(
&source,
dest.clone(),
&fs,
0,
&default_comparator(),
0,
#[cfg(zstd_any)]
None,
)?;
assert_eq!(
report.records_salvaged, 2,
"the records before the smashed frame are recovered: {report:?}",
);
assert!(
matches!(
report.dropped.first().map(|d| &d.reason),
Some(BlobDropReason::Corrupt(_))
),
"the truncated tail is recorded as a structural drop: {report:?}",
);
let recovered = scan_blob(&dest, &fs)?;
assert_eq!(recovered.len(), 2, "the salvaged copy holds the prefix");
Ok(())
}
#[test]
fn salvage_blob_file_drops_an_empty_key_frame() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("blob_source");
let dest = dir.path().join("blob_salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
build_blob(&source, &fs, &[(b"k", b"vvvv"), (b"k2", b"second")])?;
let mut bytes = std::fs::read(&source)?;
let seqno = {
let Some(b) = bytes.get(20..28) else {
panic!("seqno within the first frame");
};
u64::from_le_bytes(b.try_into().unwrap_or([0; 8]))
};
let new_hcrc = {
let mut hasher = xxhash_rust::xxh3::Xxh3::default();
hasher.update(&seqno.to_le_bytes());
hasher.update(&0u16.to_le_bytes());
hasher.update(&5u32.to_le_bytes());
hasher.update(&5u32.to_le_bytes());
#[expect(
clippy::cast_possible_truncation,
reason = "intentionally truncated to the 4-byte header CRC"
)]
{
hasher.digest() as u32
}
};
let new_checksum = {
let mut hasher = xxhash_rust::xxh3::Xxh3::default();
hasher.update(b"kvvvv");
hasher.update(&new_hcrc.to_le_bytes());
hasher.digest128()
};
let patch = |bytes: &mut Vec<u8>, range: core::ops::Range<usize>, val: &[u8]| {
let Some(slot) = bytes.get_mut(range) else {
panic!("patch range within the first frame");
};
slot.copy_from_slice(val);
};
patch(&mut bytes, 4..20, &new_checksum.to_le_bytes());
patch(&mut bytes, 28..30, &0u16.to_le_bytes());
patch(&mut bytes, 30..34, &5u32.to_le_bytes());
patch(&mut bytes, 34..38, &5u32.to_le_bytes());
patch(&mut bytes, 38..42, &new_hcrc.to_le_bytes());
std::fs::write(&source, &bytes)?;
let report = salvage_blob_file(
&source,
dest.clone(),
&fs,
0,
&default_comparator(),
0,
#[cfg(zstd_any)]
None,
)?;
assert_eq!(
report.dropped.len(),
1,
"the empty-key frame drops as corrupt: {report:?}",
);
assert!(
matches!(
report.dropped.first().map(|d| &d.reason),
Some(BlobDropReason::Corrupt(_))
),
"the drop reason names the malformed frame: {report:?}",
);
assert_eq!(
report.records_salvaged, 1,
"the record after the malformed frame is still recovered",
);
let recovered = scan_blob(&dest, &fs)?;
assert_eq!(
recovered,
vec![(b"k2".to_vec(), b"second".to_vec())],
"the salvaged copy holds exactly the healthy record",
);
Ok(())
}
#[test]
fn salvage_blob_file_drops_a_frame_with_a_forged_value_length() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("blob_source");
let dest = dir.path().join("blob_salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
build_blob(&source, &fs, &[(b"k", b"vvvv"), (b"k2", b"second")])?;
let mut bytes = std::fs::read(&source)?;
let seqno = {
let Some(b) = bytes.get(20..28) else {
panic!("seqno within the first frame");
};
u64::from_le_bytes(b.try_into().unwrap_or([0; 8]))
};
let new_hcrc = {
let mut hasher = xxhash_rust::xxh3::Xxh3::default();
hasher.update(&seqno.to_le_bytes());
hasher.update(&1u16.to_le_bytes());
hasher.update(&5u32.to_le_bytes());
hasher.update(&4u32.to_le_bytes());
#[expect(
clippy::cast_possible_truncation,
reason = "intentionally truncated to the 4-byte header CRC"
)]
{
hasher.digest() as u32
}
};
let new_checksum = {
let mut hasher = xxhash_rust::xxh3::Xxh3::default();
hasher.update(b"kvvvv");
hasher.update(&new_hcrc.to_le_bytes());
hasher.digest128()
};
let patch = |bytes: &mut Vec<u8>, range: core::ops::Range<usize>, val: &[u8]| {
let Some(slot) = bytes.get_mut(range) else {
panic!("patch range within the first frame");
};
slot.copy_from_slice(val);
};
patch(&mut bytes, 4..20, &new_checksum.to_le_bytes());
patch(&mut bytes, 30..34, &5u32.to_le_bytes());
patch(&mut bytes, 38..42, &new_hcrc.to_le_bytes());
std::fs::write(&source, &bytes)?;
let report = salvage_blob_file(
&source,
dest.clone(),
&fs,
0,
&default_comparator(),
0,
#[cfg(zstd_any)]
None,
)?;
assert_eq!(
report.dropped.len(),
1,
"the forged-length frame drops as corrupt: {report:?}",
);
assert_eq!(
report.records_salvaged, 1,
"the record after the malformed frame is still recovered",
);
let recovered = scan_blob(&dest, &fs)?;
assert_eq!(
recovered,
vec![(b"k2".to_vec(), b"second".to_vec())],
"the salvaged copy holds exactly the healthy record",
);
Ok(())
}
#[cfg(feature = "lz4")]
#[test]
#[expect(clippy::expect_used, reason = "test code")]
fn salvage_blob_file_recovers_a_compressed_source() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("blob_source");
let dest = dir.path().join("blob_salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let value = b"some compressible value aaaaaaaaaaaaaaaa";
{
let mut writer =
BlobWriter::new(&source, 0, 0, &*fs)?.use_compression(crate::CompressionType::Lz4);
writer.write(b"k0", 0, value)?;
writer.finish()?;
}
let report = salvage_blob_file(
&source,
dest.clone(),
&fs,
0,
&default_comparator(),
0,
#[cfg(zstd_any)]
None,
)?;
assert_eq!(report.records_salvaged, 1, "the record is recovered");
assert!(
report.is_complete(),
"nothing dropped: {:?}",
report.dropped
);
assert_eq!(report.salvaged_path.as_ref(), Some(&dest));
let handle = crate::vlog::recover_blob_file(
&dest,
0,
crate::Checksum::from_raw(crate::repair::compute_table_checksum(&*fs, &dest)?),
0,
&fs,
)?;
assert_eq!(
handle.compression(),
crate::CompressionType::Lz4,
"the copy keeps the source's compression descriptor",
);
let (_, relocation) = *report.offset_remap.first().expect("one remap entry");
let file = std::fs::File::open(&dest)?;
let reader = crate::vlog::blob_file::reader::Reader::new(&handle, &file);
let vhandle = crate::vlog::ValueHandle {
blob_file_id: 0,
offset: relocation.offset,
on_disk_size: relocation.on_disk_size,
};
assert_eq!(
reader.get(b"k0", &vhandle)?.as_ref(),
value,
"the salvaged record decompresses to its original value",
);
Ok(())
}
#[cfg(zstd_any)]
#[test]
#[expect(clippy::expect_used, reason = "test code")]
fn salvage_blob_file_recovers_a_dictionary_source_with_the_dictionary() -> crate::Result<()> {
use std::io::{Seek, SeekFrom, Write};
let dir = tempdir()?;
let source = dir.path().join("blob_source");
let dest = dir.path().join("blob_salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let dict = Arc::new(crate::compression::ZstdDictionary::new(
b"sample sample sample sample payload payload payload",
));
let compression = crate::CompressionType::ZstdDict {
level: 3,
dict_id: dict.id(),
};
{
let mut w = BlobWriter::new(&source, 0, 0, &*fs)?
.use_compression(compression)
.use_zstd_dictionary(Some(Arc::clone(&dict)));
w.write(b"a", 1, b"sample payload sample payload sample payload")?;
w.write(b"b", 2, b"sample payload sample payload sample sample")?;
w.finish()?;
}
let entries: Vec<_> =
crate::vlog::BlobFileScanner::new(&source, &*fs, 0)?.collect::<crate::Result<Vec<_>>>()?;
let second = entries.get(1).expect("two records");
{
let mut f = std::fs::OpenOptions::new().write(true).open(&source)?;
f.seek(SeekFrom::Start(second.frame_end - 4))?;
f.write_all(&[0xFF, 0xFF, 0xFF, 0xFF])?;
}
let report = salvage_blob_file(
&source,
dest.clone(),
&fs,
0,
&default_comparator(),
0,
#[cfg(zstd_any)]
Some(&dict),
)?;
assert_eq!(
report.records_salvaged, 1,
"the intact dictionary-compressed record is recovered: {report:?}",
);
assert_eq!(report.dropped.len(), 1, "the rotted record drops");
let salvaged = crate::vlog::BlobFileScanner::new(&dest, &*fs, 0)?
.next()
.expect("one salvaged record")?;
let value = super::decompress_blob_value(
compression,
&salvaged.value,
salvaged.uncompressed_len as usize,
#[cfg(zstd_any)]
Some(&dict),
)?;
assert_eq!(
value.as_ref(),
b"sample payload sample payload sample payload",
"the salvaged record decodes to its original value under the dictionary",
);
Ok(())
}
#[cfg(all(feature = "lz4", zstd_any))]
#[test]
fn salvage_blob_file_rejects_a_dictionary_compressed_source() -> crate::Result<()> {
let dir = tempdir()?;
let source = dir.path().join("blob_source");
let dest = dir.path().join("blob_salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
{
let mut writer =
BlobWriter::new(&source, 0, 0, &*fs)?.use_compression(crate::CompressionType::Lz4);
writer.metadata_compression_override = Some(crate::CompressionType::ZstdDict {
level: 3,
dict_id: 7,
});
writer.write(b"k0", 0, b"some compressible value aaaaaaaaaaaaaaaa")?;
writer.finish()?;
}
assert!(
matches!(
salvage_blob_file(
&source,
dest,
&fs,
0,
&default_comparator(),
0,
#[cfg(zstd_any)]
None
),
Err(crate::Error::ZstdDictMismatch {
expected: 7,
got: None
}),
),
"a dictionary-compressed blob file must be rejected rather than mis-salvaged",
);
Ok(())
}
#[test]
#[expect(clippy::expect_used, reason = "test code")]
fn blob_handle_rewrite_installs_the_relocated_size() -> crate::Result<()> {
use crate::blob_tree::handle::BlobIndirection;
use crate::coding::{Decode, Encode};
use crate::vlog::ValueHandle;
use crate::{InternalValue, ValueType};
let ind = BlobIndirection {
vhandle: ValueHandle {
blob_file_id: 7,
offset: 400,
on_disk_size: 64,
},
size: 100,
};
let mut value = Vec::new();
ind.encode_into(&mut value)?;
let entries = vec![InternalValue::from_components(
b"k".to_vec(),
value,
1,
ValueType::Indirection,
)];
let mut map = crate::HashMap::default();
map.insert(
400u64,
super::BlobRecordRelocation {
offset: 16,
on_disk_size: 61,
},
);
let mut rewrite = crate::HashMap::default();
rewrite.insert(
7u64,
super::BlobFileRewrite::Remap {
new_id: 9,
offsets: map,
},
);
let mut dropped = 0u64;
let (out, carry) = super::rewrite_block_indirections(entries, &rewrite, &mut dropped)?;
assert_eq!(dropped, 0, "the record survived, nothing drops");
assert!(carry.is_none(), "nothing was beheaded, nothing to suppress");
let entry = out.first().expect("one rewritten entry");
let rewritten = BlobIndirection::decode_from(&mut &entry.value[..])?;
assert_eq!(
rewritten.vhandle.blob_file_id, 9,
"the handle names the salvaged replacement, not the damaged original",
);
assert_eq!(rewritten.vhandle.offset, 16, "the offset is re-targeted");
assert_eq!(
rewritten.vhandle.on_disk_size, 61,
"the on-disk size is the RE-EMITTED record's, not the source's: the \
reader rejects a handle whose size disagrees with the frame header",
);
Ok(())
}
#[test]
fn blob_handle_rewrite_drops_older_versions_when_the_head_record_is_lost() -> crate::Result<()> {
use crate::blob_tree::handle::BlobIndirection;
use crate::coding::Encode;
use crate::vlog::ValueHandle;
use crate::{InternalValue, ValueType};
let indirection = |offset: u64| -> crate::Result<Vec<u8>> {
let mut value = Vec::new();
BlobIndirection {
vhandle: ValueHandle {
blob_file_id: 7,
offset,
on_disk_size: 64,
},
size: 100,
}
.encode_into(&mut value)?;
Ok(value)
};
let mut rewrite = crate::HashMap::default();
rewrite.insert(
7u64,
super::BlobFileRewrite::Remap {
new_id: 9,
offsets: crate::HashMap::default(),
},
);
let entries = vec![
InternalValue::from_components(
b"k".to_vec(),
indirection(400)?,
10,
ValueType::Indirection,
),
InternalValue::from_components(b"k".to_vec(), b"old".to_vec(), 5, ValueType::Value),
InternalValue::from_components(b"z".to_vec(), b"v".to_vec(), 1, ValueType::Value),
];
let mut dropped = 0u64;
let (out, carry) = super::rewrite_block_indirections(entries, &rewrite, &mut dropped)?;
let keys: Vec<_> = out.iter().map(|e| e.key.user_key.to_vec()).collect();
assert_eq!(
keys,
vec![b"z".to_vec()],
"the beheaded key goes entirely; the next key is untouched",
);
assert_eq!(dropped, 2, "both the lost head and the version behind it");
assert!(
carry.is_none(),
"a surviving entry with a different key ended the run",
);
let entries = vec![
InternalValue::from_components(
b"k".to_vec(),
indirection(400)?,
10,
ValueType::Indirection,
),
InternalValue::from_components(b"k".to_vec(), b"old".to_vec(), 5, ValueType::Value),
];
let mut dropped = 0u64;
let (out, carry) = super::rewrite_block_indirections(entries, &rewrite, &mut dropped)?;
assert!(out.is_empty(), "the whole block was one beheaded key");
assert_eq!(
carry.as_deref(),
Some(b"k".as_slice()),
"nothing proved the run ended, so the key keeps suppressing",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_preserves_columnar_value_subcolumns() -> crate::Result<()> {
use crate::table::columnar::{Column, TypeTag, entries_to_column_batch};
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let cmp = default_comparator();
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true);
for block in 0..2u32 {
let entries: Vec<InternalValue> = (0..4u32)
.map(|i| {
let k = format!("k{:04}", block * 4 + i);
InternalValue::from_components(k.into_bytes(), b"x".to_vec(), 0, ValueType::Value)
})
.collect();
let mut batch = entries_to_column_batch(&entries)?;
batch.columns.pop();
let mut data = Vec::new();
for i in 0..4u32 {
data.extend_from_slice(&(block * 4 + i).to_le_bytes());
}
batch.columns.push(Column {
column_id: 3,
type_tag: TypeTag::Fixed(4),
validity: None,
data: data.into(),
});
writer.write_columnar_batch(&batch, &cmp)?;
}
assert!(
writer.finish()?.is_some(),
"source columnar SST is non-empty"
);
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(
report.is_complete(),
"a healthy columnar SST drops nothing: {report:?}"
);
assert_eq!(
report.blocks_copied_verbatim, report.blocks_salvaged,
"clean columnar blocks are copied verbatim",
);
let recovered = open(dest, &fs)?;
assert!(
recovered.metadata.columnar,
"the recovered copy stays columnar"
);
let batches = recovered.columnar_scan(&[3], None)?;
let rows: u32 = batches.iter().map(|b| b.row_count).sum();
assert_eq!(rows, 8, "every row's sub-column is recovered");
assert!(
batches
.iter()
.all(|b| b.columns.iter().all(|c| c.column_id == 3)),
"the value sub-column (id 3) is preserved verbatim, not collapsed",
);
Ok(())
}
#[cfg(all(feature = "columnar", feature = "page_ecc"))]
#[test]
fn salvage_reencodes_an_ecc_recovered_columnar_block() -> crate::Result<()> {
use crate::table::block::{EccParams, Header};
use crate::table::columnar::entries_to_column_batch;
let dir = tempdir()?;
let source = dir.path().join("source");
let dest = dir.path().join("salvaged");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let cmp = default_comparator();
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_zone_map(true)
.use_ecc(Some(EccParams::RS_4_2));
for block in 0..2u32 {
let entries: Vec<InternalValue> = (0..4u32)
.map(|i| {
let k = format!("k{:04}", block * 4 + i);
InternalValue::from_components(k.into_bytes(), b"x".to_vec(), 0, ValueType::Value)
})
.collect();
let batch = entries_to_column_batch(&entries)?;
writer.write_columnar_batch(&batch, &cmp)?;
}
assert!(
writer.finish()?.is_some(),
"source columnar SST is non-empty"
);
let first_off = {
let table = open(source.clone(), &fs)?;
let Some(kh) = table.data_block_handles().find_map(Result::ok) else {
panic!("a non-empty SST has at least one data block");
};
usize::try_from(*kh.as_ref().offset()).unwrap_or(usize::MAX)
};
let pos = first_off + Header::MIN_LEN + 3;
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(pos) {
*b ^= 0x80;
}
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert!(
report.is_complete(),
"an RS-recoverable columnar block is healed, not dropped: {report:?}",
);
assert_eq!(
report.blocks_salvaged, report.blocks_total,
"every block is recovered",
);
assert!(
report.blocks_copied_verbatim < report.blocks_salvaged,
"the ECC-recovered columnar block is re-encoded, not copied verbatim: {report:?}",
);
let recovered = open(dest, &fs)?;
assert!(
recovered.metadata.columnar,
"the recovered copy stays columnar"
);
assert_eq!(recovered.metadata.item_count, 8, "every row is recovered");
Ok(())
}
#[test]
fn salvage_suppresses_only_the_shadowed_key_itself() -> crate::Result<()> {
use crate::table::Writer;
use crate::{InternalValue, ValueType};
let dir = tempfile::tempdir()?;
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let source = dir.path().join("source");
let dest = dir.path().join("dest");
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(1);
writer.write(InternalValue::new_tombstone(b"K".as_slice(), 10))?;
writer.write(InternalValue::from_components(
b"k",
b"old",
5,
ValueType::Value,
))?;
writer.write(InternalValue::from_components(
b"z",
b"v",
1,
ValueType::Value,
))?;
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let first_off = {
let table = open_with_id(source.clone(), &fs, 0)?;
let offsets: Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|kh| *kh.as_ref().offset())
.collect();
assert!(
offsets.len() >= 3,
"one block per entry, got {}",
offsets.len()
);
usize::try_from(offsets.first().copied().unwrap_or_default()).unwrap_or(usize::MAX)
};
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(first_off + crate::table::block::Header::MIN_LEN + 1) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(report.dropped.len(), 1, "one block is lost: {report:?}");
let recovered = open_with_id(dest, &fs, 0)?;
assert!(
recovered
.get(b"k", crate::SeqNo::MAX, crate::hash::hash64(b"k"))?
.is_some(),
"`k` is a different key from the lost `K`: its own newest version is \
intact and must survive",
);
Ok(())
}
#[test]
fn the_blob_rewrite_suppresses_only_the_beheaded_key() -> crate::Result<()> {
use crate::coding::Encode;
use crate::{InternalValue, ValueType};
let indirection = |offset: u64| -> crate::Result<Vec<u8>> {
let ind = crate::blob_tree::handle::BlobIndirection {
vhandle: crate::vlog::ValueHandle {
blob_file_id: 7,
offset,
on_disk_size: 16,
},
size: 16,
};
let mut buf = Vec::new();
ind.encode_into(&mut buf)?;
Ok(buf)
};
let entries = vec![
InternalValue::from_components(b"a", indirection(100)?, 10, ValueType::Indirection),
InternalValue::from_components(b"a", indirection(200)?, 5, ValueType::Indirection),
InternalValue::from_components(b"z", indirection(200)?, 1, ValueType::Indirection),
];
let mut offsets = crate::HashMap::default();
offsets.insert(
200u64,
super::BlobRecordRelocation {
offset: 300,
on_disk_size: 16,
},
);
let mut rewrite = crate::HashMap::default();
rewrite.insert(7u64, super::BlobFileRewrite::Remap { new_id: 9, offsets });
let mut dropped = 0u64;
let (kept, carry) = super::rewrite_block_indirections(entries, &rewrite, &mut dropped)?;
let keys: Vec<_> = kept.iter().map(|e| e.key.user_key.to_vec()).collect();
assert_eq!(
keys,
vec![b"z".to_vec()],
"the beheaded key's older version goes with it; a different key stays",
);
assert_eq!(
dropped, 2,
"both the lost head and its orphaned older version"
);
assert!(
carry.is_none(),
"a surviving different key ended the run inside this block",
);
Ok(())
}
#[test]
fn salvage_keeps_tied_seqno_merge_operands() -> crate::Result<()> {
use crate::table::Writer;
use crate::{InternalValue, ValueType};
let dir = tempfile::tempdir()?;
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let source = dir.path().join("source");
let dest = dir.path().join("dest");
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(100);
writer.write(InternalValue::new_merge_operand(
b"k".as_slice(),
[b'A'; 64],
10,
))?;
writer.write(InternalValue::new_merge_operand(
b"k".as_slice(),
[b'B'; 64],
10,
))?;
writer.write(InternalValue::from_components(
b"z",
[b'z'; 128],
1,
ValueType::Value,
))?;
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let last_off = {
let table = open(source.clone(), &fs)?;
let offsets: Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|kh| *kh.as_ref().offset())
.collect();
assert!(
offsets.len() >= 2,
"the fixture needs the operands and `z` in separate blocks, got {}",
offsets.len(),
);
usize::try_from(offsets.last().copied().unwrap_or_default()).unwrap_or(usize::MAX)
};
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(last_off + crate::table::block::Header::MIN_LEN + 1) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(
report.dropped.len(),
1,
"only `z`'s block is lost: {report:?}"
);
let recovered = open(dest, &fs)?;
let operands: Vec<_> = recovered
.scan()?
.filter_map(Result::ok)
.filter(|e| e.key.user_key.as_ref() == b"k")
.map(|e| e.value.to_vec())
.collect();
assert_eq!(
operands.len(),
2,
"both operands are valid input to the merge, so a clean block holding \
them must not be rejected as out of order: {operands:?}",
);
Ok(())
}
#[test]
fn salvage_keeps_tied_seqno_merge_operands_across_a_block_boundary() -> crate::Result<()> {
use crate::table::Writer;
use crate::{InternalValue, ValueType};
let dir = tempfile::tempdir()?;
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let source = dir.path().join("source");
let dest = dir.path().join("dest");
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(1);
writer.write(InternalValue::new_merge_operand(b"k".as_slice(), b"A", 10))?;
writer.write(InternalValue::new_merge_operand(b"k".as_slice(), b"B", 10))?;
writer.write(InternalValue::from_components(
b"z",
b"v",
1,
ValueType::Value,
))?;
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let last_off = {
let table = open(source.clone(), &fs)?;
let offsets: Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|kh| *kh.as_ref().offset())
.collect();
assert_eq!(offsets.len(), 3, "one block per entry");
usize::try_from(offsets.last().copied().unwrap_or_default()).unwrap_or(usize::MAX)
};
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(last_off + crate::table::block::Header::MIN_LEN + 1) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(
report.dropped.len(),
1,
"only `z`'s block is lost: {report:?}"
);
let recovered = open(dest, &fs)?;
let operands = recovered
.scan()?
.filter_map(Result::ok)
.filter(|e| e.key.user_key.as_ref() == b"k")
.count();
assert_eq!(
operands, 2,
"the tie spanning the block edge is a valid run, not a violation",
);
Ok(())
}
#[test]
fn salvage_drops_the_boundary_key_when_its_newest_version_is_lost() -> crate::Result<()> {
use crate::table::Writer;
use crate::{InternalValue, ValueType};
let dir = tempfile::tempdir()?;
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let source = dir.path().join("source");
let dest = dir.path().join("dest");
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(1);
writer.write(InternalValue::from_components(
b"a",
b"v",
1,
ValueType::Value,
))?;
writer.write(InternalValue::new_tombstone(b"k".as_slice(), 10))?;
writer.write(InternalValue::from_components(
b"k",
b"old",
5,
ValueType::Value,
))?;
writer.write(InternalValue::from_components(
b"z",
b"v",
1,
ValueType::Value,
))?;
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let second_off = {
let table = open(source.clone(), &fs)?;
let offsets: Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|kh| *kh.as_ref().offset())
.collect();
assert!(
offsets.len() >= 3,
"the fixture needs one block per entry, got {}",
offsets.len(),
);
usize::try_from(offsets.get(1).copied().unwrap_or_default()).unwrap_or(usize::MAX)
};
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(second_off + crate::table::block::Header::MIN_LEN + 1) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(report.dropped.len(), 1, "one block is lost: {report:?}");
let recovered = open(dest, &fs)?;
assert!(
recovered
.get(b"k", crate::SeqNo::MAX, crate::hash::hash64(b"k"))?
.is_none(),
"the older version of the boundary key must not be republished as \
current when the newest version was lost",
);
assert!(
recovered
.get(b"a", crate::SeqNo::MAX, crate::hash::hash64(b"a"))?
.is_some(),
"unaffected keys are still recovered",
);
Ok(())
}
#[test]
fn salvage_keeps_suppressing_a_boundary_key_across_a_whole_shadowed_block() -> crate::Result<()> {
use crate::table::Writer;
use crate::{InternalValue, ValueType};
let dir = tempfile::tempdir()?;
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let source = dir.path().join("source");
let dest = dir.path().join("dest");
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(1);
writer.write(InternalValue::new_tombstone(b"k".as_slice(), 10))?;
writer.write(InternalValue::from_components(
b"k",
b"mid",
5,
ValueType::Value,
))?;
writer.write(InternalValue::from_components(
b"k",
b"old",
1,
ValueType::Value,
))?;
writer.write(InternalValue::from_components(
b"z",
b"v",
1,
ValueType::Value,
))?;
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let first_off = {
let table = open(source.clone(), &fs)?;
let offsets: Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|kh| *kh.as_ref().offset())
.collect();
assert!(
offsets.len() >= 4,
"the fixture needs one block per entry, got {}",
offsets.len(),
);
usize::try_from(offsets.first().copied().unwrap_or_default()).unwrap_or(usize::MAX)
};
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(first_off + crate::table::block::Header::MIN_LEN + 1) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(report.dropped.len(), 1, "one block is lost: {report:?}");
let recovered = open(dest, &fs)?;
assert!(
recovered
.get(b"k", crate::SeqNo::MAX, crate::hash::hash64(b"k"))?
.is_none(),
"the chain continues past the first surviving block, so every older \
version of the boundary key has to stay suppressed",
);
assert!(
recovered
.get(b"z", crate::SeqNo::MAX, crate::hash::hash64(b"z"))?
.is_some(),
"the key that ends the run is recovered",
);
Ok(())
}
#[test]
fn salvage_arms_an_unknown_boundary_when_the_lost_block_has_no_range() -> crate::Result<()> {
use crate::table::Writer;
use crate::{InternalValue, ValueType};
let dir = tempfile::tempdir()?;
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let source = dir.path().join("source");
let dest = dir.path().join("dest");
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(1);
writer.write(InternalValue::from_components(
b"a",
b"v",
1,
ValueType::Value,
))?;
writer.write(InternalValue::new_tombstone(b"k".as_slice(), 10))?;
writer.write(InternalValue::from_components(
b"k",
b"old",
5,
ValueType::Value,
))?;
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let second_off = {
let table = open(source.clone(), &fs)?;
let offsets: Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|kh| *kh.as_ref().offset())
.collect();
assert!(
offsets.len() >= 3,
"the fixture needs one block per entry, got {}",
offsets.len(),
);
usize::try_from(offsets.get(1).copied().unwrap_or_default()).unwrap_or(usize::MAX)
};
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(second_off) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
let recovered = open(dest, &fs)?;
assert!(
recovered
.get(b"k", crate::SeqNo::MAX, crate::hash::hash64(b"k"))?
.is_none(),
"the lost region'\u{2019}s range is unknown, so the FOLLOWING block'\u{2019}s \
first key is the one it may have shadowed: {report:?}",
);
Ok(())
}
#[test]
fn salvage_does_not_byte_copy_a_block_it_suppressed_a_key_from() -> crate::Result<()> {
use crate::table::Writer;
use crate::{InternalValue, ValueType};
let dir = tempfile::tempdir()?;
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let source = dir.path().join("source");
let dest = dir.path().join("dest");
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?.use_data_block_size(100);
writer.write(InternalValue::from_components(
b"k",
[b'n'; 200],
10,
ValueType::Value,
))?;
writer.write(InternalValue::from_components(
b"k",
b"old",
5,
ValueType::Value,
))?;
writer.write(InternalValue::from_components(
b"m",
b"v",
1,
ValueType::Value,
))?;
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let first_off = {
let table = open(source.clone(), &fs)?;
let offsets: Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|kh| *kh.as_ref().offset())
.collect();
assert_eq!(
offsets.len(),
2,
"the fixture needs the lost version alone in block 1 and the \
shadowed one sharing block 2",
);
usize::try_from(offsets.first().copied().unwrap_or_default()).unwrap_or(usize::MAX)
};
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(first_off + crate::table::block::Header::MIN_LEN + 1) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(report.dropped.len(), 1, "one block is lost: {report:?}");
let recovered = open(dest, &fs)?;
assert!(
recovered
.get(b"k", crate::SeqNo::MAX, crate::hash::hash64(b"k"))?
.is_none(),
"a verbatim copy would carry the suppressed key'\u{2019}s bytes through \
untouched, republishing the version its newer one had replaced",
);
assert!(
recovered
.get(b"m", crate::SeqNo::MAX, crate::hash::hash64(b"m"))?
.is_some(),
"the block'\u{2019}s other key is still recovered",
);
Ok(())
}
#[cfg(feature = "columnar")]
#[test]
fn salvage_drops_a_columnar_boundary_key_when_its_newest_version_is_lost() -> crate::Result<()> {
use crate::table::Writer;
use crate::{InternalValue, ValueType};
let dir = tempfile::tempdir()?;
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let source = dir.path().join("source");
let dest = dir.path().join("dest");
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?
.use_columnar(true)
.use_data_block_size(1);
writer.write(InternalValue::from_components(
b"a",
b"v",
1,
ValueType::Value,
))?;
writer.write(InternalValue::new_tombstone(b"k".as_slice(), 10))?;
writer.write(InternalValue::from_components(
b"k",
b"old",
5,
ValueType::Value,
))?;
writer.write(InternalValue::from_components(
b"z",
b"v",
1,
ValueType::Value,
))?;
assert!(writer.finish()?.is_some(), "source SST is non-empty");
let second_off = {
let table = open(source.clone(), &fs)?;
let offsets: Vec<u64> = table
.data_block_handles()
.filter_map(Result::ok)
.map(|kh| *kh.as_ref().offset())
.collect();
assert!(
offsets.len() >= 3,
"the fixture needs one block per entry, got {}",
offsets.len(),
);
usize::try_from(offsets.get(1).copied().unwrap_or_default()).unwrap_or(usize::MAX)
};
let mut bytes = std::fs::read(&source)?;
if let Some(b) = bytes.get_mut(second_off + crate::table::block::Header::MIN_LEN + 1) {
*b ^= 0xFF;
}
std::fs::write(&source, &bytes)?;
let report = salvage_sst(&source, dest.clone(), &fs)?;
assert_eq!(report.dropped.len(), 1, "one block is lost: {report:?}");
let recovered = open(dest, &fs)?;
assert!(
recovered
.get(b"k", crate::SeqNo::MAX, crate::hash::hash64(b"k"))?
.is_none(),
"the columnar branch must suppress the boundary key too, or salvage \
republishes a value the deletion had removed",
);
assert!(
recovered
.get(b"a", crate::SeqNo::MAX, crate::hash::hash64(b"a"))?
.is_some(),
"unaffected keys are still recovered",
);
Ok(())
}
#[test]
fn verify_metadata_bounds_keeps_a_lone_weak_tombstone_matching_the_rt_sentinel() -> crate::Result<()>
{
use crate::UserKey;
use crate::range_tombstone::RangeTombstone;
let dir = tempdir()?;
let source = dir.path().join("source");
let fs: Arc<dyn Fs> = Arc::new(StdFs);
let mut writer = Writer::new(source.clone(), 0, 0, Arc::clone(&fs))?;
writer.write(InternalValue::new_weak_tombstone(
UserKey::from(b"key-002".as_slice()),
3,
))?;
writer.write_range_tombstone(RangeTombstone::new(
UserKey::from(b"key-002".as_slice()),
UserKey::from(b"key-005".as_slice()),
3,
));
assert!(writer.finish()?.is_some(), "source SST is non-empty");
crate::test_forge::forge_meta_value_both_mirrors(&source, b"key#max", b"key-005")?;
crate::test_forge::forge_meta_value_both_mirrors(&source, b"seqno#min", &9u64.to_le_bytes())?;
let table = open(source, &fs)?;
let err = reconcile_error(&table, crate::table::ReconcileGate::MetadataBounds, None);
assert!(
matches!(
err,
crate::Error::InvalidHeader("meta seqno#min is above the decoded minimum seqno")
),
"the rejection names the seqno#min branch specifically, got {err:?}",
);
Ok(())
}