use std::collections::HashSet;
use std::io::Write;
use std::path::Path;
use crate::content::node::NodeState;
use crate::content::provider::SegmentProvider;
use crate::error::{Error, Result};
use crate::segment::identifier::SegmentIdentifier;
use crate::segment::record::RecordIdentifier;
use crate::store::Repository;
use crate::writer::compaction::deep_copy_tree;
use crate::writer::repository_lock::RepositoryLock;
use crate::writer::segment_builder::GarbageCollectionGeneration;
use crate::writer::store_writer::WritableRepository;
pub fn backup(source_directory: &Path, target_directory: &Path) -> Result<()> {
let source = Repository::open(source_directory)?;
let target = WritableRepository::open(target_directory)?;
copy_head_between(&source, source.head_record_identifier(), &target, "b")?;
target.close()
}
pub fn restore(backup_directory: &Path, target_directory: &Path) -> Result<()> {
let backup = Repository::open(backup_directory)?;
let target = WritableRepository::open(target_directory)?;
copy_head_between(&backup, backup.head_record_identifier(), &target, "r")?;
target.close()
}
fn copy_head_between(
source: &Repository,
source_head: RecordIdentifier,
target: &WritableRepository,
writer_identifier: &str,
) -> Result<()> {
let generation = source_head_generation(source, source_head)?;
let mut writer = target.record_writer_with_identifier(generation, writer_identifier);
let (new_head, _) = deep_copy_tree(source, &mut writer, source_head)?;
writer.finish()?;
target.replace_head(new_head);
target.flush()
}
fn source_head_generation(
source: &Repository,
source_head: RecordIdentifier,
) -> Result<GarbageCollectionGeneration> {
for archive in source.archives() {
if let Some(entry) = archive.index_entry(source_head.segment) {
return Ok(GarbageCollectionGeneration {
generation: entry.generation,
full_generation: entry.full_generation,
is_compacted: entry.is_compacted,
});
}
}
let view = source.segment(source_head.segment)?;
Ok(GarbageCollectionGeneration {
generation: view.structure.generation,
full_generation: view.structure.full_generation,
is_compacted: view.structure.is_compacted,
})
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RecoveryOutcome {
pub recovered_head: RecordIdentifier,
pub candidates_examined: usize,
pub previous_journal_backup: Option<std::path::PathBuf>,
}
struct Candidate {
record: RecordIdentifier,
timestamp_milliseconds: i64,
}
pub fn recover_journal(directory: &Path) -> Result<RecoveryOutcome> {
let _repository_lock = RepositoryLock::acquire(directory)?;
let archives = crate::store::open_all_archives(directory)?;
let mut candidates: Vec<Candidate> = Vec::new();
let mut seen: HashSet<RecordIdentifier> = HashSet::new();
let provider = crate::store::ArchiveSet::new(archives);
for segment_identifier in provider.segment_identifiers() {
if segment_identifier.is_bulk_segment() {
continue;
}
let Ok(view) = provider.segment(segment_identifier) else {
continue;
};
let Some(timestamp) = read_segment_info_timestamp(&provider, segment_identifier) else {
continue;
};
let node_records: Vec<u32> = view
.structure
.record_table()
.iter()
.filter(|entry| entry.record_type() == Some(crate::segment::record::RecordType::Node))
.map(|entry| entry.record_number)
.collect();
for record_number in node_records {
let record = RecordIdentifier::new(segment_identifier, record_number);
if !seen.insert(record) {
continue;
}
if is_super_root(&provider, record) {
candidates.push(Candidate {
record,
timestamp_milliseconds: timestamp,
});
}
}
}
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();
let consistent_position = candidates
.iter()
.position(|candidate| is_fully_consistent(&provider, candidate.record, &mut corrupt_memory))
.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,
})
}
fn is_super_root(provider: &dyn SegmentProvider, record: RecordIdentifier) -> bool {
let node = NodeState::new(provider, record);
matches!(node.child_node("root"), Ok(Some(_)))
&& matches!(node.child_node("checkpoints"), Ok(Some(_)))
}
fn signed_uuid_key(record: RecordIdentifier) -> (i64, i64) {
(
record.segment.most_significant_bits as i64,
record.segment.least_significant_bits as i64,
)
}
const MAXIMUM_RECOVERY_DEPTH: usize = 4000;
fn is_fully_consistent(
provider: &dyn SegmentProvider,
record: RecordIdentifier,
corrupt_memory: &mut Vec<String>,
) -> bool {
let super_root = NodeState::new(provider, record);
let mut visited = HashSet::new();
let Ok(Some(content_root)) = super_root.child_node("root") else {
return false;
};
if !probe_corrupt_paths(provider, &content_root, corrupt_memory) {
return false;
}
if let Err(corrupt_path) = verify_tree(provider, content_root.record_identifier(), &mut visited)
{
remember_corrupt(corrupt_memory, corrupt_path);
return false;
}
let Ok(Some(checkpoints)) = super_root.child_node("checkpoints") else {
return false;
};
let Ok(checkpoint_entries) = checkpoints.child_node_entries() else {
return false;
};
for (_, checkpoint) in checkpoint_entries {
let Ok(Some(snapshot_root)) = checkpoint.child_node("root") else {
return false;
};
if !probe_corrupt_paths(provider, &snapshot_root, corrupt_memory) {
return false;
}
if let Err(corrupt_path) =
verify_tree(provider, snapshot_root.record_identifier(), &mut visited)
{
remember_corrupt(corrupt_memory, corrupt_path);
return false;
}
}
true
}
fn probe_corrupt_paths(
provider: &dyn SegmentProvider,
tree_root: &NodeState<'_>,
corrupt_memory: &[String],
) -> bool {
for corrupt_path in corrupt_memory {
let Ok(Some(corrupt_node)) = resolve_descendant(tree_root, corrupt_path) else {
return false;
};
if check_node_shallow(provider, corrupt_node.record_identifier()).is_err() {
return false;
}
}
true
}
fn resolve_descendant<'provider>(
node: &NodeState<'provider>,
relative_path: &str,
) -> Result<Option<NodeState<'provider>>> {
let mut current = *node;
for name in relative_path
.split('/')
.filter(|segment| !segment.is_empty())
{
match current.child_node(name)? {
Some(child) => current = child,
None => return Ok(None),
}
}
Ok(Some(current))
}
fn remember_corrupt(corrupt_memory: &mut Vec<String>, corrupt_path: String) {
if !corrupt_memory.contains(&corrupt_path) {
corrupt_memory.push(corrupt_path);
}
}
fn check_node_shallow(provider: &dyn SegmentProvider, record: RecordIdentifier) -> Result<()> {
let node = NodeState::new(provider, record);
for property in node.properties()? {
match &property.values {
crate::content::node::PropertyValues::Single(value) => {
verify_inline_binary(provider, value)?;
}
crate::content::node::PropertyValues::Multiple(values) => {
for value in values {
verify_inline_binary(provider, value)?;
}
}
}
}
Ok(())
}
fn verify_tree(
provider: &dyn SegmentProvider,
record: RecordIdentifier,
visited: &mut HashSet<RecordIdentifier>,
) -> std::result::Result<(), String> {
fn walk(
provider: &dyn SegmentProvider,
record: RecordIdentifier,
visited: &mut HashSet<RecordIdentifier>,
ancestors: &mut HashSet<RecordIdentifier>,
depth: usize,
path: &mut String,
) -> std::result::Result<(), String> {
if depth > MAXIMUM_RECOVERY_DEPTH {
return Err(path.clone());
}
if ancestors.contains(&record) {
return Err(path.clone());
}
if !visited.insert(record) {
return Ok(());
}
ancestors.insert(record);
check_node_shallow(provider, record).map_err(|_| path.clone())?;
let node = NodeState::new(provider, record);
let children = node.child_node_entries().map_err(|_| path.clone())?;
for (name, child) in children {
let parent_length = path.len();
path.push('/');
path.push_str(&name);
walk(
provider,
child.record_identifier(),
visited,
ancestors,
depth + 1,
path,
)?;
path.truncate(parent_length);
}
ancestors.remove(&record);
Ok(())
}
let mut ancestors = HashSet::new();
let mut path = String::new();
walk(provider, record, visited, &mut ancestors, 0, &mut path)
}
fn verify_inline_binary(
provider: &dyn SegmentProvider,
value: &crate::content::property::PropertyValue,
) -> Result<()> {
use crate::content::property::PropertyValue;
use crate::content::value::BinaryValue;
if let PropertyValue::Binary(BinaryValue::Inline {
record_identifier, ..
}) = value
{
crate::content::value::verify_binary_content(provider, *record_identifier)?;
}
Ok(())
}
fn read_segment_info_timestamp(
provider: &dyn SegmentProvider,
segment: SegmentIdentifier,
) -> Option<i64> {
let view = provider.segment(segment).ok()?;
let first_record = view.structure.record_table().first()?.record_number;
let info =
crate::content::value::read_string(provider, RecordIdentifier::new(segment, first_record))
.ok()?;
parse_info_timestamp(&info)
}
fn parse_info_timestamp(info: &str) -> Option<i64> {
let marker = "\"t\":";
let start = info.find(marker)? + marker.len();
let rest = &info[start..];
let end = rest
.find(|character: char| !character.is_ascii_digit() && character != '-')
.unwrap_or(rest.len());
rest[..end].parse().ok()
}
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(),
})
}
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::{backup, parse_info_timestamp, recover_journal, restore};
use crate::content::node::PropertyValues;
use crate::content::property::PropertyValue;
use crate::store::Repository;
use crate::writer::commit::create_checkpoint;
use crate::writer::record_writer::{ChildNodesToWrite, PropertyToWrite, PropertyValuesToWrite};
use crate::writer::store_writer::WritableRepository;
struct TestDirectory {
path: std::path::PathBuf,
}
impl TestDirectory {
fn new(name: &str) -> Self {
let path =
std::env::temp_dir().join(format!("froe-backup-{name}-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&path);
Self { path }
}
}
impl Drop for TestDirectory {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.path);
}
}
fn populate(directory: &std::path::Path) {
let store = WritableRepository::open(directory).expect("bootstrap");
let generation = store.writing_generation().expect("generation");
let mut writer = store.record_writer(generation);
let alpha = writer
.write_node(Some("nt:unstructured"), &[], &ChildNodesToWrite::Zero, &[])
.expect("alpha");
let beta = writer
.write_node(Some("nt:unstructured"), &[], &ChildNodesToWrite::Zero, &[])
.expect("beta");
let title = writer.write_string("Backup Source").expect("value");
let content = writer
.write_node(
Some("nt:unstructured"),
&[],
&ChildNodesToWrite::Many(vec![
("alpha".to_owned(), alpha),
("beta".to_owned(), beta),
]),
&[PropertyToWrite {
name: "title".to_owned(),
property_type: crate::content::property::PropertyType::String,
values: PropertyValuesToWrite::Single(title),
}],
)
.expect("content");
let root = writer
.write_node(
None,
&[],
&ChildNodesToWrite::One {
name: "content".to_owned(),
node: content,
},
&[],
)
.expect("root");
let head = writer
.write_node(
None,
&[],
&ChildNodesToWrite::One {
name: "root".to_owned(),
node: root,
},
&[],
)
.expect("super root");
writer.finish().expect("finish");
let previous = store.head();
assert!(store.set_head(previous, head));
create_checkpoint(&store, 10_000_000, &[]).expect("checkpoint");
store.close().expect("close");
}
fn assert_content(directory: &std::path::Path, expected_title: &str) {
let repository = Repository::open(directory).expect("reader");
let content = repository
.node_at_path("/content")
.expect("resolve")
.expect("present");
assert_eq!(content.child_node_count().expect("count"), 2);
assert_eq!(
content
.property("title")
.expect("read")
.expect("present")
.values,
PropertyValues::Single(PropertyValue::String(expected_title.to_owned()))
);
}
fn write_revision_with_children(directory: &std::path::Path, child_count: usize) {
let store = WritableRepository::open(directory).expect("open");
let generation = store.writing_generation().expect("generation");
let mut writer = store.record_writer(generation);
let mut children = Vec::new();
for index in 0..child_count {
let child = writer
.write_node(Some("nt:unstructured"), &[], &ChildNodesToWrite::Zero, &[])
.expect("child");
children.push((format!("child-{index}"), child));
}
let child_structure = if children.is_empty() {
ChildNodesToWrite::Zero
} else {
ChildNodesToWrite::Many(children)
};
let content = writer
.write_node(Some("nt:unstructured"), &[], &child_structure, &[])
.expect("content");
let root = writer
.write_node(
None,
&[],
&ChildNodesToWrite::One {
name: "content".to_owned(),
node: content,
},
&[],
)
.expect("root");
let head = writer
.write_node(
None,
&[],
&ChildNodesToWrite::One {
name: "root".to_owned(),
node: root,
},
&[],
)
.expect("super root");
writer.finish().expect("finish");
let previous = store.head();
assert!(store.set_head(previous, head));
create_checkpoint(&store, 10_000_000, &[]).expect("checkpoint");
store.close().expect("close");
}
#[test]
fn recovery_picks_the_newest_revision_even_when_it_has_fewer_nodes() {
let directory = TestDirectory::new("recover-newest");
write_revision_with_children(&directory.path, 8); std::thread::sleep(std::time::Duration::from_millis(5));
write_revision_with_children(&directory.path, 2);
let newest_head = {
let repository = Repository::open(&directory.path).expect("reader");
repository.head_record_identifier()
};
std::fs::remove_file(directory.path.join("journal.log")).expect("remove journal");
let outcome = recover_journal(&directory.path).expect("recover");
assert_eq!(
outcome.recovered_head, newest_head,
"recovery returns the newest revision, not the larger older one"
);
let repository = Repository::open(&directory.path).expect("reader");
let content = repository
.node_at_path("/content")
.expect("resolve")
.expect("present");
assert_eq!(content.child_node_count().expect("count"), 2);
}
#[test]
fn backup_copies_content_and_checkpoints() {
let source = TestDirectory::new("backup-source");
let target = TestDirectory::new("backup-target");
populate(&source.path);
backup(&source.path, &target.path).expect("backup");
assert_content(&target.path, "Backup Source");
let repository = Repository::open(&target.path).expect("reader");
assert_eq!(
repository.checkpoints().expect("checkpoints").len(),
1,
"checkpoints are copied with the head"
);
assert_content(&source.path, "Backup Source");
}
#[test]
fn restore_overwrites_the_target_head() {
let backup_store = TestDirectory::new("restore-backup");
let target = TestDirectory::new("restore-target");
populate(&backup_store.path);
{
let store = WritableRepository::open(&target.path).expect("bootstrap target");
store.close().expect("close");
}
restore(&backup_store.path, &target.path).expect("restore");
assert_content(&target.path, "Backup Source");
}
#[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");
}
#[test]
fn backup_stamps_the_source_head_generation_verbatim() {
let source = TestDirectory::new("backup-generation-source");
let target = TestDirectory::new("backup-generation-target");
populate(&source.path);
{
let mut store = WritableRepository::open(&source.path).expect("open");
crate::writer::compaction::compact(
&mut store,
crate::writer::compaction::CompactionKind::Full,
)
.expect("compact");
store.close().expect("close");
}
let source_head_generation = {
let repository = Repository::open(&source.path).expect("reader");
let head_segment = repository.head_record_identifier().segment;
repository
.archives()
.iter()
.find_map(|archive| archive.index_entry(head_segment))
.map(|entry| (entry.generation, entry.full_generation, entry.is_compacted))
.expect("head segment is indexed")
};
assert_ne!(
source_head_generation.0, 0,
"compaction advanced the source generation"
);
backup(&source.path, &target.path).expect("backup");
let repository = Repository::open(&target.path).expect("target reader");
let head_segment = repository.head_record_identifier().segment;
let stamped = repository
.archives()
.iter()
.find_map(|archive| archive.index_entry(head_segment))
.map(|entry| (entry.generation, entry.full_generation, entry.is_compacted))
.expect("target head segment is indexed");
assert_eq!(
stamped, source_head_generation,
"the source head's generation triple is stamped verbatim"
);
}
#[test]
fn parses_segment_info_timestamps() {
assert_eq!(
parse_info_timestamp("{\"wid\":\"froe\",\"sno\":3,\"t\":1700000000000}"),
Some(1_700_000_000_000)
);
assert_eq!(parse_info_timestamp("{\"wid\":\"x\"}"), None);
}
}