#![expect(
clippy::expect_used,
reason = "tests assert on known-present values; a panic is the failure signal"
)]
#![allow(
clippy::cast_possible_truncation,
reason = "in-file block offsets fit usize; only narrow on 32-bit targets"
)]
use super::*;
use crate::{
AbstractTree,
MAX_SEQNO,
SequenceNumberCounter,
runtime_config::EccScheme,
table::{block::Header, block_index::BlockIndex as _},
};
fn open_ecc_tree(dir: &std::path::Path) -> crate::Tree {
let crate::AnyTree::Standard(tree) = crate::Config::new(
dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.page_ecc(true)
.ecc_scheme(EccScheme::ReedSolomon {
data_shards: 8,
parity_shards: 2,
})
.open()
.expect("open ecc tree") else {
unreachable!("standard tree configured (no kv separation)");
};
tree
}
fn write_ecc_sst(dir: &std::path::Path) -> (std::path::PathBuf, crate::table::BlockHandle) {
let tree = open_ecc_tree(dir);
for i in 0u64..2_000 {
tree.insert(format!("key-{i:06}"), format!("v{i:06}"), i);
}
tree.flush_active_memtable(2_000).expect("flush");
let binding = tree.version_history.read().latest_version();
let table = binding
.version
.iter_tables()
.next()
.expect("flush produced one table");
let keyed = table
.block_index
.iter()
.next()
.expect("table has at least one data block")
.expect("block index entry decodes");
let handle = crate::table::BlockHandle::new(keyed.offset(), keyed.size());
((*table.path).clone(), handle)
}
fn write_single_block_ecc_sst(
dir: &std::path::Path,
) -> (std::path::PathBuf, crate::table::BlockHandle) {
let tree = open_ecc_tree(dir);
for i in 0u64..4 {
tree.insert(format!("key-{i:03}"), format!("v{i:03}"), i);
}
tree.flush_active_memtable(4).expect("flush");
let binding = tree.version_history.read().latest_version();
let table = binding
.version
.iter_tables()
.next()
.expect("flush produced one table");
let keyed = table
.block_index
.iter()
.next()
.expect("table has at least one data block")
.expect("block index entry decodes");
let handle = crate::table::BlockHandle::new(keyed.offset(), keyed.size());
((*table.path).clone(), handle)
}
fn write_ecc_sst_footered(
dir: &std::path::Path,
) -> (std::path::PathBuf, crate::table::BlockHandle) {
let tree = open_ecc_tree(dir);
tree.update_runtime_config(|c| {
c.kv_checksums = crate::runtime_config::KvChecksumPolicy::AllLevels;
})
.expect("enable kv checksums");
for i in 0u64..2_000 {
tree.insert(format!("key-{i:06}"), format!("v{i:06}"), i);
}
tree.flush_active_memtable(2_000).expect("flush");
let binding = tree.version_history.read().latest_version();
let table = binding
.version
.iter_tables()
.next()
.expect("flush produced one table");
let keyed = table
.block_index
.iter()
.next()
.expect("table has at least one data block")
.expect("block index entry decodes");
let handle = crate::table::BlockHandle::new(keyed.offset(), keyed.size());
((*table.path).clone(), handle)
}
fn write_ecc_sst_with_range_tombstone(
dir: &std::path::Path,
) -> (std::path::PathBuf, crate::table::BlockHandle) {
let tree = open_ecc_tree(dir);
for i in 0u64..2_000 {
tree.insert(format!("key-{i:06}"), format!("v{i:06}"), i);
}
tree.remove_range("key-000100", "key-000200", 2_000);
tree.flush_active_memtable(2_100).expect("flush");
let binding = tree.version_history.read().latest_version();
let table = binding
.version
.iter_tables()
.next()
.expect("flush produced one table");
let keyed = table
.block_index
.iter()
.next()
.expect("table has at least one data block")
.expect("block index entry decodes");
let handle = crate::table::BlockHandle::new(keyed.offset(), keyed.size());
((*table.path).clone(), handle)
}
#[test]
fn patrol_scrub_corrects_seeded_single_bit_fault_and_schedules_heal() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst(dir.path());
let corrupt_pos = block.offset().0 as usize + Header::MIN_LEN + 3;
let mut bytes = std::fs::read(&sst_path)?;
let slot = bytes
.get_mut(corrupt_pos)
.expect("corrupt_pos in range for the SST bytes");
*slot ^= 0x80;
std::fs::write(&sst_path, &bytes)?;
let tree = open_ecc_tree(dir.path());
tree.update_runtime_config(|c| c.auto_heal = true)?;
assert!(tree.heal_hints().is_empty(), "fresh tree has no hints");
let report = patrol_scrub(&tree, &PatrolScrubOptions::default());
assert!(
report.corrections_applied >= 1,
"scrub must correct the seeded fault: {report:?}",
);
assert_eq!(
report.ssts_scheduled_for_rewrite, 1,
"the corrected SST is queued for healing exactly once: {report:?}",
);
assert_eq!(report.uncorrectable_blocks, 0, "{report:?}");
assert!(
report.is_ok(),
"a fully-correctable scrub is ok: {report:?}"
);
assert!(
!tree.heal_hints().is_empty(),
"the SST is recorded in the heal queue",
);
#[cfg(feature = "metrics")]
assert_eq!(
tree.metrics().ecc_auto_heal_scheduled_count(),
1,
"the scheduled SST is counted once in metrics",
);
Ok(())
}
#[test]
fn patrol_scrub_corrects_without_scheduling_when_auto_heal_off() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst(dir.path());
let corrupt_pos = block.offset().0 as usize + Header::MIN_LEN + 3;
let mut bytes = std::fs::read(&sst_path)?;
let slot = bytes.get_mut(corrupt_pos).expect("corrupt_pos in range");
*slot ^= 0x80;
std::fs::write(&sst_path, &bytes)?;
let tree = open_ecc_tree(dir.path());
assert!(!tree.heal_hints().is_enabled(), "auto_heal defaults off");
let report = patrol_scrub(&tree, &PatrolScrubOptions::default());
assert!(
report.corrections_applied >= 1,
"correction-on-read still happens with auto_heal off: {report:?}",
);
assert_eq!(
report.ssts_scheduled_for_rewrite, 0,
"auto_heal off suppresses rewrite scheduling: {report:?}",
);
assert!(
tree.heal_hints().is_empty(),
"no SST queued when scheduling is off",
);
assert!(report.is_ok());
Ok(())
}
#[test]
fn patrol_scrub_progress_measures_physical_sizes() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _block) = write_ecc_sst(dir.path());
let physical = std::fs::metadata(&sst_path)?.len();
let tree = open_ecc_tree(dir.path());
let progress = std::sync::Arc::new(crate::RecoveryProgress::default());
let report = patrol_scrub(
&tree,
&PatrolScrubOptions {
progress: Some(std::sync::Arc::clone(&progress)),
..PatrolScrubOptions::default()
},
);
assert!(report.is_ok(), "{report:?}");
let snap = progress.snapshot();
assert_eq!(
snap.bytes_total, physical,
"the total is the SST's physical size, not the pre-section \
metadata figure: {snap:?}",
);
assert_eq!(
snap.bytes_processed, snap.bytes_total,
"a finished scrub reaches 100%: {snap:?}",
);
Ok(())
}
#[test]
fn patrol_scrub_progress_keeps_healed_within_recovered() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst(dir.path());
let corrupt_pos = block.offset().0 as usize + Header::MIN_LEN + 3;
let mut bytes = std::fs::read(&sst_path)?;
let slot = bytes.get_mut(corrupt_pos).expect("corrupt_pos in range");
*slot ^= 0x80;
std::fs::write(&sst_path, &bytes)?;
let tree = open_ecc_tree(dir.path());
let progress = std::sync::Arc::new(crate::RecoveryProgress::default());
let report = patrol_scrub(
&tree,
&PatrolScrubOptions {
progress: Some(std::sync::Arc::clone(&progress)),
..PatrolScrubOptions::default()
},
);
assert!(report.corrections_applied >= 1, "{report:?}");
let snap = progress.snapshot();
assert!(
snap.blocks_healed >= 1,
"the correction must be published: {snap:?}",
);
assert!(
snap.blocks_healed <= snap.blocks_recovered,
"healed blocks are a subset of recovered blocks: {snap:?}",
);
Ok(())
}
#[test]
fn patrol_scrub_reports_uncorrectable_block_not_silently_skipped() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst(dir.path());
let payload_start = block.offset().0 as usize + Header::MIN_LEN;
let payload_end = block.offset().0 as usize + block.size() as usize;
let mut bytes = std::fs::read(&sst_path)?;
for slot in bytes
.get_mut(payload_start..payload_end)
.expect("block payload range in bounds")
{
*slot ^= 0xFF;
}
std::fs::write(&sst_path, &bytes)?;
let tree = open_ecc_tree(dir.path());
tree.update_runtime_config(|c| c.auto_heal = true)?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default());
assert!(
report.uncorrectable_blocks >= 1,
"an unrecoverable block must be reported, not skipped: {report:?}",
);
assert!(!report.is_ok(), "uncorrectable corruption fails the scrub");
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::UncorrectableBlock { .. })),
"the finding is an UncorrectableBlock: {report:?}",
);
Ok(())
}
#[test]
fn patrol_scrub_clean_ecc_tree_reports_no_corrections() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let _ = write_ecc_sst(dir.path());
let tree = open_ecc_tree(dir.path());
let report = patrol_scrub(&tree, &PatrolScrubOptions::default());
assert_eq!(report.sst_files_scanned, 1);
assert!(report.blocks_scanned >= 1);
assert_eq!(report.corrections_applied, 0, "no fault → no correction");
assert_eq!(report.uncorrectable_blocks, 0);
assert!(report.is_ok());
let got = tree.get(b"key-000000", MAX_SEQNO)?.expect("key present");
assert_eq!(&*got, b"v000000");
Ok(())
}
#[test]
fn patrol_scrub_heals_in_place_restoring_the_block_byte_for_byte() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst(dir.path());
let original = std::fs::read(&sst_path)?;
let corrupt_pos = block.offset().0 as usize + Header::MIN_LEN + 3;
let mut bytes = original.clone();
let slot = bytes
.get_mut(corrupt_pos)
.expect("corrupt_pos in range for the SST bytes");
*slot ^= 0x80;
std::fs::write(&sst_path, &bytes)?;
assert_ne!(bytes, original, "the seeded fault changed the file");
let tree = open_ecc_tree(dir.path());
let opts = PatrolScrubOptions::default().heal_in_place(true);
let report = patrol_scrub(&tree, &opts);
assert_eq!(
report.blocks_healed_in_place, 1,
"exactly the corrupted block is healed in place: {report:?}",
);
assert_eq!(report.corrections_applied, 1, "{report:?}");
assert_eq!(
report.ssts_scheduled_for_rewrite, 0,
"in-place heal schedules no full-file rewrite: {report:?}",
);
assert_eq!(report.uncorrectable_blocks, 0, "{report:?}");
assert!(report.is_ok(), "{report:?}");
let healed = std::fs::read(&sst_path)?;
assert_eq!(
healed, original,
"in-place heal restores the SST byte-for-byte (O(damage), nothing else moved)",
);
drop(tree);
let tree2 = open_ecc_tree(dir.path());
let report2 = patrol_scrub(&tree2, &PatrolScrubOptions::default().heal_in_place(true));
assert_eq!(
report2.blocks_healed_in_place, 0,
"nothing left to heal after a clean heal: {report2:?}",
);
assert_eq!(report2.corrections_applied, 0, "{report2:?}");
Ok(())
}
#[test]
fn patrol_scrub_heal_in_place_leaves_an_uncorrectable_block_for_salvage() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst(dir.path());
let payload_start = block.offset().0 as usize + Header::MIN_LEN;
let payload_end = block.offset().0 as usize + block.size() as usize;
let mut bytes = std::fs::read(&sst_path)?;
for slot in bytes
.get_mut(payload_start..payload_end)
.expect("block payload range in bounds")
{
*slot ^= 0xFF;
}
std::fs::write(&sst_path, &bytes)?;
let corrupted = std::fs::read(&sst_path)?;
let tree = open_ecc_tree(dir.path());
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert_eq!(
report.blocks_healed_in_place, 0,
"an uncorrectable block is not healed in place: {report:?}",
);
assert!(
report.uncorrectable_blocks >= 1,
"the uncorrectable block is reported, not silently skipped: {report:?}",
);
assert!(
!report.is_ok(),
"uncorrectable corruption fails the heal pass"
);
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::UncorrectableBlock { .. })),
"the finding is an UncorrectableBlock: {report:?}",
);
let after = std::fs::read(&sst_path)?;
assert_eq!(
after, corrupted,
"an uncorrectable block is left untouched in place for salvage",
);
Ok(())
}
#[test]
fn patrol_scrub_heal_in_place_still_checks_a_non_ecc_table() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let sst_path;
let block_off;
{
let crate::AnyTree::Standard(tree) = crate::Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.open()
.expect("open plain tree") else {
unreachable!("standard tree configured (no kv separation)");
};
for i in 0u64..2_000 {
tree.insert(format!("key-{i:06}"), format!("v{i:06}"), i);
}
tree.flush_active_memtable(2_000).expect("flush");
let binding = tree.version_history.read().latest_version();
let table = binding
.version
.iter_tables()
.next()
.expect("flush produced one table");
let keyed = table
.block_index
.iter()
.next()
.expect("table has a data block")
.expect("index entry decodes");
sst_path = (*table.path).clone();
block_off = keyed.offset().0 as usize;
}
let mut bytes = std::fs::read(&sst_path)?;
let slot = bytes
.get_mut(block_off + Header::MIN_LEN + 3)
.expect("corrupt position in range");
*slot ^= 0x80;
std::fs::write(&sst_path, &bytes)?;
let crate::AnyTree::Standard(tree) = crate::Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.open()
.expect("reopen plain tree") else {
unreachable!("standard tree configured (no kv separation)");
};
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert_eq!(
report.blocks_healed_in_place, 0,
"a non-ECC table has nothing to heal in place: {report:?}",
);
assert!(
report.uncorrectable_blocks >= 1,
"a corrupt block in a non-ECC table is reported, not silently clean: {report:?}",
);
assert!(!report.is_ok(), "uncorrectable corruption fails the pass");
Ok(())
}
#[test]
fn heal_in_place_restores_a_rotted_parity_trailer() -> crate::Result<()> {
use crate::coding::Decode;
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst(dir.path());
let mut bytes = std::fs::read(&sst_path)?;
let base = block.offset().0 as usize;
let Some(mut cursor) = bytes.get(base..) else {
panic!("first data block within the file");
};
let header = Header::decode_from(&mut cursor)?;
let trailer_pos = base + 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");
};
let original = *slot;
*slot = original ^ 0xFF;
std::fs::write(&sst_path, &bytes)?;
let tree = open_ecc_tree(dir.path());
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report.blocks_healed_in_place >= 1,
"the rotted parity trailer is rebuilt and persisted: {report:?}",
);
assert_eq!(report.uncorrectable_blocks, 0, "{report:?}");
assert!(report.is_ok(), "a parity rebuild is a heal, not a finding");
let healed = std::fs::read(&sst_path)?;
let Some(&now) = healed.get(trailer_pos) else {
panic!("parity trailer within the healed file");
};
assert_eq!(now, original, "the original parity byte was restored");
Ok(())
}
fn open_ecc_tree_on(dir: &std::path::Path, fs: std::sync::Arc<dyn crate::fs::Fs>) -> crate::Tree {
let crate::AnyTree::Standard(tree) = crate::Config::new(
dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.page_ecc(true)
.ecc_scheme(EccScheme::ReedSolomon {
data_shards: 8,
parity_shards: 2,
})
.with_shared_fs(fs)
.open()
.expect("open ecc tree on injected fs") else {
unreachable!("standard tree configured (no kv separation)");
};
tree
}
fn corrupt_parity_trailer_byte(
path: &std::path::Path,
block: &crate::table::BlockHandle,
) -> crate::Result<()> {
use crate::coding::Decode;
let mut bytes = std::fs::read(path)?;
let base = block.offset().0 as usize;
let Some(mut cursor) = bytes.get(base..) else {
panic!("data block within the file");
};
let header = Header::decode_from(&mut cursor)?;
let trailer_pos = base + 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(path, &bytes)?;
Ok(())
}
fn rebuild_manifest_over_current_bytes(dir: &std::path::Path) -> crate::Result<()> {
use crate::version::{Level, Run, Version};
use std::sync::Arc;
let fs: Arc<dyn crate::fs::Fs> = Arc::new(crate::fs::StdFs);
let sst_path = dir.join("tables").join("0");
let checksum =
crate::Checksum::from_raw(crate::repair::compute_table_checksum(&*fs, &sst_path)?);
#[cfg(feature = "metrics")]
let metrics = Arc::new(crate::Metrics::default());
let table = {
#[cfg_attr(not(feature = "metrics"), expect(unused_mut))]
let mut params = crate::table::RecoverParams::new(
sst_path,
checksum,
0,
Arc::clone(&fs),
crate::comparator::default_comparator(),
Arc::new(crate::Cache::with_capacity_bytes(1_000_000)),
);
#[cfg(feature = "metrics")]
{
params.metrics = metrics;
}
crate::table::Table::recover(params)?
};
let mut next_version_id = 0u64;
for entry in fs.read_dir(dir)? {
if let Some(rest) = entry.file_name.strip_prefix('v')
&& let Ok(n) = rest.parse::<u64>()
{
next_version_id = next_version_id.max(n + 1);
}
}
let run = Run::new(alloc::vec![table]).expect("a non-empty run");
let mut levels = alloc::vec![Level::from_runs(alloc::vec![Arc::new(run)])];
for _ in 1..7 {
levels.push(Level::empty());
}
let version = Version::from_levels(
next_version_id,
crate::config::TreeType::Standard,
levels,
crate::version::BlobFileList::new(crate::HashMap::default()),
crate::blob_tree::FragmentationMap::default(),
);
crate::version::persist_version(
dir,
&version,
crate::comparator::default_comparator().name(),
&*fs,
Arc::new(crate::runtime_config::types::RuntimeConfig::default()),
None,
crate::fs::SyncMode::Full,
)?;
Ok(())
}
fn open_ecc_tree_with_failing_edit_log(
dir: &std::path::Path,
) -> (crate::Tree, std::sync::Arc<crate::fs::FaultInjector>) {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let fault = FaultFs::new(crate::fs::StdFs);
let injector = fault.injector();
let tree = open_ecc_tree_on(dir, std::sync::Arc::new(fault));
injector.arm(
FaultRule::new(FaultOp::Open, Fault::Error(ErrorKind::Other))
.on_path("edits")
.once(),
);
(tree, injector)
}
#[test]
fn heal_in_place_reports_a_failed_parity_reread_as_uncorrectable() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let (_sst_path, _block) = write_single_block_ecc_sst(dir.path());
let fault = FaultFs::new(crate::fs::StdFs);
let injector = fault.injector();
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(fault));
injector.arm(
FaultRule::new(FaultOp::ReadAt, Fault::Error(ErrorKind::Other))
.on_path("tables")
.skip(3)
.once(),
);
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert_eq!(
report.uncorrectable_blocks, 1,
"the unverifiable trailer is a finding: {report:?}",
);
assert!(
format!("{report:?}").contains("parity re-read failed"),
"the finding names the failed re-read: {report:?}",
);
assert_eq!(
report.blocks_healed_in_place, 0,
"nothing was persisted for the failed block: {report:?}",
);
Ok(())
}
#[test]
fn heal_in_place_reports_a_failed_trailer_writeback_as_uncorrectable() -> crate::Result<()> {
use crate::coding::Decode;
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst(dir.path());
let mut bytes = std::fs::read(&sst_path)?;
let base = block.offset().0 as usize;
let Some(mut cursor) = bytes.get(base..) else {
panic!("first data block within the file");
};
let header = Header::decode_from(&mut cursor)?;
let trailer_pos = base + 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");
};
let rotted = *slot ^ 0xFF;
*slot = rotted;
std::fs::write(&sst_path, &bytes)?;
let fault = FaultFs::new(crate::fs::StdFs);
let injector = fault.injector();
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(fault));
injector.arm(
FaultRule::new(FaultOp::Write, Fault::Error(ErrorKind::Other))
.on_path("tables")
.skip(1)
.once(),
);
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert_eq!(
report.blocks_healed_in_place, 0,
"a write-back that failed is not counted as a heal: {report:?}",
);
assert_eq!(report.uncorrectable_blocks, 1, "{report:?}");
assert!(
format!("{report:?}").contains("in-place parity rebuild"),
"the finding names the failed rebuild: {report:?}",
);
let after = std::fs::read(&sst_path)?;
assert_eq!(
after.get(trailer_pos).copied(),
Some(rotted),
"the rotted trailer byte is untouched after the failed write-back",
);
Ok(())
}
#[cfg(feature = "page_ecc")]
#[test]
fn heal_in_place_keeps_the_marker_when_a_write_back_fails() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst_footered(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let fault = FaultFs::new(crate::fs::StdFs);
let injector = fault.injector();
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(fault));
injector.arm(
FaultRule::new(FaultOp::Write, Fault::Error(ErrorKind::Other))
.on_path("tables")
.skip(1)
.once(),
);
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert_eq!(
report.blocks_healed_in_place, 0,
"a failed write-back heals no block: {report:?}",
);
assert!(
report.uncorrectable_blocks >= 1
&& format!("{report:?}").contains("in-place parity rebuild"),
"the failed trailer write-back must be recorded, proving the write was \
attempted: {report:?}",
);
assert!(
heal_attest_path(&sst_path).exists(),
"the marker must be KEPT after a failed write-back — the file may already be \
partially modified, so a later patrol still needs the attribution",
);
Ok(())
}
#[cfg(feature = "page_ecc")]
#[test]
fn heal_in_place_keeps_the_marker_when_a_restorative_write_back_fails() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst_footered(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
let fault = FaultFs::new(crate::fs::StdFs);
let injector = fault.injector();
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(fault));
injector.arm(
FaultRule::new(FaultOp::Write, Fault::Error(ErrorKind::Other))
.on_path("tables")
.skip(1)
.once(),
);
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert_eq!(
report.blocks_healed_in_place, 0,
"a failed write-back heals no block: {report:?}",
);
assert!(
heal_attest_path(&sst_path).exists(),
"the restorative heal must persist its marker BEFORE the first write-back, so a \
crash cannot expose unattested intermediate bytes to a checkpoint hard-link",
);
Ok(())
}
#[cfg(feature = "page_ecc")]
#[test]
fn heal_in_place_removes_the_marker_when_no_write_is_attempted() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst_footered(dir.path());
let start = block.offset().0 as usize + Header::MIN_LEN;
let end = block.offset().0 as usize + block.size() as usize;
let mut bytes = std::fs::read(&sst_path)?;
for off in start..end {
if let Some(b) = bytes.get_mut(off) {
*b ^= 0xFF;
}
}
std::fs::write(&sst_path, &bytes)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(crate::fs::StdFs));
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert_eq!(
report.blocks_healed_in_place, 0,
"an uncorrectable block heals nothing: {report:?}",
);
assert!(
report.uncorrectable_blocks >= 1,
"the uncorrectable block is recorded: {report:?}",
);
assert!(
!heal_attest_path(&sst_path).exists(),
"the marker must be removed when no write was attempted, so it cannot authorize a \
later unrelated mismatch",
);
Ok(())
}
#[cfg(feature = "page_ecc")]
#[test]
fn heal_in_place_keeps_the_marker_when_the_reconcile_walk_read_fails_transiently()
-> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let (sst_path, block) = {
let tree = open_ecc_tree(dir.path());
for i in 0u64..4 {
tree.insert(format!("key-{i:03}"), format!("v{i:03}"), i);
}
tree.flush_active_memtable(4).expect("flush");
let binding = tree.version_history.read().latest_version();
let table = binding
.version
.iter_tables()
.next()
.expect("flush produced one table");
let keyed = table
.block_index
.iter()
.next()
.expect("table has at least one data block")
.expect("block index entry decodes");
(
(*table.path).clone(),
crate::table::BlockHandle::new(keyed.offset(), keyed.size()),
)
};
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let fault = FaultFs::new(crate::fs::StdFs);
let injector = fault.injector();
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(fault));
injector.arm(
FaultRule::new(FaultOp::ReadAt, Fault::Error(ErrorKind::Other))
.on_path("tables")
.skip(4),
);
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert_eq!(
report.blocks_healed_in_place, 1,
"the trailer rebuild must land before the walk fault: {report:?}",
);
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"the transient walk read must refuse the digest refresh: {report:?}",
);
assert!(
heal_attest_path(&sst_path).exists(),
"a transient walk failure must KEEP the marker for retry, not strand the \
healed SST under the stale manifest digest",
);
Ok(())
}
#[cfg(feature = "page_ecc")]
#[test]
fn heal_in_place_propagates_a_transient_read_during_correction_prediction() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let (sst_path, block) = {
let tree = open_ecc_tree(dir.path());
for i in 0u64..4 {
tree.insert(format!("key-{i:03}"), format!("v{i:03}"), i);
}
tree.flush_active_memtable(4).expect("flush");
let binding = tree.version_history.read().latest_version();
let table = binding
.version
.iter_tables()
.next()
.expect("flush produced one table");
let keyed = table
.block_index
.iter()
.next()
.expect("table has at least one data block")
.expect("block index entry decodes");
(
(*table.path).clone(),
crate::table::BlockHandle::new(keyed.offset(), keyed.size()),
)
};
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let fault = FaultFs::new(crate::fs::StdFs);
let injector = fault.injector();
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(fault));
injector.arm(
FaultRule::new(FaultOp::ReadAt, Fault::Error(ErrorKind::Other))
.on_path("tables")
.once(),
);
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert!(
!report.is_ok(),
"a transient prediction read must surface as an error, not report a clean pass \
over the un-healed fault: {report:?}",
);
Ok(())
}
#[cfg(feature = "page_ecc")]
#[test]
fn heal_in_place_refuses_the_refresh_when_the_reconcile_attestation_write_fails()
-> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let (sst_path, block) = {
let tree = open_ecc_tree(dir.path());
for i in 0u64..4 {
tree.insert(format!("key-{i:03}"), format!("v{i:03}"), i);
}
tree.flush_active_memtable(4).expect("flush");
let binding = tree.version_history.read().latest_version();
let table = binding
.version
.iter_tables()
.next()
.expect("flush produced one table");
let keyed = table
.block_index
.iter()
.next()
.expect("table has at least one data block")
.expect("block index entry decodes");
(
(*table.path).clone(),
crate::table::BlockHandle::new(keyed.offset(), keyed.size()),
)
};
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let fault = FaultFs::new(crate::fs::StdFs);
let injector = fault.injector();
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(fault));
injector.arm(
FaultRule::new(FaultOp::Write, Fault::Error(ErrorKind::Other))
.on_path(".heal-attest")
.skip(1)
.once(),
);
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert_eq!(
report.blocks_healed_in_place, 1,
"the trailer rebuild lands before the reconcile: {report:?}",
);
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"a failed reconcile attestation write must refuse the digest refresh, not \
install it without a durable marker: {report:?}",
);
assert!(
heal_attest_path(&sst_path).exists(),
"the marker is kept for the next patrol to retry the reconcile",
);
Ok(())
}
#[cfg(feature = "page_ecc")]
#[test]
fn heal_skips_blocks_below_the_restriction_bound() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, first_block) = write_ecc_sst(dir.path());
let start = first_block.offset().0 as usize + Header::MIN_LEN;
let mut bytes = std::fs::read(&sst_path)?;
for off in start..start + 256 {
if let Some(b) = bytes.get_mut(off) {
*b ^= 0xFF;
}
}
std::fs::write(&sst_path, &bytes)?;
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(crate::fs::StdFs));
let binding = tree.version_history.read().latest_version();
let table = binding
.version
.iter_tables()
.next()
.expect("flush produced one table");
let restricted = table.reopen_restricted(crate::UserKey::from(b"key-001000".as_slice()))?;
let (report, _healed) =
restricted.heal_data_blocks_in_place(crate::fs::SyncMode::Full, restricted.checksum());
assert!(
report.errors.is_empty()
&& report.uncorrectable_blocks == 0
&& report.blocks_healed_in_place == 0,
"the wrecked block below the restriction bound must be skipped entirely, \
neither healed nor reported: {report:?}",
);
Ok(())
}
#[cfg(feature = "page_ecc")]
#[test]
fn every_published_table_carries_the_heal_hint_sink() -> crate::Result<()> {
use crate::AbstractTree;
let dir = tempfile::tempdir()?;
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(crate::fs::StdFs));
for round in 0..2u64 {
for i in 0..32u64 {
tree.insert(
format!("key-{i:06}").as_bytes(),
format!("v{round}").as_bytes(),
round * 100 + i,
);
}
tree.flush_active_memtable(0)?;
}
tree.major_compact(u64::MAX, 1_000)?;
let mut ingestion = crate::tree::ingest::Ingestion::new(&tree)?;
for i in 0..16u64 {
ingestion.write(
format!("ingested-{i:06}").as_bytes().into(),
format!("i{i}").as_bytes().into(),
)?;
}
ingestion.finish()?;
let live: Vec<_> = {
let binding = tree.version_history.read().latest_version();
binding.version.iter_tables().cloned().collect()
};
assert!(
live.len() >= 2,
"flush, compaction and ingest all published"
);
for table in &live {
assert!(
table
.heal_hints_for_test()
.is_some_and(|sink| std::sync::Arc::ptr_eq(&sink, &tree.heal_hints)),
"live table {} carries no heal-hint sink: a correctable fault in \
it could never schedule a durable heal",
table.id(),
);
}
Ok(())
}
#[cfg(feature = "page_ecc")]
#[test]
fn a_compaction_output_carries_the_heal_hint_sink() -> crate::Result<()> {
use crate::AbstractTree;
let dir = tempfile::tempdir()?;
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(crate::fs::StdFs));
for round in 0..2u64 {
for i in 0..64u64 {
tree.insert(
format!("key-{i:06}").as_bytes(),
format!("v{round}").as_bytes(),
round * 100 + i,
);
}
tree.flush_active_memtable(0)?;
}
tree.major_compact(u64::MAX, 1_000)?;
let outputs: Vec<_> = {
let binding = tree.version_history.read().latest_version();
binding.version.iter_tables().cloned().collect()
};
assert!(!outputs.is_empty(), "the compaction produced an output");
for table in &outputs {
assert!(
table
.heal_hints_for_test()
.is_some_and(|sink| { std::sync::Arc::ptr_eq(&sink, &tree.heal_hints) }),
"compaction output {} must carry the tree's heal-hint sink, or a \
correctable read from it can never schedule a durable heal",
table.id(),
);
}
Ok(())
}
#[cfg(feature = "page_ecc")]
#[test]
fn restricted_reopen_carries_the_heal_hint_sink_forward() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (_sst_path, _block) = write_ecc_sst(dir.path());
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(crate::fs::StdFs));
let table = {
let binding = tree.version_history.read().latest_version();
binding
.version
.iter_tables()
.next()
.expect("flush produced one table")
.clone()
};
table.install_heal_hints(crate::heal_hints::HealHints::new_shared(true));
let installed = table.heal_hints_for_test().expect("the sink is installed");
let restricted = table.reopen_restricted(crate::UserKey::from(b"key-001000".as_slice()))?;
let carried = restricted.heal_hints_for_test();
assert!(
carried.is_some_and(|c| std::sync::Arc::ptr_eq(&c, &installed)),
"the restricted reopen must carry the SAME heal-hint sink forward, or \
correctable reads from the restricted view stop queueing the table \
for a durable healing recompaction",
);
Ok(())
}
#[cfg(feature = "page_ecc")]
#[test]
fn heal_reports_a_contended_checksum_refresh_as_a_finding() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_single_block_ecc_sst(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(crate::fs::StdFs));
let _compaction_guard = tree.compaction_state.lock();
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert_eq!(
report.blocks_healed_in_place, 1,
"the heal itself lands; only the digest install is contended: {report:?}",
);
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"install-lock contention must surface as a finding, not a clean pass: {report:?}",
);
assert!(!report.is_ok(), "{report:?}");
assert!(
heal_attest_path(&sst_path).exists(),
"the marker is kept for the next patrol to reconcile",
);
Ok(())
}
#[cfg(feature = "page_ecc")]
#[test]
fn heal_scans_the_current_view_when_the_captured_one_went_stale() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _first_block) = write_ecc_sst(dir.path());
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(crate::fs::StdFs));
let captured = {
let binding = tree.version_history.read().latest_version();
binding
.version
.iter_tables()
.next()
.expect("flush produced one table")
.clone()
};
let restricted = captured.reopen_restricted(crate::UserKey::from(b"key-001000".as_slice()))?;
tree.version_history.write().upgrade_version(
&tree.config.path,
|current| {
let mut copy = current.clone();
let ctx = crate::version::TransformContext::new(tree.config.comparator.as_ref());
copy.version = copy.version.with_tight_slice(
&[(captured.id(), restricted.clone())],
&[],
&[],
vec![],
None,
0,
&ctx,
);
Ok(copy)
},
&tree.config.seqno,
&tree.config.visible_seqno,
&*tree.config.fs,
tree.runtime_config.load_full(),
tree.config.encryption.clone(),
crate::version::RetentionEffect::Keep,
)?;
let last_block = {
let mut last = None;
for handle in restricted.block_index.iter() {
let handle = handle?;
last = Some(crate::table::BlockHandle::new(
handle.offset(),
handle.size(),
));
}
last.expect("table has data blocks")
};
let corrupt_pos = last_block.offset().0 as usize + Header::MIN_LEN + 3;
let mut bytes = std::fs::read(&sst_path)?;
let Some(slot) = bytes.get_mut(corrupt_pos) else {
panic!("corrupt_pos in range for the SST bytes");
};
*slot ^= 0x80;
std::fs::write(&sst_path, &bytes)?;
let report = super::scan_and_reconcile(
&tree,
&captured,
&PatrolScrubOptions::default().heal_in_place(true),
);
assert!(
report.blocks_healed_in_place >= 1,
"the known correctable fault must be healed through the current view, \
not silently skipped by the divergent-heal guard: {report:?}",
);
assert!(report.is_ok(), "{report:?}");
Ok(())
}
#[test]
fn heal_in_place_reports_a_failed_heal_reread_as_uncorrectable() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_single_block_ecc_sst(dir.path());
let corrupt_pos = block.offset().0 as usize + Header::MIN_LEN + 3;
let mut bytes = std::fs::read(&sst_path)?;
let Some(slot) = bytes.get_mut(corrupt_pos) else {
panic!("corrupt_pos in range for the SST bytes");
};
*slot ^= 0x80;
std::fs::write(&sst_path, &bytes)?;
let fault = FaultFs::new(crate::fs::StdFs);
let injector = fault.injector();
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(fault));
injector.arm(
FaultRule::new(FaultOp::ReadAt, Fault::Error(ErrorKind::Other))
.on_path("tables")
.skip(3)
.once(),
);
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert_eq!(
report.blocks_healed_in_place, 0,
"nothing was persisted for the failed block: {report:?}",
);
assert_eq!(report.uncorrectable_blocks, 1, "{report:?}");
assert!(
format!("{report:?}").contains("heal re-read failed"),
"the finding names the failed heal re-read: {report:?}",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_mutate_a_hard_linked_checkpoint_inode() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst(dir.path());
let corrupt_pos = block.offset().0 as usize + Header::MIN_LEN + 3;
let mut bytes = std::fs::read(&sst_path)?;
let slot = bytes.get_mut(corrupt_pos).expect("corrupt_pos in range");
*slot ^= 0x80;
std::fs::write(&sst_path, &bytes)?;
let cp_dir = tempfile::tempdir_in(dir.path().parent().expect("tempdir has a parent"))?;
let link_path = cp_dir.path().join("checkpoint.sst");
std::fs::hard_link(&sst_path, &link_path)?;
let snapshot = std::fs::read(&link_path)?;
let tree = open_ecc_tree(dir.path());
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report.blocks_healed_in_place >= 1,
"the live file's fault is healed: {report:?}",
);
assert!(report.is_ok(), "{report:?}");
let live = std::fs::read(&sst_path)?;
assert_ne!(
live, snapshot,
"the live path must expose the healed bytes after the scrub",
);
let checkpoint = std::fs::read(&link_path)?;
assert_eq!(
checkpoint, snapshot,
"the checkpoint's hard-linked inode must keep its snapshot bytes: \
healing through a shared inode desynchronizes the checkpoint from \
its own manifest digest",
);
Ok(())
}
#[test]
fn heal_in_place_rebinds_the_descriptor_cache_after_an_unshare() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
let cp_dir = tempfile::tempdir_in(dir.path().parent().expect("tempdir has a parent"))?;
std::fs::hard_link(&sst_path, cp_dir.path().join("checkpoint.sst"))?;
let tree = open_ecc_tree(dir.path());
assert!(tree.get("key-000000", crate::SeqNo::MAX)?.is_some());
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(report.is_ok(), "the trailer rebuild succeeds: {report:?}");
assert!(
report.blocks_healed_in_place >= 1,
"the first pass must write, so the unshare runs: {report:?}",
);
let corrupt_pos = block.offset().0 as usize + Header::MIN_LEN + 3;
let mut bytes = std::fs::read(&sst_path)?;
let slot = bytes
.get_mut(corrupt_pos)
.expect("corrupt_pos in range for the SST bytes");
*slot ^= 0x80;
std::fs::write(&sst_path, &bytes)?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report.blocks_healed_in_place >= 1,
"the live fault is ECC-recovered and healed in place: {report:?}",
);
assert!(report.is_ok(), "{report:?}");
Ok(())
}
#[test]
fn heal_in_place_rebinds_the_descriptor_cache_when_the_directory_sync_fails() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
let cp_dir = tempfile::tempdir_in(dir.path().parent().expect("tempdir has a parent"))?;
std::fs::hard_link(&sst_path, cp_dir.path().join("checkpoint.sst"))?;
let fault = FaultFs::new(crate::fs::StdFs);
let injector = fault.injector();
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(fault));
assert!(tree.get("key-000000", crate::SeqNo::MAX)?.is_some());
injector.arm(
FaultRule::new(FaultOp::SyncDirectory, Fault::Error(ErrorKind::Other))
.on_path("tables")
.once(),
);
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert!(
!report.is_ok(),
"the failed unshare is a finding: {report:?}"
);
let corrupt_pos = block.offset().0 as usize + Header::MIN_LEN + 3;
let mut bytes = std::fs::read(&sst_path)?;
let slot = bytes
.get_mut(corrupt_pos)
.expect("corrupt_pos in range for the SST bytes");
*slot ^= 0x80;
std::fs::write(&sst_path, &bytes)?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report.blocks_healed_in_place >= 1,
"the live fault is ECC-recovered and healed in place: {report:?}",
);
assert!(report.is_ok(), "{report:?}");
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_a_renamed_section() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst(dir.path());
let crate::AnyTree::Standard(tree) = crate::Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.page_ecc(true)
.ecc_scheme(EccScheme::ReedSolomon {
data_shards: 8,
parity_shards: 2,
})
.filter_block_pinning_policy(crate::config::PinningPolicy::new([false]))
.open()
.expect("open ecc tree with lazy filters") else {
unreachable!("standard tree configured (no kv separation)");
};
crate::test_forge::forge_section_name(&sst_path, b"filter", b"filtex")?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"an unknown section name must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the renamed-section SST must keep failing verify_integrity: \
restamping its digest would legitimize the vanished section",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_diverged_meta_mirrors() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst(dir.path());
let tree = open_ecc_tree(dir.path());
crate::test_forge::forge_tail_meta_value(&sst_path, b"compression#data", &[1])?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"diverged meta mirrors must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: restamping its \
digest would mask a forge only the mirror comparison detects",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_a_forged_zone_map() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let sst_path = {
let tree = open_ecc_tree(dir.path());
tree.update_runtime_config(|c| c.zone_map = true)?;
for i in 0u64..2_000 {
tree.insert(format!("key-{i:06}"), format!("v{i:06}"), i);
}
tree.flush_active_memtable(2_000).expect("flush");
let binding = tree.version_history.read().latest_version();
let table = binding
.version
.iter_tables()
.next()
.expect("flush produced one table");
(*table.path).clone()
};
let tree = open_ecc_tree(dir.path());
crate::test_forge::forge_flip_section_last_payload_byte(&sst_path, b"zone_map", Some((8, 2)))?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"a forged zone_map must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: restamping its \
digest would let predicate scans silently skip matching blocks",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_altered_range_tombstones() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst_with_range_tombstone(dir.path());
let tree = open_ecc_tree(dir.path());
crate::test_forge::forge_flip_section_last_payload_byte(
&sst_path,
b"range_tombstones",
Some((8, 2)),
)?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"altered range tombstones must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the altered SST must keep failing verify_integrity: restamping its \
digest would let reads resurrect deleted data or hide live data",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_a_forged_footerless_value() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst(dir.path());
let tree = open_ecc_tree(dir.path());
crate::test_forge::forge_value_byte_in_first_data_block(&sst_path, Some((8, 2)))?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"a forged footer-less value must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: restamping its \
digest would erase the only record of the original value bytes",
);
Ok(())
}
fn write_ecc_zstd_multiblock_sst(dir: &std::path::Path) -> std::path::PathBuf {
let crate::AnyTree::Standard(tree) = crate::Config::new(
dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.page_ecc(true)
.ecc_scheme(EccScheme::ReedSolomon {
data_shards: 8,
parity_shards: 2,
})
.data_block_size_policy(crate::config::BlockSizePolicy::all(256 * 1024))
.data_block_compression_policy(crate::config::CompressionPolicy::all(
crate::CompressionType::Zstd(19),
))
.open()
.expect("open ecc zstd tree") else {
unreachable!("standard tree configured (no kv separation)");
};
tree.update_runtime_config(|c| {
c.kv_checksums = crate::runtime_config::KvChecksumPolicy::AllLevels;
})
.expect("enable kv checksums");
for i in 0u64..20_000 {
tree.insert(format!("key-{i:012}"), format!("value-{i:08}-payload"), i);
}
tree.flush_active_memtable(20_000).expect("flush");
let binding = tree.version_history.read().latest_version();
let table = binding
.version
.iter_tables()
.next()
.expect("flush produced one table");
assert!(
table.regions.block_layout.is_some(),
"the multi-inner-block fixture must carry a block_layout section",
);
(*table.path).clone()
}
#[test]
fn salvage_load_block_re_encodes_a_multi_inner_block() -> crate::Result<()> {
use crate::table::BlockHandle;
use crate::table::block::BlockType;
let dir = tempfile::tempdir()?;
let crate::AnyTree::Standard(tree) = crate::Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.page_ecc(true)
.ecc_scheme(EccScheme::ReedSolomon {
data_shards: 8,
parity_shards: 2,
})
.data_block_size_policy(crate::config::BlockSizePolicy::all(256 * 1024))
.data_block_compression_policy(crate::config::CompressionPolicy::all(
crate::CompressionType::Zstd(19),
))
.open()
.expect("open ecc zstd tree") else {
unreachable!("standard tree configured (no kv separation)");
};
tree.update_runtime_config(|c| {
c.kv_checksums = crate::runtime_config::KvChecksumPolicy::AllLevels;
})
.expect("enable kv checksums");
for i in 0u64..20_000 {
tree.insert(format!("key-{i:012}"), format!("value-{i:08}-payload"), i);
}
tree.flush_active_memtable(20_000).expect("flush");
let binding = tree.version_history.read().latest_version();
let table = binding
.version
.iter_tables()
.next()
.expect("flush produced one table");
let multi_inner_offsets = table.block_layout.offsets();
assert!(
!multi_inner_offsets.is_empty(),
"the fixture must carry at least one multi-inner-block frame",
);
let multi = table
.block_index
.iter()
.filter_map(Result::ok)
.find(|kh| multi_inner_offsets.contains(&kh.offset().0))
.map(|kh| BlockHandle::new(kh.offset(), kh.size()))
.expect("a block index handle for the multi-inner offset");
let sb = table.salvage_load_block(&multi, BlockType::Data)?;
assert!(
sb.verbatim.is_none(),
"a multi-inner block must re-encode from the verified payload, not byte-copy its \
unauthenticated recorded layout",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_a_forged_block_layout() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let sst_path = write_ecc_zstd_multiblock_sst(dir.path());
let crate::AnyTree::Standard(tree) = crate::Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.page_ecc(true)
.ecc_scheme(EccScheme::ReedSolomon {
data_shards: 8,
parity_shards: 2,
})
.data_block_size_policy(crate::config::BlockSizePolicy::all(256 * 1024))
.data_block_compression_policy(crate::config::CompressionPolicy::all(
crate::CompressionType::Zstd(19),
))
.open()
.expect("reopen ecc zstd tree") else {
unreachable!("standard tree configured (no kv separation)");
};
crate::test_forge::forge_block_layout_shift_middle_end(&sst_path, Some((8, 2)))?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"a forged block_layout must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: restamping its \
digest would let partial range reads silently omit keys",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_consistently_forged_tli_mirrors() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst_footered(dir.path());
let tree = open_ecc_tree(dir.path());
crate::test_forge::forge_tli_mirrors_truncated(
&sst_path,
0,
Some(crate::table::block::EccParams::try_new(8, 2)?),
)?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"consistently forged TLI mirrors must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: restamping its \
digest would let the next recovery hide the dropped block",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_forged_meta_key_bounds() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst_footered(dir.path());
let tree = open_ecc_tree(dir.path());
crate::test_forge::forge_meta_value_both_mirrors(&sst_path, b"key#max", b"key-000999")?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"forged meta key bounds must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: restamping its \
digest would let run selection silently skip real keys",
);
Ok(())
}
fn open_ecc_hashed_tree(dir: &std::path::Path) -> crate::Tree {
let crate::AnyTree::Standard(tree) = crate::Config::new(
dir,
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.page_ecc(true)
.ecc_scheme(EccScheme::ReedSolomon {
data_shards: 8,
parity_shards: 2,
})
.data_block_compression_policy(crate::config::CompressionPolicy::all(
crate::CompressionType::None,
))
.data_block_hash_ratio_policy(crate::config::HashRatioPolicy::all(2.0))
.open()
.expect("open ecc hashed tree") else {
unreachable!("standard tree configured (no kv separation)");
};
tree
}
fn write_ecc_sst_footered_hashed(dir: &std::path::Path) -> std::path::PathBuf {
let tree = open_ecc_hashed_tree(dir);
tree.update_runtime_config(|c| {
c.kv_checksums = crate::runtime_config::KvChecksumPolicy::AllLevels;
})
.expect("enable kv checksums");
for i in 0u64..2_000 {
tree.insert(format!("key-{i:06}"), format!("v{i:06}"), i);
}
tree.flush_active_memtable(2_000).expect("flush");
let binding = tree.version_history.read().latest_version();
let table = binding
.version
.iter_tables()
.next()
.expect("flush produced one table");
(*table.path).clone()
}
#[test]
fn heal_in_place_does_not_restamp_over_a_forged_hash_index() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let sst_path = write_ecc_sst_footered_hashed(dir.path());
let tree = open_ecc_hashed_tree(dir.path());
crate::test_forge::forge_hash_index_all_free(&sst_path, Some((8, 2)))?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"a forged hash index must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: restamping its \
digest would let point reads miss existing keys",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_trust_cached_blocks_over_the_disk_bytes() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let sst_path = write_ecc_sst_footered_hashed(dir.path());
let tree = open_ecc_hashed_tree(dir.path());
assert!(
tree.get("key-000000", MAX_SEQNO)?.is_some(),
"pre-forge read warms the cache",
);
crate::test_forge::forge_hash_index_all_free(&sst_path, Some((8, 2)))?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"a forged hash index must refuse the digest refresh even when the \
pristine block is cached: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: restamping its \
digest would let point reads miss existing keys after the cache cools",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_forged_meta_seqno_bounds() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst_footered(dir.path());
let tree = open_ecc_tree(dir.path());
let recorded_max = {
let binding = tree.version_history.read().latest_version();
let table = binding.version.iter_tables().next().expect("one table");
table.get_highest_seqno()
};
crate::test_forge::forge_meta_value_both_mirrors(
&sst_path,
b"seqno#min",
&recorded_max.to_le_bytes(),
)?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"forged metadata seqno bounds must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: restamping its \
digest would let snapshot reads silently skip older visible versions",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_a_forged_created_at() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst_footered(dir.path());
let tree = open_ecc_tree(dir.path());
crate::test_forge::forge_meta_value_both_mirrors(
&sst_path,
b"created_at",
&1u128.to_le_bytes(),
)?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"a forged created_at must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: restamping its \
digest would let FIFO compaction drop the live SST as TTL-expired",
);
Ok(())
}
#[test]
fn heal_in_place_rejects_a_created_at_restamped_before_open() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst_footered(dir.path());
crate::test_forge::forge_meta_value_both_mirrors(
&sst_path,
b"created_at",
&1u128.to_le_bytes(),
)?;
let tree = open_ecc_tree(dir.path());
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report.errors.iter().any(|e| matches!(
e,
ScrubError::ChecksumRefreshFailed { reason, .. }
if reason.contains("not attributable to this pass's heal")
)),
"an unattributed mismatch on a footer-bearing table must not reconcile a \
pre-open created_at restamp: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: restamping its digest \
would let FIFO compaction drop the live SST as TTL-expired",
);
Ok(())
}
#[cfg(feature = "page_ecc")]
#[test]
fn heal_in_place_ignores_a_stale_in_progress_marker() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst_footered(dir.path());
crate::test_forge::forge_meta_value_both_mirrors(
&sst_path,
b"created_at",
&1u128.to_le_bytes(),
)?;
let tree = open_ecc_tree(dir.path());
let manifest_checksum = {
let binding = tree.version_history.read().latest_version();
binding
.version
.iter_tables()
.next()
.expect("flush produced one table")
.checksum()
};
crate::scrub::heal_attest::write_in_progress(
&crate::fs::StdFs,
&sst_path,
None,
0,
manifest_checksum,
)?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report.errors.iter().any(|e| matches!(
e,
ScrubError::ChecksumRefreshFailed { reason, .. }
if reason.contains("not attributable to this pass's heal")
)),
"a stale in-progress marker must not make an unattributed mismatch \
attributable: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: a stale in-progress \
marker must not authorize restamping its digest",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_a_forged_kv_checksum_descriptor() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst_footered(dir.path());
let tree = open_ecc_tree(dir.path());
crate::test_forge::forge_meta_value_both_mirrors(&sst_path, b"descriptor#kv_checksum", &[0])?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"a forged kv-checksum descriptor must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: restamping its \
digest would let point reads misread footer bytes as the trailer",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_a_forged_data_block_count() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst_footered(dir.path());
let tree = open_ecc_tree(dir.path());
crate::test_forge::forge_meta_value_both_mirrors(
&sst_path,
b"block_count#data",
&1u64.to_le_bytes(),
)?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"a forged data-block count must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: restamping its \
digest would let compaction scans silently drop the omitted blocks",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_a_forged_tli_separator() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst_footered(dir.path());
let tree = open_ecc_tree(dir.path());
crate::test_forge::forge_tli_mirrors_lower_first_separator(
&sst_path,
0,
Some(crate::table::block::EccParams::try_new(8, 2)?),
)?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"a forged TLI separator must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: restamping its \
digest would let point reads miss keys routed to the wrong block",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_a_forged_tli_binary_index() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst_footered(dir.path());
let tree = open_ecc_tree(dir.path());
crate::test_forge::forge_tli_binary_index_pointer(
&sst_path,
0,
Some(crate::table::block::EccParams::try_new(8, 2)?),
)?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"a forged TLI binary index must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: restamping its \
digest would let index seeks start at the wrong restart head",
);
Ok(())
}
#[test]
fn heal_in_place_reconciles_a_tombstone_bearing_table_after_a_legit_heal() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst_with_range_tombstone(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let tree = open_ecc_tree(dir.path());
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report.blocks_healed_in_place >= 1,
"the rotted trailer is rebuilt in place: {report:?}",
);
assert!(
report.is_ok(),
"an attributable heal reconciles the digest despite the deletion \
metadata: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
integrity.is_ok(),
"the healed table verifies clean against the refreshed digest, got {:?}",
integrity.errors,
);
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_a_forged_filter() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst(dir.path());
let tree = open_ecc_tree(dir.path());
crate::test_forge::forge_filter_false_negative(
&sst_path,
crate::hash::hash64(b"key-000000"),
Some((8, 2)),
)?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"a forged filter must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: restamping its \
digest would let point reads silently miss existing keys",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_a_forged_seqno_bounds() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let sst_path = {
let tree = open_ecc_tree(dir.path());
tree.update_runtime_config(|c| c.seqno_in_index = true)?;
for i in 0u64..2_000 {
tree.insert(format!("key-{i:06}"), format!("v{i:06}"), i);
}
tree.flush_active_memtable(2_000).expect("flush");
let binding = tree.version_history.read().latest_version();
let table = binding
.version
.iter_tables()
.next()
.expect("flush produced one table");
(*table.path).clone()
};
let tree = open_ecc_tree(dir.path());
crate::test_forge::forge_seqno_bounds_zeroed_entry(&sst_path, Some((8, 2)))?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"a forged seqno_bounds map must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: restamping its \
digest would let scans silently skip live blocks",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_a_forged_tli_tail() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst(dir.path());
let tree = open_ecc_tree(dir.path());
crate::test_forge::forge_tli_tail_truncated(
&sst_path,
0,
Some(crate::table::block::EccParams::try_new(8, 2)?),
)?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"diverged TLI mirrors must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the forged SST must keep failing verify_integrity: restamping its \
digest would hide a mirror only the decoded comparison detects",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_a_relabeled_section_block() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst(dir.path());
crate::test_forge::forge_section_block_role(
&sst_path,
b"filter",
crate::table::block::BlockType::Data,
)?;
let crate::AnyTree::Standard(tree) = crate::Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.page_ecc(true)
.ecc_scheme(EccScheme::ReedSolomon {
data_shards: 8,
parity_shards: 2,
})
.filter_block_pinning_policy(crate::config::PinningPolicy::new([false]))
.open()
.expect("open ecc tree with lazy filters") else {
unreachable!("standard tree configured (no kv separation)");
};
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"the relabeled block must refuse the digest refresh: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the relabeled SST must keep failing verify_integrity: restamping \
its digest would mask a forge only the role cross-check detects",
);
Ok(())
}
#[cfg(unix)]
#[test]
fn heal_in_place_keeps_a_healthy_sst_hard_linked() -> crate::Result<()> {
use std::os::unix::fs::MetadataExt as _;
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst(dir.path());
let cp_dir = tempfile::tempdir_in(dir.path().parent().expect("tempdir has a parent"))?;
std::fs::hard_link(&sst_path, cp_dir.path().join("checkpoint.sst"))?;
let tree = open_ecc_tree(dir.path());
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(report.is_ok(), "{report:?}");
assert_eq!(report.blocks_healed_in_place, 0, "nothing to heal");
assert_eq!(
std::fs::metadata(&sst_path)?.nlink(),
2,
"a clean scan must leave the checkpoint link in place: detaching \
without a write to protect it from costs a full-file copy and \
doubles the SST's disk usage",
);
Ok(())
}
#[test]
fn heal_in_place_treats_an_unknown_link_count_as_shared() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
let fault = FaultFs::new(crate::fs::StdFs);
let injector = fault.injector();
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(fault));
injector.arm(
FaultRule::new(FaultOp::HardLinkCount, Fault::Error(ErrorKind::Other))
.on_path("tables")
.once(),
);
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert!(
report.blocks_healed_in_place >= 1,
"fail-closed still heals (through the detached copy): {report:?}",
);
assert!(report.is_ok(), "{report:?}");
Ok(())
}
#[test]
fn heal_in_place_cleans_up_the_temp_copy_when_the_unshare_fails() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
let cp_dir = tempfile::tempdir_in(dir.path().parent().expect("tempdir has a parent"))?;
std::fs::hard_link(&sst_path, cp_dir.path().join("checkpoint.sst"))?;
let fault = FaultFs::new(crate::fs::StdFs);
let injector = fault.injector();
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(fault));
injector.arm(
FaultRule::new(FaultOp::SyncAll, Fault::Error(ErrorKind::Other))
.on_path("healtmp")
.once(),
);
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert!(
!report.is_ok(),
"the failed unshare is a finding: {report:?}"
);
let leftovers: Vec<_> = std::fs::read_dir(sst_path.parent().expect("sst in tables dir"))?
.filter_map(Result::ok)
.filter(|e| e.file_name().to_string_lossy().contains("healtmp"))
.collect();
assert!(
leftovers.is_empty(),
"a failed unshare must remove its temp copy: {leftovers:?}",
);
drop(tree);
let tree = open_ecc_tree(dir.path());
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the refused heal leaves the rot in place and visible",
);
Ok(())
}
#[test]
fn heal_in_place_does_not_restamp_over_side_section_rot() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _) = write_ecc_sst(dir.path());
let (pos, len) = {
let mut f = std::fs::File::open(&sst_path)?;
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") else {
panic!("the SST must carry a filter section");
};
(entry.pos(), entry.len())
};
assert!(len > 128, "filter section large enough to rot: {len}");
let start = usize::try_from(pos).expect("filter offset fits usize") + 40;
let mut bytes = std::fs::read(&sst_path)?;
let Some(run) = bytes.get_mut(start..start + 64) else {
panic!("filter payload within the file");
};
for b in run {
*b ^= 0xFF;
}
std::fs::write(&sst_path, &bytes)?;
let tree = open_ecc_tree(dir.path());
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
!report.is_ok(),
"a digest mismatch the scan cannot attribute to a heal must be a \
finding, not silently restamped: {report:?}",
);
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"the finding must be the refused digest refresh: {report:?}",
);
assert_eq!(
report.uncorrectable_blocks, 0,
"the data blocks themselves stay clean: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the corrupted file must keep failing verify_integrity: restamping \
its digest over unverified side sections would mask the rot",
);
Ok(())
}
fn heal_attest_path(sst_path: &std::path::Path) -> std::path::PathBuf {
let mut name = sst_path.as_os_str().to_os_string();
name.push(".heal-attest");
std::path::PathBuf::from(name)
}
#[test]
fn heal_in_place_reconciles_a_crashed_refresh_via_the_attestation() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst_footered(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let (tree, injector) = open_ecc_tree_with_failing_edit_log(dir.path());
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert!(report.blocks_healed_in_place >= 1, "{report:?}");
assert!(
!report.is_ok(),
"the failed refresh is a finding: {report:?}"
);
assert!(
heal_attest_path(&sst_path).exists(),
"the heal must leave an attestation for the crashed refresh",
);
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report.is_ok(),
"a crashed refresh must reconcile via the attestation on the next scrub: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
integrity.is_ok(),
"the reconciled digest matches the healed file, got {:?}",
integrity.errors,
);
assert!(
!heal_attest_path(&sst_path).exists(),
"the attestation is consumed once its reconciliation lands",
);
Ok(())
}
#[test]
fn heal_in_place_reconciles_via_a_completed_marker() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst_footered(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let (tree, injector) = open_ecc_tree_with_failing_edit_log(dir.path());
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert!(report.blocks_healed_in_place >= 1, "{report:?}");
let (table_id, manifest_digest) = {
let binding = tree.version_history.read().latest_version();
let table = binding
.version
.iter_tables()
.next()
.expect("the healed table is still in the manifest");
(table.id(), table.checksum())
};
let healed_digest = crate::Checksum::from_raw(crate::repair::compute_table_checksum(
&crate::fs::StdFs,
&sst_path,
)?);
std::fs::remove_file(heal_attest_path(&sst_path))?;
crate::scrub::heal_attest::write(
&crate::fs::StdFs,
&sst_path,
None,
table_id,
manifest_digest,
healed_digest,
)?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert_eq!(
report.blocks_healed_in_place, 0,
"the second pass must find every block clean, so ONLY the marker attributes the \
mismatch (a re-heal would make attribution direct and skip the marker path): {report:?}",
);
assert!(
report.is_ok(),
"a crash before the manifest refresh must reconcile via the completed \
marker on the next scrub: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
integrity.is_ok(),
"the reconciled digest matches the healed file, got {:?}",
integrity.errors,
);
assert!(
!heal_attest_path(&sst_path).exists(),
"the marker is consumed once its reconciliation lands",
);
Ok(())
}
#[test]
fn attests_post_matches_only_the_recorded_completed_post() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let path = dir.path().join("t.sst");
std::fs::write(&path, b"x")?;
let pre = crate::Checksum::from_raw(0x1111);
let post = crate::Checksum::from_raw(0x2222);
crate::scrub::heal_attest::write(&crate::fs::StdFs, &path, None, 7, pre, post)?;
use crate::scrub::heal_attest::AttestResult;
assert!(matches!(
crate::scrub::heal_attest::attests_post(&crate::fs::StdFs, &path, None, 7, post),
AttestResult::Attests,
));
assert!(matches!(
crate::scrub::heal_attest::attests_post(
&crate::fs::StdFs,
&path,
None,
7,
crate::Checksum::from_raw(0x3333),
),
AttestResult::Absent,
));
assert!(matches!(
crate::scrub::heal_attest::attests_post(&crate::fs::StdFs, &path, None, 8, post),
AttestResult::Absent,
));
Ok(())
}
#[test]
fn attests_post_is_inconclusive_on_a_transient_sidecar_read() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
use crate::scrub::heal_attest::AttestResult;
let dir = tempfile::tempdir()?;
let path = dir.path().join("t.sst");
std::fs::write(&path, b"x")?;
let post = crate::Checksum::from_raw(0x2222);
crate::scrub::heal_attest::write(
&crate::fs::StdFs,
&path,
None,
7,
crate::Checksum::from_raw(0x1111),
post,
)?;
let fault = FaultFs::new(crate::fs::StdFs);
fault.injector().arm(
FaultRule::new(FaultOp::Open, Fault::Error(ErrorKind::Interrupted)).on_path("heal-attest"),
);
assert!(
matches!(
crate::scrub::heal_attest::attests_post(&fault, &path, None, 7, post),
AttestResult::Inconclusive,
),
"a transient sidecar read must be Inconclusive, not Absent",
);
Ok(())
}
#[test]
fn heal_in_place_fails_closed_without_a_heal_attestation() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst_footered(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let (tree, injector) = open_ecc_tree_with_failing_edit_log(dir.path());
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert!(report.blocks_healed_in_place >= 1, "{report:?}");
std::fs::remove_file(heal_attest_path(&sst_path))?;
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report.errors.iter().any(|e| matches!(
e,
ScrubError::ChecksumRefreshFailed { reason, .. }
if reason.contains("not attributable to this pass's heal")
)),
"without an attestation the stale digest must fail closed: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"the unattested stale digest keeps flagging, got {:?}",
integrity.errors,
);
Ok(())
}
#[test]
fn reconcile_keeps_the_marker_when_the_sidecar_read_is_inconclusive() -> crate::Result<()> {
use crate::fs::{Fault, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst_footered(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let (tree, injector) = open_ecc_tree_with_failing_edit_log(dir.path());
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert!(report.blocks_healed_in_place >= 1, "{report:?}");
assert!(
heal_attest_path(&sst_path).exists(),
"the heal left a marker"
);
injector
.arm(FaultRule::new(FaultOp::Open, Fault::Error(ErrorKind::Other)).on_path(".heal-attest"));
let _ = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert!(
heal_attest_path(&sst_path).exists(),
"an inconclusive sidecar read must keep the marker, not delete it",
);
Ok(())
}
#[test]
fn recovery_sweeps_an_orphaned_heal_attestation_but_keeps_a_live_one() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _block) = write_ecc_sst(dir.path());
let tables_dir = sst_path.parent().expect("the SST lives in a tables folder");
let orphan = tables_dir.join("99999.heal-attest");
std::fs::write(&orphan, b"orphaned attestation")?;
let live = heal_attest_path(&sst_path);
std::fs::write(&live, b"live attestation")?;
let _tree = open_ecc_tree(dir.path());
assert!(
!orphan.exists(),
"recovery must sweep a sidecar whose table id is absent from the manifest",
);
assert!(
live.exists(),
"recovery must preserve a live table's pending attestation",
);
Ok(())
}
#[test]
fn reconcile_reclaims_an_obsolete_marker_when_the_digest_already_matches() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _block) = write_ecc_sst(dir.path());
let marker = heal_attest_path(&sst_path);
std::fs::write(&marker, b"obsolete marker")?;
assert!(marker.exists());
let tree = open_ecc_tree(dir.path());
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report.is_ok(),
"a clean ECC table scrubs cleanly: {report:?}"
);
assert!(
!marker.exists(),
"an obsolete marker must be reclaimed once the digest already matches the manifest",
);
Ok(())
}
#[test]
fn abort_checkpoint_ignores_an_obsolete_marker_but_aborts_on_a_stale_digest() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, _block) = write_ecc_sst(dir.path());
let marker = heal_attest_path(&sst_path);
std::fs::write(&marker, b"obsolete marker")?;
let tree = open_ecc_tree(dir.path());
crate::scrub::abort_checkpoint_if_pending_heals(&tree, "obsolete-marker case")?;
assert!(
!marker.exists(),
"an obsolete marker matching the manifest must be reclaimed, not wedge the checkpoint",
);
std::fs::write(&marker, b"pending marker")?;
let mut bytes = std::fs::read(&sst_path)?;
if let Some(b) = bytes.get_mut(64) {
*b ^= 0xFF;
}
std::fs::write(&sst_path, &bytes)?;
let err = crate::scrub::abort_checkpoint_if_pending_heals(&tree, "stale-digest case")
.expect_err("a pending heal whose digest does not match the manifest must abort");
assert!(
matches!(err, crate::Error::Io(_)),
"the abort is surfaced as an Io error, got {err:?}",
);
assert!(
marker.exists(),
"a genuine pending marker is kept for the next reconciliation",
);
Ok(())
}
#[cfg(feature = "page_ecc")]
#[test]
fn checkpoint_flush_does_not_deadlock_with_patrol_and_compaction() -> crate::Result<()> {
use crate::AbstractTree;
use std::sync::Arc;
use std::time::Duration;
let dir = tempfile::tempdir()?;
write_ecc_sst(dir.path());
let tree = Arc::new(open_ecc_tree_on(dir.path(), Arc::new(crate::fs::StdFs)));
tree.insert("pending-row", "v", 5_000);
let table = {
let binding = tree.version_history.read().latest_version();
binding
.version
.iter_tables()
.next()
.expect("flush produced one table")
.clone()
};
let state_guard = tree.compaction_state.lock();
let checkpoint = {
let tree = Arc::clone(&tree);
let dst = dir.path().join("checkpoint");
std::thread::spawn(move || tree.create_checkpoint(&dst))
};
std::thread::sleep(Duration::from_millis(300));
let patrol = {
let tree = Arc::clone(&tree);
std::thread::spawn(move || {
patrol_scrub(&*tree, &PatrolScrubOptions::default().heal_in_place(true))
})
};
std::thread::sleep(Duration::from_millis(300));
drop(table.heal_lock_arc().lock());
drop(state_guard);
let report = patrol.join().expect("patrol thread must not panic");
assert!(report.is_ok(), "{report:?}");
checkpoint
.join()
.expect("checkpoint thread must not panic")?;
Ok(())
}
#[cfg(feature = "page_ecc")]
#[test]
fn checkpoint_reconciles_a_pending_heal_before_snapshotting() -> crate::Result<()> {
use crate::AbstractTree;
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst_footered(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let (tree, injector) = open_ecc_tree_with_failing_edit_log(dir.path());
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert!(report.blocks_healed_in_place >= 1, "{report:?}");
assert!(
heal_attest_path(&sst_path).exists(),
"the failed refresh must leave a pending attestation",
);
let dst = dir.path().join("checkpoint");
tree.create_checkpoint(&dst)?;
let checkpoint = open_ecc_tree(&dst);
let integrity = crate::verify::verify_integrity(&checkpoint);
assert!(
integrity.is_ok(),
"the checkpoint must not capture a table's stale pre-heal digest, got {:?}",
integrity.errors,
);
Ok(())
}
mod exists_fail_fs {
use crate::fs::{Fs, FsDirEntry, FsFile, FsMetadata, FsOpenOptions, StdFs};
use crate::io;
use std::path::Path;
pub(super) struct ExistsFailFs;
impl Fs for ExistsFailFs {
fn open(&self, path: &Path, opts: &FsOpenOptions) -> io::Result<Box<dyn FsFile>> {
StdFs.open(path, opts)
}
fn create_dir_all(&self, path: &Path) -> io::Result<()> {
StdFs.create_dir_all(path)
}
fn read_dir(&self, path: &Path) -> io::Result<Vec<FsDirEntry>> {
StdFs.read_dir(path)
}
fn remove_file(&self, path: &Path) -> io::Result<()> {
StdFs.remove_file(path)
}
fn remove_dir_all(&self, path: &Path) -> io::Result<()> {
StdFs.remove_dir_all(path)
}
fn rename(&self, from: &Path, to: &Path) -> io::Result<()> {
StdFs.rename(from, to)
}
fn metadata(&self, path: &Path) -> io::Result<FsMetadata> {
StdFs.metadata(path)
}
fn sync_directory(&self, path: &Path) -> io::Result<()> {
StdFs.sync_directory(path)
}
fn exists(&self, path: &Path) -> io::Result<bool> {
if path.to_string_lossy().ends_with(".heal-attest") {
return Err(io::Error::other("injected attestation probe failure"));
}
StdFs.exists(path)
}
fn backend_id(&self) -> Option<u64> {
StdFs.backend_id()
}
fn volume_id(&self, path: &Path) -> Option<u64> {
StdFs.volume_id(path)
}
}
}
#[test]
fn reconcile_pending_heals_aborts_when_the_attestation_probe_fails() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let _ = write_ecc_sst_footered(dir.path());
let tree = open_ecc_tree_on(
dir.path(),
std::sync::Arc::new(exists_fail_fs::ExistsFailFs),
);
let result = crate::scrub::reconcile_pending_heals(&tree);
assert!(
result.is_err(),
"a failed attestation probe must abort the reconciliation, not silently \
skip the table (pre-fix the probe error was swallowed as 'no pending heal'): \
{result:?}",
);
Ok(())
}
#[test]
fn heal_in_place_aborts_when_the_attestation_cannot_be_persisted() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst_footered(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let before = std::fs::read(&sst_path)?;
let fault = FaultFs::new(crate::fs::StdFs);
let injector = fault.injector();
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(fault));
injector
.arm(FaultRule::new(FaultOp::Open, Fault::Error(ErrorKind::Other)).on_path(".heal-attest"));
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert_eq!(
report.blocks_healed_in_place, 0,
"the heal must not mutate any block when its attestation cannot be persisted: {report:?}",
);
let after = std::fs::read(&sst_path)?;
assert_eq!(
after, before,
"the corrupt block must stay untouched so the table stays reconcilable",
);
assert!(
!report.is_ok(),
"aborting the heal must surface a finding: {report:?}",
);
assert!(
!heal_attest_path(&sst_path).exists(),
"no partial attestation marker after a failed write",
);
Ok(())
}
#[test]
fn heal_reconcile_refuses_to_refresh_when_the_sst_data_cannot_be_synced() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst_footered(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let manifest_digest = |tree: &crate::Tree| {
let binding = tree.version_history.read().latest_version();
binding
.version
.iter_tables()
.next()
.map(crate::table::Table::checksum)
.expect("the table is in the manifest")
};
let fault = FaultFs::new(crate::fs::StdFs);
let injector = fault.injector();
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(fault));
let digest_before = manifest_digest(&tree);
injector
.arm(FaultRule::new(FaultOp::SyncData, Fault::Error(ErrorKind::Other)).on_path("tables"));
let _ = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
heal_attest_path(&sst_path).exists(),
"a write-attempted heal keeps its attestation for a later patrol",
);
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert_eq!(
manifest_digest(&tree),
digest_before,
"the manifest digest must not be refreshed over bytes that were never \
synced: a power loss would discard the healed block: {report:?}",
);
assert!(
heal_attest_path(&sst_path).exists(),
"the marker survives so a later, syncable patrol can still reconcile",
);
Ok(())
}
#[test]
fn checkpoint_aborts_when_a_pending_marker_survives_the_pre_window_reconcile() -> crate::Result<()>
{
use crate::AbstractTree;
let dir = tempfile::tempdir()?;
let crate::AnyTree::Standard(tree) = crate::Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.open()?
else {
unreachable!("standard tree configured");
};
for i in 0u64..4 {
tree.insert(format!("key-{i:03}"), format!("v{i:03}"), i);
}
tree.flush_active_memtable(4)?;
let sst_path = {
let binding = tree.version_history.read().latest_version();
let table = binding
.version
.iter_tables()
.next()
.expect("flush produced one table");
(*table.path).clone()
};
let mut bytes = std::fs::read(&sst_path)?;
if let Some(b) = bytes.get_mut(50) {
*b ^= 0xFF;
}
std::fs::write(&sst_path, &bytes)?;
std::fs::write(heal_attest_path(&sst_path), b"pending marker")?;
let dst = dir.path().join("checkpoint");
let result = tree.create_checkpoint(&dst);
assert!(
result.is_err(),
"a pending marker with a stale digest the pre-window reconcile cannot consume must \
abort the checkpoint at the post-link-window guard: {result:?}",
);
Ok(())
}
#[test]
fn heal_in_place_aborts_when_the_pre_heal_digest_probe_fails() -> crate::Result<()> {
use crate::fs::{Fault, FaultFs, FaultOp, FaultRule};
use crate::io::ErrorKind;
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst_footered(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let before = std::fs::read(&sst_path)?;
let fault = FaultFs::new(crate::fs::StdFs);
let injector = fault.injector();
let tree = open_ecc_tree_on(dir.path(), std::sync::Arc::new(fault));
injector.arm(FaultRule::new(FaultOp::Read, Fault::Error(ErrorKind::Other)).on_path("tables"));
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert_eq!(
report.blocks_healed_in_place, 0,
"a failed pre-heal digest probe must abort before any block is mutated: {report:?}",
);
let after = std::fs::read(&sst_path)?;
assert_eq!(
after, before,
"the corrupt block must stay untouched so the table stays reconcilable",
);
assert!(
!report.is_ok(),
"aborting the heal must surface a finding: {report:?}",
);
Ok(())
}
#[test]
fn heal_in_place_refreshes_the_manifest_checksum() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let tree = open_ecc_tree(dir.path());
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report.blocks_healed_in_place >= 1,
"the rotted trailer is rebuilt in place: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
integrity.is_ok(),
"a healed table must verify clean against a refreshed manifest \
checksum, got {:?}",
integrity.errors,
);
drop(tree);
let tree = open_ecc_tree(dir.path());
let integrity = crate::verify::verify_integrity(&tree);
assert!(
integrity.is_ok(),
"the refreshed checksum is durable across reopen, got {:?}",
integrity.errors,
);
Ok(())
}
#[test]
fn checksum_refresh_is_idempotent_against_a_stale_table_view() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let tree = open_ecc_tree(dir.path());
let stale_view = {
let binding = tree.version_history.read().latest_version();
binding
.version
.iter_tables()
.next()
.expect("one table")
.clone()
};
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report.blocks_healed_in_place >= 1,
"the rotted trailer is rebuilt in place: {report:?}",
);
assert!(report.is_ok(), "the first patrol reconciles: {report:?}");
let finding = refresh_healed_checksum(&tree, &stale_view, false);
assert!(
finding.is_none(),
"an already-reconciled file must not be reported: {finding:?}",
);
Ok(())
}
#[test]
fn heal_in_place_refuses_to_detach_under_a_pinned_descriptor() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let link = dir.path().join("checkpoint-link");
std::fs::hard_link(&sst_path, &link)?;
let crate::AnyTree::Standard(tree) = crate::Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.page_ecc(true)
.ecc_scheme(EccScheme::ReedSolomon {
data_shards: 8,
parity_shards: 2,
})
.use_descriptor_table(None)
.open()
.expect("open pinned ecc tree") else {
unreachable!("standard tree configured (no kv separation)");
};
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report.errors.iter().any(|e| matches!(
e,
ScrubError::UncorrectableBlock { reason, .. }
if reason.starts_with("unshare hard-linked SST for heal: ")
)),
"a refused detach must surface as an unshare finding: {report:?}",
);
assert_eq!(
report.blocks_healed_in_place, 0,
"no block may be healed when the detach is refused: {report:?}",
);
assert_eq!(
std::fs::read(&sst_path)?,
std::fs::read(&link)?,
"the live path must stay byte-identical to its checkpoint link (no detach)",
);
Ok(())
}
#[test]
fn heal_attribution_compares_against_the_current_manifest_digest() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let stale_view = {
let tree = open_ecc_tree(dir.path());
let stale = {
let binding = tree.version_history.read().latest_version();
binding
.version
.iter_tables()
.next()
.expect("one table")
.clone()
};
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(report.is_ok(), "the first heal reconciles: {report:?}");
stale
};
let second_block = {
let keyed = stale_view
.block_index
.iter()
.nth(1)
.expect("the fixture writes several data blocks")
.expect("block index entry decodes");
crate::table::BlockHandle::new(keyed.offset(), keyed.size())
};
corrupt_parity_trailer_byte(&sst_path, &second_block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let tree = open_ecc_tree(dir.path());
let report = scan_and_reconcile(
&tree,
&stale_view,
&PatrolScrubOptions::default().heal_in_place(true),
);
assert!(
report.blocks_healed_in_place >= 1,
"the second fault is rebuilt in place: {report:?}",
);
assert!(
report.is_ok(),
"a heal whose pre-write digest matches the CURRENT manifest is \
attributable and must reconcile: {report:?}",
);
Ok(())
}
#[test]
fn heal_in_place_reports_a_failed_checksum_refresh() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, block) = write_ecc_sst(dir.path());
corrupt_parity_trailer_byte(&sst_path, &block)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let (tree, injector) = open_ecc_tree_with_failing_edit_log(dir.path());
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
injector.clear();
assert!(
report.blocks_healed_in_place >= 1,
"the rotted trailer is rebuilt in place: {report:?}",
);
assert!(
report
.errors
.iter()
.any(|e| matches!(e, ScrubError::ChecksumRefreshFailed { .. })),
"a failed manifest-digest refresh must be a scrub finding, not a \
swallowed log line: {report:?}",
);
assert!(
!report.is_ok(),
"a scrub whose findings include a failed checksum refresh is not ok",
);
Ok(())
}
#[test]
fn heal_in_place_skips_the_checksum_refresh_while_corruption_remains() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let (sst_path, first) = write_ecc_sst(dir.path());
let second = {
let tree = open_ecc_tree(dir.path());
let binding = tree.version_history.read().latest_version();
let table = binding
.version
.iter_tables()
.next()
.expect("one table recovered");
let mut it = table.block_index.iter();
let _ = it.next().expect("first block").expect("decodes");
let keyed = it.next().expect("second block").expect("decodes");
crate::table::BlockHandle::new(keyed.offset(), keyed.size())
};
corrupt_parity_trailer_byte(&sst_path, &first)?;
let mut bytes = std::fs::read(&sst_path)?;
let payload_start = second.offset().0 as usize + Header::MIN_LEN;
let payload_end = second.offset().0 as usize + second.size() as usize;
for slot in bytes
.get_mut(payload_start..payload_end)
.expect("second block payload range in bounds")
{
*slot ^= 0xFF;
}
std::fs::write(&sst_path, &bytes)?;
rebuild_manifest_over_current_bytes(dir.path())?;
let tree = open_ecc_tree(dir.path());
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report.blocks_healed_in_place >= 1,
"the rotted trailer is rebuilt in place: {report:?}",
);
assert!(
report.uncorrectable_blocks >= 1,
"the wrecked block is reported uncorrectable: {report:?}",
);
let integrity = crate::verify::verify_integrity(&tree);
assert!(
!integrity.is_ok(),
"an SST with a known uncorrectable block must keep failing \
verify_integrity — restamping its manifest checksum would mask the \
corruption",
);
Ok(())
}
#[cfg(all(feature = "columnar", feature = "encryption", zstd_any))]
#[test]
fn heal_in_place_leaves_a_clean_encrypted_columnar_sst_with_no_findings() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let enc: std::sync::Arc<dyn crate::encryption::EncryptionProvider> =
std::sync::Arc::new(crate::Aes256GcmProvider::new(&[0x51; 32]));
let crate::AnyTree::Standard(tree) = crate::Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.page_ecc(true)
.ecc_scheme(EccScheme::ReedSolomon {
data_shards: 8,
parity_shards: 2,
})
.with_encryption(Some(enc))
.open()
.expect("open encrypted ecc tree") else {
unreachable!("standard tree configured (no kv separation)");
};
tree.update_runtime_config(|cfg| cfg.columnar = true)?;
for i in 0u64..2_000 {
tree.insert(format!("key-{i:06}"), format!("v{i:06}"), i);
}
tree.flush_active_memtable(2_000).expect("flush");
{
let binding = tree.version_history.read().latest_version();
let table = binding
.version
.iter_tables()
.next()
.expect("flush produced one table");
assert!(table.metadata.columnar, "the SST is columnar");
assert!(table.metadata.ecc_params.is_some(), "the SST carries ECC");
assert!(table.encryption.is_some(), "the SST is encrypted");
}
let report = patrol_scrub(&tree, &PatrolScrubOptions::default().heal_in_place(true));
assert!(
report.blocks_scanned >= 1,
"the columnar SST has at least one data block to scrub: {report:?}",
);
assert_eq!(
report.uncorrectable_blocks, 0,
"a clean encrypted columnar block must decrypt and verify cleanly, \
not be reported uncorrectable: {report:?}",
);
assert!(
report.is_ok(),
"a clean encrypted columnar SST heals with no findings: {report:?}",
);
Ok(())
}
#[test]
fn tree_open_skips_a_lingering_heal_attestation() -> crate::Result<()> {
let dir = tempfile::tempdir()?;
let sst_path = {
let crate::AnyTree::Standard(tree) = crate::Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.open()?
else {
unreachable!("standard tree configured");
};
for i in 0u64..50 {
tree.insert(format!("k{i:04}"), b"v", i);
}
tree.flush_active_memtable(50)?;
let binding = tree.version_history.read().latest_version();
let Some(table) = binding.version.iter_tables().next() else {
panic!("flush produced one table");
};
(*table.path).clone()
};
std::fs::write(heal_attest_path(&sst_path), b"pending attestation")?;
let reopened = crate::Config::new(
dir.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.open();
assert!(
reopened.is_ok(),
"a lingering heal-attest sidecar must not break open: {:?}",
reopened.err(),
);
assert!(
heal_attest_path(&sst_path).exists(),
"open must leave the pending attestation for the next scrub to consume",
);
Ok(())
}