use super::archive_certificate::certify_active_archives;
use super::providers::CertifiedReclaimSources;
use super::reclaim::{
ArchiveRewritePolicy, ReclaimRule, analyze_standalone_segment_cleanup,
reject_duplicate_active_segments,
};
use crate::error::{Error, Result};
use crate::segment::identifier::SegmentIdentifier;
use std::collections::HashMap;
use std::path::Path;
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) enum PlannedArchiveSweep {
Remove {
file_name: String,
segment_count: usize,
file_bytes: u64,
},
Rewrite {
file_name: String,
replacement_name: String,
segment_count: usize,
eligible_entry_bytes: u64,
},
DeferredBySavings {
file_name: String,
segment_count: usize,
eligible_entry_bytes: u64,
},
DeferredAtLastGeneration {
file_name: String,
segment_count: usize,
eligible_entry_bytes: u64,
},
BlockedByOccupiedGeneration {
file_name: String,
occupied_name: String,
segment_count: usize,
eligible_entry_bytes: u64,
},
}
impl PlannedArchiveSweep {
pub(crate) fn file_name(&self) -> &str {
match self {
Self::Remove { file_name, .. }
| Self::Rewrite { file_name, .. }
| Self::DeferredBySavings { file_name, .. }
| Self::DeferredAtLastGeneration { file_name, .. }
| Self::BlockedByOccupiedGeneration { file_name, .. } => file_name,
}
}
pub(crate) fn changes_disk(&self) -> bool {
matches!(self, Self::Remove { .. } | Self::Rewrite { .. })
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct StandaloneSegmentCompactionPlan {
pub(crate) archives: Vec<PlannedArchiveSweep>,
pub(crate) marked_segments: usize,
pub(super) reclaimable: std::collections::HashSet<SegmentIdentifier>,
}
impl StandaloneSegmentCompactionPlan {
pub(crate) fn reclaimable_segments(&self) -> &std::collections::HashSet<SegmentIdentifier> {
&self.reclaimable
}
}
pub(super) fn sorted_sweep_plan(
planned: &HashMap<String, PlannedArchiveSweep>,
reclaimable: &std::collections::HashSet<SegmentIdentifier>,
) -> StandaloneSegmentCompactionPlan {
let mut archives: Vec<PlannedArchiveSweep> = planned.values().cloned().collect();
archives.sort_by(|left, right| left.file_name().cmp(right.file_name()));
StandaloneSegmentCompactionPlan {
archives,
marked_segments: reclaimable.len(),
reclaimable: reclaimable.clone(),
}
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub(crate) struct StandaloneSegmentCompactionOutcome {
pub(crate) rewritten_archives: usize,
pub(crate) removed_archives: usize,
pub(crate) removed_segments: usize,
pub(crate) archive_bytes_before: u64,
pub(crate) archive_bytes_after: u64,
pub(crate) deletion_failures: Vec<DeferredFileDeletion>,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct DeferredFileDeletion {
pub(crate) file_name: String,
pub(crate) error: String,
pub(crate) target_was_already_absent: bool,
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub(super) enum ArchiveSweepDisposition {
#[default]
Unchanged,
Removed,
Rewritten,
}
#[derive(Debug, Default)]
pub(super) struct ArchiveSweepOutcome {
pub(super) disposition: ArchiveSweepDisposition,
pub(super) deletion_failures: Vec<DeferredFileDeletion>,
pub(super) newly_unavailable: std::collections::HashSet<SegmentIdentifier>,
}
#[derive(Clone, Copy)]
pub(crate) struct GenerationReclaimRequest<'sources> {
pub(crate) rule: ReclaimRule,
pub(crate) rewrite_policy: ArchiveRewritePolicy,
pub(crate) certified_sources: Option<&'sources CertifiedReclaimSources>,
pub(crate) expected: Option<&'sources StandaloneSegmentCompactionPlan>,
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub(crate) struct SegmentSweepOutcome {
pub(crate) removed_archives: usize,
pub(crate) rewritten_archives: usize,
pub(crate) removed_segments: usize,
pub(crate) deletion_failures: Vec<DeferredFileDeletion>,
}
pub(crate) const RETAINED_GENERATIONS: i32 = 1;
pub(crate) fn plan_standalone_segment_cleanup(
directory: &Path,
repository: &crate::store::Repository,
rule: ReclaimRule,
current_head_segment: SegmentIdentifier,
protected: &std::collections::HashSet<SegmentIdentifier>,
rewrite_policy: ArchiveRewritePolicy,
observer: &mut dyn crate::progress::ProgressObserver,
) -> Result<StandaloneSegmentCompactionPlan> {
reject_duplicate_active_segments(repository.archives())?;
certify_active_archives(repository, repository.archives())?;
analyze_standalone_segment_cleanup(
directory,
repository.archives(),
rule,
current_head_segment,
protected,
rewrite_policy,
observer,
)
}
pub(crate) fn plan_reclaimed_totals(plan: &StandaloneSegmentCompactionPlan) -> (usize, u64) {
let mut segments = 0usize;
let mut bytes = 0u64;
for archive in &plan.archives {
let (archive_segments, archive_bytes) = match archive {
PlannedArchiveSweep::Remove {
segment_count,
file_bytes,
..
} => (*segment_count, *file_bytes),
PlannedArchiveSweep::Rewrite {
segment_count,
eligible_entry_bytes,
..
} => (*segment_count, *eligible_entry_bytes),
PlannedArchiveSweep::DeferredBySavings { .. }
| PlannedArchiveSweep::DeferredAtLastGeneration { .. }
| PlannedArchiveSweep::BlockedByOccupiedGeneration { .. } => continue,
};
segments = segments.saturating_add(archive_segments);
bytes = bytes.saturating_add(archive_bytes);
}
(segments, bytes)
}
pub(crate) fn measure_unvetoed_reclamation(
directory: &Path,
repository: &crate::store::Repository,
rule: ReclaimRule,
current_head_segment: SegmentIdentifier,
rewrite_policy: ArchiveRewritePolicy,
observer: &mut dyn crate::progress::ProgressObserver,
) -> Result<(usize, u64)> {
let unvetoed = analyze_standalone_segment_cleanup(
directory,
repository.archives(),
rule,
current_head_segment,
&std::collections::HashSet::new(),
rewrite_policy,
observer,
)?;
Ok(plan_reclaimed_totals(&unvetoed))
}
pub(crate) fn planned_unavailable_segments(
directory: &Path,
plan: &StandaloneSegmentCompactionPlan,
) -> Result<std::collections::HashSet<SegmentIdentifier>> {
let actionable: std::collections::HashSet<&str> = plan
.archives
.iter()
.filter(|archive| archive.changes_disk())
.map(PlannedArchiveSweep::file_name)
.collect();
let archives = crate::store::open_all_archives(directory)?;
let mut unavailable = std::collections::HashSet::new();
for archive in archives {
if !actionable.contains(archive.file_name()) {
continue;
}
let Some(index) = archive.index() else {
return Err(Error::InvalidFormat {
details: format!(
"cleanup planned to mutate recovered archive {}, which has no valid index",
archive.file_name()
),
});
};
unavailable.extend(
index
.entries()
.iter()
.map(|entry| entry.segment_identifier)
.filter(|identifier| plan.reclaimable.contains(identifier)),
);
}
Ok(unavailable)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::writer::compaction::CompactionKind;
use crate::writer::store_writer::reclaim::*;
use crate::writer::store_writer::repository::*;
use crate::writer::store_writer::test_support::*;
use std::collections::HashSet;
#[test]
fn dangling_future_boundary_is_shared_across_archive_files() {
let directory = TestDirectory::new("dangling-future-cross-archive");
let older = data_identifier(50);
let root = data_identifier(51);
let future = data_identifier(52);
let current = generation(9, 9, true);
write_test_archive(
&directory,
"data00000a.tar",
&[
TestArchiveEntry::new(older, 1, current),
TestArchiveEntry::new(root, 1, current),
],
);
write_test_archive(
&directory,
"data00001a.tar",
&[TestArchiveEntry::new(future, 1, current)],
);
write_manifest(&directory);
let plan = plan_cleanup_from_directory(&directory.path, current, root, &HashSet::new())
.expect("plan across archives");
assert_eq!(plan.marked_segments, 1);
assert_eq!(plan.reclaimable_segments(), &HashSet::from([future]));
assert!(matches!(
plan.archives.as_slice(),
[PlannedArchiveSweep::Remove {
file_name,
segment_count: 1,
..
}] if file_name == "data00001a.tar"
));
assert!(!plan.reclaimable_segments().contains(&root));
assert!(!plan.reclaimable_segments().contains(&older));
}
#[test]
fn a_post_compaction_sweep_refuses_an_unconfirmed_disposition_before_it_unlinks() {
let (directory, live) = reclaimable_base_fixture("sweep-authorization");
let before: Vec<String> = crate::store::list_archive_file_names(&directory.path)
.expect("list the archives before");
assert!(!before.is_empty());
let impossible = StandaloneSegmentCompactionPlan {
archives: vec![PlannedArchiveSweep::Remove {
file_name: "data09999a.tar".to_owned(),
segment_count: 1,
file_bytes: 1,
}],
marked_segments: 1,
reclaimable: std::collections::HashSet::new(),
};
let mut store = WritableRepository::open(&directory.path).expect("open for the sweep");
let error = store
.reclaim_old_generations_with(GenerationReclaimRequest {
rule: ReclaimRule {
reference: live,
kind: CompactionKind::Full,
retained_generations: RETAINED_GENERATIONS,
},
rewrite_policy: ArchiveRewritePolicy::EveryReclaimableArchive,
certified_sources: None,
expected: Some(&impossible),
})
.expect_err("an unconfirmed disposition must be refused");
store.close().expect("close after the refusal");
assert!(
error.to_string().contains("changed after confirmation"),
"the refusal names its reason: {error}"
);
assert_eq!(
crate::store::list_archive_file_names(&directory.path)
.expect("list the archives after"),
before,
"a refused sweep unlinks nothing"
);
}
#[test]
fn the_savings_gate_reaches_the_post_compaction_sweep() {
for (policy, expect_unchanged) in [
(ArchiveRewritePolicy::OakSavingsGate, true),
(ArchiveRewritePolicy::EveryReclaimableArchive, false),
] {
let (directory, live) = reclaimable_base_fixture(&format!("sweep-policy-{policy:?}"));
let before = crate::store::list_archive_file_names(&directory.path).expect("list");
let mut store = WritableRepository::open(&directory.path).expect("open for the sweep");
let outcome = store
.reclaim_old_generations_with(GenerationReclaimRequest {
rule: ReclaimRule {
reference: live,
kind: CompactionKind::Full,
retained_generations: RETAINED_GENERATIONS,
},
rewrite_policy: policy,
certified_sources: None,
expected: None,
})
.expect("the sweep runs");
store.close().expect("close after the sweep");
let after = crate::store::list_archive_file_names(&directory.path).expect("list");
if expect_unchanged {
assert_eq!(
after, before,
"Oak's heuristic leaves an archive it will not repay: {outcome:?}"
);
assert_eq!(outcome, SegmentSweepOutcome::default());
} else {
assert_ne!(
after, before,
"the default policy reclaims what the heuristic declined"
);
assert!(
outcome.removed_archives + outcome.rewritten_archives > 0,
"and reports what it did: {outcome:?}"
);
}
}
}
}