use crate::content::provider::SegmentProvider as _;
use super::{
BTreeSet, CheckpointPlan, CompactionKind, Error, GarbageCollectionGeneration, HashSet,
HistoryProtection, JournalAnalysis, MaintenanceTask, NodeTreeVerifier, Path,
PlannedArchiveSweep, ProgressObserver, ReclaimRule, RecordIdentifier, RecordType, Repository,
Result, SegmentIdentifier, StaleArchive, StandaloneSegmentCompactionPlan, Step, WorkUnit,
active_index_generations, compaction_target_generation, extend_segment_closure,
generation_from_header, plan_stale_archives, plan_standalone_segment_cleanup,
predict_shared_bulk_segments, prospective_retained_roots, segments_ahead_of_the_head,
validate_prospective_segment_plan, validate_reclaim_reference_invariant,
};
pub(crate) struct HeadClosureInputs<'inputs> {
pub(crate) repository: &'inputs Repository,
pub(crate) options: &'inputs crate::writer::maintenance::options::CompactionOptions,
pub(crate) checkpoints: &'inputs CheckpointPlan,
pub(crate) current_head: RecordIdentifier,
pub(crate) reference_generation: GarbageCollectionGeneration,
pub(crate) reclaim_rule: ReclaimRule,
pub(crate) index_available: bool,
}
pub(crate) struct HeadClosure {
pub(crate) segments: HashSet<SegmentIdentifier>,
pub(crate) residue_segments: usize,
pub(crate) indexed_bytes: Option<(u64, u64)>,
pub(crate) uniformly_compacted_at_reference: bool,
}
pub(crate) fn trace_head_closure(
inputs: &HeadClosureInputs<'_>,
observer: &mut dyn ProgressObserver,
) -> Result<HeadClosure> {
let &HeadClosureInputs {
repository,
options,
checkpoints,
current_head,
reference_generation,
reclaim_rule,
index_available,
} = inputs;
let mut current_closure = HashSet::new();
let mut residue_segments = 0usize;
let mut indexed_bytes = None;
let mut uniformly_compacted_at_reference = false;
if index_available
&& (options.contains(MaintenanceTask::Segments) || !checkpoints.names.is_empty())
{
let active_index_generations = active_index_generations(repository)?;
residue_segments =
segments_ahead_of_the_head(&active_index_generations, reference_generation);
crate::progress::observe(
observer,
&Step::new(
"tracing segments reachable from the head",
WorkUnit::Segments,
),
|observer| {
extend_segment_closure(
repository,
[current_head.segment],
&mut current_closure,
observer,
)?;
let (data_bytes, bulk_bytes) =
crate::writer::maintenance::reclamation::indexed_bytes_by_kind(
repository,
¤t_closure,
);
indexed_bytes = Some((data_bytes, bulk_bytes));
if bulk_bytes == 0 {
observer.step_concluded(&format!(
"the head reaches {} of node data",
crate::units::format_byte_size(data_bytes),
));
} else {
observer.step_concluded(&format!(
"the head reaches {} of node data and {} of shared binary blocks",
crate::units::format_byte_size(data_bytes),
crate::units::format_byte_size(bulk_bytes),
));
}
Ok::<(), Error>(())
},
)?;
validate_reclaim_reference_invariant(
repository,
¤t_closure,
&active_index_generations,
reclaim_rule,
)?;
uniformly_compacted_at_reference = reference_generation.is_compacted
&& current_closure
.iter()
.filter(|identifier| identifier.is_data_segment())
.all(|identifier| {
active_index_generations.get(identifier) == Some(&reference_generation)
});
}
Ok(HeadClosure {
segments: current_closure,
residue_segments,
indexed_bytes,
uniformly_compacted_at_reference,
})
}
pub(crate) struct PredictedSweepInputs<'inputs> {
pub(crate) directory: &'inputs Path,
pub(crate) repository: &'inputs Repository,
pub(crate) options: &'inputs crate::writer::maintenance::options::CompactionOptions,
pub(crate) compaction_kind: Option<CompactionKind>,
pub(crate) index_available: bool,
pub(crate) current_head: RecordIdentifier,
pub(crate) reference_generation: GarbageCollectionGeneration,
pub(crate) residue_segments: usize,
pub(crate) checkpoints: &'inputs CheckpointPlan,
pub(crate) stale_archives: &'inputs [StaleArchive],
}
pub(crate) fn predict_sweeps(
inputs: &PredictedSweepInputs<'_>,
observer: &mut dyn ProgressObserver,
) -> Result<(
Option<StandaloneSegmentCompactionPlan>,
Option<StandaloneSegmentCompactionPlan>,
)> {
let &PredictedSweepInputs {
directory,
repository,
options,
compaction_kind,
index_available,
current_head,
reference_generation,
residue_segments,
checkpoints,
stale_archives,
} = inputs;
let residue_sweep = if index_available && residue_segments != 0 {
Some(crate::progress::observe(
observer,
&Step::new("planning residue retirement", WorkUnit::Archives)
.with_total(crate::progress::count(repository.archives().len())),
|observer| {
plan_standalone_segment_cleanup(
directory,
repository,
ReclaimRule {
reference: reference_generation,
kind: CompactionKind::Full,
retained_generations: i32::MAX,
},
current_head.segment,
&HashSet::new(),
options.archive_rewrite_policy,
observer,
)
},
)?)
} else {
None
};
let predicted_sweep = match compaction_kind {
Some(kind) if index_available => {
let omitted: BTreeSet<String> = checkpoints.names.iter().cloned().collect();
let shared =
predict_shared_bulk_segments(repository, current_head, &omitted, observer)?;
let target = compaction_target_generation(reference_generation, kind);
let absent: HashSet<String> = stale_archives
.iter()
.map(|stale| stale.file_name.clone())
.collect();
Some(crate::progress::observe(
observer,
&Step::new("predicting the reclamation", WorkUnit::Archives)
.with_total(crate::progress::count(repository.archives().len())),
|observer| {
let predicted =
crate::writer::store_writer::predict_post_compaction_reclamation(
directory,
repository,
ReclaimRule {
reference: target,
kind,
retained_generations:
crate::writer::store_writer::RETAINED_GENERATIONS,
},
&shared,
options.archive_rewrite_policy,
&absent,
)?;
observer.step_concluded(&predicted_sweep_conclusion(&predicted));
Ok::<_, Error>(predicted)
},
)?)
}
_ => None,
};
Ok((residue_sweep, predicted_sweep))
}
fn predicted_sweep_conclusion(predicted: &StandaloneSegmentCompactionPlan) -> String {
let mut removed_archives = 0u64;
let mut removed_bytes = 0u64;
let mut rewritten_archives = 0u64;
let mut rewritten_entry_bytes = 0u64;
for archive in &predicted.archives {
match archive {
PlannedArchiveSweep::Remove { file_bytes, .. } => {
removed_archives += 1;
removed_bytes = removed_bytes.saturating_add(*file_bytes);
}
PlannedArchiveSweep::Rewrite {
eligible_entry_bytes,
..
} => {
rewritten_archives += 1;
rewritten_entry_bytes = rewritten_entry_bytes.saturating_add(*eligible_entry_bytes);
}
PlannedArchiveSweep::DeferredBySavings { .. }
| PlannedArchiveSweep::DeferredAtLastGeneration { .. }
| PlannedArchiveSweep::BlockedByOccupiedGeneration { .. } => {}
}
}
let archives_phrase = |count: u64| {
let noun = if count == 1 { "archive" } else { "archives" };
format!("{} {noun}", crate::units::format_count(count))
};
match (removed_archives, rewritten_archives) {
(0, 0) => "the sweep has nothing to reclaim".to_owned(),
(_, 0) => format!(
"the sweep removes {} ({})",
archives_phrase(removed_archives),
crate::units::format_byte_size(removed_bytes),
),
(0, _) => format!(
"the sweep rewrites {} ({} of entries)",
archives_phrase(rewritten_archives),
crate::units::format_byte_size(rewritten_entry_bytes),
),
(_, _) => format!(
"the sweep removes {} ({}) and rewrites {} ({} of entries)",
archives_phrase(removed_archives),
crate::units::format_byte_size(removed_bytes),
archives_phrase(rewritten_archives),
crate::units::format_byte_size(rewritten_entry_bytes),
),
}
}
pub(crate) struct SegmentPlanInputs<'inputs> {
pub(crate) directory: &'inputs Path,
pub(crate) repository: &'inputs Repository,
pub(crate) options: &'inputs crate::writer::maintenance::options::CompactionOptions,
pub(crate) effective_compaction_kind: Option<CompactionKind>,
pub(crate) index_available: bool,
pub(crate) current_head: RecordIdentifier,
pub(crate) reclaim_rule: ReclaimRule,
pub(crate) journal_analysis: &'inputs JournalAnalysis,
}
pub(crate) fn trace_history_closure(
repository: &Repository,
journal_analysis: &JournalAnalysis,
retained_closure: &mut HashSet<SegmentIdentifier>,
observer: &mut dyn ProgressObserver,
) -> Result<()> {
crate::progress::observe(
observer,
&Step::new(
"tracing segments reachable from history",
WorkUnit::Segments,
),
|observer| {
extend_segment_closure(
repository,
journal_analysis
.retained_record_ids
.iter()
.map(|record| record.segment),
retained_closure,
observer,
)
},
)?;
Ok(())
}
pub(crate) fn plan_standalone_segments(
inputs: &SegmentPlanInputs<'_>,
current_closure: HashSet<SegmentIdentifier>,
protected_history_segments: &mut HashSet<SegmentIdentifier>,
history_protection: &mut HistoryProtection,
observer: &mut dyn ProgressObserver,
) -> Result<Option<StandaloneSegmentCompactionPlan>> {
let &SegmentPlanInputs {
directory,
repository,
options,
effective_compaction_kind,
index_available,
current_head,
reclaim_rule,
journal_analysis,
} = inputs;
let segment_plan = if index_available
&& options.contains(MaintenanceTask::Segments)
&& effective_compaction_kind.is_none()
{
let head_data_segments: HashSet<SegmentIdentifier> = current_closure
.iter()
.copied()
.filter(|identifier| identifier.is_data_segment())
.collect();
let mut retained_closure = current_closure;
trace_history_closure(
repository,
journal_analysis,
&mut retained_closure,
observer,
)?;
protected_history_segments.extend(
retained_closure
.into_iter()
.filter(|identifier| identifier.is_data_segment()),
);
history_protection.history_only_segments = protected_history_segments
.iter()
.filter(|identifier| !head_data_segments.contains(identifier))
.count();
let plan = crate::progress::observe(
observer,
&Step::new("planning segment reclamation", WorkUnit::Archives)
.with_total(crate::progress::count(repository.archives().len())),
|observer| {
plan_standalone_segment_cleanup(
directory,
repository,
reclaim_rule,
current_head.segment,
protected_history_segments,
options.archive_rewrite_policy,
observer,
)
},
)?;
if history_protection.history_only_segments != 0 {
let (unvetoed_segments, unvetoed_bytes) = crate::progress::observe(
observer,
&Step::new("pricing the journal-history protection", WorkUnit::Archives)
.with_total(crate::progress::count(repository.archives().len())),
|observer| {
crate::writer::store_writer::measure_unvetoed_reclamation(
directory,
repository,
reclaim_rule,
current_head.segment,
options.archive_rewrite_policy,
observer,
)
},
)?;
let (vetoed_segments, vetoed_bytes) =
crate::writer::store_writer::plan_reclaimed_totals(&plan);
history_protection.would_be_reclaimable_segments =
unvetoed_segments.saturating_sub(vetoed_segments);
history_protection.would_be_reclaimable_bytes =
unvetoed_bytes.saturating_sub(vetoed_bytes);
}
let retained_roots = prospective_retained_roots(
directory,
repository,
&plan,
&journal_analysis.retained_record_ids,
);
crate::progress::observe(
observer,
&Step::new("validating the prospective plan", WorkUnit::Nodes),
|observer| {
validate_prospective_segment_plan(
directory,
repository,
&plan,
&retained_roots,
observer,
)
},
)?;
Some(plan)
} else {
None
};
Ok(segment_plan)
}
pub(crate) struct SegmentWork {
pub(crate) stale_archives: Vec<StaleArchive>,
pub(crate) reference_generation: GarbageCollectionGeneration,
pub(crate) segment_plan: Option<StandaloneSegmentCompactionPlan>,
pub(crate) residue_sweep: Option<StandaloneSegmentCompactionPlan>,
pub(crate) predicted_sweep: Option<StandaloneSegmentCompactionPlan>,
pub(crate) protected_history_segments: HashSet<SegmentIdentifier>,
pub(crate) history_protection: HistoryProtection,
pub(crate) residue_segments: usize,
pub(crate) closure_indexed_bytes: Option<(u64, u64)>,
pub(crate) effective_compaction_kind: Option<CompactionKind>,
pub(crate) already_fully_compacted: bool,
}
pub(crate) struct JournalConvergence {
pub(crate) single_line_naming_head: bool,
}
fn convergence_gate(
options: &crate::writer::maintenance::options::CompactionOptions,
closure: &HeadClosure,
checkpoints: &CheckpointPlan,
journal: &JournalConvergence,
purge_selected: bool,
) -> (Option<CompactionKind>, bool) {
let Some(kind) = options.compaction_kind else {
return (None, false);
};
let already_fully_compacted = closure.uniformly_compacted_at_reference
&& checkpoints.names.is_empty()
&& journal.single_line_naming_head
&& !purge_selected;
if already_fully_compacted && !options.always_copy {
(None, true)
} else {
(Some(kind), already_fully_compacted)
}
}
pub(crate) struct SegmentWorkInputs<'inputs> {
pub(crate) directory: &'inputs Path,
pub(crate) repository: &'inputs Repository,
pub(crate) options: &'inputs crate::writer::maintenance::options::CompactionOptions,
pub(crate) current_head: RecordIdentifier,
pub(crate) index_available: bool,
pub(crate) checkpoints: &'inputs CheckpointPlan,
pub(crate) journal_analysis: &'inputs JournalAnalysis,
pub(crate) journal_convergence: &'inputs JournalConvergence,
pub(crate) purge_selected: bool,
}
pub(crate) fn plan_segment_work(
inputs: &SegmentWorkInputs<'_>,
warnings: &mut Vec<String>,
observer: &mut dyn ProgressObserver,
) -> Result<SegmentWork> {
let &SegmentWorkInputs {
directory,
repository,
options,
current_head,
index_available,
checkpoints,
journal_analysis,
journal_convergence,
purge_selected,
} = inputs;
let stale_archives = if index_available && options.contains(MaintenanceTask::StaleArchives) {
crate::progress::observe(
observer,
&Step::new("scanning for stale archives", WorkUnit::Archives)
.with_total(crate::progress::count(repository.archives().len())),
|observer| plan_stale_archives(directory, repository, warnings, observer),
)?
} else {
Vec::new()
};
let reference_generation = generation_from_header(repository, current_head.segment)?;
let reclaim_rule = ReclaimRule {
reference: reference_generation,
kind: CompactionKind::Full,
retained_generations: crate::writer::store_writer::RETAINED_GENERATIONS,
};
let closure = trace_head_closure(
&HeadClosureInputs {
repository,
options,
checkpoints,
current_head,
reference_generation,
reclaim_rule,
index_available,
},
observer,
)?;
let (effective_compaction_kind, already_fully_compacted) = convergence_gate(
options,
&closure,
checkpoints,
journal_convergence,
purge_selected,
);
let HeadClosure {
segments: current_closure,
residue_segments,
indexed_bytes: closure_indexed_bytes,
..
} = closure;
let (residue_sweep, predicted_sweep) = predict_sweeps(
&PredictedSweepInputs {
directory,
repository,
options,
compaction_kind: effective_compaction_kind,
index_available,
current_head,
reference_generation,
residue_segments,
checkpoints,
stale_archives: &stale_archives,
},
observer,
)?;
let mut protected_history_segments = HashSet::new();
let mut history_protection = HistoryProtection::default();
let segment_plan = plan_standalone_segments(
&SegmentPlanInputs {
directory,
repository,
options,
effective_compaction_kind,
index_available,
current_head,
reclaim_rule,
journal_analysis,
},
current_closure,
&mut protected_history_segments,
&mut history_protection,
observer,
)?;
Ok(SegmentWork {
stale_archives,
reference_generation,
segment_plan,
residue_sweep,
predicted_sweep,
protected_history_segments,
history_protection,
residue_segments,
closure_indexed_bytes,
effective_compaction_kind,
already_fully_compacted,
})
}
pub(in crate::writer::maintenance) fn verify_exact_super_root(
repository: &Repository,
head: RecordIdentifier,
observer: &mut dyn ProgressObserver,
) -> Result<()> {
verify_exact_super_root_counting_nodes(repository, head, observer).map(|_| ())
}
pub(in crate::writer::maintenance) fn verify_head_for_planning(
repository: &Repository,
head: RecordIdentifier,
census: &mut super::PlanningContentCensus,
collect_internal_identifiers: bool,
observer: &mut dyn ProgressObserver,
) -> Result<u64> {
let view = repository.segment(head.segment)?;
if view.structure.record_type(head.record_number) != Some(RecordType::Node) {
return Err(Error::InvalidFormat {
details: format!("current journal head {head} is not a node record"),
});
}
let super_root = repository.node(head);
let content_root = super_root
.child_node("root")?
.ok_or_else(|| Error::InvalidFormat {
details: format!("journal root {head} has no content \"root\" child node"),
})?;
let mut verifier = NodeTreeVerifier::new(repository);
let version_storage = content_root
.child_node("jcr:system")?
.map(|system| system.child_node("jcr:versionStorage"))
.transpose()?
.flatten();
crate::progress::observe(
observer,
&Step::new("verifying the current head", WorkUnit::Nodes),
|observer| {
census.match_live_identifiers = true;
let outcome = verifier.verify_collecting_pruned_with_progress(
content_root.record_identifier(),
census,
&["/jcr:system/jcr:versionStorage"],
observer,
);
census.match_live_identifiers = false;
outcome
},
)?;
if let Some(version_storage) = &version_storage {
crate::progress::observe(
observer,
&Step::new("scanning version storage", WorkUnit::Nodes),
|observer| {
let mut scan = super::VersionStoragePreScan::new(
&mut census.version_storage,
&mut census.external_binaries,
collect_internal_identifiers,
);
verifier.verify_collecting_with_progress(
version_storage.record_identifier(),
&mut scan,
observer,
)
},
)?;
}
census.version_storage.resolve_live_matches();
crate::progress::observe(
observer,
&Step::new("verifying the checkpoints", WorkUnit::Nodes),
|observer| verifier.verify_collecting_with_progress(head, census, observer),
)?;
if let Some(checkpoints) = super_root.child_node("checkpoints")? {
for (name, checkpoint) in checkpoints.child_node_entries()? {
checkpoint
.child_node("root")?
.ok_or_else(|| Error::InvalidFormat {
details: format!(
"checkpoint {name} under journal root {head} has no snapshot \"root\" child node"
),
})?;
}
}
Ok(verifier.verified_nodes())
}
pub(in crate::writer::maintenance) fn verify_exact_super_root_counting_nodes(
repository: &Repository,
head: RecordIdentifier,
observer: &mut dyn ProgressObserver,
) -> Result<u64> {
verify_exact_super_root_collecting(
repository,
head,
&mut crate::tooling::DiscardedVerifiedContent,
observer,
)
}
pub(in crate::writer::maintenance) fn verify_exact_super_root_collecting(
repository: &Repository,
head: RecordIdentifier,
content: &mut dyn crate::tooling::VerifiedContentObserver,
observer: &mut dyn ProgressObserver,
) -> Result<u64> {
crate::progress::observe(
observer,
&Step::new("verifying the current head", WorkUnit::Nodes),
|observer| {
let mut verifier = NodeTreeVerifier::new(repository);
verify_exact_super_root_collecting_with_verifier(
repository,
head,
&mut verifier,
content,
observer,
)?;
Ok(verifier.verified_nodes())
},
)
}
pub(in crate::writer::maintenance) fn verify_exact_super_root_with_verifier(
repository: &Repository,
head: RecordIdentifier,
verifier: &mut NodeTreeVerifier<'_>,
observer: &mut dyn ProgressObserver,
) -> Result<()> {
verify_exact_super_root_collecting_with_verifier(
repository,
head,
verifier,
&mut crate::tooling::DiscardedVerifiedContent,
observer,
)
}
pub(in crate::writer::maintenance) fn verify_exact_super_root_collecting_with_verifier(
provider: &dyn crate::content::provider::SegmentProvider,
head: RecordIdentifier,
verifier: &mut NodeTreeVerifier<'_>,
content: &mut dyn crate::tooling::VerifiedContentObserver,
observer: &mut dyn ProgressObserver,
) -> Result<()> {
let view = provider.segment(head.segment)?;
if view.structure.record_type(head.record_number) != Some(RecordType::Node) {
return Err(Error::InvalidFormat {
details: format!("current journal head {head} is not a node record"),
});
}
verifier.verify_collecting_with_progress(head, content, observer)?;
let super_root = crate::content::node::NodeState::new(provider, head);
super_root
.child_node("root")?
.ok_or_else(|| Error::InvalidFormat {
details: format!("journal root {head} has no content \"root\" child node"),
})?;
if let Some(checkpoints) = super_root.child_node("checkpoints")? {
for (name, checkpoint) in checkpoints.child_node_entries()? {
checkpoint
.child_node("root")?
.ok_or_else(|| Error::InvalidFormat {
details: format!(
"checkpoint {name} under journal root {head} has no snapshot \"root\" child node"
),
})?;
}
}
Ok(())
}