use std::collections::HashMap;
use std::io::Write;
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::segment::record::RecordIdentifier;
use crate::writer::record_writer::{
ChildNodesToWrite, PropertyToWrite, PropertyValuesToWrite, RecordWriter, SegmentSink,
sort_properties_for_template,
};
use crate::writer::segment_builder::GarbageCollectionGeneration;
use crate::writer::store_writer::WritableRepository;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CompactionKind {
Full,
Tail,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub 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)> {
let mut copier = Compactor {
source,
writer,
rewritten_nodes: HashMap::new(),
compacted_nodes: 0,
};
let root = copier.compact_node(source_root, 0)?;
Ok((root, copier.compacted_nodes))
}
const MAXIMUM_COMPACTION_DEPTH: usize = 4000;
struct Compactor<'writer, Sink: SegmentSink> {
source: &'writer dyn SegmentProvider,
writer: &'writer mut RecordWriter<Sink>,
rewritten_nodes: HashMap<RecordIdentifier, RecordIdentifier>,
compacted_nodes: u64,
}
impl<Sink: SegmentSink> Compactor<'_, Sink> {
fn compact_node(
&mut self,
source_node: RecordIdentifier,
depth: usize,
) -> Result<RecordIdentifier> {
if let Some(&rewritten) = self.rewritten_nodes.get(&source_node) {
return Ok(rewritten);
}
if depth > MAXIMUM_COMPACTION_DEPTH {
return Err(Error::InvalidFormat {
details: format!(
"node tree exceeds depth {MAXIMUM_COMPACTION_DEPTH}; \
the source records probably form a cycle"
),
});
}
let node = NodeState::new(self.source, source_node);
let template = node.template()?;
let stable_identifier = node.stable_identifier_bytes()?;
let mut child_entries = Vec::new();
for (name, child) in node.child_node_entries()? {
child_entries.push((
name,
self.compact_node(child.record_identifier(), depth + 1)?,
));
}
let children = match child_entries.len() {
0 => ChildNodesToWrite::Zero,
1 => {
let (name, node) = child_entries.into_iter().next().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.rewritten_nodes.insert(source_node, rewritten);
self.compacted_nodes += 1;
Ok(rewritten)
}
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,
})
}
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)
}
_ => 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)
}
}
pub fn compact(store: &mut WritableRepository, kind: CompactionKind) -> 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 mut writer = store.record_writer_with_identifier(target_generation, "c");
let (new_head, compacted_nodes) = deep_copy_tree(store, &mut writer, head)?;
writer.finish()?;
if !store.set_head(head, new_head) {
return Err(Error::InvalidFormat {
details: "the head moved during compaction".to_owned(),
});
}
store.flush()?;
store.reclaim_old_generations(target_generation, kind == CompactionKind::Full)?;
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,
})
}
fn append_gc_log(
store: &WritableRepository,
repository_size: u64,
reclaimed_size: u64,
generation: GarbageCollectionGeneration,
compacted_nodes: u64,
root: RecordIdentifier,
) -> Result<()> {
use std::io::Write;
let timestamp = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |duration| duration.as_millis());
let line = format!(
"{repository_size},{reclaimed_size},{timestamp},{},{},{compacted_nodes},{}:{}\n",
generation.generation, generation.full_generation, root.segment, root.record_number as i32,
);
let mut file = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(store.directory().join("gc.log"))?;
file.write_all(line.as_bytes())?;
file.sync_data()?;
Ok(())
}
fn rewrite_journal_to_head(store: &WritableRepository, head: RecordIdentifier) -> Result<()> {
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)]
mod tests {
use super::{CompactionKind, compact};
use crate::content::node::PropertyValues;
use crate::content::property::PropertyValue;
use crate::content::provider::SegmentProvider;
use crate::store::Repository;
use crate::writer::commit::{create_checkpoint, list_checkpoints};
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-compaction-{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 build_populated_store(directory: &TestDirectory) {
let store = WritableRepository::open(&directory.path).expect("bootstrap");
let generation = store.writing_generation().expect("generation");
let mut writer = store.record_writer(generation);
let first_child = writer
.write_node(Some("nt:unstructured"), &[], &ChildNodesToWrite::Zero, &[])
.expect("child");
let second_child = writer
.write_node(Some("nt:unstructured"), &[], &ChildNodesToWrite::Zero, &[])
.expect("child");
let title = writer.write_string("Compaction Test").expect("value");
let content = writer
.write_node(
Some("nt:unstructured"),
&[],
&ChildNodesToWrite::Many(vec![
("alpha".to_owned(), first_child),
("beta".to_owned(), second_child),
]),
&[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,
&[("purpose".to_owned(), "test".to_owned())],
)
.expect("checkpoint");
store.close().expect("close");
}
fn assert_content_intact(directory: &TestDirectory) {
let repository = Repository::open(&directory.path).expect("reader opens");
let content = repository
.node_at_path("/content")
.expect("resolve")
.expect("present");
assert_eq!(content.child_node_count().expect("count"), 2);
assert!(content.child_node("alpha").expect("read").is_some());
assert!(content.child_node("beta").expect("read").is_some());
let title = content.property("title").expect("read").expect("present");
assert_eq!(
title.values,
PropertyValues::Single(PropertyValue::String("Compaction Test".to_owned()))
);
let checkpoints = repository.checkpoints().expect("checkpoints");
assert_eq!(checkpoints.len(), 1, "the checkpoint survives compaction");
let (_, checkpoint) = &checkpoints[0];
let snapshot = checkpoint
.child_node("root")
.expect("read")
.expect("snapshot");
assert!(snapshot.child_node("content").expect("read").is_some());
}
#[test]
fn full_compaction_preserves_content_and_checkpoints() {
let directory = TestDirectory::new("full");
build_populated_store(&directory);
let outcome = {
let mut store = WritableRepository::open(&directory.path).expect("open for compaction");
let before_generation = store
.segment_generation(store.head().segment)
.expect("generation");
let outcome = compact(&mut store, CompactionKind::Full).expect("compact");
let after_generation = store
.segment_generation(store.head().segment)
.expect("generation");
assert_eq!(
after_generation.generation,
before_generation.generation + 1
);
assert_eq!(
after_generation.full_generation,
before_generation.full_generation + 1
);
assert!(after_generation.is_compacted);
store.close().expect("close");
outcome
};
assert!(outcome.compacted_nodes > 0);
assert_content_intact(&directory);
let journal = std::fs::read_to_string(directory.path.join("journal.log")).expect("journal");
assert_eq!(journal.lines().count(), 1, "journal rewritten to one line");
let gc_log = std::fs::read_to_string(directory.path.join("gc.log")).expect("gc.log");
assert_eq!(gc_log.lines().count(), 1);
assert_eq!(gc_log.split(',').count(), 7, "seven gc.log fields");
}
#[test]
fn compaction_preserves_stable_identifiers() {
let directory = TestDirectory::new("stable-ids");
build_populated_store(&directory);
let before = {
let repository = Repository::open(&directory.path).expect("reader");
repository
.node_at_path("/content")
.expect("resolve")
.expect("present")
.stable_identifier()
.expect("stable id")
};
{
let mut store = WritableRepository::open(&directory.path).expect("open");
compact(&mut store, CompactionKind::Full).expect("compact");
store.close().expect("close");
}
let after = {
let repository = Repository::open(&directory.path).expect("reader");
repository
.node_at_path("/content")
.expect("resolve")
.expect("present")
.stable_identifier()
.expect("stable id")
};
assert_eq!(
before, after,
"the stable identifier survives compaction so Oak's fast path keeps matching"
);
}
#[test]
fn compaction_preserves_infinite_doubles_and_type_named_properties() {
let directory = TestDirectory::new("edge-values");
{
let store = WritableRepository::open(&directory.path).expect("open");
let generation = store.writing_generation().expect("generation");
let mut writer = store.record_writer(generation);
let infinity_value = writer.write_string("Infinity").expect("value");
let odd_name_value = writer.write_string("literal").expect("value");
let content = writer
.write_node(
None,
&[],
&ChildNodesToWrite::Zero,
&[
PropertyToWrite {
name: "ratio".to_owned(),
property_type: crate::content::property::PropertyType::Double,
values: PropertyValuesToWrite::Single(infinity_value),
},
PropertyToWrite {
name: "jcr:primaryType".to_owned(),
property_type: crate::content::property::PropertyType::String,
values: PropertyValuesToWrite::Single(odd_name_value),
},
],
)
.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));
store.close().expect("close");
}
{
let mut store = WritableRepository::open(&directory.path).expect("open");
compact(&mut store, CompactionKind::Full).expect("compact");
store.close().expect("close");
}
let repository = Repository::open(&directory.path).expect("reader");
let content = repository
.node_at_path("/content")
.expect("resolve")
.expect("present");
let ratio = content.property("ratio").expect("read").expect("present");
assert_eq!(
ratio.values,
PropertyValues::Single(PropertyValue::Double(f64::INFINITY))
);
let odd = content
.property("jcr:primaryType")
.expect("read")
.expect("present");
assert_eq!(
odd.property_type,
crate::content::property::PropertyType::String
);
assert_eq!(
odd.values,
PropertyValues::Single(PropertyValue::String("literal".to_owned()))
);
}
#[test]
fn compaction_streams_long_binaries_through_bulk_segments() {
let directory = TestDirectory::new("long-binary");
let content: Vec<u8> = (0..300 * 1024).map(|index| (index % 251) as u8).collect();
{
let store = WritableRepository::open(&directory.path).expect("open");
let generation = store.writing_generation().expect("generation");
let mut writer = store.record_writer(generation);
let binary_value = writer.write_binary_content(&content).expect("binary");
let content_node = writer
.write_node(
Some("nt:file"),
&[],
&ChildNodesToWrite::Zero,
&[PropertyToWrite {
name: "data".to_owned(),
property_type: crate::content::property::PropertyType::Binary,
values: PropertyValuesToWrite::Single(binary_value),
}],
)
.expect("content");
let root = writer
.write_node(
None,
&[],
&ChildNodesToWrite::One {
name: "content".to_owned(),
node: content_node,
},
&[],
)
.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));
store.close().expect("close");
}
{
let mut store = WritableRepository::open(&directory.path).expect("open");
compact(&mut store, CompactionKind::Full).expect("compact");
store.close().expect("close");
}
let repository = Repository::open(&directory.path).expect("reader");
let content_node = repository
.node_at_path("/content")
.expect("resolve")
.expect("present");
let data = content_node
.property("data")
.expect("read")
.expect("present");
let record = match &data.values {
PropertyValues::Single(PropertyValue::Binary(
crate::content::value::BinaryValue::Inline {
record_identifier, ..
},
)) => *record_identifier,
other => panic!("expected an inline binary, got {other:?}"),
};
let read_back =
crate::content::value::read_binary_content(&repository, record).expect("content");
assert_eq!(
read_back, content,
"the long binary round-trips through compaction"
);
}
#[test]
fn committing_after_compaction_in_one_session_persists_the_journal() {
let directory = TestDirectory::new("commit-after-compact");
build_populated_store(&directory);
{
let mut store = WritableRepository::open(&directory.path).expect("open");
compact(&mut store, CompactionKind::Full).expect("compact");
create_checkpoint(&store, 10_000_000, &[]).expect("checkpoint");
store.close().expect("close");
}
let repository = Repository::open(&directory.path).expect("reader");
assert_eq!(
repository.checkpoints().expect("checkpoints").len(),
2,
"the checkpoint created after compaction is visible in the journal"
);
}
#[test]
fn tail_compaction_keeps_the_full_generation() {
let directory = TestDirectory::new("tail");
build_populated_store(&directory);
{
let mut store = WritableRepository::open(&directory.path).expect("open");
let before = store
.segment_generation(store.head().segment)
.expect("generation");
compact(&mut store, CompactionKind::Tail).expect("compact");
let after = store
.segment_generation(store.head().segment)
.expect("generation");
assert_eq!(after.generation, before.generation + 1);
assert_eq!(
after.full_generation, before.full_generation,
"tail compaction keeps the full generation"
);
store.close().expect("close");
}
assert_content_intact(&directory);
}
#[test]
fn compaction_reclaims_disk_space_from_garbage() {
let directory = TestDirectory::new("reclaim");
{
let store = WritableRepository::open(&directory.path).expect("bootstrap");
for revision in 0..30 {
let generation = store.writing_generation().expect("generation");
let mut writer = store.record_writer(generation);
let value = writer
.write_string(&format!("revision-{revision}").repeat(2000))
.expect("value");
let content = writer
.write_node(
Some("nt:unstructured"),
&[],
&ChildNodesToWrite::Zero,
&[PropertyToWrite {
name: "data".to_owned(),
property_type: crate::content::property::PropertyType::String,
values: PropertyValuesToWrite::Single(value),
}],
)
.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));
store.flush().expect("flush");
}
store.close().expect("close");
}
let mut store = WritableRepository::open(&directory.path).expect("open");
let outcome = compact(&mut store, CompactionKind::Full).expect("compact");
store.close().expect("close");
assert!(
outcome.size_after < outcome.size_before,
"compaction reclaims garbage: {} -> {}",
outcome.size_before,
outcome.size_after
);
let repository = Repository::open(&directory.path).expect("reader");
let content = repository
.node_at_path("/content")
.expect("resolve")
.expect("present");
let data = content.property("data").expect("read").expect("present");
assert_eq!(
data.values,
PropertyValues::Single(PropertyValue::String("revision-29".repeat(2000)))
);
}
#[test]
fn compacted_stores_survive_a_second_compaction() {
let directory = TestDirectory::new("twice");
build_populated_store(&directory);
for _ in 0..2 {
let mut store = WritableRepository::open(&directory.path).expect("open");
compact(&mut store, CompactionKind::Full).expect("compact");
store.close().expect("close");
assert_content_intact(&directory);
}
let store = WritableRepository::open(&directory.path).expect("open");
assert_eq!(list_checkpoints(&store).expect("list").len(), 1);
store.close().expect("close");
}
fn assert_no_dangling_segment_references(directory: &TestDirectory) {
let repository = Repository::open(&directory.path).expect("reader opens");
for segment_identifier in repository.segment_identifiers() {
if segment_identifier.is_bulk_segment() {
continue;
}
let view = repository
.segment(segment_identifier)
.expect("data segment readable");
for referenced in &view.structure.referenced_segments {
assert!(
repository.contains_segment(*referenced),
"kept data segment {segment_identifier} references missing segment \
{referenced}"
);
}
}
}
#[test]
fn tail_compaction_keeps_bulk_segments_referenced_by_retained_data_segments() {
let directory = TestDirectory::new("tail-bulk-mark");
build_populated_store(&directory);
{
let store = WritableRepository::open(&directory.path).expect("open");
let generation = store.writing_generation().expect("generation");
let mut writer = store.record_writer(generation);
let large = writer
.write_string(&"bulk-backed-value ".repeat(20_000))
.expect("large value");
let content = writer
.write_node(
Some("nt:unstructured"),
&[],
&ChildNodesToWrite::Zero,
&[PropertyToWrite {
name: "data".to_owned(),
property_type: crate::content::property::PropertyType::String,
values: PropertyValuesToWrite::Single(large),
}],
)
.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));
store.close().expect("close");
}
{
let mut store = WritableRepository::open(&directory.path).expect("open");
compact(&mut store, CompactionKind::Full).expect("full compact");
store.close().expect("close");
}
assert_no_dangling_segment_references(&directory);
{
let mut store = WritableRepository::open(&directory.path).expect("open");
compact(&mut store, CompactionKind::Tail).expect("tail compact");
store.close().expect("close");
}
assert_no_dangling_segment_references(&directory);
let repository = Repository::open(&directory.path).expect("reader opens");
let content = repository
.node_at_path("/content")
.expect("resolve")
.expect("present");
let data = content.property("data").expect("read").expect("present");
assert_eq!(
data.values,
PropertyValues::Single(PropertyValue::String("bulk-backed-value ".repeat(20_000)))
);
}
}