use super::file_identity::preserve_file_metadata;
use super::providers::ArchiveSegmentsProvider;
use super::providers::read_blob_identifiers;
use super::repair::{AuthorizeVersionTwoWrite, VersionTwoAlreadyEstablished};
use super::startup::install_target_generation;
use crate::content::provider::SegmentProvider as _;
use crate::error::{Error, Result};
use crate::segment::identifier::SegmentIdentifier;
use crate::segment::parsed_segment::ParsedSegment;
use crate::tar_archive::archive::TarArchiveReader;
use crate::tar_archive::file_name::ArchiveFileName;
use crate::writer::segment_builder::GarbageCollectionGeneration;
use crate::writer::tar_writer::TarArchiveWriter;
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::Arc;
pub(super) fn select_writable_generation(
directory: &Path,
generations: &[ArchiveFileName],
) -> (Option<TarArchiveReader>, bool) {
let mut any_nonempty = false;
for candidate in generations.iter().rev() {
let path = directory.join(&candidate.file_name);
if std::fs::metadata(&path).is_ok_and(|metadata| metadata.len() == 0) {
continue;
}
any_nonempty = true;
if let Ok(reader) = TarArchiveReader::open(&path)
&& !reader.is_recovered()
{
return (Some(reader), any_nonempty);
}
}
(None, any_nonempty)
}
pub(super) fn open_archive_numbers_for_writing(
directory: &Path,
by_number: std::collections::BTreeMap<u32, Vec<ArchiveFileName>>,
observer: &mut dyn crate::progress::ProgressObserver,
) -> Result<Vec<TarArchiveReader>> {
let archive_numbers = by_number.len();
let mut archives = Vec::new();
for (opened, (_, mut generations)) in by_number.into_iter().enumerate() {
observer.step_advanced(crate::progress::count(opened));
generations.sort_by_key(|name| name.file_generation);
let (winner, any_nonempty) = select_writable_generation(directory, &generations);
match winner {
Some(reader) => {
for stale in &generations {
if stale.file_name != reader.file_name() {
std::fs::remove_file(directory.join(&stale.file_name))?;
}
}
archives.push(reader);
}
None if !any_nonempty => {}
None => {
archives.push(recover_archive_number(
directory,
&generations,
&mut VersionTwoAlreadyEstablished,
)?);
}
}
}
observer.step_advanced(crate::progress::count(archive_numbers));
Ok(archives)
}
pub(super) fn unrecoverable_archive_number_refusal(generations: &[ArchiveFileName]) -> Error {
let names: Vec<&str> = generations
.iter()
.map(|generation| generation.file_name.as_str())
.collect();
Error::InvalidFormat {
details: format!(
"archive number {} has no valid index and no recoverable segment in {}; \
refusing to replace it with an empty archive. Cleanup preserves this file \
rather than removing it, so opening the store for writing needs it moved \
aside — keep it, it is the only copy of whatever it holds",
generations.first().map_or(0, |first| first.archive_number),
names.join(", ")
),
}
}
pub(super) fn inherit_replaced_archive_metadata(
directory: &Path,
target_name: &str,
temporary_path: &Path,
) -> Result<()> {
let Ok(source_metadata) = std::fs::metadata(directory.join(target_name)) else {
return Ok(());
};
let staged = std::fs::OpenOptions::new()
.write(true)
.open(temporary_path)?;
preserve_file_metadata(&staged, &source_metadata)
}
pub(super) fn authorize_before_install(
authorize: &mut dyn AuthorizeVersionTwoWrite,
temporary_path: &Path,
) -> Result<()> {
if let Err(error) = authorize.authorize() {
let _ = std::fs::remove_file(temporary_path);
return Err(error);
}
Ok(())
}
pub(super) fn recover_archive_number(
directory: &Path,
generations: &[ArchiveFileName],
authorize: &mut dyn AuthorizeVersionTwoWrite,
) -> Result<TarArchiveReader> {
let recovered = scan_recoverable_segments(directory, generations);
if recovered.is_empty() {
return Err(unrecoverable_archive_number_refusal(generations));
}
let mut parsed_segments: HashMap<SegmentIdentifier, Arc<ParsedSegment>> = HashMap::new();
for (identifier, bytes) in &recovered {
parsed_segments.insert(
*identifier,
Arc::new(ParsedSegment::parse(*identifier, bytes)?),
);
}
let provider = ArchiveSegmentsProvider {
segments: recovered
.iter()
.filter_map(|(identifier, bytes)| {
parsed_segments
.get(identifier)
.map(|parsed| (*identifier, (Arc::clone(parsed), bytes.as_slice())))
})
.collect(),
};
let target_name = &install_target_generation(directory, generations).file_name;
let temporary_name = format!("{target_name}.recovering");
let temporary_path = directory.join(&temporary_name);
let _ = std::fs::remove_file(&temporary_path);
let write_replacement =
|| -> Result<()> {
let mut writer = TarArchiveWriter::new(directory, &temporary_name);
for (identifier, bytes) in &recovered {
let (generation, references, binary_references) =
if let Some(parsed) = parsed_segments.get(identifier) {
let segment = provider.segment(*identifier)?;
let binary_references = read_blob_identifiers(&provider, &segment)
.map_err(|error| Error::InvalidFormat {
details: format!(
"cannot rebuild the binary references catalog while \
recovering {target_name}: an external blob identifier in \
segment {identifier} does not resolve within the recovered \
segments ({error}); refusing to publish an incomplete \
catalog, which could let blob garbage collection delete \
referenced binaries"
),
})?;
(
GarbageCollectionGeneration {
generation: parsed.generation,
full_generation: parsed.full_generation,
is_compacted: parsed.is_compacted,
},
parsed.referenced_segments.clone(),
binary_references,
)
} else {
(
GarbageCollectionGeneration {
generation: 0,
full_generation: 0,
is_compacted: false,
},
Vec::new(),
Vec::new(),
)
};
writer.write_segment(
*identifier,
bytes,
generation,
&references,
&binary_references,
)?;
}
writer.close()?;
Ok(())
};
if let Err(error) = write_replacement() {
let _ = std::fs::remove_file(&temporary_path);
return Err(error);
}
if let Err(error) = inherit_replaced_archive_metadata(directory, target_name, &temporary_path) {
let _ = std::fs::remove_file(&temporary_path);
return Err(error);
}
crate::writer::compaction::fsync_directory(directory);
match TarArchiveReader::open(&temporary_path) {
Ok(validated) if !validated.is_recovered() => drop(validated),
Ok(_) => {
let _ = std::fs::remove_file(&temporary_path);
return Err(Error::InvalidFormat {
details: format!("the rebuilt archive {temporary_name} failed index validation"),
});
}
Err(error) => {
let _ = std::fs::remove_file(&temporary_path);
return Err(error);
}
}
authorize_before_install(authorize, &temporary_path)?;
install_recovered_archive(directory, generations, target_name, &temporary_path)
}
pub(super) fn any_recoverable_segment(directory: &Path, generations: &[ArchiveFileName]) -> bool {
generations.iter().any(|generation| {
TarArchiveReader::open(&directory.join(&generation.file_name))
.is_ok_and(|reader| reader.segment_count() > 0)
})
}
pub(super) fn scan_recoverable_segments(
directory: &Path,
generations: &[ArchiveFileName],
) -> Vec<(SegmentIdentifier, Vec<u8>)> {
let mut recovered: Vec<(SegmentIdentifier, Vec<u8>)> = Vec::new();
let mut positions: HashMap<SegmentIdentifier, usize> = HashMap::new();
for generation in generations {
let path = directory.join(&generation.file_name);
if let Ok(reader) = TarArchiveReader::open(&path) {
for identifier in reader.segment_identifiers() {
if let Some(bytes) = reader.segment_data(identifier) {
if let Some(&position) = positions.get(&identifier) {
recovered[position].1 = bytes.to_vec();
} else {
positions.insert(identifier, recovered.len());
recovered.push((identifier, bytes.to_vec()));
}
}
}
}
}
recovered
}
pub(super) fn install_recovered_archive(
directory: &Path,
generations: &[ArchiveFileName],
target_name: &str,
temporary_path: &Path,
) -> Result<TarArchiveReader> {
let mut renamed: Vec<(PathBuf, PathBuf)> = Vec::new();
let mut target_backup: Option<PathBuf> = None;
let roll_back = |renamed: &[(PathBuf, PathBuf)]| {
for (original, backup) in renamed.iter().rev() {
let _ = std::fs::rename(backup, original);
}
};
for generation in generations {
let path = directory.join(&generation.file_name);
if std::fs::metadata(&path).is_ok_and(|metadata| metadata.len() == 0) {
continue;
}
let backup = backup_path(directory, &generation.file_name);
if generation.file_name == *target_name {
if std::fs::hard_link(&path, &backup).is_err()
&& let Err(error) = std::fs::copy(&path, &backup)
{
roll_back(&renamed);
return Err(error.into());
}
target_backup = Some(backup);
} else if let Err(error) = std::fs::rename(&path, &backup) {
roll_back(&renamed);
return Err(error.into());
} else {
renamed.push((path, backup));
}
}
let target_path = directory.join(target_name);
if let Err(error) = std::fs::rename(temporary_path, &target_path) {
if let Some(backup) = &target_backup {
let _ = std::fs::remove_file(backup);
}
roll_back(&renamed);
return Err(error.into());
}
crate::writer::compaction::fsync_directory(directory);
match TarArchiveReader::open(&target_path) {
Ok(reader) => Ok(reader),
Err(error) => {
if let Some(backup) = &target_backup {
let _ = std::fs::rename(backup, &target_path);
}
roll_back(&renamed);
crate::writer::compaction::fsync_directory(directory);
Err(error)
}
}
}
pub(super) fn backup_path(directory: &Path, file_name: &str) -> PathBuf {
let first = directory.join(format!("{file_name}.bak"));
if !first.exists() {
return first;
}
let mut counter = 2u32;
loop {
let candidate = directory.join(format!("{file_name}.{counter}.bak"));
if !candidate.exists() {
return candidate;
}
counter += 1;
}
}
#[cfg(test)]
mod tests {
use crate::content::provider::SegmentProvider;
use crate::store::Repository;
use crate::writer::store_writer::repository::*;
use crate::writer::store_writer::test_support::*;
#[test]
fn an_empty_archive_file_does_not_break_opening_for_writing() {
let directory = TestDirectory::new("empty-archive-open");
{
let store = WritableRepository::open(&directory.path).expect("bootstrap");
store.close().expect("close");
}
let empty = directory.path.join("data00009a.tar");
std::fs::write(&empty, b"").expect("create the empty archive");
let store = WritableRepository::open(&directory.path).expect("write open must succeed");
let head = store.head();
assert!(store.segment(head.segment).is_ok(), "head still resolves");
store.close().expect("close");
Repository::open(&directory.path).expect("the reader still opens");
}
#[test]
fn an_unrecoverable_archive_is_refused_with_its_file_name() {
let directory = TestDirectory::new("unrecoverable-archive");
{
let store = WritableRepository::open(&directory.path).expect("bootstrap");
store.close().expect("close");
}
let junk = directory.path.join("data00009a.tar");
std::fs::write(&junk, vec![0x5au8; 4096]).expect("write junk archive");
let message = match WritableRepository::open(&directory.path) {
Ok(_) => panic!("opening for writing must refuse an unrecoverable archive"),
Err(error) => error.to_string(),
};
assert!(
message.contains("data00009a.tar"),
"the refusal names the unusable file: {message}"
);
assert!(
message.contains("no recoverable segment"),
"the refusal states why: {message}"
);
assert!(
!message.contains("No such file or directory"),
"the refusal must not surface a bare errno: {message}"
);
assert!(junk.exists(), "the refusal leaves the file in place");
}
#[test]
fn archives_without_an_index_are_recovered_with_backups() {
let directory = TestDirectory::new("write-recovery");
{
let store = WritableRepository::open(&directory.path).expect("bootstrap");
store.close().expect("close");
}
let path = directory.path.join("data00000a.tar");
let full = std::fs::read(&path).expect("read");
let trailer_start = full
.windows(4)
.position(|window| window == b".brf")
.map(|position| (position / 512) * 512)
.expect("brf trailer present");
let mut truncated = full[..trailer_start].to_vec();
truncated.extend_from_slice(&[0u8; 1024]);
std::fs::write(&path, &truncated).expect("truncate");
{
let store = WritableRepository::open(&directory.path).expect("recovering open");
let head = store.head();
assert!(
store.segment(head.segment).is_ok(),
"head segment recovered"
);
store.close().expect("close");
}
assert!(
directory.path.join("data00000a.tar.bak").exists(),
"the damaged archive is backed up"
);
let repository = Repository::open(&directory.path).expect("reader opens");
assert!(
!repository
.archives()
.iter()
.any(crate::tar_archive::archive::TarArchiveReader::is_recovered),
"the regenerated archive has a valid index"
);
repository.content_root().expect("content root resolves");
}
}