use super::*;
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) fn trace_head_closure(
inputs: &HeadClosureInputs<'_>,
observer: &mut dyn ProgressObserver,
) -> Result<(HashSet<SegmentIdentifier>, usize)> {
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;
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,
)
},
)?;
validate_reclaim_reference_invariant(
repository,
¤t_closure,
&active_index_generations,
reclaim_rule,
)?;
}
Ok((current_closure, residue_segments))
}
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) 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,
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 options.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| {
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,
)
},
)?)
}
_ => None,
};
Ok((residue_sweep, predicted_sweep))
}
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) 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,
index_available,
current_head,
reclaim_rule,
journal_analysis,
} = inputs;
let segment_plan = if index_available
&& options.contains(MaintenanceTask::Segments)
&& options.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) struct TracedSweeps {
pub(crate) current_closure: HashSet<SegmentIdentifier>,
pub(crate) residue_sweep: Option<StandaloneSegmentCompactionPlan>,
pub(crate) predicted_sweep: Option<StandaloneSegmentCompactionPlan>,
pub(crate) residue_segments: usize,
}
#[allow(
clippy::too_many_arguments,
reason = "the trace and the prediction judge the same store against the same rule"
)]
pub(crate) fn trace_and_predict(
directory: &Path,
repository: &Repository,
options: &crate::writer::maintenance::options::CompactionOptions,
current_head: RecordIdentifier,
index_available: bool,
checkpoints: &CheckpointPlan,
stale_archives: &[StaleArchive],
reference_generation: GarbageCollectionGeneration,
reclaim_rule: ReclaimRule,
observer: &mut dyn ProgressObserver,
) -> Result<TracedSweeps> {
let (current_closure, residue_segments) = trace_head_closure(
&HeadClosureInputs {
repository,
options,
checkpoints,
current_head,
reference_generation,
reclaim_rule,
index_available,
},
observer,
)?;
let (residue_sweep, predicted_sweep) = predict_sweeps(
&PredictedSweepInputs {
directory,
repository,
options,
index_available,
current_head,
reference_generation,
residue_segments,
checkpoints,
stale_archives,
},
observer,
)?;
Ok(TracedSweeps {
current_closure,
residue_sweep,
predicted_sweep,
residue_segments,
})
}
#[allow(
clippy::too_many_arguments,
reason = "every phase here judges the same store against the same head and options"
)]
pub(crate) fn plan_segment_work(
directory: &Path,
repository: &Repository,
options: &crate::writer::maintenance::options::CompactionOptions,
current_head: RecordIdentifier,
index_available: bool,
checkpoints: &CheckpointPlan,
journal_analysis: &JournalAnalysis,
warnings: &mut Vec<String>,
observer: &mut dyn ProgressObserver,
) -> Result<SegmentWork> {
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 TracedSweeps {
current_closure,
residue_sweep,
predicted_sweep,
residue_segments,
} = trace_and_predict(
directory,
repository,
options,
current_head,
index_available,
checkpoints,
&stale_archives,
reference_generation,
reclaim_rule,
observer,
)?;
let mut protected_history_segments = HashSet::new();
let mut history_protection = HistoryProtection::default();
let segment_plan = plan_standalone_segments(
&SegmentPlanInputs {
directory,
repository,
options,
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,
})
}
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_exact_super_root_counting_nodes(
repository: &Repository,
head: RecordIdentifier,
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_with_verifier(repository, head, &mut verifier, 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<()> {
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"),
});
}
verifier.verify_with_progress(head, observer)?;
let super_root = repository.node(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(())
}