mod accessor;
pub mod blob_file;
mod handle;
pub use {
accessor::Accessor, blob_file::BlobFile,
blob_file::merge::MergeScanner as BlobFileMergeScanner,
blob_file::multi_writer::MultiWriter as BlobFileWriter,
blob_file::scanner::Scanner as BlobFileScanner, handle::ValueHandle,
};
use crate::path::{Path, PathBuf};
use crate::{
Checksum, DescriptorTable, TreeId,
file_accessor::FileAccessor,
fs::Fs,
vlog::blob_file::{Inner as BlobFileInner, Metadata},
};
use alloc::sync::Arc;
#[cfg(not(feature = "std"))]
use alloc::vec::Vec;
use core::sync::atomic::AtomicBool;
fn dedupe_blob_sightings(
fs: &Arc<dyn Fs>,
blob_files: &mut Vec<BlobFile>,
orphaned: &mut Vec<PathBuf>,
) -> crate::Result<()> {
let mut seen: crate::HashMap<BlobFileId, usize> = crate::HashMap::default();
let mut duplicate_ids: crate::HashSet<BlobFileId> = crate::HashSet::default();
for (idx, bf) in blob_files.iter().enumerate() {
if seen.insert(bf.id(), idx).is_some() {
duplicate_ids.insert(bf.id());
}
}
if duplicate_ids.is_empty() {
return Ok(());
}
let rank = |bf: &BlobFile| -> crate::Result<u8> {
let digest = match crate::file::checksum_from_with_overrides(
&**fs,
&bf.0.path,
bf.0.live_data_start,
&[],
) {
Ok(d) => Some(d),
Err(e) if e.is_environmental() => return Err(e),
Err(_) => None,
};
if digest.is_some_and(|d| Checksum::from_raw(d) == bf.0.checksum) {
return Ok(0);
}
let canonical = bf.0.path.parent().is_some_and(|dir| {
dir.join(alloc::string::ToString::to_string(&bf.id())) == *bf.0.path
});
Ok(if canonical { 1 } else { 2 })
};
let mut best: crate::HashMap<BlobFileId, (u8, usize)> = crate::HashMap::default();
for (idx, bf) in blob_files.iter().enumerate() {
let score = rank(bf)?;
match best.get(&bf.id()) {
Some(&(prev, _)) if prev <= score => {}
_ => {
best.insert(bf.id(), (score, idx));
}
}
}
if let Some(id) = best
.iter()
.find(|(id, (score, _))| *score > 0 && duplicate_ids.contains(*id))
.map(|(id, _)| *id)
{
log::error!(
"blob file {id} has multiple copies and none reproduces the manifest's \
checksum; refusing to pick one. Run a repair, which can salvage the \
authoritative copy's intact frames",
);
return Err(crate::Error::Unrecoverable);
}
let winners: crate::HashSet<usize> = best.values().map(|&(_, idx)| idx).collect();
let mut idx = 0;
blob_files.retain(|bf| {
let keep = winners.contains(&idx);
if !keep {
log::warn!(
"blob file {} duplicates id {}; the manifest names another copy, so this one \
is an orphan",
bf.0.path.display(),
bf.id(),
);
orphaned.push(bf.0.path.clone());
}
idx += 1;
keep
});
Ok(())
}
pub fn recover_blob_files(
folder: &Path,
ids: &[(BlobFileId, Checksum, u64)],
tree_id: TreeId,
descriptor_table: Option<&Arc<DescriptorTable>>,
fs: &Arc<dyn Fs>,
) -> crate::Result<(Vec<BlobFile>, Vec<PathBuf>)> {
let entries = match fs.read_dir(folder) {
Ok(entries) => entries,
Err(e) if e.kind() == crate::io::ErrorKind::NotFound => {
if ids.is_empty() {
return Ok((vec![], vec![]));
}
return Err(crate::Error::Unrecoverable);
}
Err(e) => return Err(e.into()),
};
let cnt = ids.len();
let progress_mod = match cnt {
_ if cnt <= 20 => 1,
_ if cnt <= 100 => 10,
_ => 100,
};
log::debug!("Recovering {cnt} blob files from {:?}", folder.display());
let mut blob_files = Vec::with_capacity(ids.len());
let mut orphaned_blob_files = vec![];
let mut pending_cache_inserts = Vec::new();
for (idx, dirent) in entries.into_iter().enumerate() {
let file_name = &dirent.file_name;
if dirent.is_dir {
continue;
}
let blob_file_id = match crate::file::BlobDirEntry::classify(file_name) {
crate::file::BlobDirEntry::Blob(id) => id,
crate::file::BlobDirEntry::SalvageTmp(_) => {
orphaned_blob_files.push(dirent.path.clone());
continue;
}
crate::file::BlobDirEntry::Foreign => {
log::debug!("Ignoring {file_name:?} in the blobs folder: not an engine file");
continue;
}
};
let blob_file_path = &dirent.path;
if let Some(&(_, checksum, live_data_start)) =
ids.iter().find(|(id, _, _)| id == &blob_file_id)
{
log::trace!(
"Recovering blob file #{blob_file_id:?} from {}",
blob_file_path.display(),
);
let mut file = fs.open(blob_file_path, &crate::fs::FsOpenOptions::new().read(true))?;
let meta = {
let reader = crate::sfa::Reader::from_reader(&mut file)?;
let toc = reader.toc();
let metadata_section = toc.section(b"meta")
.ok_or(crate::Error::Unrecoverable)
.inspect_err(|_| {
log::error!("meta section in blob file #{blob_file_id} is missing - maybe the file is corrupted?");
})?;
let metadata_len = usize::try_from(metadata_section.len())
.map_err(|_| crate::Error::Unrecoverable)?;
let metadata_slice =
crate::file::read_exact(&*file, metadata_section.pos(), metadata_len)?;
Metadata::from_slice(&metadata_slice)?
};
let file: Arc<dyn crate::fs::FsFile> = Arc::from(file);
let file_accessor = if let Some(dt) = descriptor_table.cloned() {
let global_id = (tree_id, blob_file_id).into();
pending_cache_inserts.push((
dt.clone(),
global_id,
blob_file_path.clone(),
file.clone(),
));
FileAccessor::DescriptorTable {
table: dt,
fs: fs.clone(),
}
} else {
FileAccessor::File(file)
};
blob_files.push(BlobFile(Arc::new(BlobFileInner {
id: blob_file_id,
path: blob_file_path.clone(),
meta,
is_deleted: AtomicBool::new(false),
punch_on_drop: portable_atomic::AtomicU64::new(u64::MAX),
checksum,
live_data_start,
file_accessor,
tree_id,
fs: fs.clone(),
deletion_pause: once_cell::race::OnceBox::new(),
#[cfg(feature = "std")]
background_deleter: once_cell::race::OnceBox::new(),
})));
if idx % progress_mod == 0 {
log::debug!("Recovered {idx}/{cnt} blob files");
}
} else {
orphaned_blob_files.push(blob_file_path.clone());
}
}
dedupe_blob_sightings(fs, &mut blob_files, &mut orphaned_blob_files)?;
if blob_files.len() < ids.len() {
return Err(crate::Error::Unrecoverable);
}
let retained: crate::HashSet<PathBuf> = blob_files.iter().map(|bf| bf.0.path.clone()).collect();
for (dt, global_id, path, file) in pending_cache_inserts {
if retained.contains(&path) {
dt.insert_for_blob_file(global_id, file);
}
}
log::debug!("Successfully recovered {} blob files", blob_files.len());
Ok((blob_files, orphaned_blob_files))
}
#[cfg_attr(
not(feature = "std"),
allow(
dead_code,
reason = "single-file blob recovery for the std-gated repair surface; the no_std open path uses recover_blob_files"
)
)]
pub fn recover_blob_file(
path: &Path,
id: BlobFileId,
checksum: Checksum,
tree_id: TreeId,
fs: &Arc<dyn Fs>,
) -> crate::Result<BlobFile> {
recover_blob_file_from(path, id, checksum, tree_id, fs, 0)
}
pub fn recover_blob_file_from(
path: &Path,
id: BlobFileId,
checksum: Checksum,
tree_id: TreeId,
fs: &Arc<dyn Fs>,
live_data_start: u64,
) -> crate::Result<BlobFile> {
let mut file = fs.open(path, &crate::fs::FsOpenOptions::new().read(true))?;
let meta = {
let reader = crate::sfa::Reader::from_reader(&mut file)?;
let toc = reader.toc();
let metadata_section = toc.section(b"meta").ok_or_else(|| {
log::error!("meta section in blob file #{id} is missing (file may be corrupted)");
crate::Error::Unrecoverable
})?;
let metadata_len =
usize::try_from(metadata_section.len()).map_err(|_| crate::Error::Unrecoverable)?;
let metadata_slice = crate::file::read_exact(&*file, metadata_section.pos(), metadata_len)?;
Metadata::from_slice(&metadata_slice)?
};
let file: Arc<dyn crate::fs::FsFile> = Arc::from(file);
Ok(BlobFile(Arc::new(BlobFileInner {
id,
path: path.to_path_buf(),
meta,
is_deleted: AtomicBool::new(false),
punch_on_drop: portable_atomic::AtomicU64::new(u64::MAX),
checksum,
live_data_start,
file_accessor: FileAccessor::File(file),
tree_id,
fs: fs.clone(),
deletion_pause: once_cell::race::OnceBox::new(),
#[cfg(feature = "std")]
background_deleter: once_cell::race::OnceBox::new(),
})))
}
pub type BlobFileId = u64;
#[cfg(test)]
#[expect(clippy::unwrap_used, reason = "test code")]
mod tests;