use super::{
ArchiveSet, Candidate, DiscardedProgress, Error, Path, ProgressObserver, RecordIdentifier,
RepositoryLock, Result, Step, WorkUnit, Write, collect_super_root_candidates,
is_fully_consistent, signed_uuid_key,
};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RecoveryOutcome {
pub recovered_head: RecordIdentifier,
pub candidates_examined: usize,
pub previous_journal_backup: Option<std::path::PathBuf>,
}
pub fn recover_journal(directory: &Path) -> Result<RecoveryOutcome> {
recover_journal_with_progress(directory, &mut DiscardedProgress)
}
pub fn recover_journal_with_progress(
directory: &Path,
observer: &mut dyn ProgressObserver,
) -> Result<RecoveryOutcome> {
let _repository_lock = RepositoryLock::acquire(directory)?;
let archives = crate::store::open_all_archives_with_progress(directory, observer)?;
let provider = ArchiveSet::new(archives);
let mut candidates = collect_super_root_candidates(&provider, observer);
let candidates_examined = candidates.len();
candidates.sort_by(|first, second| {
second
.timestamp_milliseconds
.cmp(&first.timestamp_milliseconds)
.then_with(|| signed_uuid_key(second.record).cmp(&signed_uuid_key(first.record)))
.then_with(|| second.record.record_number.cmp(&first.record.record_number))
});
let mut corrupt_memory: Vec<String> = Vec::new();
observer.step_began(
&Step::new("probing candidates for consistency", WorkUnit::Revisions)
.with_total(crate::progress::count(candidates_examined)),
);
let consistent_position = candidates
.iter()
.enumerate()
.position(|(probed, candidate)| {
observer.step_advanced(crate::progress::count(probed));
is_fully_consistent(&provider, candidate.record, &mut corrupt_memory)
});
observer.step_advanced(crate::progress::count(
consistent_position.map_or(candidates_examined, |position| position + 1),
));
observer.step_ended();
let consistent_position = consistent_position.ok_or_else(|| Error::InvalidFormat {
details: format!(
"no consistent super-root found among {candidates_examined} candidates in {}",
directory.display()
),
})?;
let survivors = &candidates[consistent_position..];
let recovered_head = survivors[0].record;
let previous_journal_backup = back_up_existing_journal(directory)?;
if let Err(error) = write_recovered_journal(directory, survivors) {
let temporary_path = directory.join("journal.log.recovered");
if temporary_path.exists() {
let _ = std::fs::remove_file(&temporary_path);
if let Some(backup_path) = &previous_journal_backup {
let _ = std::fs::remove_file(backup_path);
}
}
return Err(error);
}
Ok(RecoveryOutcome {
recovered_head,
candidates_examined,
previous_journal_backup,
})
}
pub(crate) fn back_up_existing_journal(directory: &Path) -> Result<Option<std::path::PathBuf>> {
let journal_path = directory.join("journal.log");
if !journal_path.exists() {
return Ok(None);
}
for counter in 0..1000 {
let backup = directory.join(format!("journal.log.bak.{counter:03}"));
if !backup.exists() {
std::fs::copy(&journal_path, &backup)?;
std::fs::File::open(&backup)?.sync_all()?;
return Ok(Some(backup));
}
}
Err(Error::InvalidFormat {
details: "all journal backup names (000-999) are taken".to_owned(),
})
}
pub(crate) fn write_recovered_journal(directory: &Path, survivors: &[Candidate]) -> Result<()> {
let temporary_path = directory.join("journal.log.recovered");
{
let mut file = std::fs::File::create(&temporary_path)?;
for candidate in survivors.iter().rev() {
let line = format!(
"{}:{} root {}\n",
candidate.record.segment,
candidate.record.record_number as i32,
candidate.timestamp_milliseconds
);
file.write_all(line.as_bytes())?;
}
file.sync_all()?;
}
std::fs::rename(&temporary_path, directory.join("journal.log"))?;
std::fs::File::open(directory)?.sync_all()?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::recover_journal;
use crate::writer::backup::test_support::{TestDirectory, assert_content, populate};
#[test]
fn recover_journal_rebuilds_a_deleted_journal() {
let directory = TestDirectory::new("recover");
populate(&directory.path);
std::fs::remove_file(directory.path.join("journal.log")).expect("remove journal");
let outcome = recover_journal(&directory.path).expect("recover");
assert!(outcome.candidates_examined >= 1);
assert_content(&directory.path, "Backup Source");
}
#[test]
fn recover_journal_backs_up_a_corrupt_journal() {
let directory = TestDirectory::new("recover-backup");
populate(&directory.path);
std::fs::write(
directory.path.join("journal.log"),
"garbage-with-no-space\n",
)
.expect("corrupt journal");
let outcome = recover_journal(&directory.path).expect("recover");
assert!(
outcome.previous_journal_backup.is_some(),
"the corrupt journal is backed up"
);
assert!(directory.path.join("journal.log.bak.000").exists());
assert_content(&directory.path, "Backup Source");
}
#[test]
fn recover_journal_requires_the_repository_lock() {
let directory = TestDirectory::new("recover-locked");
populate(&directory.path);
let held_lock =
crate::writer::repository_lock::RepositoryLock::acquire(&directory.path).expect("lock");
assert!(
recover_journal(&directory.path).is_err(),
"recovery must refuse to run while another process holds repo.lock"
);
drop(held_lock);
recover_journal(&directory.path).expect("recovery succeeds once the lock is free");
}
}