use crate::UserKey;
use crate::encryption::EncryptionProvider;
use alloc::string::String;
use alloc::sync::Arc;
use alloc::vec::Vec;
use std::path::PathBuf;
#[derive(Clone, Default)]
pub struct SalvageOptions {
pub encryption: Option<Arc<dyn EncryptionProvider>>,
#[cfg(zstd_any)]
pub zstd_dictionary: Option<Arc<crate::compression::ZstdDictionary>>,
pub table_id: crate::TableId,
pub expected_stored_id: Option<crate::TableId>,
pub output_id: Option<crate::TableId>,
pub allow_delete_resurrection: bool,
pub sync_mode: crate::fs::SyncMode,
pub prefix_extractor: Option<Arc<dyn crate::prefix::PrefixExtractor>>,
pub blob_rewrite: Option<Arc<crate::HashMap<crate::vlog::BlobFileId, BlobFileRewrite>>>,
pub progress: Option<Arc<crate::RecoveryProgress>>,
}
#[derive(Debug, Clone)]
pub enum BlobFileRewrite {
Remap {
new_id: crate::vlog::BlobFileId,
offsets: crate::HashMap<u64, BlobRecordRelocation>,
},
DropBelow(u64),
}
#[derive(Debug, Clone)]
pub enum DropReason {
HeaderCorrupted(String),
ChecksumMismatch,
ReadError(String),
DecodeError(String),
}
#[derive(Debug, Clone)]
pub struct DroppedBlock {
pub offset: u64,
pub section: Vec<u8>,
pub reason: DropReason,
pub key_range: Option<(UserKey, UserKey)>,
}
#[derive(Debug)]
pub struct SalvageReport {
pub salvaged_path: Option<PathBuf>,
pub blocks_total: usize,
pub blocks_salvaged: usize,
pub blocks_copied_verbatim: usize,
pub entries_salvaged: u64,
pub entries_dropped_by_rewrite: u64,
pub dropped: Vec<DroppedBlock>,
pub delete_rows_resurrected: bool,
}
impl SalvageReport {
#[must_use]
pub fn is_complete(&self) -> bool {
self.dropped.is_empty()
}
}
pub fn salvage_sst(
source: &std::path::Path,
dest: std::path::PathBuf,
fs: &alloc::sync::Arc<dyn crate::fs::Fs>,
) -> crate::Result<SalvageReport> {
salvage_sst_with_options(source, dest, fs, &SalvageOptions::default())
}
pub fn salvage_sst_with_options(
source: &std::path::Path,
dest: std::path::PathBuf,
fs: &alloc::sync::Arc<dyn crate::fs::Fs>,
options: &SalvageOptions,
) -> crate::Result<SalvageReport> {
salvage_with_context(
source,
dest,
fs,
&crate::comparator::default_comparator(),
options,
)
}
pub(crate) fn salvage_with_context(
source: &std::path::Path,
dest: std::path::PathBuf,
fs: &alloc::sync::Arc<dyn crate::fs::Fs>,
comparator: &crate::comparator::SharedComparator,
options: &SalvageOptions,
) -> crate::Result<SalvageReport> {
let diverged = meta_mirrors_diverge(source, fs, options)?;
if diverged == MirrorDivergence::Agree {
return salvage_attempt(source, dest, fs, comparator, options, false, true);
}
if diverged == MirrorDivergence::NonDerivable {
log::error!(
"{}: metadata mirrors disagree in fields no entry derives \
(bulk-ingest provenance, L0 recency, or compaction lineage); \
neither copy can be authenticated, so salvage will not choose",
source.display(),
);
return Err(crate::Error::Unrecoverable);
}
let tail_dest = next_arb_temp(fs, &dest)?;
let tail = salvage_attempt(
source,
tail_dest.clone(),
fs,
comparator,
options,
false,
false,
);
if let Ok(r) = &tail
&& r.blocks_total > 0
&& r.dropped.is_empty()
{
return publish_from_temp(fs, tail, &tail_dest, &dest, options);
}
let mid_dest = next_arb_temp(fs, &dest)?;
let mid = salvage_attempt(
source,
mid_dest.clone(),
fs,
comparator,
options,
true,
false,
);
match arbitrate_mirrors(&tail, &mid) {
MirrorArbitration::Propagate => {
if attempt_owns_temp(&tail) {
discard_partial(fs, &tail_dest);
}
if attempt_owns_temp(&mid) {
discard_partial(fs, &mid_dest);
}
if matches!(&tail, Err(e) if e.is_environmental()) {
tail
} else {
mid
}
}
MirrorArbitration::PublishMid => {
if attempt_owns_temp(&tail) {
discard_partial(fs, &tail_dest);
}
publish_from_temp(fs, mid, &mid_dest, &dest, options)
}
MirrorArbitration::PublishTail => {
if attempt_owns_temp(&mid) {
discard_partial(fs, &mid_dest);
}
publish_from_temp(fs, tail, &tail_dest, &dest, options)
}
}
}
fn attempt_owns_temp(result: &crate::Result<SalvageReport>) -> bool {
matches!(result, Ok(r) if r.salvaged_path.is_some())
}
#[derive(Debug, PartialEq, Eq)]
enum MirrorArbitration {
PublishTail,
PublishMid,
Propagate,
}
fn arbitrate_mirrors(
tail: &crate::Result<SalvageReport>,
mid: &crate::Result<SalvageReport>,
) -> MirrorArbitration {
let is_retryable =
|r: &crate::Result<SalvageReport>| matches!(r, Err(e) if e.is_environmental());
let complete = |r: &crate::Result<SalvageReport>| matches!(r, Ok(rep) if rep.is_complete());
if (is_retryable(tail) && !complete(mid)) || (is_retryable(mid) && !complete(tail)) {
return MirrorArbitration::Propagate;
}
let completeness = |rep: &SalvageReport| {
(
rep.blocks_salvaged,
rep.entries_salvaged,
core::cmp::Reverse(rep.dropped.len()),
)
};
match (tail, mid) {
(Err(_), Ok(_)) => MirrorArbitration::PublishMid,
(_, Err(_)) => MirrorArbitration::PublishTail,
(Ok(t), Ok(m)) => {
if completeness(m) > completeness(t) {
MirrorArbitration::PublishMid
} else {
MirrorArbitration::PublishTail
}
}
}
}
fn next_arb_temp(
fs: &alloc::sync::Arc<dyn crate::fs::Fs>,
dest: &std::path::Path,
) -> crate::Result<std::path::PathBuf> {
static ARB_TMP_SEQ: core::sync::atomic::AtomicU64 = core::sync::atomic::AtomicU64::new(0);
let pid_hi = u64::from(std::process::id()) << 32;
loop {
let counter = ARB_TMP_SEQ.fetch_add(1, core::sync::atomic::Ordering::Relaxed);
let seq = pid_hi | (counter & 0xFFFF_FFFF);
let candidate = dest.with_extension(alloc::format!("healtmp-{seq}"));
match fs.exists(&candidate) {
Ok(false) => return Ok(candidate),
Ok(true) => {}
Err(e) => return Err(e.into()),
}
}
}
fn publish_from_temp(
fs: &alloc::sync::Arc<dyn crate::fs::Fs>,
result: crate::Result<SalvageReport>,
temp: &std::path::Path,
dest: &std::path::Path,
options: &SalvageOptions,
) -> crate::Result<SalvageReport> {
let already_exists = || {
crate::Error::Io(crate::io::Error::new(
crate::io::ErrorKind::AlreadyExists,
"salvage destination already exists",
))
};
let mut rep = match result {
Ok(rep) => rep,
Err(e) => {
return Err(e);
}
};
if rep.salvaged_path.is_none() {
return Ok(rep);
}
match fs.hard_link(temp, dest) {
Ok(()) => {}
Err(e) if e.kind() == crate::io::ErrorKind::AlreadyExists => {
discard_partial(fs, temp);
return Err(already_exists());
}
Err(e) if e.kind() == crate::io::ErrorKind::Unsupported => {
match fs.open(
dest,
&crate::fs::FsOpenOptions::new().write(true).create_new(true),
) {
Ok(claim) => {
drop(claim);
if let Err(e) = fs.rename(temp, dest) {
discard_partial(fs, dest);
discard_partial(fs, temp);
return Err(e.into());
}
}
Err(e) if e.kind() == crate::io::ErrorKind::AlreadyExists => {
discard_partial(fs, temp);
return Err(already_exists());
}
Err(e) => {
discard_partial(fs, temp);
return Err(e.into());
}
}
}
Err(e) => {
discard_partial(fs, temp);
return Err(e.into());
}
}
match fs.remove_file(temp) {
Ok(()) => {}
Err(e) if e.kind() == crate::io::ErrorKind::NotFound => {}
Err(e) => {
discard_partial(fs, dest);
return Err(e.into());
}
}
if let Err(e) = fs.sync_directory_with(entry_directory(dest), options.sync_mode) {
discard_partial(fs, dest);
return Err(e.into());
}
rep.salvaged_path = Some(dest.to_path_buf());
Ok(rep)
}
#[derive(Debug, PartialEq, Eq)]
enum MirrorDivergence {
Agree,
Derivable,
NonDerivable,
}
fn meta_mirrors_diverge(
source: &std::path::Path,
fs: &alloc::sync::Arc<dyn crate::fs::Fs>,
options: &SalvageOptions,
) -> crate::Result<MirrorDivergence> {
let mut file = match fs.open(source, &crate::fs::FsOpenOptions::new().read(true)) {
Ok(f) => f,
Err(e) => return Err(crate::Error::Io(e)),
};
let trailer = match crate::sfa::Reader::from_reader(&mut file) {
Ok(t) => t,
Err(crate::sfa::Error::Io(e)) => return Err(crate::Error::Io(e)),
Err(_) => return Ok(MirrorDivergence::Agree),
};
let Ok(regions) = crate::table::regions::ParsedRegions::parse_from_toc(trailer.toc()) else {
return Ok(MirrorDivergence::Agree);
};
let Some(mid_handle) = regions.metadata_mid else {
return Ok(MirrorDivergence::Agree);
};
let expected_id = if options.encryption.is_some() {
Some(options.table_id)
} else {
options.expected_stored_id
};
let tail = crate::table::meta::ParsedMeta::load_with_handle(
&*file,
®ions.metadata,
expected_id,
options.encryption.as_deref(),
);
let mid = crate::table::meta::ParsedMeta::load_with_handle(
&*file,
&mid_handle,
expected_id,
options.encryption.as_deref(),
);
match (tail, mid) {
(Ok(t), Ok(m)) => {
if t == m {
return Ok(MirrorDivergence::Agree);
}
let non_derivable_disagree = t.bulk_ingested != m.bulk_ingested
|| t.recency != m.recency
|| t.lineage != m.lineage
|| t.lineage_prev != m.lineage_prev
|| t.lineage_transformed != m.lineage_transformed
|| t.lineage_last != m.lineage_last;
Ok(if non_derivable_disagree {
MirrorDivergence::NonDerivable
} else {
MirrorDivergence::Derivable
})
}
(Err(crate::Error::Io(io)), _) | (_, Err(crate::Error::Io(io)))
if io.kind().is_environmental() =>
{
Err(crate::Error::Io(io))
}
_ => Ok(MirrorDivergence::Agree),
}
}
fn salvage_attempt(
source: &std::path::Path,
dest: std::path::PathBuf,
fs: &alloc::sync::Arc<dyn crate::fs::Fs>,
comparator: &crate::comparator::SharedComparator,
options: &SalvageOptions,
prefer_mid_meta: bool,
allow_verbatim: bool,
) -> crate::Result<SalvageReport> {
let checksum = match crate::repair::compute_table_checksum(&**fs, source) {
Ok(c) => crate::Checksum::from_raw(c),
Err(_) => crate::Checksum::from_raw(0),
};
let table = {
let mut params = crate::table::RecoverParams::new(
source.to_path_buf(),
checksum,
options.table_id,
Arc::clone(fs),
comparator.clone(),
Arc::new(crate::cache::Cache::with_capacity_bytes(8 * 1024 * 1024)),
);
params.descriptor_table = Some(Arc::new(crate::descriptor_table::DescriptorTable::new(64)));
params.encryption.clone_from(&options.encryption);
#[cfg(zstd_any)]
{
params.zstd_dictionary.clone_from(&options.zstd_dictionary);
}
crate::table::Table::recover_inner(
params,
crate::table::RecoveryMode::Salvage {
expected_id: options.expected_stored_id,
prefer_mid_meta,
},
)?
};
if !table.range_tombstones().is_empty() || table.metadata.range_tombstone_count > 0 {
return Err(crate::Error::FeatureUnsupported(
"salvage of an SST with range tombstones",
));
}
#[cfg(feature = "columnar")]
let has_visible_deletion = table.has_delete_bitmap_section();
#[cfg(not(feature = "columnar"))]
let has_visible_deletion = false;
if !has_visible_deletion && table.salvage_degraded_a_rebuildable_section() {
return Err(crate::Error::FeatureUnsupported(
"salvage of an SST with a degraded rebuildable section that may hide \
a relabeled deletion",
));
}
if !has_visible_deletion && crate::repair::toc_may_hide_deletions(fs, source)? {
return Err(crate::Error::FeatureUnsupported(
"salvage of an SST whose TOC may hide a deletion section \
(an omitted, renamed, or shadowed entry)",
));
}
#[cfg(feature = "columnar")]
let delete_mask_unpositionable = table.delete_bitmap_degraded
|| (table.has_delete_bitmap_section() && table.delete_bitmap().is_empty())
|| (!table.delete_bitmap().is_empty() && !table.delete_positions_verified()?)
|| !table.delete_bitmap_authenticated();
#[cfg(not(feature = "columnar"))]
let delete_mask_unpositionable = table.delete_bitmap_degraded;
if delete_mask_unpositionable && !options.allow_delete_resurrection {
return Err(crate::Error::InvalidHeader(
"salvage: the delete bitmap cannot be applied; recovering would resurrect deleted \
rows (opt in with allow_delete_resurrection)",
));
}
let writer = crate::table::Writer::new(
dest.clone(),
options.output_id.unwrap_or(table.metadata.id),
0,
Arc::clone(fs),
)?
.mirror_from(
&table.metadata,
table.has_zone_map(),
table.has_seqno_bounds(),
)
.use_recency(Some(table.l0_recency()))
.use_lineage(table.metadata.lineage.clone())
.use_lineage_prev(table.metadata.lineage_prev)
.use_lineage_transformed(table.metadata.lineage_transformed)
.use_lineage_last(table.metadata.lineage_last)
.use_sync_mode(options.sync_mode)
.use_prefix_extractor(options.prefix_extractor.clone())
.use_encryption(options.encryption.clone());
let writer = if options.prefix_extractor.is_none() {
writer.use_bloom_policy(crate::config::BloomConstructionPolicy::BitsPerKey(0.0))
} else {
writer
};
#[cfg(zstd_any)]
let writer = writer.use_zstd_dictionary(options.zstd_dictionary.clone());
let walk = match salvage_blocks(
&table,
writer,
comparator,
!delete_mask_unpositionable,
allow_verbatim,
options.blob_rewrite.as_deref(),
options.progress.as_deref(),
) {
Ok(walk) => walk,
Err(e) => {
discard_partial(fs, &dest);
return Err(e);
}
};
let salvaged_path = if walk.wrote {
Some(dest)
} else {
discard_partial(fs, &dest);
None
};
Ok(SalvageReport {
salvaged_path,
blocks_total: walk.blocks_total,
blocks_salvaged: walk.blocks_salvaged,
blocks_copied_verbatim: walk.blocks_copied_verbatim,
entries_salvaged: walk.entries_salvaged,
entries_dropped_by_rewrite: walk.entries_dropped_by_rewrite,
dropped: walk.dropped,
delete_rows_resurrected: delete_mask_unpositionable && options.allow_delete_resurrection,
})
}
struct SalvageWalk {
blocks_total: usize,
blocks_salvaged: usize,
blocks_copied_verbatim: usize,
entries_salvaged: u64,
entries_dropped_by_rewrite: u64,
dropped: Vec<DroppedBlock>,
wrote: bool,
}
fn discard_partial(fs: &alloc::sync::Arc<dyn crate::fs::Fs>, dest: &std::path::Path) {
if let Err(e) = fs.remove_file(dest) {
log::warn!(
"salvage: could not remove the incomplete destination {}: {e}",
dest.display(),
);
}
}
fn entry_directory(path: &std::path::Path) -> &std::path::Path {
match path.parent() {
Some(parent) if !parent.as_os_str().is_empty() => parent,
_ => std::path::Path::new("."),
}
}
fn suppress_shadowed_boundary(
entries: Vec<crate::InternalValue>,
boundary: Option<&UserKey>,
) -> (Vec<crate::InternalValue>, Option<UserKey>) {
let Some(shadowed) = boundary
.cloned()
.or_else(|| entries.first().map(|e| e.key.user_key.clone()))
else {
return (entries, None);
};
let kept: Vec<_> = entries
.into_iter()
.filter(|e| !crate::comparator::same_user_key(&e.key.user_key, &shadowed))
.collect();
let carry = kept.is_empty().then_some(shadowed);
(kept, carry)
}
#[cfg(feature = "columnar")]
enum BoundarySuppression {
Unchanged,
Emptied(UserKey),
Rebuilt(
crate::table::columnar::ColumnBatch,
Vec<crate::InternalValue>,
),
}
#[cfg(feature = "columnar")]
fn suppress_columnar_boundary(
batch: &crate::table::columnar::ColumnBatch,
boundary: Option<&UserKey>,
) -> crate::Result<BoundarySuppression> {
let entries = crate::table::columnar::column_batch_to_entries(batch)?;
let before = entries.len();
let (kept, carry) = suppress_shadowed_boundary(entries, boundary);
if kept.len() == before {
return Ok(BoundarySuppression::Unchanged);
}
if let Some(key) = carry {
return Ok(BoundarySuppression::Emptied(key));
}
if batch.columns.len() > 4 {
return Err(crate::Error::FeatureUnsupported(
"boundary-key suppression in a columnar block with value sub-columns",
));
}
let rebuilt = crate::table::columnar::entries_to_column_batch(&kept)?;
Ok(BoundarySuppression::Rebuilt(rebuilt, kept))
}
fn classify_drop(
e: &crate::Error,
offset: u64,
prev_end: Option<&UserKey>,
end_key: Option<&UserKey>,
) -> DroppedBlock {
use alloc::format;
let reason = match e {
crate::Error::ChecksumMismatch { .. } => DropReason::ChecksumMismatch,
crate::Error::InvalidHeader(_) | crate::Error::InvalidTag(_) => {
DropReason::DecodeError(format!("{e:?}"))
}
_ => DropReason::ReadError(format!("{e:?}")),
};
DroppedBlock {
offset,
section: b"data".to_vec(),
reason,
key_range: end_key.map(|ek| (prev_end.cloned().unwrap_or_else(UserKey::empty), ek.clone())),
}
}
fn rewrite_block_indirections(
entries: Vec<crate::InternalValue>,
rewrite: &crate::HashMap<crate::vlog::BlobFileId, BlobFileRewrite>,
dropped_entries: &mut u64,
) -> crate::Result<(Vec<crate::InternalValue>, Option<UserKey>)> {
use crate::coding::{Decode, Encode};
let mut out = Vec::with_capacity(entries.len());
let mut headless: Option<UserKey> = None;
for mut entry in entries {
if headless
.as_ref()
.is_some_and(|h| crate::comparator::same_user_key(&entry.key.user_key, h))
{
*dropped_entries += 1;
continue;
}
headless = None;
if entry.key.value_type != crate::ValueType::Indirection {
out.push(entry);
continue;
}
let mut cursor = &entry.value[..];
let mut ind = crate::blob_tree::handle::BlobIndirection::decode_from(&mut cursor)?;
match rewrite.get(&ind.vhandle.blob_file_id) {
None => out.push(entry),
Some(BlobFileRewrite::Remap { new_id, offsets }) => {
if let Some(&relocation) = offsets.get(&ind.vhandle.offset) {
ind.vhandle.blob_file_id = *new_id;
ind.vhandle.offset = relocation.offset;
ind.vhandle.on_disk_size = relocation.on_disk_size;
let mut buf = Vec::new();
ind.encode_into(&mut buf)?;
entry.value = buf.into();
out.push(entry);
} else {
*dropped_entries += 1;
headless = Some(entry.key.user_key);
}
}
Some(BlobFileRewrite::DropBelow(frontier)) => {
if ind.vhandle.offset < *frontier {
*dropped_entries += 1;
headless = Some(entry.key.user_key);
} else {
out.push(entry);
}
}
}
}
Ok((out, headless))
}
fn collect_indirections(
entries: &[crate::InternalValue],
) -> crate::Result<Vec<crate::blob_tree::handle::BlobIndirection>> {
use crate::coding::Decode;
let mut out = Vec::new();
for entry in entries {
if entry.key.value_type == crate::ValueType::Indirection {
let mut cursor = &entry.value[..];
out.push(crate::blob_tree::handle::BlobIndirection::decode_from(
&mut cursor,
)?);
}
}
Ok(out)
}
#[cfg(feature = "columnar")]
fn collect_columnar_indirections(
batch: &crate::table::columnar::ColumnBatch,
) -> crate::Result<Vec<crate::blob_tree::handle::BlobIndirection>> {
let tag = u8::from(crate::ValueType::Indirection);
let has_indirections = batch.columns.get(2).is_some_and(|c| c.data.contains(&tag));
if !has_indirections {
return Ok(Vec::new());
}
let entries = crate::table::columnar::column_batch_to_entries(batch)?;
collect_indirections(&entries)
}
fn fold_blob_links(
derived: &mut crate::HashMap<crate::vlog::BlobFileId, crate::table::writer::LinkedFile>,
indirections: &[crate::blob_tree::handle::BlobIndirection],
) {
for ind in indirections {
derived
.entry(ind.vhandle.blob_file_id)
.and_modify(|link| {
link.bytes += u64::from(ind.size);
link.on_disk_bytes += u64::from(ind.vhandle.on_disk_size);
link.len += 1;
})
.or_insert_with(|| crate::table::writer::LinkedFile {
blob_file_id: ind.vhandle.blob_file_id,
bytes: u64::from(ind.size),
on_disk_bytes: u64::from(ind.vhandle.on_disk_size),
len: 1,
});
}
}
#[derive(Default)]
struct PublishedProgress {
blocks_scanned: usize,
blocks_recovered: usize,
blocks_dropped: usize,
blocks_healed: u64,
kvs: u64,
columns: u64,
}
#[expect(
clippy::too_many_arguments,
reason = "a plain projection of the walk's running totals; bundling them into a struct would only rename the call sites"
)]
fn publish_progress(
progress: Option<&crate::RecoveryProgress>,
published: &mut PublishedProgress,
blocks_scanned: usize,
blocks_recovered: usize,
blocks_dropped: usize,
blocks_healed: u64,
kvs: u64,
columns: u64,
) {
let Some(p) = progress else { return };
p.add_blocks(
(blocks_scanned - published.blocks_scanned) as u64,
(blocks_recovered - published.blocks_recovered) as u64,
(blocks_dropped - published.blocks_dropped) as u64,
blocks_healed - published.blocks_healed,
);
p.add_rows(kvs - published.kvs, columns - published.columns);
*published = PublishedProgress {
blocks_scanned,
blocks_recovered,
blocks_dropped,
blocks_healed,
kvs,
columns,
};
}
#[cfg_attr(
not(feature = "columnar"),
expect(
unused_variables,
reason = "the delete mask exists only for columnar sources; without the feature the flag has no consumer"
)
)]
fn salvage_blocks(
table: &crate::table::Table,
mut writer: crate::table::Writer,
comparator: &crate::comparator::SharedComparator,
apply_delete_mask: bool,
allow_verbatim: bool,
blob_rewrite: Option<&crate::HashMap<crate::vlog::BlobFileId, BlobFileRewrite>>,
progress: Option<&crate::RecoveryProgress>,
) -> crate::Result<SalvageWalk> {
use crate::table::block::ParsedItem;
use alloc::format;
let allow_verbatim = allow_verbatim && blob_rewrite.is_none();
let mut blocks_total = 0usize;
let mut blocks_salvaged = 0usize;
let mut blocks_copied_verbatim = 0usize;
let mut entries_salvaged = 0u64;
let mut blocks_healed = 0u64;
#[cfg_attr(
not(feature = "columnar"),
expect(unused_mut, reason = "only the columnar walk arms count columns")
)]
let mut columns_salvaged = 0u64;
let mut entries_dropped_by_rewrite = 0u64;
let mut published = PublishedProgress::default();
let mut dropped: Vec<DroppedBlock> = Vec::new();
let mut derived_blob_links: crate::HashMap<
crate::vlog::BlobFileId,
crate::table::writer::LinkedFile,
> = crate::HashMap::default();
let mut prev_end: Option<UserKey> = None;
let mut lost_boundary: Option<Option<UserKey>> = None;
let mut prev_block_end: Option<UserKey> = None;
let mut dropped_seen;
let mut indexed: Vec<crate::table::KeyedBlockHandle> = Vec::new();
let mut index_enum_error: Option<String> = None;
for handle in table.data_block_handles() {
match handle {
Ok(k) => indexed.push(k),
Err(e) => {
index_enum_error = Some(format!("{e:?}"));
break;
}
}
}
let mut items: Vec<(crate::table::BlockHandle, Option<UserKey>)> = Vec::new();
let data_section = {
let mut file = table
.fs
.open(&table.path, &crate::fs::FsOpenOptions::new().read(true))?;
match crate::sfa::Reader::from_reader(&mut file) {
Ok(t) => {
let toc_pos = t.toc_pos();
t.toc().section(b"data").and_then(|s| {
let end = s.pos().checked_add(s.len())?;
(end <= toc_pos).then_some((s.pos(), end))
})
}
Err(crate::sfa::Error::Io(e)) if e.kind().is_environmental() => {
return Err(crate::Error::Io(e));
}
Err(_) => None,
}
};
if let Some((section_pos, section_end)) = data_section {
let probe_file = table
.fs
.open(&table.path, &crate::fs::FsOpenOptions::new().read(true))?;
let probe_block_type = {
#[cfg(feature = "columnar")]
{
if table.metadata.columnar {
crate::table::block::BlockType::Columnar
} else {
crate::table::block::BlockType::Data
}
}
#[cfg(not(feature = "columnar"))]
{
crate::table::block::BlockType::Data
}
};
let frames_and_loads =
|at: u64, to: u64| -> crate::Result<Option<crate::table::BlockHandle>> {
match table.probe_block_handle_in(&*probe_file, at, to) {
Ok(h) => match table.salvage_load_block(&h, probe_block_type) {
Ok(_) => Ok(Some(h)),
Err(crate::Error::Io(io)) if io.kind().is_environmental() => {
Err(crate::Error::Io(io))
}
Err(_) => Ok(None),
},
Err(crate::Error::Io(io)) if io.kind().is_environmental() => {
Err(crate::Error::Io(io))
}
Err(_) => Ok(None),
}
};
let probe_gap = |from: u64,
to: u64,
items: &mut Vec<(crate::table::BlockHandle, Option<UserKey>)>,
dropped: &mut Vec<DroppedBlock>|
-> crate::Result<bool> {
let mut at = from;
while at < to {
if let Some(h) = frames_and_loads(at, to)? {
let next = at + u64::from(h.size());
if next > at {
items.push((h, None));
at = next;
continue;
}
}
dropped.push(DroppedBlock {
offset: at,
section: b"data".to_vec(),
reason: DropReason::HeaderCorrupted(
"unanchored bytes after a broken block boundary; the tail \
cannot be proven to be original block starts"
.to_owned(),
),
key_range: None,
});
return Ok(false);
}
Ok(true)
};
let drop_unanchored_handle = |off: u64, dropped: &mut Vec<DroppedBlock>| {
dropped.push(DroppedBlock {
offset: off,
section: b"data".to_vec(),
reason: DropReason::HeaderCorrupted(
"indexed block after a broken physical chain has no \
authenticated boundary"
.to_owned(),
),
key_range: None,
});
};
let mut cursor = section_pos;
indexed.sort_unstable_by_key(|k| *k.as_ref().offset());
let tli_trusted = table.tli_structure_authenticated()?;
let mut chain_anchored = true;
for keyed in indexed {
let off = *keyed.as_ref().offset();
if off >= section_end {
continue;
}
if off < cursor {
continue;
}
if !chain_anchored {
drop_unanchored_handle(off, &mut dropped);
continue;
}
if off > cursor {
let reached = probe_gap(cursor, off, &mut items, &mut dropped)?;
if !reached && !tli_trusted {
chain_anchored = false;
drop_unanchored_handle(off, &mut dropped);
continue;
}
cursor = off;
}
let (handle, end_key) =
match table.probe_block_handle_in(&*probe_file, off, section_end) {
Ok(probed) if probed.size() == keyed.as_ref().size() => {
(*keyed.as_ref(), Some(keyed.end_key().clone()))
}
Ok(probed) => (probed, None),
Err(crate::Error::Io(io)) if io.kind().is_environmental() => {
return Err(crate::Error::Io(io));
}
Err(_) => {
if !tli_trusted {
chain_anchored = false;
drop_unanchored_handle(off, &mut dropped);
}
continue;
}
};
let next = (off + u64::from(handle.size())).min(section_end);
items.push((handle, end_key));
cursor = cursor.max(next);
}
if chain_anchored && cursor < section_end {
probe_gap(cursor, section_end, &mut items, &mut dropped)?;
}
} else {
if let Some(reason) = index_enum_error {
dropped.push(DroppedBlock {
offset: 0,
section: b"index".to_vec(),
reason: DropReason::HeaderCorrupted(reason),
key_range: None,
});
}
for keyed in indexed {
let handle = *keyed.as_ref();
items.push((handle, Some(keyed.end_key().clone())));
}
}
blocks_total += dropped.len();
let mut lost_regions: Vec<u64> = dropped
.iter()
.filter(|d| d.section == b"data")
.map(|d| d.offset)
.collect();
lost_regions.sort_unstable();
let mut lost_regions = lost_regions.into_iter().peekable();
dropped_seen = dropped.len();
for (block_handle, end_key) in items {
publish_progress(
progress,
&mut published,
blocks_total,
blocks_salvaged,
dropped.len(),
blocks_healed,
entries_salvaged,
columns_salvaged,
);
blocks_total += 1;
if dropped.len() > dropped_seen {
lost_boundary = Some(prev_block_end.clone());
}
dropped_seen = dropped.len();
let offset = *block_handle.offset();
let mut passed_lost_region = false;
while lost_regions.next_if(|start| *start < offset).is_some() {
passed_lost_region = true;
}
if passed_lost_region {
lost_boundary = Some(None);
}
prev_block_end.clone_from(&end_key);
#[cfg(feature = "columnar")]
if table.metadata.columnar {
if table.has_delete_bitmap_section() && apply_delete_mask {
if end_key.is_none() {
dropped.push(classify_drop(
&crate::Error::InvalidHeader(
"delete positions unverifiable for an index-omitted block",
),
offset,
prev_end.as_ref(),
None,
));
continue;
}
match table.load_columnar_block_masked(&block_handle) {
Ok(Some(batch)) => {
let mut rewrite_carry: Option<UserKey> = None;
let batch = match blob_rewrite {
Some(rw) => {
if batch.columns.len() > 4 {
return Err(crate::Error::FeatureUnsupported(
"blob-handle rewrite of a columnar block \
with value sub-columns",
));
}
let step = crate::table::columnar::column_batch_to_entries(&batch)
.and_then(|entries| {
rewrite_block_indirections(
entries,
rw,
&mut entries_dropped_by_rewrite,
)
});
match step {
Ok((entries, carry)) if entries.is_empty() => {
if let Some(key) = carry {
lost_boundary = Some(Some(key));
}
prev_end = end_key.or(prev_end);
continue;
}
Ok((entries, carry)) => {
rewrite_carry = carry;
match crate::table::columnar::entries_to_column_batch(
&entries,
) {
Ok(batch) => batch,
Err(e) => {
dropped.push(classify_drop(
&e,
offset,
prev_end.as_ref(),
end_key.as_ref(),
));
prev_end = end_key.or(prev_end);
continue;
}
}
}
Err(e) => {
dropped.push(classify_drop(
&e,
offset,
prev_end.as_ref(),
end_key.as_ref(),
));
prev_end = end_key.or(prev_end);
continue;
}
}
}
None => batch,
};
let batch = match lost_boundary.take() {
Some(boundary) => {
match suppress_columnar_boundary(&batch, boundary.as_ref()) {
Ok(BoundarySuppression::Unchanged) => batch,
Ok(BoundarySuppression::Emptied(key)) => {
lost_boundary = Some(Some(key));
prev_end = end_key.or(prev_end);
continue;
}
Ok(BoundarySuppression::Rebuilt(batch, _)) => batch,
Err(e @ crate::Error::FeatureUnsupported(_)) => {
return Err(e);
}
Err(e) => {
dropped.push(classify_drop(
&e,
offset,
prev_end.as_ref(),
end_key.as_ref(),
));
prev_end = end_key.or(prev_end);
continue;
}
}
}
None => batch,
};
if let Some(key) = rewrite_carry {
lost_boundary = Some(Some(key));
}
let rows = u64::from(batch.row_count);
let block_links = match collect_columnar_indirections(&batch) {
Ok(links) => links,
Err(e) => {
dropped.push(classify_drop(
&e,
offset,
prev_end.as_ref(),
end_key.as_ref(),
));
prev_end = end_key.or(prev_end);
continue;
}
};
match writer.write_columnar_block_verbatim(&batch, comparator) {
Ok(_) => {
entries_salvaged += rows;
blocks_salvaged += 1;
columns_salvaged += batch.columns.len() as u64;
fold_blob_links(&mut derived_blob_links, &block_links);
}
Err(
e @ (crate::Error::InvalidHeader(_) | crate::Error::InvalidTag(_)),
) => {
dropped.push(classify_drop(
&e,
offset,
prev_end.as_ref(),
end_key.as_ref(),
));
}
Err(e) => return Err(e),
}
}
Ok(None) => {}
Err(crate::Error::Io(io)) if io.kind().is_environmental() => {
return Err(crate::Error::Io(io));
}
Err(e) => dropped.push(classify_drop(
&e,
offset,
prev_end.as_ref(),
end_key.as_ref(),
)),
}
} else {
match table
.salvage_load_block(&block_handle, crate::table::block::BlockType::Columnar)
{
Ok(mut sb) => {
let ecc_healed = sb.ecc_recovered;
if !allow_verbatim {
sb.verbatim = None;
}
match crate::table::columnar::ColumnBatch::decode(&sb.block.data).and_then(
|batch| {
crate::table::columnar::column_batch_to_entries(&batch)
.map(|entries| (batch, entries))
},
) {
Ok((batch, _)) if batch.row_count == 0 => {
dropped.push(classify_drop(
&crate::Error::InvalidHeader("columnar: zero-row data block"),
offset,
prev_end.as_ref(),
end_key.as_ref(),
));
}
Ok((batch, entries)) => {
let mut rewrite_carry: Option<UserKey> = None;
let (batch, entries) = match blob_rewrite {
Some(rw) => {
if batch.columns.len() > 4 {
return Err(crate::Error::FeatureUnsupported(
"blob-handle rewrite of a columnar block \
with value sub-columns",
));
}
let entries = match rewrite_block_indirections(
entries,
rw,
&mut entries_dropped_by_rewrite,
) {
Ok((entries, carry)) => {
rewrite_carry = carry;
entries
}
Err(e) => {
dropped.push(classify_drop(
&e,
offset,
prev_end.as_ref(),
end_key.as_ref(),
));
prev_end = end_key.or(prev_end);
continue;
}
};
if entries.is_empty() {
if let Some(key) = rewrite_carry {
lost_boundary = Some(Some(key));
}
prev_end = end_key.or(prev_end);
continue;
}
match crate::table::columnar::entries_to_column_batch(
&entries,
) {
Ok(batch) => (batch, entries),
Err(e) => {
dropped.push(classify_drop(
&e,
offset,
prev_end.as_ref(),
end_key.as_ref(),
));
prev_end = end_key.or(prev_end);
continue;
}
}
}
None => (batch, entries),
};
let mut rebuilt_by_suppression = false;
let (batch, entries) = match lost_boundary.take() {
Some(boundary) => {
match suppress_columnar_boundary(&batch, boundary.as_ref())
{
Ok(BoundarySuppression::Unchanged) => (batch, entries),
Ok(BoundarySuppression::Emptied(key)) => {
lost_boundary = Some(Some(key));
prev_end = end_key.or(prev_end);
continue;
}
Ok(BoundarySuppression::Rebuilt(batch, kept)) => {
rebuilt_by_suppression = true;
(batch, kept)
}
Err(e @ crate::Error::FeatureUnsupported(_)) => {
return Err(e);
}
Err(e) => {
dropped.push(classify_drop(
&e,
offset,
prev_end.as_ref(),
end_key.as_ref(),
));
prev_end = end_key.or(prev_end);
continue;
}
}
}
None => (batch, entries),
};
if let Some(key) = rewrite_carry {
lost_boundary = Some(Some(key));
}
let rows = u64::from(batch.row_count);
let block_links = match collect_indirections(&entries) {
Ok(links) => links,
Err(e) => {
dropped.push(classify_drop(
&e,
offset,
prev_end.as_ref(),
end_key.as_ref(),
));
prev_end = end_key.or(prev_end);
continue;
}
};
let verbatim_source = if table.has_delete_bitmap_section()
|| rebuilt_by_suppression
{
None
} else {
sb.verbatim
};
let emitted = match verbatim_source {
Some((raw, header, layout)) => writer
.append_verbatim_data_block(
&raw,
header,
layout,
&entries,
Some(batch.zone_stats()),
comparator,
)
.map(|_| true),
None => writer
.write_columnar_block_verbatim(&batch, comparator)
.map(|_| false),
};
match emitted {
Ok(verbatim) => {
if verbatim {
blocks_copied_verbatim += 1;
}
if ecc_healed {
blocks_healed += 1;
}
entries_salvaged += rows;
blocks_salvaged += 1;
columns_salvaged += batch.columns.len() as u64;
fold_blob_links(&mut derived_blob_links, &block_links);
}
Err(
e @ (crate::Error::InvalidHeader(_)
| crate::Error::InvalidTag(_)),
) => {
dropped.push(classify_drop(
&e,
offset,
prev_end.as_ref(),
end_key.as_ref(),
));
}
Err(e) => return Err(e),
}
}
Err(crate::Error::Io(io)) if io.kind().is_environmental() => {
return Err(crate::Error::Io(io));
}
Err(e) => {
dropped.push(classify_drop(
&e,
offset,
prev_end.as_ref(),
end_key.as_ref(),
));
}
}
}
Err(crate::Error::Io(io)) if io.kind().is_environmental() => {
return Err(crate::Error::Io(io));
}
Err(e) => dropped.push(classify_drop(
&e,
offset,
prev_end.as_ref(),
end_key.as_ref(),
)),
}
}
prev_end = end_key.or(prev_end);
continue;
}
match table.salvage_load_block(&block_handle, crate::table::block::BlockType::Data) {
Ok(mut sb) => {
let ecc_healed = sb.ecc_recovered;
if !allow_verbatim {
sb.verbatim = None;
}
let has_kv_footer = table.metadata.kv_checksum_algo.is_some();
if has_kv_footer
&& let Err(e) = crate::table::DataBlock::verify_kv_checked(
&sb.block.data,
sb.block.header,
comparator.clone(),
table.metadata.kv_checksum_algo,
)
{
dropped.push(classify_drop(
&e,
offset,
prev_end.as_ref(),
end_key.as_ref(),
));
prev_end = end_key.or(prev_end);
continue;
}
match crate::table::DataBlock::from_loaded(sb.block, has_kv_footer) {
Ok(data_block) => match data_block.try_iter(comparator.clone()) {
Ok(iter) => {
let entries: Vec<crate::InternalValue> =
iter.map(|p| p.materialize(data_block.as_slice())).collect();
if entries.is_empty() {
dropped.push(classify_drop(
&crate::Error::InvalidHeader(
"row block decodes to zero entries",
),
offset,
prev_end.as_ref(),
end_key.as_ref(),
));
prev_end = end_key.or(prev_end);
continue;
}
if entries.len() != data_block.len() {
dropped.push(classify_drop(
&crate::Error::InvalidHeader(
"row block iterates to fewer entries than its \
trailer declares",
),
offset,
prev_end.as_ref(),
end_key.as_ref(),
));
prev_end = end_key.or(prev_end);
continue;
}
let (entries, rewrite_carry) = match blob_rewrite {
Some(rw) => match rewrite_block_indirections(
entries,
rw,
&mut entries_dropped_by_rewrite,
) {
Ok(pair) => pair,
Err(e) => {
dropped.push(classify_drop(
&e,
offset,
prev_end.as_ref(),
end_key.as_ref(),
));
prev_end = end_key.or(prev_end);
continue;
}
},
None => (entries, None),
};
let entries = match lost_boundary.take() {
Some(boundary) => {
let before = entries.len();
let (kept, carry) =
suppress_shadowed_boundary(entries, boundary.as_ref());
if let Some(key) = carry {
lost_boundary = Some(Some(key));
}
if kept.len() != before {
sb.verbatim = None;
}
kept
}
None => entries,
};
if let Some(key) = rewrite_carry {
lost_boundary = Some(Some(key));
}
if entries.is_empty() {
prev_end = end_key.or(prev_end);
continue;
}
let count = entries.len() as u64;
let block_links = match collect_indirections(&entries) {
Ok(links) => links,
Err(e) => {
dropped.push(classify_drop(
&e,
offset,
prev_end.as_ref(),
end_key.as_ref(),
));
prev_end = end_key.or(prev_end);
continue;
}
};
let emitted = writer
.validate_direct_block_order(&entries, comparator)
.and_then(|()| {
if let Some((raw, header, layout)) = sb.verbatim {
writer
.append_verbatim_data_block(
&raw, header, layout, &entries, None, comparator,
)
.map(|_| true)
} else {
for e in entries {
writer.write(e)?;
}
Ok(false)
}
});
match emitted {
Ok(verbatim) => {
if verbatim {
blocks_copied_verbatim += 1;
}
if ecc_healed {
blocks_healed += 1;
}
entries_salvaged += count;
blocks_salvaged += 1;
fold_blob_links(&mut derived_blob_links, &block_links);
}
Err(
e @ (crate::Error::InvalidHeader(_)
| crate::Error::InvalidTag(_)),
) => {
dropped.push(classify_drop(
&e,
offset,
prev_end.as_ref(),
end_key.as_ref(),
));
}
Err(e) => return Err(e),
}
}
Err(e) => dropped.push(DroppedBlock {
offset,
section: b"data".to_vec(),
reason: DropReason::DecodeError(format!("{e:?}")),
key_range: end_key.as_ref().map(|ek| {
(prev_end.clone().unwrap_or_else(UserKey::empty), ek.clone())
}),
}),
},
Err(e) => dropped.push(DroppedBlock {
offset,
section: b"data".to_vec(),
reason: DropReason::DecodeError(format!("{e:?}")),
key_range: end_key.as_ref().map(|ek| {
(prev_end.clone().unwrap_or_else(UserKey::empty), ek.clone())
}),
}),
}
}
Err(crate::Error::Io(io)) if io.kind().is_environmental() => {
return Err(crate::Error::Io(io));
}
Err(e) => dropped.push(classify_drop(
&e,
offset,
prev_end.as_ref(),
end_key.as_ref(),
)),
}
prev_end = end_key.or(prev_end);
}
publish_progress(
progress,
&mut published,
blocks_total,
blocks_salvaged,
dropped.len(),
blocks_healed,
entries_salvaged,
columns_salvaged,
);
let wrote = blocks_salvaged > 0;
if wrote {
let mut links: Vec<crate::table::writer::LinkedFile> =
derived_blob_links.into_values().collect();
links.sort_unstable_by_key(|l| l.blob_file_id);
for link in links {
writer.link_blob_file(link.blob_file_id, link.len, link.bytes, link.on_disk_bytes);
}
writer.finish()?;
} else {
drop(writer);
}
Ok(SalvageWalk {
blocks_total,
blocks_salvaged,
blocks_copied_verbatim,
entries_salvaged,
entries_dropped_by_rewrite,
dropped,
wrote,
})
}
#[derive(Debug, Clone)]
pub enum BlobDropReason {
ChecksumMismatch,
Corrupt(String),
}
#[derive(Debug, Clone)]
pub struct DroppedBlob {
pub reason: BlobDropReason,
}
#[derive(Debug)]
pub struct BlobSalvageReport {
pub salvaged_path: Option<PathBuf>,
pub records_total: usize,
pub records_salvaged: usize,
pub offset_remap: Vec<(u64, BlobRecordRelocation)>,
pub dropped: Vec<DroppedBlob>,
}
impl BlobSalvageReport {
#[must_use]
pub fn is_complete(&self) -> bool {
self.dropped.is_empty()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct BlobRecordRelocation {
pub offset: u64,
pub on_disk_size: u32,
}
#[cfg_attr(
not(zstd_any),
expect(
clippy::elidable_lifetime_names,
reason = "named to stay valid in the zstd build, where a second reference parameter exists"
)
)]
pub(crate) fn decompress_blob_value<'a>(
compression: crate::CompressionType,
on_disk: &'a [u8],
real_len: usize,
#[cfg(zstd_any)] zstd_dictionary: Option<&crate::compression::ZstdDictionary>,
) -> crate::Result<alloc::borrow::Cow<'a, [u8]>> {
match compression {
crate::CompressionType::None => {
if on_disk.len() != real_len {
return Err(crate::Error::InvalidHeader("Blob"));
}
Ok(alloc::borrow::Cow::Borrowed(on_disk))
}
#[cfg(feature = "lz4")]
crate::CompressionType::Lz4 => {
let mut buf = alloc::vec![0u8; real_len];
let written = lz4_flex::block::decompress_into(on_disk, &mut buf)
.map_err(|_| crate::Error::Decompress(compression))?;
if written != real_len {
return Err(crate::Error::Decompress(compression));
}
Ok(alloc::borrow::Cow::Owned(buf))
}
#[cfg(zstd_any)]
crate::CompressionType::Zstd(_) => {
use crate::compression::CompressionProvider as _;
let decompressed = crate::compression::ZstdBackend::decompress(on_disk, real_len)
.map_err(|_| crate::Error::Decompress(compression))?;
if decompressed.len() != real_len {
return Err(crate::Error::Decompress(compression));
}
Ok(alloc::borrow::Cow::Owned(decompressed))
}
#[cfg(zstd_any)]
crate::CompressionType::ZstdDict { dict_id, .. } => {
use crate::compression::CompressionProvider as _;
let Some(dict) = zstd_dictionary else {
return Err(crate::Error::FeatureUnsupported(
"salvage of a dictionary-compressed blob file",
));
};
if dict.id() != dict_id {
return Err(crate::Error::ZstdDictMismatch {
expected: dict_id,
got: Some(dict.id()),
});
}
let decompressed =
crate::compression::ZstdBackend::decompress_with_dict(on_disk, dict, real_len)
.map_err(|_| crate::Error::Decompress(compression))?;
if decompressed.len() != real_len {
return Err(crate::Error::Decompress(compression));
}
Ok(alloc::borrow::Cow::Owned(decompressed))
}
}
}
pub(crate) fn blob_key_regresses(
comparator: &crate::comparator::SharedComparator,
prev: Option<&(crate::UserKey, crate::SeqNo)>,
entry: &crate::vlog::blob_file::scanner::ScanEntry,
) -> bool {
let Some((prev_key, prev_seqno)) = prev else {
return false;
};
match comparator.compare(entry.key.as_ref(), prev_key.as_ref()) {
core::cmp::Ordering::Less => true,
core::cmp::Ordering::Equal => entry.seqno > *prev_seqno,
core::cmp::Ordering::Greater => false,
}
}
pub fn salvage_blob_file(
source: &std::path::Path,
dest: std::path::PathBuf,
fs: &alloc::sync::Arc<dyn crate::fs::Fs>,
blob_file_id: crate::vlog::BlobFileId,
comparator: &crate::comparator::SharedComparator,
live_data_start: u64,
#[cfg(zstd_any)] zstd_dictionary: Option<&alloc::sync::Arc<crate::compression::ZstdDictionary>>,
) -> crate::Result<BlobSalvageReport> {
use crate::vlog::blob_file::{scanner::Scanner, writer::Writer as BlobWriter};
use alloc::format;
let source_handle =
crate::vlog::recover_blob_file(source, blob_file_id, crate::Checksum::from_raw(0), 0, fs)?;
let compression = source_handle.compression();
#[cfg(zstd_any)]
if let crate::CompressionType::ZstdDict { dict_id, .. } = compression {
let Some(dict) = zstd_dictionary else {
return Err(crate::Error::ZstdDictMismatch {
expected: dict_id,
got: None,
});
};
if dict.id() != dict_id {
return Err(crate::Error::ZstdDictMismatch {
expected: dict_id,
got: Some(dict.id()),
});
}
}
let scanner = if live_data_start > 0 {
Scanner::resume(source, &**fs, blob_file_id, live_data_start)?
} else {
Scanner::new(source, &**fs, blob_file_id)?
};
let sync_mode = crate::fs::SyncMode::Full;
let mut writer = BlobWriter::new(&dest, blob_file_id, 0, &**fs)?
.use_sync_mode(sync_mode)
.use_compression(compression);
#[cfg(zstd_any)]
{
writer = writer.use_zstd_dictionary(zstd_dictionary.cloned());
}
let mut records_total = 0usize;
let mut records_salvaged = 0usize;
let mut offset_remap: Vec<(u64, BlobRecordRelocation)> = Vec::new();
let mut dropped: Vec<DroppedBlob> = Vec::new();
let mut prev_written: Option<(crate::UserKey, crate::SeqNo)> = None;
let walk = (|| -> crate::Result<()> {
for item in scanner {
records_total += 1;
match item {
Ok(entry) if entry.resynced => {
let _ = entry;
dropped.push(DroppedBlob {
reason: BlobDropReason::Corrupt(
"tail surrendered at the first resync: every frame past a \
damaged frame has an unprovable boundary, dropped as one"
.to_string(),
),
});
break;
}
Ok(entry) if entry.key.is_empty() => {
dropped.push(DroppedBlob {
reason: BlobDropReason::Corrupt("frame carries an empty key".to_string()),
});
}
Ok(entry)
if compression == crate::CompressionType::None
&& entry.uncompressed_len as usize != entry.value.len() =>
{
dropped.push(DroppedBlob {
reason: BlobDropReason::Corrupt(
"frame's declared value length disagrees with its stored bytes"
.to_string(),
),
});
}
Ok(entry) if blob_key_regresses(comparator, prev_written.as_ref(), &entry) => {
dropped.push(DroppedBlob {
reason: BlobDropReason::Corrupt(
"frame's internal key regresses below the previous salvaged \
record; re-emitting it would violate the blob writer's \
sorted-input contract"
.to_string(),
),
});
}
Ok(entry) => {
let value = match decompress_blob_value(
compression,
&entry.value,
entry.uncompressed_len as usize,
#[cfg(zstd_any)]
zstd_dictionary.map(alloc::sync::Arc::as_ref),
) {
Ok(value) => value,
#[cfg(zstd_any)]
Err(e @ crate::Error::ZstdDictMismatch { .. }) => return Err(e),
Err(e) => {
dropped.push(DroppedBlob {
reason: BlobDropReason::Corrupt(format!(
"frame's value does not decompress: {e:?}"
)),
});
continue;
}
};
let salvaged_offset = writer.offset();
let on_disk_size = writer.write(&entry.key, entry.seqno, &value)?;
offset_remap.push((
entry.offset,
BlobRecordRelocation {
offset: salvaged_offset,
on_disk_size,
},
));
prev_written = Some((entry.key.clone(), entry.seqno));
records_salvaged += 1;
}
Err(crate::Error::ChecksumMismatch { .. }) => dropped.push(DroppedBlob {
reason: BlobDropReason::ChecksumMismatch,
}),
Err(
e @ (crate::Error::HeaderCrcMismatch { .. } | crate::Error::InvalidHeader(_)),
) => {
dropped.push(DroppedBlob {
reason: BlobDropReason::Corrupt(format!("{e:?}")),
});
}
Err(crate::Error::Io(io)) if io.kind().is_environmental() => {
return Err(crate::Error::Io(io));
}
Err(e) => {
dropped.push(DroppedBlob {
reason: BlobDropReason::Corrupt(format!("{e:?}")),
});
break;
}
}
}
Ok(())
})();
let salvaged_path = match walk {
Err(e) => {
drop(writer);
discard_partial(fs, &dest);
return Err(e);
}
Ok(()) if records_salvaged > 0 => {
if let Err(e) = writer.finish() {
discard_partial(fs, &dest);
return Err(e);
}
if let Err(e) = fs.sync_directory_with(entry_directory(&dest), sync_mode) {
discard_partial(fs, &dest);
return Err(e.into());
}
Some(dest)
}
Ok(()) => {
drop(writer);
discard_partial(fs, &dest);
None
}
};
Ok(BlobSalvageReport {
salvaged_path,
records_total,
records_salvaged,
offset_remap,
dropped,
})
}
#[cfg(test)]
mod tests;