const SEGMENT_TRACE_REPORT_STRIDE: u64 = 256;
use super::indexless_refusal::indexless_active_archive_refusal;
use super::plan::RetainedReclaimable;
use super::planning::add_estimate;
use super::stale_archives::generation_from_header;
use crate::content::provider::SegmentProvider;
use crate::content::template::{Template, read_template};
use crate::content::value::read_string;
use crate::error::{Error, Result};
use crate::progress::{ProgressObserver, Step, WorkUnit};
use crate::segment::identifier::SegmentIdentifier;
use crate::segment::record::RecordIdentifier;
#[cfg(test)]
use crate::segment::record::RecordType;
use crate::segment::view::SegmentView;
use crate::store::Repository;
use crate::tooling::NodeTreeVerifier;
use crate::writer::compaction::CompactionKind;
use crate::writer::segment_builder::GarbageCollectionGeneration;
use crate::writer::store_writer::{
PlannedArchiveSweep, ReclaimRule, StandaloneSegmentCompactionPlan, is_reclaimable,
planned_unavailable_segments,
};
use std::borrow::Cow;
use std::collections::{BTreeSet, HashMap, HashSet, VecDeque};
use std::path::Path;
use std::sync::Arc;
pub(super) fn compaction_target_generation(
base: GarbageCollectionGeneration,
kind: CompactionKind,
) -> GarbageCollectionGeneration {
match kind {
CompactionKind::Full => GarbageCollectionGeneration {
generation: base.generation.wrapping_add(1),
full_generation: base.full_generation.wrapping_add(1),
is_compacted: true,
},
CompactionKind::Tail => GarbageCollectionGeneration {
generation: base.generation.wrapping_add(1),
full_generation: base.full_generation,
is_compacted: true,
},
}
}
pub(super) fn retained_reclaimable_from(
planned: &PlannedArchiveSweep,
retained: &mut RetainedReclaimable,
warnings: &mut Vec<String>,
) -> Result<()> {
match planned {
PlannedArchiveSweep::DeferredBySavings {
file_name,
segment_count,
eligible_entry_bytes,
} => {
retained.below_savings_gate += segment_count;
add_estimate(&mut retained.bytes, *eligible_entry_bytes)?;
warnings.push(format!(
"{file_name}: {segment_count} reclaimable segments ({}) retained because savings do not exceed Oak's 25% rewrite gate",
crate::units::format_byte_size(*eligible_entry_bytes)
));
}
PlannedArchiveSweep::DeferredAtLastGeneration {
file_name,
segment_count,
eligible_entry_bytes,
} => {
retained.at_last_generation += segment_count;
add_estimate(&mut retained.bytes, *eligible_entry_bytes)?;
warnings.push(format!(
"{file_name}: {segment_count} reclaimable segments ({}) retained because archive generation z cannot be rewritten",
crate::units::format_byte_size(*eligible_entry_bytes)
));
}
PlannedArchiveSweep::BlockedByOccupiedGeneration {
file_name,
occupied_name,
segment_count,
eligible_entry_bytes,
} => {
retained.blocked_by_occupied_generation += segment_count;
add_estimate(&mut retained.bytes, *eligible_entry_bytes)?;
warnings.push(format!(
"{file_name}: {segment_count} reclaimable segments ({}) retained because {occupied_name} already exists",
crate::units::format_byte_size(*eligible_entry_bytes)
));
}
PlannedArchiveSweep::Remove { .. } | PlannedArchiveSweep::Rewrite { .. } => {}
}
Ok(())
}
pub(super) fn predict_shared_bulk_segments(
repository: &Repository,
head: RecordIdentifier,
omitted_checkpoints: &BTreeSet<String>,
observer: &mut dyn ProgressObserver,
) -> Result<HashSet<SegmentIdentifier>> {
crate::progress::observe(
observer,
&Step::new("predicting the shared binary content", WorkUnit::Nodes),
|observer| {
let mut shared = HashSet::new();
let mut visited: HashSet<RecordIdentifier> = HashSet::new();
let mut pending = vec![head];
let mut traced = crate::progress::StrideCounter::new(SEGMENT_TRACE_REPORT_STRIDE);
while let Some(record) = pending.pop() {
if !visited.insert(record) {
continue;
}
traced.advance(observer);
let node = repository.node(record);
for property in node.properties()? {
if property.property_type != crate::content::property::PropertyType::Binary {
continue;
}
let values = match &property.values {
crate::content::node::PropertyValues::Single(value) => {
std::slice::from_ref(value)
}
crate::content::node::PropertyValues::Multiple(values) => values.as_slice(),
};
for value in values {
let crate::content::property::PropertyValue::Binary(
crate::content::value::BinaryValue::Inline {
length,
record_identifier,
},
) = value
else {
continue;
};
if *length < crate::writer::record_writer::MEDIUM_VALUE_LIMIT as u64 {
continue;
}
collect_shared_bulk_blocks(
repository,
*record_identifier,
*length,
&mut shared,
)?;
}
}
let is_checkpoint_container = record == head;
for (name, child) in node.child_node_entries()? {
if !is_checkpoint_container
&& !omitted_checkpoints.is_empty()
&& omitted_checkpoints.contains(&name)
{
continue;
}
pending.push(child.record_identifier());
}
}
traced.finish(observer);
if shared.is_empty() {
observer.step_concluded("no pre-existing bulk segments will be shared in place");
} else {
let (_, bulk_bytes) = indexed_bytes_by_kind(repository, &shared);
observer.step_concluded(&format!(
"{} pre-existing bulk segments ({}) will be shared in place and retained",
crate::units::format_count(crate::progress::count(shared.len())),
crate::units::format_byte_size(bulk_bytes),
));
}
Ok(shared)
},
)
}
pub(super) fn collect_shared_bulk_blocks(
repository: &Repository,
value: RecordIdentifier,
length: u64,
shared: &mut HashSet<SegmentIdentifier>,
) -> Result<()> {
let block_count = length.div_ceil(crate::content::value::BLOCK_SIZE);
let view = repository.segment(value.segment)?;
let list_identifier = view.read_record_identifier(value.record_number, 8, 0)?;
for block in
crate::content::list::uncounted_list_entries(repository, list_identifier, block_count)?
{
if block.segment.is_bulk_segment() {
shared.insert(block.segment);
}
}
Ok(())
}
pub(super) fn segments_ahead_of_the_head(
active_index_generations: &HashMap<SegmentIdentifier, GarbageCollectionGeneration>,
reference: GarbageCollectionGeneration,
) -> usize {
active_index_generations
.iter()
.filter(|(identifier, generation)| {
identifier.is_data_segment() && generation.generation > reference.generation
})
.count()
}
pub(super) fn active_index_generations(
repository: &Repository,
) -> Result<HashMap<SegmentIdentifier, GarbageCollectionGeneration>> {
if let Some(details) = indexless_active_archive_refusal(repository) {
return Err(Error::InvalidFormat { details });
}
let mut generations = HashMap::new();
for archive in repository.archives() {
for identifier in archive.segment_identifiers() {
let entry = archive
.index_entry(identifier)
.ok_or_else(|| Error::InvalidFormat {
details: format!(
"active archive {} has no index metadata for its own segment {identifier}; refusing generation cleanup",
archive.file_name()
),
})?;
generations.insert(
identifier,
GarbageCollectionGeneration {
generation: entry.generation,
full_generation: entry.full_generation,
is_compacted: entry.is_compacted,
},
);
}
}
Ok(generations)
}
pub(in crate::writer::maintenance) fn indexed_bytes_by_kind(
repository: &Repository,
segments: &HashSet<SegmentIdentifier>,
) -> (u64, u64) {
let mut data_bytes = 0u64;
let mut bulk_bytes = 0u64;
for archive in repository.archives() {
for identifier in archive.segment_identifiers() {
if !segments.contains(&identifier) {
continue;
}
let Some(entry) = archive.index_entry(identifier) else {
continue;
};
if identifier.is_data_segment() {
data_bytes = data_bytes.saturating_add(u64::from(entry.size));
} else {
bulk_bytes = bulk_bytes.saturating_add(u64::from(entry.size));
}
}
}
(data_bytes, bulk_bytes)
}
pub(super) fn extend_segment_closure(
provider: &dyn SegmentProvider,
roots: impl IntoIterator<Item = SegmentIdentifier>,
seen: &mut HashSet<SegmentIdentifier>,
observer: &mut dyn ProgressObserver,
) -> Result<()> {
let mut pending: VecDeque<_> = roots.into_iter().collect();
let mut traced = crate::progress::StrideCounter::new(SEGMENT_TRACE_REPORT_STRIDE);
while let Some(identifier) = pending.pop_front() {
if !seen.insert(identifier) {
continue;
}
let segment = provider.segment(identifier)?;
traced.advance(observer);
pending.extend(segment.structure.referenced_segments.iter().copied());
}
traced.finish(observer);
Ok(())
}
pub(super) fn validate_reclaim_reference_invariant(
repository: &Repository,
current_closure: &HashSet<SegmentIdentifier>,
active_index_generations: &HashMap<SegmentIdentifier, GarbageCollectionGeneration>,
rule: ReclaimRule,
) -> Result<()> {
for &identifier in current_closure {
if !identifier.is_data_segment() {
continue;
}
let header = generation_from_header(repository, identifier)?;
let indexed = active_index_generations.get(&identifier).ok_or_else(|| {
Error::InvalidFormat {
details: format!(
"current head reaches data segment {identifier}, but no active archive index describes it"
),
}
})?;
if *indexed != header {
return Err(Error::InvalidFormat {
details: format!(
"segment {identifier} has index generation {indexed:?}, but its header says {header:?}"
),
});
}
if is_reclaimable(rule.reference, header, rule.kind, rule.retained_generations) {
return Err(Error::InvalidFormat {
details: format!(
"current head reaches data segment {identifier} in reclaimable generation {header:?}; refusing to trust generation cleanup"
),
});
}
}
Ok(())
}
pub(super) struct ExcludingProvider<'repository> {
pub(super) repository: &'repository Repository,
pub(super) unavailable: &'repository HashSet<SegmentIdentifier>,
}
impl SegmentProvider for ExcludingProvider<'_> {
fn segment(&self, identifier: SegmentIdentifier) -> Result<SegmentView<'_>> {
if self.unavailable.contains(&identifier) {
return Err(Error::SegmentNotFound {
segment_identifier: identifier,
});
}
self.repository.segment(identifier)
}
fn string(&self, identifier: RecordIdentifier) -> Result<Arc<str>> {
read_string(self, identifier).map(Arc::from)
}
fn template(&self, identifier: RecordIdentifier) -> Result<Arc<Template>> {
read_template(self, identifier).map(Arc::new)
}
}
pub(super) fn validate_prospective_segment_plan(
directory: &Path,
repository: &Repository,
plan: &StandaloneSegmentCompactionPlan,
retained_roots: &[RecordIdentifier],
observer: &mut dyn ProgressObserver,
) -> Result<()> {
let unavailable = planned_unavailable_segments(directory, plan)?;
if unavailable.is_empty() {
return Ok(());
}
let provider = ExcludingProvider {
repository,
unavailable: &unavailable,
};
let mut verifier = NodeTreeVerifier::new(&provider);
for &root in retained_roots {
verifier
.verify_with_progress(root, observer)
.map_err(|error| Error::InvalidFormat {
details: format!(
"segment cleanup would make retained journal root {root} unreadable: {error}"
),
})?;
}
let reclaimable = plan.reclaimable_segments();
for identifier in repository.distinct_segment_identifiers() {
if !identifier.is_data_segment()
|| unavailable.contains(&identifier)
|| reclaimable.contains(&identifier)
{
continue;
}
let segment = repository.segment(identifier)?;
if let Some(target) = segment
.structure
.referenced_segments
.iter()
.find(|target| unavailable.contains(target))
{
return Err(Error::InvalidFormat {
details: format!(
"surviving data segment {identifier} references segment {target}, which the cleanup plan would remove"
),
});
}
}
Ok(())
}
pub(super) fn prospective_retained_roots<'roots>(
directory: &Path,
repository: &Repository,
plan: &StandaloneSegmentCompactionPlan,
retained_roots: &'roots [RecordIdentifier],
) -> Cow<'roots, [RecordIdentifier]> {
#[cfg(test)]
if crate::writer::fault_injection::is_substitution_armed(
"cleanup.before-prospective-retained-root-verification",
) {
let unavailable = planned_unavailable_segments(directory, plan)
.expect("armed prospective-root fixture must have a valid physical plan");
for identifier in unavailable {
let segment = repository
.segment(identifier)
.expect("armed prospective-root fixture segment must be readable");
if let Some(entry) = segment
.structure
.record_table()
.iter()
.find(|entry| entry.record_type() == Some(RecordType::Node))
{
let mut injected = retained_roots.to_vec();
injected.push(RecordIdentifier::new(identifier, entry.record_number));
return Cow::Owned(injected);
}
}
panic!("prospective retained-root fault fixture has no removable node record");
}
let _ = (directory, repository, plan);
Cow::Borrowed(retained_roots)
}
#[cfg(test)]
mod tests {
use crate::store::Repository;
use crate::writer::maintenance::options::*;
use crate::writer::maintenance::plan::*;
use crate::writer::maintenance::prepared::*;
use crate::writer::maintenance::test_support::*;
use crate::writer::record_writer::ChildNodesToWrite;
use crate::writer::segment_builder::GarbageCollectionGeneration;
use crate::writer::store_writer::WritableRepository;
#[test]
fn prospective_plan_refuses_a_survivor_that_references_a_planned_removal() {
let directory = TestDirectory::repository("prospective-survivor-reference");
let old_generation = GarbageCollectionGeneration {
generation: 0,
full_generation: 0,
is_compacted: false,
};
let target = {
let store = WritableRepository::open(&directory.path).expect("open old-target writer");
let target = write_empty_node_segment(&store, old_generation);
store.close().expect("close old-target archive");
target
};
let store = WritableRepository::open(&directory.path).expect("open survivor writer");
let current_generation = GarbageCollectionGeneration {
generation: 2,
full_generation: 2,
is_compacted: false,
};
let mut survivor_writer = store.record_writer(current_generation);
let survivor = survivor_writer
.write_node(
None,
&[],
&ChildNodesToWrite::One {
name: "old-target".to_owned(),
node: target,
},
&[],
)
.expect("write unjournaled newer-generation survivor");
survivor_writer.finish().expect("finish survivor segment");
let content_root = write_empty_node_segment(&store, current_generation);
let mut head_writer = store.record_writer(current_generation);
let head = head_writer
.write_node(
None,
&[],
&ChildNodesToWrite::One {
name: "root".to_owned(),
node: content_root,
},
&[],
)
.expect("write unrelated current head");
head_writer.finish().expect("finish current head segment");
assert!(store.compare_and_set_head(store.head(), head));
store.close().expect("close fixture writer");
let before = file_bytes(&directory.path);
let error = plan_compaction(
&directory.path,
&CompactionOptions::default().with_tasks([MaintenanceTask::Segments]),
)
.expect_err("prospective deletion must reject a surviving cross-reference");
assert_eq!(
error.to_string(),
format!(
"invalid segment-tar data: surviving data segment {} references segment {}, which the cleanup plan would remove",
survivor.segment, target.segment
)
);
assert_eq!(
file_bytes(&directory.path),
before,
"prospective validation remains read-only"
);
Repository::open(&directory.path).expect("refused fixture remains readable");
}
#[test]
fn current_head_reaching_a_one_generation_old_segment_fails_closed() {
let directory = TestDirectory::repository("generation-invariant");
let store = WritableRepository::open(&directory.path).expect("open writer");
let mut root_writer = store.record_writer(GarbageCollectionGeneration {
generation: 1,
full_generation: 1,
is_compacted: false,
});
let old_root = root_writer
.write_node(None, &[], &ChildNodesToWrite::Zero, &[])
.expect("one-generation-old content root");
root_writer.finish().expect("finish the older generation");
let mut writer = store.record_writer(GarbageCollectionGeneration {
generation: 2,
full_generation: 2,
is_compacted: false,
});
let new_head = writer
.write_node(
None,
&[],
&ChildNodesToWrite::One {
name: "root".to_owned(),
node: old_root,
},
&[],
)
.expect("new super root");
writer.finish().expect("finish");
assert!(store.compare_and_set_head(store.head(), new_head));
store.close().expect("close");
let before = file_bytes(&directory.path);
let options = CompactionOptions::default().with_tasks([MaintenanceTask::Segments]);
let error = plan_compaction(&directory.path, &options)
.expect_err("a live one-generation-old child is unsafe at one retained generation");
assert!(
error
.to_string()
.contains("current head reaches data segment")
);
assert_eq!(file_bytes(&directory.path), before);
Repository::open(&directory.path).expect("refusal leaves repository healthy");
}
#[test]
fn gate_deferred_garbage_is_counted_rather_than_reported_as_nothing() {
let (directory, new_head) = sub_gate_garbage_fixture("retained-reclaimable");
let options = CompactionOptions::default()
.with_tasks([MaintenanceTask::Segments])
.with_oak_savings_gate();
let plan = plan_compaction(&directory.path, &options).expect("segment plan");
assert_eq!(plan.current_head(), new_head);
let deferred = plan
.warnings()
.iter()
.any(|warning| warning.contains("25% rewrite gate"));
assert!(
deferred,
"expected a savings deferral: {:?}",
plan.warnings()
);
assert!(
plan.retained_reclaimable_segments() != 0,
"declined garbage must be counted, not silently dropped"
);
assert!(plan.retained_reclaimable_bytes() != 0);
assert_eq!(plan.estimated_reclaimable_bytes(), 0);
}
#[test]
fn the_default_policy_reclaims_what_the_oak_gate_declines() {
let (directory, new_head) = sub_gate_garbage_fixture("default-policy-reclaims");
let options = CompactionOptions::default().with_tasks([MaintenanceTask::Segments]);
let plan = plan_compaction(&directory.path, &options).expect("segment plan");
assert_eq!(plan.current_head(), new_head);
assert!(
!plan
.warnings()
.iter()
.any(|warning| warning.contains("rewrite gate")),
"the default policy never defers for savings: {:?}",
plan.warnings()
);
assert_eq!(
plan.retained_reclaimable_segments(),
0,
"nothing identified may be left behind on an archive with letters to spare"
);
assert_eq!(plan.retained_reclaimable_bytes(), 0);
assert!(
plan.estimated_reclaimable_bytes() != 0,
"the garbage the gate declined is now actually reclaimed"
);
assert!(
plan.actions().iter().any(|action| matches!(
action,
CompactionAction::RewriteArchive { file_name, replacement_name, .. }
if file_name == "data00001a.tar" && replacement_name == "data00001b.tar"
)),
"the sub-gate archive is rewritten to its next generation: {:?}",
plan.actions()
);
}
#[test]
fn the_plan_prices_the_history_veto_against_oaks_own_predicate() {
let (directory, _old_head, _new_head) = history_veto_fixture("history-veto-price");
let options = CompactionOptions::default().with_tasks([MaintenanceTask::Segments]);
let plan = plan_compaction(&directory.path, &options).expect("segment plan");
assert!(
plan.history_protected_segments() != 0,
"the bootstrap revision must be counted as history-only"
);
let (reclaimable_segments, reclaimable_bytes) = plan.history_protected_reclaimable();
assert!(
reclaimable_segments != 0,
"generation-zero history must be priced as reclaimable-but-protected"
);
assert!(reclaimable_bytes != 0);
assert!(reclaimable_segments <= plan.history_protected_segments());
}
#[test]
fn segment_cleanup_removes_old_unjournaled_archive_but_preserves_history() {
let directory = TestDirectory::repository("orphan-segment-history");
let old_head = Repository::open(&directory.path)
.expect("old repository")
.head_record_identifier();
{
let store = WritableRepository::open(&directory.path).expect("open orphan writer");
let mut writer = store.record_writer(GarbageCollectionGeneration {
generation: 0,
full_generation: 0,
is_compacted: false,
});
writer
.write_node(None, &[], &ChildNodesToWrite::Zero, &[])
.expect("orphan node");
writer.finish().expect("finish orphan segment");
store.close().expect("close orphan writer");
}
assert!(directory.path.join("data00001a.tar").is_file());
let new_head = {
let store = WritableRepository::open(&directory.path).expect("open new head writer");
let generation = GarbageCollectionGeneration {
generation: 2,
full_generation: 2,
is_compacted: false,
};
let mut writer = store.record_writer(generation);
let root = writer
.write_node(None, &[], &ChildNodesToWrite::Zero, &[])
.expect("new content root");
let head = writer
.write_node(
None,
&[],
&ChildNodesToWrite::One {
name: "root".to_owned(),
node: root,
},
&[],
)
.expect("new super root");
writer.finish().expect("finish new head");
assert!(store.compare_and_set_head(store.head(), head));
store.close().expect("close new head writer");
head
};
assert!(directory.path.join("data00002a.tar").is_file());
let options = CompactionOptions::default().with_tasks([MaintenanceTask::Segments]);
let plan = plan_compaction(&directory.path, &options).expect("segment plan");
assert!(plan.actions().iter().any(|action| matches!(
action,
CompactionAction::RemoveReclaimableArchive { file_name, .. }
if file_name == "data00001a.tar"
)));
assert!(!plan.actions().iter().any(|action| matches!(
action,
CompactionAction::RemoveReclaimableArchive { file_name, .. }
if file_name == "data00000a.tar"
)));
let planned_removed_segments: usize = plan
.actions()
.iter()
.filter_map(|action| match action {
CompactionAction::RemoveReclaimableArchive { segments, .. }
| CompactionAction::RewriteArchive { segments, .. } => Some(*segments),
_ => None,
})
.sum();
assert!(planned_removed_segments != 0);
let outcome = compact(&directory.path, options).expect("segment cleanup");
assert_eq!(outcome.head_after, new_head);
assert_eq!(outcome.removed_segments(), planned_removed_segments);
assert!(!directory.path.join("data00001a.tar").exists());
let repository = Repository::open(&directory.path).expect("healthy final repository");
assert_eq!(repository.head_record_identifier(), new_head);
crate::tooling::verify_node_tree(&repository, old_head)
.expect("historical root remains readable");
}
#[test]
fn a_dead_survivor_pointing_at_a_removed_segment_is_handled() {
let directory = TestDirectory::repository("dead-survivor-reference");
let dead_generation = GarbageCollectionGeneration {
generation: 0,
full_generation: 0,
is_compacted: false,
};
let live_generation = GarbageCollectionGeneration {
generation: 2,
full_generation: 2,
is_compacted: false,
};
let target = {
let store = WritableRepository::open(&directory.path).expect("open target writer");
let mut writer = store.record_writer(dead_generation);
let node = writer
.write_node(None, &[], &ChildNodesToWrite::Zero, &[])
.expect("dead target node");
writer.finish().expect("finish target");
store.close().expect("close target writer");
node
};
let new_head = {
let store = WritableRepository::open(&directory.path).expect("open mixed writer");
let mut referencing = store.record_writer(dead_generation);
referencing
.write_node(
None,
&[],
&ChildNodesToWrite::One {
name: "target".to_owned(),
node: target,
},
&[],
)
.expect("dead referencing node");
referencing.finish().expect("finish referencing");
for _ in 0..8 {
let mut filler = store.record_writer(live_generation);
filler
.write_node(None, &[], &ChildNodesToWrite::Zero, &[])
.expect("live filler");
filler.finish().expect("finish filler");
}
let mut writer = store.record_writer(live_generation);
let root = writer
.write_node(None, &[], &ChildNodesToWrite::Zero, &[])
.expect("content root");
let head = writer
.write_node(
None,
&[],
&ChildNodesToWrite::One {
name: "root".to_owned(),
node: root,
},
&[],
)
.expect("super root");
writer.finish().expect("finish head");
assert!(store.compare_and_set_head(store.head(), head));
store.close().expect("close mixed writer");
head
};
let options = CompactionOptions::default().with_tasks([MaintenanceTask::Segments]);
let plan = plan_compaction(&directory.path, &options)
.expect("a dead survivor pointing at removed garbage must not refuse the plan");
let removed: usize = plan
.actions()
.iter()
.filter_map(|action| match action {
CompactionAction::RemoveReclaimableArchive { segments, .. }
| CompactionAction::RewriteArchive { segments, .. } => Some(*segments),
_ => None,
})
.sum();
assert!(
removed != 0,
"the fixture must actually remove something, or it proves nothing"
);
let outcome = compact(&directory.path, options).expect("apply the plan");
assert_eq!(outcome.head_after, new_head);
assert!(outcome.removed_segments() != 0);
let repository = Repository::open(&directory.path).expect("healthy store");
assert_eq!(repository.head_record_identifier(), new_head);
crate::tooling::verify_node_tree(&repository, new_head).expect("head verifies");
}
}