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)> {
let source_root = super_root;
let mut copier = Compactor {
source,
writer,
omitted_checkpoints,
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();
assert_eq!(
copier.compacted_nodes, memoized 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 {
use super::*;
use super::{CompactionKind, compact};
use crate::content::node::PropertyValues;
use crate::content::property::PropertyValue;
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;
#[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.compare_and_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.compare_and_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.compare_and_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");
}
#[test]
fn compaction_certifies_base_archives_before_writing_a_retry_copy() {
let directory = TestDirectory::new("preflight-base-certificate");
build_populated_store(&directory);
let repository = Repository::open(&directory.path).expect("open healthy repository");
let archive_name = repository.archives()[0].file_name().to_owned();
drop(repository);
corrupt_graph_checksum(&directory.path.join(&archive_name));
let journal_before =
std::fs::read(directory.path.join("journal.log")).expect("read journal before");
let archives_before =
crate::store::list_archive_file_names(&directory.path).expect("list archives before");
let bytes_before: Vec<_> = archives_before
.iter()
.map(|name| {
(
name.clone(),
std::fs::read(directory.path.join(name)).expect("read archive before"),
)
})
.collect();
for attempt in 1..=2 {
let mut store = WritableRepository::open(&directory.path)
.expect("ordinary read path tolerates an invalid optional graph");
let error = compact(&mut store, CompactionKind::Full)
.expect_err("strict reclaim source preflight must refuse the graph");
assert!(error.to_string().contains("segment graph"), "{error}");
drop(store);
assert_eq!(
crate::store::list_archive_file_names(&directory.path)
.expect("list archives after refused attempt"),
archives_before,
"refused retry {attempt} must not allocate another compacted TAR"
);
}
assert_eq!(
crate::store::list_archive_file_names(&directory.path).expect("list archives after"),
archives_before,
"preflight refusal must not allocate a compacted TAR"
);
for (name, expected) in bytes_before {
assert_eq!(
std::fs::read(directory.path.join(name)).expect("read archive after"),
expected
);
}
assert_eq!(
std::fs::read(directory.path.join("journal.log")).expect("read journal after"),
journal_before,
"preflight refusal must not publish another head"
);
}
#[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.compare_and_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)))
);
}
}