use super::apply_identity::append_apply_identity_preview_warning;
use super::checkpoints::plan_checkpoints;
use super::journal::JournalAnalysis;
use super::journal::analyze_journal;
use super::manifest::upgrade_manifest_atomically;
use super::options::{CompactionOptions, MaintenanceTask};
use super::plan::{
CompactionAction, CompactionPlan, HistoryProtection, JournalLineRemoval, RetainedReclaimable,
StaleArchiveReason,
};
use super::reclamation::{
active_index_generations, compaction_target_generation, extend_segment_closure,
predict_shared_bulk_segments, prospective_retained_roots, retained_reclaimable_from,
segments_ahead_of_the_head, validate_prospective_segment_plan,
validate_reclaim_reference_invariant,
};
use super::recovery_backups::{plan_recovery_backups, recovery_backup_target};
use super::stale_archives::{
generation_from_header, plan_stale_archives, planned_archive_repairs,
reject_cross_number_duplicate_active_segments, reject_duplicate_active_segments,
unrepairable_archive_names,
};
use super::temporaries::{plan_stale_temporaries, temporary_kind};
use crate::error::{Error, Result};
use crate::progress::{ProgressObserver, Step, WorkUnit};
use crate::segment::identifier::SegmentIdentifier;
use crate::segment::record::{RecordIdentifier, RecordType};
use crate::store::Repository;
use crate::tar_archive::file_name::ArchiveFileName;
use crate::tooling::NodeTreeVerifier;
use crate::writer::compaction::CompactionKind;
use crate::writer::maintenance::journal::RawJournal;
use crate::writer::maintenance::journal::scan_raw_journal;
use crate::writer::segment_builder::GarbageCollectionGeneration;
use crate::writer::store_writer::StandaloneSegmentCompactionPlan;
use crate::writer::store_writer::{
PlannedArchiveSweep, ReclaimRule, next_cleanup_archive_number, plan_standalone_segment_cleanup,
};
use std::collections::{BTreeSet, HashSet};
use std::ffi::{OsStr, OsString};
use std::fmt::Write as _;
use std::fs::Metadata;
#[cfg(unix)]
use std::os::unix::fs::MetadataExt;
use std::path::{Path, PathBuf};
use std::time::SystemTime;
mod content_census;
mod listing;
mod segments;
mod shape;
mod version_history_plan;
mod version_storage;
pub(crate) use content_census::*;
pub(crate) use listing::*;
pub(crate) use segments::*;
pub(in crate::writer::maintenance) use shape::*;
pub(crate) use shape::{
DirectoryFingerprint, FileFingerprint, canonical_repository_directory, directory_fingerprint,
validate_repository_shape,
};
use version_history_plan::{select_version_history_purge, version_history_plan_parts};
pub(crate) use version_storage::*;
#[derive(Clone, Debug, PartialEq, Eq)]
pub(super) struct JournalPlan {
pub(super) retained_record_ids: Vec<RecordIdentifier>,
pub(super) retained_raw_lines: Vec<Vec<u8>>,
pub(super) removals: Vec<JournalLineRemoval>,
pub(super) removed_lines: usize,
pub(super) parser_ignored: usize,
pub(super) missing_segments: usize,
pub(super) unreadable_revisions: usize,
pub(super) beyond_retention: usize,
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub(super) struct CheckpointPlan {
pub(super) names: Vec<String>,
pub(super) expired: usize,
pub(super) unreferenced: usize,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(super) struct StaleArchive {
pub(super) file_name: String,
pub(super) reason: StaleArchiveReason,
pub(super) bytes: u64,
pub(super) fingerprint: FileFingerprint,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(super) struct PlannedFileRemoval {
pub(super) file_name: String,
pub(super) bytes: u64,
pub(super) fingerprint: FileFingerprint,
}
pub(super) fn build_plan(
directory: &Path,
options: &CompactionOptions,
now: SystemTime,
observer: &mut dyn ProgressObserver,
) -> Result<CompactionPlan> {
let mut warnings = Vec::new();
build_plan_collecting(directory, options, now, observer, &mut warnings)
.map_err(|error| attach_planning_warnings(error, &warnings))
}
pub(super) fn attach_completed_repairs(
error: Error,
repaired: &[crate::writer::store_writer::RepairedArchive],
) -> Error {
if repaired.is_empty() {
return error;
}
let names: Vec<&str> = repaired
.iter()
.map(|archive| archive.file_name.as_str())
.collect();
let rebuilds = if repaired.len() == 1 {
"1 archive index rebuild, which is already durable".to_owned()
} else {
format!(
"{} archive index rebuilds, which are already durable",
repaired.len()
)
};
Error::InvalidFormat {
details: format!(
"{error} This refusal came after {rebuilds}: {}. The originals are retained under \
`.bak` names, and those archives need no second attempt.",
names.join(", ")
),
}
}
pub(super) fn attach_planning_warnings(error: Error, warnings: &[String]) -> Error {
const WARNINGS_SHOWN: usize = 3;
match error {
Error::InvalidFormat { details } if !warnings.is_empty() => {
let mut attached = warnings
.iter()
.take(WARNINGS_SHOWN)
.map(String::as_str)
.collect::<Vec<_>>()
.join("; ");
let remaining = warnings.len() - warnings.len().min(WARNINGS_SHOWN);
if remaining == 1 {
attached.push_str(", and 1 further warning");
} else if remaining > 1 {
let _ = write!(attached, ", and {remaining} further warnings");
}
Error::InvalidFormat {
details: format!("{details} Also established before the refusal: {attached}."),
}
}
other => other,
}
}
pub(super) struct ManifestUpgradeOnFirstInstall<'directory> {
pub(super) directory: &'directory Path,
pub(super) done: bool,
}
impl<'directory> ManifestUpgradeOnFirstInstall<'directory> {
pub(super) fn new(directory: &'directory Path) -> Self {
Self {
directory,
done: false,
}
}
}
impl crate::writer::store_writer::AuthorizeVersionTwoWrite for ManifestUpgradeOnFirstInstall<'_> {
fn authorize(&mut self) -> Result<()> {
if self.done {
return Ok(());
}
if crate::store::read_manifest_store_version(&self.directory.join("manifest"))? == 1 {
ensure_numbered_name_available(self.directory, "manifest.cleaning")?;
upgrade_manifest_atomically(self.directory)?;
}
self.done = true;
Ok(())
}
}
pub(crate) struct RepositoryState {
pub(crate) raw_journal: RawJournal,
pub(crate) journal_analysis: JournalAnalysis,
pub(crate) pending_repairs: Vec<CompactionAction>,
pub(crate) index_available: bool,
pub(crate) checkpoints: CheckpointPlan,
pub(crate) checkpoint_archive_number: Option<u32>,
}
pub(crate) fn survey_repository_state(
directory: &Path,
repository: &Repository,
options: &super::options::CompactionOptions,
current_head: RecordIdentifier,
now: SystemTime,
warnings: &mut Vec<String>,
observer: &mut dyn ProgressObserver,
) -> Result<RepositoryState> {
let raw_journal = scan_raw_journal(directory)?;
let journal_analysis = analyze_journal(
repository,
&raw_journal,
current_head,
options.journal_revision_retention,
observer,
)?;
let pending_repairs = if options.contains(MaintenanceTask::RepairArchives) {
let unrepairable = unrepairable_archive_names(repository);
if !unrepairable.is_empty() {
return Err(crate::writer::store_writer::unrepairable_archives_refusal(
&unrepairable,
));
}
planned_archive_repairs(repository)
} else {
Vec::new()
};
let index_available = pending_repairs.is_empty();
let checkpoints = if index_available {
plan_checkpoints(repository, options, now, warnings)?
} else {
CheckpointPlan::default()
};
if options.journal_revision_retention.is_some() && !checkpoints.names.is_empty() {
return Err(Error::InvalidFormat {
details: format!(
"a journal revision retention bound cannot run beside the removal of {} checkpoint(s): \
checkpoint removal moves the head and appends a journal line, which would put the \
bound's newest revision beyond the one this plan retained — remove the checkpoints \
first, then bound the journal in a second run",
checkpoints.names.len()
),
});
}
let checkpoint_archive_number =
if checkpoints.names.is_empty() && options.compaction_kind.is_none() {
None
} else {
Some(next_cleanup_archive_number(directory)?)
};
if options.contains(MaintenanceTask::Segments)
|| options.contains(MaintenanceTask::StaleArchives)
|| !checkpoints.names.is_empty()
{
if index_available {
reject_duplicate_active_segments(repository)?;
} else {
reject_cross_number_duplicate_active_segments(repository)?;
}
}
Ok(RepositoryState {
raw_journal,
journal_analysis,
pending_repairs,
index_available,
checkpoints,
checkpoint_archive_number,
})
}
#[allow(
clippy::too_many_arguments,
reason = "both scans read the same store, options, clock and reporter"
)]
pub(crate) fn plan_leftover_files(
directory: &Path,
repository: &Repository,
options: &super::options::CompactionOptions,
index_available: bool,
raw_journal: &RawJournal,
now: SystemTime,
warnings: &mut Vec<String>,
observer: &mut dyn ProgressObserver,
) -> Result<(Vec<PlannedFileRemoval>, Vec<PlannedFileRemoval>)> {
let temporaries = if index_available && options.contains(MaintenanceTask::StaleTemporaries) {
crate::progress::observe(
observer,
&Step::new("scanning for stale temporary files", WorkUnit::Files),
|observer| {
plan_stale_temporaries(directory, repository, raw_journal, warnings, observer)
},
)?
} else {
Vec::new()
};
let recovery_backups = if options.contains(MaintenanceTask::RecoveryBackups) {
plan_recovery_backups(
directory,
now,
options
.recovery_backup_policy
.expect("validated recovery backup policy"),
)?
} else {
Vec::new()
};
Ok((temporaries, recovery_backups))
}
#[allow(
clippy::too_many_arguments,
reason = "every task that writes version-2 data has a say in this one decision"
)]
pub(crate) fn decide_manifest_upgrade(
directory: &Path,
_options: &super::options::CompactionOptions,
pending_repairs: &[CompactionAction],
checkpoints: &CheckpointPlan,
segment_plan: Option<&StandaloneSegmentCompactionPlan>,
) -> Result<bool> {
let writes_v2 = !pending_repairs.is_empty()
|| !checkpoints.names.is_empty()
|| segment_plan.as_ref().is_some_and(|plan| {
plan.archives
.iter()
.any(|archive| matches!(archive, PlannedArchiveSweep::Rewrite { .. }))
});
let manifest_upgrade =
writes_v2 && crate::store::read_manifest_store_version(&directory.join("manifest"))? < 2;
if manifest_upgrade {
ensure_numbered_name_available(directory, "manifest.cleaning")?;
}
Ok(manifest_upgrade)
}
fn require_directory_unchanged(
directory: &Path,
fingerprint_before: &DirectoryFingerprint,
) -> Result<DirectoryFingerprint> {
let fingerprint_after = directory_fingerprint(directory)?;
if *fingerprint_before != fingerprint_after {
return Err(Error::InvalidFormat {
details:
"the repository changed while cleanup was planning; retry against a quiescent store"
.to_owned(),
});
}
Ok(fingerprint_after)
}
fn decide_manifest_upgrade_and_reserve_journal_names(
directory: &Path,
options: &CompactionOptions,
state: &RepositoryState,
segment_plan: Option<&StandaloneSegmentCompactionPlan>,
) -> Result<bool> {
let manifest_upgrade = decide_manifest_upgrade(
directory,
options,
&state.pending_repairs,
&state.checkpoints,
segment_plan,
)?;
if options.contains(MaintenanceTask::Journal) && state.journal_analysis.plan.removed_lines != 0
{
ensure_numbered_name_available(directory, "journal.log.cleaning")?;
ensure_numbered_name_available(directory, "journal.log.bak")?;
}
Ok(manifest_upgrade)
}
fn journal_convergence_of(
raw_journal: &RawJournal,
current_head: RecordIdentifier,
) -> JournalConvergence {
JournalConvergence {
single_line_naming_head: raw_journal.lines().len() == 1
&& matches!(
raw_journal.lines()[0].classification(),
crate::writer::maintenance::journal::RawJournalLineClassification::Record(record)
if record.record_identifier == current_head
),
}
}
fn verify_head_and_census(
repository: &Repository,
current_head: RecordIdentifier,
options: &CompactionOptions,
observer: &mut dyn ProgressObserver,
) -> Result<(u64, PlanningContentCensus)> {
let mut content_census = PlanningContentCensus::default();
let head_nodes = verify_head_for_planning(
repository,
current_head,
&mut content_census,
options.purges_orphaned_version_histories(),
observer,
)?;
Ok((head_nodes, content_census))
}
struct GatheredPlanFacts {
repository: Repository,
current_head: RecordIdentifier,
head_nodes: u64,
content_census: PlanningContentCensus,
state: RepositoryState,
version_history_purge_selection: Option<PurgeSelection>,
segment_work: SegmentWork,
temporaries: Vec<PlannedFileRemoval>,
recovery_backups: Vec<PlannedFileRemoval>,
manifest_upgrade: bool,
}
fn gather_plan_facts(
directory: &Path,
options: &CompactionOptions,
now: SystemTime,
observer: &mut dyn ProgressObserver,
warnings: &mut Vec<String>,
) -> Result<GatheredPlanFacts> {
let repository = Repository::open_with_progress(directory, observer)?;
let current_head = repository.head_record_identifier();
if !options.contains(MaintenanceTask::RepairArchives)
&& (options.compaction_kind().is_some() || options.contains(MaintenanceTask::Segments))
&& let Some(details) =
super::indexless_refusal::indexless_active_archive_refusal(&repository)
{
return Err(Error::InvalidFormat { details });
}
let (head_nodes, content_census) =
verify_head_and_census(&repository, current_head, options, observer)?;
let state = survey_repository_state(
directory,
&repository,
options,
current_head,
now,
warnings,
observer,
)?;
let version_history_purge_selection = select_version_history_purge(
&repository,
options,
current_head,
&content_census,
now,
warnings,
observer,
)?;
let journal_convergence = journal_convergence_of(&state.raw_journal, current_head);
let segment_work = plan_segment_work(
&SegmentWorkInputs {
directory,
repository: &repository,
options,
current_head,
index_available: state.index_available,
checkpoints: &state.checkpoints,
journal_analysis: &state.journal_analysis,
journal_convergence: &journal_convergence,
purge_selected: version_history_purge_selection
.as_ref()
.is_some_and(|selection| selection.histories != 0),
},
warnings,
observer,
)?;
let (temporaries, recovery_backups) = plan_leftover_files(
directory,
&repository,
options,
state.index_available,
&state.raw_journal,
now,
warnings,
observer,
)?;
let manifest_upgrade = decide_manifest_upgrade_and_reserve_journal_names(
directory,
options,
&state,
segment_work.segment_plan.as_ref(),
)?;
Ok(GatheredPlanFacts {
repository,
current_head,
head_nodes,
content_census,
state,
version_history_purge_selection,
segment_work,
temporaries,
recovery_backups,
manifest_upgrade,
})
}
fn warn_about_pending_reindexes(repository: &crate::store::Repository, warnings: &mut Vec<String>) {
let Ok(Some(oak_index)) =
repository.node_at_path(&format!("/{}", crate::index::INDEX_DEFINITIONS_NAME))
else {
return;
};
let Ok(entries) = oak_index.child_node_entries() else {
return;
};
for (name, node) in entries {
let path = format!("/{}/{name}", crate::index::INDEX_DEFINITIONS_NAME);
let Ok(definition) = crate::index::IndexDefinition::read(&node, &path) else {
continue;
};
if !definition.reindex.flagged {
continue;
}
let stored_type = node.property("type").ok().flatten();
let index_type =
crate::index::strict_string(stored_type.as_ref()).unwrap_or("no readable type");
let rebuildable = matches!(
definition.index_type,
Some(
crate::index::IndexType::Property
| crate::index::IndexType::Reference
| crate::index::IndexType::Counter
)
);
warnings.push(if rebuildable {
format!(
"pending reindex: {path} ({index_type}; froe index reindex rebuilds it offline)"
)
} else {
format!(
"pending reindex: {path} ({index_type}, which froe index reindex does not \
rebuild; Oak rebuilds it at startup)"
)
});
}
}
pub(super) fn build_plan_collecting(
directory: &Path,
options: &CompactionOptions,
now: SystemTime,
observer: &mut dyn ProgressObserver,
warnings: &mut Vec<String>,
) -> Result<CompactionPlan> {
let fingerprint_before = directory_fingerprint(directory)?;
let GatheredPlanFacts {
repository,
current_head,
head_nodes,
content_census,
state,
version_history_purge_selection,
segment_work,
temporaries,
recovery_backups,
manifest_upgrade,
} = gather_plan_facts(directory, options, now, observer, warnings)?;
warn_about_pending_reindexes(&repository, warnings);
let (orphaned_version_histories, version_history_purge) = version_history_plan_parts(
&repository,
current_head,
&content_census,
head_nodes,
segment_work.closure_indexed_bytes,
crate::progress::count(state.checkpoints.names.len()),
version_history_purge_selection,
)?;
let PlanListing {
actions,
estimated_reclaimable_bytes,
estimated_archive_rewrite_source_bytes,
retained_reclaimable,
} = list_planned_actions(
&PlanFindings {
directory,
options,
state: &state,
segment_work: &segment_work,
temporaries: &temporaries,
recovery_backups: &recovery_backups,
manifest_upgrade,
head_nodes,
version_history_purge: version_history_purge
.as_ref()
.map(|purge| (purge.histories, purge.nodes, purge.retained_checkpoints)),
},
warnings,
)?;
let fingerprint_after = require_directory_unchanged(directory, &fingerprint_before)?;
let mut plan = CompactionPlan {
directory: directory.to_owned(),
tasks: options.tasks().collect(),
current_head,
actions,
warnings: warnings.clone(),
estimated_reclaimable_bytes,
estimated_archive_rewrite_source_bytes,
retained_reclaimable,
history_protection: segment_work.history_protection,
fingerprint: fingerprint_after,
journal: state.journal_analysis.plan,
checkpoints: state.checkpoints,
checkpoint_archive_number: state.checkpoint_archive_number,
stale_archives: segment_work.stale_archives,
temporaries,
recovery_backups,
segment_plan: segment_work.segment_plan,
residue_sweep: segment_work.residue_sweep,
reference_generation: segment_work.reference_generation,
protected_history_segments: segment_work.protected_history_segments,
manifest_upgrade,
predicted_copy_output_bytes: match (
segment_work.effective_compaction_kind,
segment_work.closure_indexed_bytes,
) {
(Some(_), Some((data_bytes, _))) => Some(data_bytes.saturating_sub(
if version_history_purge.is_some() {
orphaned_version_histories.node_record_bytes_estimate
} else {
0
},
)),
_ => None,
},
external_binary_footprint: content_census.external_binaries.footprint(),
effective_compaction_kind: segment_work.effective_compaction_kind,
already_fully_compacted: segment_work.already_fully_compacted,
orphaned_version_histories,
version_history_purge,
};
append_apply_identity_preview_warning(directory, &mut plan);
plan.warnings.sort();
plan.warnings.dedup();
Ok(plan)
}
pub(super) fn add_estimate(total: &mut u64, amount: u64) -> Result<()> {
*total = total
.checked_add(amount)
.ok_or_else(|| Error::InvalidFormat {
details: "cleanup byte estimate overflow".to_owned(),
})?;
Ok(())
}
pub(super) fn ensure_numbered_name_available(directory: &Path, stem: &str) -> Result<()> {
for counter in 0..1000u16 {
let path = directory.join(format!("{stem}.{counter:03}"));
match std::fs::symlink_metadata(path) {
Ok(_) => {}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(()),
Err(error) => return Err(error.into()),
}
}
Err(Error::InvalidFormat {
details: format!("all numbered names for {stem} (000-999) are occupied"),
})
}
pub(crate) fn available_filesystem_bytes(directory: &Path) -> Option<u64> {
#[cfg(unix)]
{
use std::os::unix::ffi::OsStrExt as _;
let path = std::ffi::CString::new(directory.as_os_str().as_bytes()).ok()?;
let mut statistics = std::mem::MaybeUninit::<libc::statvfs>::uninit();
if unsafe { libc::statvfs(path.as_ptr(), statistics.as_mut_ptr()) } != 0 {
return None;
}
let statistics = unsafe { statistics.assume_init() };
let fragment_size = if statistics.f_frsize == 0 {
statistics.f_bsize
} else {
statistics.f_frsize
};
let bytes = u128::from(statistics.f_bavail).checked_mul(u128::from(fragment_size))?;
u64::try_from(bytes).ok()
}
#[cfg(not(unix))]
{
let _ = directory;
None
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn planning_warnings_survive_a_refusal_and_stay_bounded() {
let untouched = attach_planning_warnings(
crate::error::Error::InvalidFormat {
details: "refused".to_owned(),
},
&[],
);
let crate::error::Error::InvalidFormat { details } = untouched else {
panic!("variant must be preserved");
};
assert_eq!(details, "refused", "no warnings means no tail");
let warnings: Vec<String> = (0..5).map(|index| format!("warning {index}")).collect();
let attached = attach_planning_warnings(
crate::error::Error::InvalidFormat {
details: "refused.".to_owned(),
},
&warnings,
);
let crate::error::Error::InvalidFormat { details } = attached else {
panic!("variant must be preserved");
};
assert!(
details.contains("warning 0") && details.contains("warning 2"),
"established warnings reach the operator: {details}"
);
assert!(
!details.contains("warning 3"),
"the tail is bounded so it cannot bury the refusal: {details}"
);
assert!(
details.contains("and 2 further warnings"),
"the omitted warnings are counted: {details}"
);
let input_output = attach_planning_warnings(
crate::error::Error::InputOutput(std::io::Error::other("disk")),
&warnings,
);
assert!(
matches!(input_output, crate::error::Error::InputOutput(_)),
"only format refusals carry the tail"
);
}
}