use crate::content::node::{NodeState, PropertyState, PropertyValues};
use crate::content::property::{PropertyType, PropertyValue};
use crate::content::provider::SegmentProvider;
use crate::content::value::BinaryValue;
use crate::error::{Error, Result};
use crate::packed_records::SegmentInterner;
use crate::progress::{DiscardedProgress, ProgressObserver};
#[cfg(test)]
use crate::progress::{Step, WorkUnit};
use crate::segment::record::RecordIdentifier;
use crate::writer::record_writer::{
BulkBlockSharing, ChildNodesToWrite, PropertyToWrite, PropertyValuesToWrite, RecordWriter,
SegmentSink, sort_properties_for_template,
};
use crate::writer::segment_builder::GarbageCollectionGeneration;
#[cfg(test)]
use crate::writer::store_writer::{
ArchiveRewritePolicy, GenerationReclaimRequest, RETAINED_GENERATIONS, ReclaimRule,
WritableRepository,
};
mod gc_log;
mod memo;
#[cfg(test)]
mod test_support;
mod walk;
pub(crate) use gc_log::*;
pub(crate) use memo::*;
#[cfg(test)]
pub(crate) use test_support::*;
pub(crate) use walk::*;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CompactionKind {
Full,
Tail,
}
#[cfg(test)]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct CompactionOutcome {
pub size_before: u64,
pub size_after: u64,
pub compacted_nodes: u64,
}
pub fn deep_copy_tree<Sink: SegmentSink>(
source: &dyn SegmentProvider,
writer: &mut RecordWriter<Sink>,
source_root: RecordIdentifier,
) -> Result<(RecordIdentifier, u64)> {
deep_copy_tree_with_progress(source, writer, source_root, &mut DiscardedProgress)
}
pub fn deep_copy_tree_with_progress<Sink: SegmentSink>(
source: &dyn SegmentProvider,
writer: &mut RecordWriter<Sink>,
source_root: RecordIdentifier,
observer: &mut dyn ProgressObserver,
) -> Result<(RecordIdentifier, u64)> {
deep_copy_super_root_with_progress(
source,
writer,
source_root,
&std::collections::BTreeSet::new(),
observer,
)
}
pub fn deep_copy_tree_across_stores_with_progress<Sink: SegmentSink>(
source: &dyn SegmentProvider,
writer: &mut RecordWriter<Sink>,
source_root: RecordIdentifier,
observer: &mut dyn ProgressObserver,
) -> Result<(RecordIdentifier, u64)> {
deep_copy_super_root_sharing(
source,
writer,
source_root,
&std::collections::BTreeSet::new(),
BulkBlockSharing::AcrossStores,
observer,
)
}
pub fn deep_copy_super_root_with_progress<Sink: SegmentSink>(
source: &dyn SegmentProvider,
writer: &mut RecordWriter<Sink>,
super_root: RecordIdentifier,
omitted_checkpoints: &std::collections::BTreeSet<String>,
observer: &mut dyn ProgressObserver,
) -> Result<(RecordIdentifier, u64)> {
deep_copy_super_root_sharing(
source,
writer,
super_root,
omitted_checkpoints,
BulkBlockSharing::WithinOneStore,
observer,
)
}
pub fn deep_copy_super_root_sharing<Sink: SegmentSink>(
source: &dyn SegmentProvider,
writer: &mut RecordWriter<Sink>,
super_root: RecordIdentifier,
omitted_checkpoints: &std::collections::BTreeSet<String>,
bulk_sharing: BulkBlockSharing,
observer: &mut dyn ProgressObserver,
) -> Result<(RecordIdentifier, u64)> {
deep_copy_super_root_omitting_subtrees(
source,
writer,
super_root,
omitted_checkpoints,
&SubtreeOmissions {
omitted_subtree_records: &std::collections::HashSet::new(),
context_dependent_records: &std::collections::HashSet::new(),
},
bulk_sharing,
observer,
)
}
pub struct SubtreeOmissions<'omissions> {
pub omitted_subtree_records: &'omissions std::collections::HashSet<RecordIdentifier>,
pub context_dependent_records: &'omissions std::collections::HashSet<RecordIdentifier>,
}
pub fn deep_copy_super_root_omitting_subtrees<Sink: SegmentSink>(
source: &dyn SegmentProvider,
writer: &mut RecordWriter<Sink>,
super_root: RecordIdentifier,
omitted_checkpoints: &std::collections::BTreeSet<String>,
omissions: &SubtreeOmissions<'_>,
bulk_sharing: BulkBlockSharing,
observer: &mut dyn ProgressObserver,
) -> Result<(RecordIdentifier, u64)> {
let source_root = super_root;
let mut copier = Compactor {
source,
writer,
omitted_checkpoints,
omitted_subtree_records: omissions.omitted_subtree_records,
context_dependent_records: omissions.context_dependent_records,
scoped_rewrites: [
std::collections::HashMap::new(),
std::collections::HashMap::new(),
],
bulk_sharing,
segments: SegmentInterner::new(),
rewritten_nodes: RewrittenNodes::new(),
nodes_on_path: std::collections::HashSet::new(),
compacted_nodes: 0,
reported_nodes: 0,
observer,
};
let root = copier.compact_tree(source_root)?;
let memoized = copier.rewritten_nodes.occupied_slots();
let scoped: usize = copier
.scoped_rewrites
.iter()
.map(std::collections::HashMap::len)
.sum();
assert_eq!(
copier.compacted_nodes,
(memoized + scoped) as u64,
"copied node count diverged from the number of memoized nodes"
);
assert_eq!(
copier.rewritten_nodes.len, memoized,
"the memo's entry count diverged from its occupancy"
);
copier.observer.step_advanced(copier.compacted_nodes);
Ok((root, copier.compacted_nodes))
}
#[cfg(test)]
pub(crate) fn compact(
store: &mut WritableRepository,
kind: CompactionKind,
) -> Result<CompactionOutcome> {
compact_with_progress(store, kind, &mut DiscardedProgress)
}
#[cfg(test)]
pub(crate) fn compact_with_progress(
store: &mut WritableRepository,
kind: CompactionKind,
observer: &mut dyn ProgressObserver,
) -> Result<CompactionOutcome> {
let size_before = store.archive_size_on_disk()?;
let head = store.head();
let base_generation = store
.segment_generation(head.segment)
.ok_or(Error::SegmentNotFound {
segment_identifier: head.segment,
})?;
let target_generation = match kind {
CompactionKind::Full => GarbageCollectionGeneration {
generation: base_generation.generation.wrapping_add(1),
full_generation: base_generation.full_generation.wrapping_add(1),
is_compacted: true,
},
CompactionKind::Tail => GarbageCollectionGeneration {
generation: base_generation.generation.wrapping_add(1),
full_generation: base_generation.full_generation,
is_compacted: true,
},
};
let certified_sources = store.preflight_reclaim_sources_with_progress(observer)?;
let mut writer = store.record_writer_with_identifier(target_generation, "c");
let (new_head, compacted_nodes) = crate::progress::observe(
observer,
&Step::new("copying nodes into a fresh generation", WorkUnit::Nodes),
|observer| deep_copy_tree_with_progress(store, &mut writer, head, observer),
)?;
writer.finish()?;
if !store.compare_and_set_head(head, new_head) {
return Err(Error::InvalidFormat {
details: "the head moved during compaction".to_owned(),
});
}
store.flush()?;
crate::progress::observe(
observer,
&Step::new("reclaiming old generations", WorkUnit::Archives),
|_observer| {
store.reclaim_old_generations_with(GenerationReclaimRequest {
rule: ReclaimRule {
reference: target_generation,
kind,
retained_generations: RETAINED_GENERATIONS,
},
rewrite_policy: ArchiveRewritePolicy::EveryReclaimableArchive,
certified_sources: Some(&certified_sources),
expected: None,
})
},
)?;
rewrite_journal_to_head(store, new_head)?;
let size_after = store.archive_size_on_disk()?;
append_gc_log(
store,
size_after,
size_before.saturating_sub(size_after),
target_generation,
compacted_nodes,
new_head,
)?;
Ok(CompactionOutcome {
size_before,
size_after,
compacted_nodes,
})
}
#[cfg(test)]
pub(crate) fn rewrite_journal_to_head(
store: &WritableRepository,
head: RecordIdentifier,
) -> Result<()> {
use std::io::Write as _;
let timestamp = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |duration| duration.as_millis());
let line = format!(
"{}:{} root {timestamp}\n",
head.segment, head.record_number as i32
);
let journal_path = store.directory().join("journal.log");
let temporary_path = store.directory().join("journal.log.compacting");
{
let mut file = std::fs::File::create(&temporary_path)?;
file.write_all(line.as_bytes())?;
file.sync_all()?;
}
std::fs::rename(&temporary_path, &journal_path)?;
fsync_directory(store.directory());
store.reset_persisted_head(head)?;
Ok(())
}
pub(crate) fn fsync_directory(directory: &std::path::Path) {
if let Ok(handle) = std::fs::File::open(directory) {
let _ = handle.sync_all();
}
}
#[cfg(test)]
pub(crate) fn append_gc_log(
store: &WritableRepository,
repository_size: u64,
reclaimed_size: u64,
generation: GarbageCollectionGeneration,
compacted_nodes: u64,
root: RecordIdentifier,
) -> Result<()> {
let line = garbage_collection_log_entry(
repository_size,
reclaimed_size,
generation,
compacted_nodes,
root,
);
append_garbage_collection_log_entry(store.directory(), &line)
}
#[cfg(test)]
mod tests;