use super::{
BinaryValue, BulkBlockSharing, ChildNodesToWrite, Error, NodeState, ProgressObserver,
PropertyState, PropertyToWrite, PropertyType, PropertyValue, PropertyValues,
PropertyValuesToWrite, RecordIdentifier, RecordWriter, Result, RewrittenNodes, SegmentInterner,
SegmentProvider, SegmentSink, sort_properties_for_template,
};
pub(crate) const COPIED_NODE_REPORT_STRIDE: u64 = 512;
pub(crate) struct Compactor<'writer, Sink: SegmentSink> {
pub(crate) source: &'writer dyn SegmentProvider,
pub(crate) writer: &'writer mut RecordWriter<Sink>,
pub(crate) bulk_sharing: BulkBlockSharing,
pub(crate) omitted_checkpoints: &'writer std::collections::BTreeSet<String>,
pub(crate) omitted_subtree_records:
&'writer std::collections::HashSet<crate::segment::record::RecordIdentifier>,
pub(crate) context_dependent_records:
&'writer std::collections::HashSet<crate::segment::record::RecordIdentifier>,
pub(crate) scoped_rewrites: [std::collections::HashMap<u64, u64>; 2],
pub(crate) segments: SegmentInterner,
pub(crate) rewritten_nodes: RewrittenNodes,
pub(crate) nodes_on_path: std::collections::HashSet<u64>,
pub(crate) compacted_nodes: u64,
pub(crate) reported_nodes: u64,
pub(crate) observer: &'writer mut dyn ProgressObserver,
}
pub(crate) enum Entered {
Fresh(CompactionFrame),
Memoized(RecordIdentifier),
}
pub(crate) struct CompactionFrame {
pub(crate) source: RecordIdentifier,
pub(crate) packed: u64,
pub(crate) name_in_parent: Option<String>,
pub(crate) within_checkpoints: bool,
pub(crate) pending_children: Vec<(String, RecordIdentifier)>,
pub(crate) rewritten_children: Vec<(String, RecordIdentifier)>,
}
impl<Sink: SegmentSink> Compactor<'_, Sink> {
pub(crate) fn compact_tree(
&mut self,
source_root: RecordIdentifier,
) -> Result<RecordIdentifier> {
let mut stack = match self.enter(source_root, false)? {
Entered::Fresh(root) => vec![root],
Entered::Memoized(rewritten) => return Ok(rewritten),
};
loop {
let next = stack
.last_mut()
.expect("the loop returns before the stack empties")
.pending_children
.pop();
if let Some((name, child)) = next {
let parent = stack.last().expect("the parent frame is on the stack");
let within_checkpoints =
parent.within_checkpoints || (stack.len() == 1 && name == "checkpoints");
if !within_checkpoints && self.omitted_subtree_records.contains(&child) {
continue;
}
match self.enter(child, within_checkpoints)? {
Entered::Fresh(mut frame) => {
if stack.len() == 1
&& name == "checkpoints"
&& !self.omitted_checkpoints.is_empty()
{
frame.pending_children.retain(|(checkpoint, _)| {
!self.omitted_checkpoints.contains(checkpoint)
});
}
frame.name_in_parent = Some(name);
frame.within_checkpoints = within_checkpoints;
stack.push(frame);
}
Entered::Memoized(rewritten) => stack
.last_mut()
.expect("the parent frame is still on the stack")
.rewritten_children
.push((name, rewritten)),
}
continue;
}
let finished = stack.pop().expect("a frame was just inspected");
let rewritten = self.emit(
finished.source,
finished.packed,
finished.within_checkpoints,
finished.rewritten_children,
)?;
match stack.last_mut() {
Some(parent) => parent.rewritten_children.push((
finished
.name_in_parent
.expect("only the root frame has no name"),
rewritten,
)),
None => return Ok(rewritten),
}
}
}
fn scope_index(within_checkpoints: bool) -> usize {
usize::from(within_checkpoints)
}
pub(crate) fn enter(
&mut self,
source_node: RecordIdentifier,
within_checkpoints: bool,
) -> Result<Entered> {
let packed = self.segments.pack(source_node);
if self.context_dependent_records.contains(&source_node) {
if let Some(rewritten) = self.scoped_rewrites[Self::scope_index(within_checkpoints)]
.get(&packed)
.copied()
{
return Ok(Entered::Memoized(self.segments.unpack(rewritten)));
}
} else if let Some(rewritten) = self.rewritten_nodes.get(packed) {
return Ok(Entered::Memoized(self.segments.unpack(rewritten)));
}
if !self.nodes_on_path.insert(packed) {
return Err(Error::InvalidFormat {
details: format!(
"node record {source_node} is contained in its own subtree; \
the source records form a cycle"
),
});
}
let node = NodeState::new(self.source, source_node);
let mut pending_children: Vec<(String, RecordIdentifier)> = node
.child_node_entries()?
.into_iter()
.map(|(name, child)| (name, child.record_identifier()))
.collect();
pending_children.reverse();
Ok(Entered::Fresh(CompactionFrame {
source: source_node,
packed,
name_in_parent: None,
within_checkpoints: false,
pending_children,
rewritten_children: Vec::new(),
}))
}
pub(crate) fn emit(
&mut self,
source_node: RecordIdentifier,
packed_source: u64,
within_checkpoints: bool,
mut child_entries: Vec<(String, RecordIdentifier)>,
) -> Result<RecordIdentifier> {
let node = NodeState::new(self.source, source_node);
let template = node.template()?;
let stable_identifier = node.stable_identifier_bytes()?;
let children = match child_entries.len() {
0 => ChildNodesToWrite::Zero,
1 => {
let (name, node) = child_entries.pop().expect("one child");
ChildNodesToWrite::One { name, node }
}
_ => ChildNodesToWrite::Many(child_entries),
};
let mut properties = Vec::new();
for property in node.stored_properties()? {
properties.push(self.rewrite_property(&property)?);
}
sort_properties_for_template(&mut properties);
let rewritten = self.writer.write_node_with_stable_identifier(
template.primary_type.as_deref(),
&template.mixin_types,
&children,
&properties,
Some(stable_identifier),
)?;
self.nodes_on_path.remove(&packed_source);
let packed_rewritten = self.segments.pack(rewritten);
if self.context_dependent_records.contains(&source_node) {
self.scoped_rewrites[Self::scope_index(within_checkpoints)]
.insert(packed_source, packed_rewritten);
} else {
self.rewritten_nodes.insert(packed_source, packed_rewritten);
}
self.compacted_nodes += 1;
if self.compacted_nodes - self.reported_nodes >= COPIED_NODE_REPORT_STRIDE {
self.reported_nodes = self.compacted_nodes;
self.observer.step_advanced(self.compacted_nodes);
}
Ok(rewritten)
}
pub(crate) fn rewrite_property(&mut self, property: &PropertyState) -> Result<PropertyToWrite> {
let values = match &property.values {
PropertyValues::Single(value) => {
PropertyValuesToWrite::Single(self.rewrite_value(property.property_type, value)?)
}
PropertyValues::Multiple(values) => {
let mut rewritten = Vec::with_capacity(values.len());
for value in values {
rewritten.push(self.rewrite_value(property.property_type, value)?);
}
PropertyValuesToWrite::Multiple(rewritten)
}
};
Ok(PropertyToWrite {
name: property.name.clone(),
property_type: property.property_type,
values,
})
}
pub(crate) fn rewrite_value(
&mut self,
property_type: PropertyType,
value: &PropertyValue,
) -> Result<RecordIdentifier> {
if property_type == PropertyType::Binary {
return match value {
PropertyValue::Binary(BinaryValue::External { blob_identifier }) => self
.writer
.write_external_binary_identifier(blob_identifier),
PropertyValue::Binary(BinaryValue::Inline {
record_identifier, ..
}) => {
self.writer.copy_binary_value(
self.source,
*record_identifier,
self.bulk_sharing,
)
}
_ => Err(Error::InvalidFormat {
details: "binary property did not decode to a binary value".to_owned(),
}),
};
}
let text = value.as_text().ok_or_else(|| Error::InvalidFormat {
details: format!("property value {value:?} has no string form"),
})?;
self.writer.write_string(&text)
}
}
#[cfg(test)]
mod tests {
use crate::content::provider::SegmentProvider;
use crate::store::Repository;
use crate::writer::compaction::test_support::*;
use crate::writer::compaction::{
deep_copy_super_root_with_progress, deep_copy_tree_with_progress,
};
use crate::writer::record_writer::ChildNodesToWrite;
use crate::writer::store_writer::WritableRepository;
#[test]
fn a_copy_that_omits_a_checkpoint_reproduces_every_other_child_exactly() {
let source = TestDirectory::new("omit-checkpoint-source");
let names = build_store_with_checkpoints(&source);
let dropped = names[1].clone();
let omitting = TestDirectory::new("omit-checkpoint-copy");
copy_repository(&source.path, &omitting.path);
let removing = TestDirectory::new("omit-checkpoint-removal");
copy_repository(&source.path, &removing.path);
let (omitted_head, omitted_nodes) = {
let store = WritableRepository::open(&omitting.path).expect("open the omitting store");
let generation = store.writing_generation().expect("generation");
let mut writer = store.record_writer(generation);
let (head, nodes) = deep_copy_super_root_with_progress(
&store,
&mut writer,
store.head(),
&std::collections::BTreeSet::from([dropped.clone()]),
&mut crate::progress::DiscardedProgress,
)
.expect("copy while omitting one checkpoint");
writer.finish().expect("finish");
assert!(store.compare_and_set_head(store.head(), head));
store.flush().expect("flush");
store.close().expect("close the omitting store");
(head, nodes)
};
let (removed_head, removed_nodes) = {
let store = WritableRepository::open(&removing.path).expect("open the removing store");
crate::writer::commit::remove_checkpoints(&store, std::slice::from_ref(&dropped))
.expect("remove the checkpoint from the live head");
let generation = store.writing_generation().expect("generation");
let mut writer = store.record_writer(generation);
let (head, nodes) = deep_copy_tree_with_progress(
&store,
&mut writer,
store.head(),
&mut crate::progress::DiscardedProgress,
)
.expect("copy the already-reduced head");
writer.finish().expect("finish");
assert!(store.compare_and_set_head(store.head(), head));
store.flush().expect("flush");
store.close().expect("close the removing store");
(head, nodes)
};
let omitted = Repository::open(&omitting.path).expect("reopen the omitting store");
let removed = Repository::open(&removing.path).expect("reopen the removing store");
let omitted_root = omitted.node(omitted_head);
let removed_root = removed.node(removed_head);
assert_eq!(
child_names(&omitted_root),
child_names(&removed_root),
"the super-root carries the same children either way"
);
let omitted_checkpoints = omitted_root
.child_node("checkpoints")
.expect("read")
.expect("present");
let removed_checkpoints = removed_root
.child_node("checkpoints")
.expect("read")
.expect("present");
let mut surviving = child_names(&omitted_checkpoints);
let mut expected = child_names(&removed_checkpoints);
surviving.sort();
expected.sort();
assert_eq!(
surviving, expected,
"exactly the same checkpoints survive both mechanisms"
);
assert!(
!surviving.contains(&dropped),
"the retired checkpoint is gone: {surviving:?}"
);
assert_eq!(
omitted_nodes, removed_nodes,
"both mechanisms copy the same tree, so the same number of distinct nodes"
);
for name in &surviving {
for (label, container) in [
("omitted", &omitted_checkpoints),
("removed", &removed_checkpoints),
] {
let checkpoint = container
.child_node(name)
.expect("read the checkpoint")
.unwrap_or_else(|| panic!("{label} store keeps checkpoint {name}"));
assert!(
checkpoint
.child_node("root")
.expect("read the snapshot")
.is_some(),
"{label} store resolves the snapshot of checkpoint {name}"
);
}
}
}
#[test]
fn an_omitted_checkpoint_leaves_its_exclusive_records_uncopied() {
let directory = TestDirectory::new("omit-exclusive-records");
let names = build_store_with_checkpoints(&directory);
let store = WritableRepository::open(&directory.path).expect("open");
let generation = store.writing_generation().expect("generation");
let mut whole_writer = store.record_writer(generation);
let (_, whole) = deep_copy_super_root_with_progress(
&store,
&mut whole_writer,
store.head(),
&std::collections::BTreeSet::new(),
&mut crate::progress::DiscardedProgress,
)
.expect("copy everything");
whole_writer.finish().expect("finish");
let mut reduced_writer = store.record_writer(generation);
let (_, reduced) = deep_copy_super_root_with_progress(
&store,
&mut reduced_writer,
store.head(),
&std::collections::BTreeSet::from([names[0].clone()]),
&mut crate::progress::DiscardedProgress,
)
.expect("copy while omitting one checkpoint");
reduced_writer.finish().expect("finish");
assert!(
reduced < whole,
"omitting a checkpoint copies strictly fewer nodes: {reduced} against {whole}"
);
store.close().expect("close");
}
#[test]
fn omitting_a_shared_subtree_still_copies_it_through_the_live_root() {
let directory = TestDirectory::new("omit-shared-subtree");
let names = build_store_with_checkpoints(&directory);
let store = WritableRepository::open(&directory.path).expect("open");
let generation = store.writing_generation().expect("generation");
let mut writer = store.record_writer(generation);
let (head, _) = deep_copy_super_root_with_progress(
&store,
&mut writer,
store.head(),
&std::collections::BTreeSet::from([names[0].clone()]),
&mut crate::progress::DiscardedProgress,
)
.expect("copy while omitting one checkpoint");
writer.finish().expect("finish");
assert!(store.compare_and_set_head(store.head(), head));
store.flush().expect("flush");
store.close().expect("close");
let repository = Repository::open(&directory.path).expect("reopen");
let content = repository
.node_at_path("/content")
.expect("resolve the content")
.expect("the shared content survives the omission");
assert!(
content.child_node_count().expect("count") > 0,
"the shared subtree is intact"
);
}
#[test]
fn a_deep_copy_copies_each_distinct_node_exactly_once() {
let directory = TestDirectory::new("memo-exact");
build_populated_store(&directory);
let (copied, distinct) = {
let store = WritableRepository::open(&directory.path).expect("open");
let head = store.head();
let generation = store.writing_generation().expect("generation");
let mut writer = store.record_writer(generation);
let (root, copied) = deep_copy_tree_with_progress(
&store,
&mut writer,
head,
&mut crate::progress::DiscardedProgress,
)
.expect("deep copy");
let distinct = distinct_reachable_nodes(&store, head);
writer.finish().expect("finish");
assert!(store.compare_and_set_head(head, root));
store.close().expect("close");
(copied, distinct)
};
assert_eq!(
copied as usize, distinct,
"every distinct node is copied exactly once"
);
assert_content_intact(&directory);
}
#[test]
fn a_shared_subtree_is_copied_once_however_deep_the_sharing_nests() {
for (levels, ballast) in [(4usize, 0usize), (14, 0), (14, 32), (24, 4)] {
let directory = TestDirectory::new(&format!("diamond-{levels}-{ballast}"));
build_diamond_chain(&directory, levels, ballast);
let store = WritableRepository::open(&directory.path).expect("open");
let head = store.head();
let distinct = distinct_reachable_nodes(&store, head);
let generation = store.writing_generation().expect("generation");
let mut writer = store.record_writer(generation);
let (_root, copied) = deep_copy_tree_with_progress(
&store,
&mut writer,
head,
&mut crate::progress::DiscardedProgress,
)
.expect("deep copy");
writer.finish().expect("finish");
store.close().expect("close");
assert_eq!(
copied as usize, distinct,
"levels={levels} ballast={ballast}: copied must equal the distinct node count"
);
}
}
#[test]
fn a_tree_deeper_than_any_call_stack_copies_whole() {
let handle = std::thread::Builder::new()
.stack_size(2 * 1024 * 1024)
.spawn(|| {
let directory = TestDirectory::new("deep-chain");
build_diamond_chain(&directory, 100_000, 0);
let store = WritableRepository::open(&directory.path).expect("open");
let head = store.head();
let distinct = distinct_reachable_nodes(&store, head);
let generation = store.writing_generation().expect("generation");
let mut writer = store.record_writer(generation);
let (_root, copied) = deep_copy_tree_with_progress(
&store,
&mut writer,
head,
&mut crate::progress::DiscardedProgress,
)
.expect("a deep tree copies rather than aborting");
writer.finish().expect("finish");
assert_eq!(copied as usize, distinct);
assert!(distinct > 100_000, "the chain really is that deep");
crate::tooling::verify_node_tree(&store, head)
.expect("the verifier has no depth limit either");
store.close().expect("close");
})
.expect("spawn");
handle.join().expect("the walk stays off the call stack");
}
#[test]
#[ignore = "measurement, not an assertion"]
fn measure_deep_chain_walk_footprint() {
for levels in [100_000usize, 400_000] {
let directory = TestDirectory::new(&format!("deep-footprint-{levels}"));
build_diamond_chain(&directory, levels, 0);
let store = WritableRepository::open(&directory.path).expect("open");
let head = store.head();
let before = resident_bytes();
crate::tooling::verify_node_tree(&store, head).expect("verifies");
let after_verify = resident_bytes();
store.close().expect("close");
println!(
"levels={levels:>7} verify_rss_delta={:>6} MiB = {:>4} B/level",
after_verify.saturating_sub(before) / 1024 / 1024,
after_verify.saturating_sub(before) / levels,
);
}
}
#[test]
#[ignore = "measurement, not an assertion"]
fn measure_copy_throughput() {
let directory = TestDirectory::new("throughput");
let build_started = std::time::Instant::now();
build_wide_store(&directory, 1000);
let built = build_started.elapsed();
let store = WritableRepository::open(&directory.path).expect("open");
let head = store.head();
let generation = store.writing_generation().expect("generation");
let mut writer = store.record_writer(generation);
let started = std::time::Instant::now();
let (_root, copied) = deep_copy_tree_with_progress(
&store,
&mut writer,
head,
&mut crate::progress::DiscardedProgress,
)
.expect("deep copy");
let elapsed = started.elapsed();
writer.finish().expect("finish");
store.close().expect("close");
let per_second =
u32::try_from(copied).map_or(f64::INFINITY, f64::from) / elapsed.as_secs_f64();
println!(
"built {copied} nodes in {:.1}s; copied in {:.2}s = {per_second:.0} nodes/s; \
18.8M nodes extrapolates to {:.1} min",
built.as_secs_f64(),
elapsed.as_secs_f64(),
18_800_000.0 / per_second / 60.0,
);
}
#[test]
fn every_random_shape_copies_each_distinct_node_exactly_once() {
for seed in 1..=200u64 {
let directory = TestDirectory::new(&format!("random-dag-{seed}"));
build_random_dag(&directory, seed);
let store = WritableRepository::open(&directory.path).expect("open");
let head = store.head();
let distinct = distinct_reachable_nodes(&store, head);
let generation = store.writing_generation().expect("generation");
let mut writer = store.record_writer(generation);
let (_root, copied) = deep_copy_tree_with_progress(
&store,
&mut writer,
head,
&mut crate::progress::DiscardedProgress,
)
.expect("deep copy");
writer.finish().expect("finish");
store.close().expect("close");
assert_eq!(
copied as usize, distinct,
"seed {seed}: copied {copied} but {distinct} distinct nodes are reachable"
);
}
}
#[test]
fn a_cyclic_source_is_refused_at_the_record_that_closes_the_cycle() {
use crate::content::provider::tests::MemorySegmentProvider;
use crate::error::Error;
use crate::writer::segment_builder::SegmentBufferBuilder;
let directory = TestDirectory::new("compaction-cycle");
let store = WritableRepository::open(&directory.path).expect("open");
let generation = store.writing_generation().expect("generation");
let mut writer = store.record_writer(generation);
let original_child = writer
.write_node(None, &[], &ChildNodesToWrite::Zero, &[])
.expect("original child");
let root = writer
.write_node(
None,
&[],
&ChildNodesToWrite::One {
name: "loop".to_owned(),
node: original_child,
},
&[],
)
.expect("root");
writer.finish().expect("finish segment");
let view = store.segment(root.segment).expect("root segment");
let root_position = view
.record_position(root.record_number)
.expect("root position");
let mut cyclic_bytes = view.bytes.to_vec();
let child_slot: &mut [u8; 6] = (&mut cyclic_bytes[root_position + 12..root_position + 18])
.try_into()
.expect("one child identifier slot");
SegmentBufferBuilder::write_record_identifier_bytes(0, root.record_number, child_slot);
let mut memory = MemorySegmentProvider::default();
memory.insert(root.segment, cyclic_bytes);
let mut sink_writer = store.record_writer(generation);
let error = deep_copy_tree_with_progress(
&memory,
&mut sink_writer,
root,
&mut crate::progress::DiscardedProgress,
)
.expect_err("a cyclic source is refused");
let Error::InvalidFormat { details } = &error else {
panic!("a cycle is a format error, got {error:?}");
};
assert!(
details.contains("contained in its own subtree"),
"unexpected detail: {details}"
);
assert!(
details.contains(&root.to_string()),
"the error names the offending record: {details}"
);
store.close().expect("close");
}
}