use super::options::MaintenanceTask;
use super::planning::{
CheckpointPlan, DirectoryFingerprint, JournalPlan, PlannedFileRemoval, StaleArchive,
};
use crate::segment::identifier::SegmentIdentifier;
use crate::segment::record::RecordIdentifier;
use crate::writer::compaction::CompactionKind;
use crate::writer::segment_builder::GarbageCollectionGeneration;
use crate::writer::store_writer::StandaloneSegmentCompactionPlan;
use std::collections::HashSet;
use std::path::{Path, PathBuf};
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub enum StaleArchiveReason {
Superseded,
EmptyIncomplete,
}
impl std::fmt::Display for StaleArchiveReason {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str(match self {
Self::Superseded => "superseded by the active archive generation",
Self::EmptyIncomplete => "empty incomplete archive",
})
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub enum JournalRemovalReason {
ParserSkippedNoSpace,
InvalidRecordIdentifier,
MissingSegment,
UnreadableRevision,
BeyondRetention,
}
impl std::fmt::Display for JournalRemovalReason {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str(match self {
Self::ParserSkippedNoSpace => "parser-skipped (no ASCII space)",
Self::InvalidRecordIdentifier => "invalid record identifier",
Self::MissingSegment => "missing segment",
Self::UnreadableRevision => "unreadable historical revision",
Self::BeyondRetention => "beyond the journal retention bound",
})
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub struct JournalLineRemoval {
pub(super) line_number: usize,
pub(super) record_identifier: Option<RecordIdentifier>,
pub(super) reason: JournalRemovalReason,
pub(super) preview: Vec<u8>,
pub(super) preview_truncated: bool,
}
impl JournalLineRemoval {
#[must_use]
pub fn line_number(&self) -> usize {
self.line_number
}
#[must_use]
pub fn record_identifier(&self) -> Option<RecordIdentifier> {
self.record_identifier
}
#[must_use]
pub fn reason(&self) -> JournalRemovalReason {
self.reason
}
#[must_use]
pub fn preview_bytes(&self) -> &[u8] {
&self.preview
}
#[must_use]
pub fn preview_truncated(&self) -> bool {
self.preview_truncated
}
}
pub(super) const ALREADY_ABSENT_DELETION_DETAIL: &str =
"file was already absent when deletion was attempted";
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(super) enum FileDeletionFailureKind {
Retained,
AlreadyAbsent,
}
#[derive(Clone, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub struct FileDeletionFailure {
pub(super) file_name: String,
pub(super) error: String,
pub(super) kind: FileDeletionFailureKind,
}
impl FileDeletionFailure {
pub(super) fn retained(file_name: String, error: impl Into<String>) -> Self {
Self {
file_name,
error: error.into(),
kind: FileDeletionFailureKind::Retained,
}
}
pub(super) fn already_absent(file_name: String, error: impl Into<String>) -> Self {
Self {
file_name,
error: error.into(),
kind: FileDeletionFailureKind::AlreadyAbsent,
}
}
#[must_use]
pub fn file_name(&self) -> &str {
&self.file_name
}
#[must_use]
pub fn error(&self) -> &str {
&self.error
}
#[must_use]
pub fn target_was_already_absent(&self) -> bool {
self.kind == FileDeletionFailureKind::AlreadyAbsent
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub enum CompactionAction {
RepairArchiveIndex {
file_name: String,
retired_file_names: Vec<String>,
reason: String,
bytes: u64,
},
PruneJournal {
lines: usize,
parser_ignored: usize,
missing_segments: usize,
unreadable_revisions: usize,
beyond_retention: usize,
},
UpgradeManifest,
RemoveCheckpoints {
names: Vec<String>,
expired: usize,
unreferenced: usize,
},
RemoveReclaimableArchive {
file_name: String,
segments: usize,
bytes: u64,
},
RewriteArchive {
file_name: String,
replacement_name: String,
segments: usize,
eligible_bytes: u64,
},
RemoveStaleArchive {
file_name: String,
reason: StaleArchiveReason,
bytes: u64,
},
RemoveTemporary {
file_name: String,
bytes: u64,
},
RetireJournalHistory {
revisions: usize,
},
RetireInterruptedCompactionResidue {
segments: usize,
},
CopyHeadIntoFreshGeneration {
head_nodes: u64,
target_generation: GarbageCollectionGeneration,
kind: CompactionKind,
},
RemoveRecoveryBackup {
file_name: String,
bytes: u64,
},
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub(super) struct RetainedReclaimable {
pub(super) below_savings_gate: usize,
pub(super) at_last_generation: usize,
pub(super) blocked_by_occupied_generation: usize,
pub(super) bytes: u64,
}
impl RetainedReclaimable {
fn segments(self) -> usize {
self.below_savings_gate
.saturating_add(self.at_last_generation)
.saturating_add(self.blocked_by_occupied_generation)
}
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub(super) struct HistoryProtection {
pub(super) history_only_segments: usize,
pub(super) would_be_reclaimable_segments: usize,
pub(super) would_be_reclaimable_bytes: u64,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct CompactionPlan {
pub(super) directory: PathBuf,
pub(super) tasks: Vec<MaintenanceTask>,
pub(super) current_head: RecordIdentifier,
pub(super) actions: Vec<CompactionAction>,
pub(super) warnings: Vec<String>,
pub(super) estimated_reclaimable_bytes: u64,
pub(super) estimated_archive_rewrite_source_bytes: u64,
pub(super) retained_reclaimable: RetainedReclaimable,
pub(super) history_protection: HistoryProtection,
pub(super) fingerprint: DirectoryFingerprint,
pub(super) journal: JournalPlan,
pub(super) checkpoints: CheckpointPlan,
pub(super) checkpoint_archive_number: Option<u32>,
pub(super) stale_archives: Vec<StaleArchive>,
pub(super) temporaries: Vec<PlannedFileRemoval>,
pub(super) recovery_backups: Vec<PlannedFileRemoval>,
pub(super) segment_plan: Option<StandaloneSegmentCompactionPlan>,
pub(super) residue_sweep: Option<StandaloneSegmentCompactionPlan>,
pub(super) reference_generation: GarbageCollectionGeneration,
pub(super) protected_history_segments: HashSet<SegmentIdentifier>,
pub(super) manifest_upgrade: bool,
}
impl CompactionPlan {
#[must_use]
pub fn directory(&self) -> &Path {
&self.directory
}
#[must_use]
#[cfg(test)]
pub(crate) fn tasks(&self) -> &[MaintenanceTask] {
&self.tasks
}
#[must_use]
pub fn current_head(&self) -> RecordIdentifier {
self.current_head
}
#[must_use]
pub fn actions(&self) -> &[CompactionAction] {
&self.actions
}
#[must_use]
pub fn journal_line_removals(&self) -> &[JournalLineRemoval] {
if self.tasks.contains(&MaintenanceTask::Journal) {
&self.journal.removals
} else {
&[]
}
}
#[must_use]
pub fn warnings(&self) -> &[String] {
&self.warnings
}
#[must_use]
pub fn estimated_reclaimable_bytes(&self) -> u64 {
self.estimated_reclaimable_bytes
}
#[must_use]
pub fn estimated_archive_rewrite_source_bytes(&self) -> u64 {
self.estimated_archive_rewrite_source_bytes
}
#[must_use]
pub fn retained_reclaimable_segments(&self) -> usize {
self.retained_reclaimable.segments()
}
#[must_use]
pub fn retained_reclaimable_bytes(&self) -> u64 {
self.retained_reclaimable.bytes
}
#[must_use]
pub fn history_protected_segments(&self) -> usize {
self.history_protection.history_only_segments
}
#[must_use]
pub fn history_protected_reclaimable(&self) -> (usize, u64) {
(
self.history_protection.would_be_reclaimable_segments,
self.history_protection.would_be_reclaimable_bytes,
)
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.actions.is_empty()
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub struct CompactedGeneration {
pub nodes: u64,
pub generation: GarbageCollectionGeneration,
}
#[derive(Clone, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub struct CompactionOutcome {
pub head_before: RecordIdentifier,
pub head_after: RecordIdentifier,
pub removed_checkpoints: u64,
pub removed_journal_lines: usize,
pub rewritten_archives: usize,
pub removed_reclaimable_archives: usize,
pub removed_stale_archives: usize,
pub removed_temporaries: usize,
pub removed_recovery_backups: usize,
pub repaired_archives: usize,
pub files_not_deleted: Vec<String>,
pub archive_bytes_before: u64,
pub archive_bytes_after: u64,
pub retained_recovery_backup_bytes: u64,
pub compacted: Option<CompactedGeneration>,
pub(super) removed_segments: usize,
pub(super) journal_backup_path: Option<PathBuf>,
pub(super) deletion_failures: Vec<FileDeletionFailure>,
}
impl CompactionOutcome {
#[must_use]
pub fn removed_segments(&self) -> usize {
self.removed_segments
}
#[must_use]
pub fn journal_backup_path(&self) -> Option<&Path> {
self.journal_backup_path.as_deref()
}
#[must_use]
pub fn deletion_failures(&self) -> &[FileDeletionFailure] {
&self.deletion_failures
}
#[must_use]
pub fn is_complete(&self) -> bool {
self.deletion_failures.is_empty()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::store::Repository;
use crate::writer::maintenance::options::*;
use crate::writer::maintenance::prepared::*;
use crate::writer::maintenance::test_support::*;
use std::io::Write as _;
use std::num::NonZeroUsize;
#[test]
fn dangling_journal_line_is_pruned_with_backup_and_archives_untouched() {
let directory = TestDirectory::repository("dangling-journal");
let missing = SegmentIdentifier::new(7, 0xA000_0000_0000_0007);
let journal_path = directory.path.join("journal.log");
let retained_journal = std::fs::read(&journal_path).expect("read retained journal");
let mut journal = std::fs::OpenOptions::new()
.append(true)
.open(&journal_path)
.expect("open journal");
writeln!(journal, "{missing}:0 root 123").expect("append dangling line");
drop(journal);
let archive_before =
std::fs::read(directory.path.join("data00000a.tar")).expect("read archive");
std::fs::write(
directory.path.join("manifest"),
b"custom.property=untouched\nstore.version=1\n",
)
.expect("version-one manifest");
let manifest_before = std::fs::read(directory.path.join("manifest")).expect("manifest");
let options = CompactionOptions::default().with_tasks([MaintenanceTask::Journal]);
let plan = plan_compaction(&directory.path, &options).expect("plan");
assert_eq!(plan.tasks(), &[MaintenanceTask::Journal]);
assert_eq!(plan.journal_line_removals().len(), 1);
let removal = &plan.journal_line_removals()[0];
assert_eq!(
removal.record_identifier().map(|record| record.segment),
Some(missing)
);
assert_eq!(removal.reason(), JournalRemovalReason::MissingSegment);
assert!(
removal
.preview_bytes()
.starts_with(missing.to_string().as_bytes())
);
assert!(!removal.preview_truncated());
assert!(plan.actions().iter().any(|action| matches!(
action,
CompactionAction::PruneJournal {
missing_segments: 1,
..
}
)));
let outcome = compact(&directory.path, options).expect("apply");
assert_eq!(outcome.removed_journal_lines, 1);
let expected_backup =
canonical_fixture_directory(&directory.path).join("journal.log.bak.000");
assert_eq!(
outcome.journal_backup_path(),
Some(expected_backup.as_path())
);
assert!(outcome.is_complete());
assert!(directory.path.join("journal.log.bak.000").is_file());
assert!(
!std::fs::read_to_string(&journal_path)
.expect("journal")
.contains(&missing.to_string())
);
assert_eq!(
std::fs::read(&journal_path).expect("rewritten journal"),
retained_journal,
"the retained physical journal line must be byte-exact"
);
assert_eq!(
std::fs::read(directory.path.join("data00000a.tar")).expect("archive"),
archive_before
);
assert_eq!(
std::fs::read(directory.path.join("manifest")).expect("manifest"),
manifest_before
);
Repository::open(&directory.path).expect("healthy repository");
}
#[test]
fn deletion_absence_state_does_not_depend_on_diagnostic_text() {
let retained = super::FileDeletionFailure::retained(
"data00000a.tar".to_owned(),
ALREADY_ABSENT_DELETION_DETAIL,
);
let absent = super::FileDeletionFailure::already_absent(
"data00001a.tar".to_owned(),
"a deliberately different ENOENT diagnostic",
);
assert!(!retained.target_was_already_absent());
assert!(absent.target_was_already_absent());
}
#[test]
fn a_journal_retention_bound_retires_the_history_the_veto_protects() {
let (directory, old_head, new_head) = history_veto_fixture("history-veto-retention");
let protected = plan_compaction(
&directory.path,
&CompactionOptions::default().with_tasks([MaintenanceTask::Segments]),
)
.expect("unbounded plan");
assert!(protected.history_protected_reclaimable().0 != 0);
assert!(
!protected.actions().iter().any(|action| matches!(
action,
CompactionAction::RemoveReclaimableArchive { file_name, .. }
if file_name == "data00000a.tar"
)),
"without a bound the veto must keep the bootstrap archive"
);
let bounded = CompactionOptions::default()
.with_tasks([MaintenanceTask::Segments, MaintenanceTask::Journal])
.with_journal_revision_retention(NonZeroUsize::new(1).expect("one revision"));
let plan = plan_compaction(&directory.path, &bounded).expect("bounded plan");
assert!(
plan.journal_line_removals().iter().any(|removal| {
removal.reason() == JournalRemovalReason::BeyondRetention
&& removal.record_identifier() == Some(old_head)
}),
"the superseded revision must be removed as beyond retention"
);
assert!(
plan.actions().iter().any(|action| matches!(
action,
CompactionAction::RemoveReclaimableArchive { file_name, .. }
if file_name == "data00000a.tar"
)),
"the bound must release the bootstrap archive to Oak's predicate"
);
assert!(plan.estimated_reclaimable_bytes() != 0);
let outcome = compact(&directory.path, bounded).expect("bounded cleanup");
assert_eq!(outcome.head_after, new_head);
assert!(!directory.path.join("data00000a.tar").exists());
let repository = Repository::open(&directory.path).expect("healthy final repository");
assert_eq!(repository.head_record_identifier(), new_head);
let journal =
std::fs::read_to_string(directory.path.join("journal.log")).expect("read journal");
assert_eq!(
journal.lines().count(),
1,
"a bound of one leaves one journal line"
);
assert!(
crate::tooling::verify_node_tree(&repository, old_head).is_err(),
"the retired revision must no longer resolve"
);
}
#[test]
fn a_bound_counts_only_revisions_that_actually_resolve() {
let (directory, old_head, new_head) = history_veto_fixture("retention-counts-readable");
let unreadable = RecordIdentifier::new(new_head.segment, new_head.record_number + 1);
let journal_path = directory.path.join("journal.log");
let journal = std::fs::read_to_string(&journal_path).expect("read journal");
let mut lines: Vec<&str> = journal.lines().collect();
let head_line = lines.pop().expect("a head line");
let unreadable_line = format!("{unreadable} root 0");
lines.push(&unreadable_line);
lines.push(head_line);
std::fs::write(&journal_path, format!("{}\n", lines.join("\n")))
.expect("insert unreadable line");
let options = CompactionOptions::default()
.with_tasks([MaintenanceTask::Segments, MaintenanceTask::Journal])
.with_journal_revision_retention(NonZeroUsize::new(2).expect("two revisions"));
let plan = plan_compaction(&directory.path, &options).expect("bounded plan");
assert!(
!plan.journal_line_removals().iter().any(|removal| {
removal.reason() == JournalRemovalReason::BeyondRetention
&& removal.record_identifier() == Some(old_head)
}),
"a readable revision was retired to make room for an unreadable one: {:?}",
plan.journal_line_removals()
);
}
#[test]
fn a_bound_larger_than_the_journal_removes_nothing() {
let (directory, _old_head, _new_head) = history_veto_fixture("history-veto-wide-bound");
let options = CompactionOptions::default()
.with_tasks([MaintenanceTask::Segments, MaintenanceTask::Journal])
.with_journal_revision_retention(NonZeroUsize::new(64).expect("wide bound"));
let plan = plan_compaction(&directory.path, &options).expect("wide plan");
assert!(
!plan
.journal_line_removals()
.iter()
.any(|removal| { removal.reason() == JournalRemovalReason::BeyondRetention }),
"a bound wider than the journal must retire nothing"
);
assert!(plan.history_protected_reclaimable().0 != 0);
}
}