use super::archive_certificate::certify_active_archives_with_progress;
use super::file_identity::{archive_file_bytes, sync_directory_strict};
use super::reclaim::{
ArchiveRewritePolicy, ReclaimRule, analyze_standalone_segment_cleanup,
reject_duplicate_active_segments,
};
#[cfg(test)]
use super::sweep::probe_archive_sweep_phase_boundary;
use super::sweep::sweep_one_archive;
use super::sweep_plan::DeferredFileDeletion;
use super::sweep_plan::{
ArchiveSweepDisposition, PlannedArchiveSweep, StandaloneSegmentCompactionOutcome,
StandaloneSegmentCompactionPlan,
};
use crate::content::provider::SegmentProvider;
use crate::error::{Error, Result};
use crate::segment::identifier::SegmentIdentifier;
use crate::tar_archive::archive::TarArchiveReader;
use std::collections::HashMap;
use std::path::Path;
#[allow(
clippy::too_many_lines,
reason = "replanning, mutation, and exact partial-outcome accounting form one locked application sequence"
)]
pub(crate) fn apply_standalone_segment_cleanup(
directory: &Path,
rule: ReclaimRule,
current_head_segment: SegmentIdentifier,
protected: &std::collections::HashSet<SegmentIdentifier>,
rewrite_policy: ArchiveRewritePolicy,
expected: Option<&StandaloneSegmentCompactionPlan>,
observer: &mut dyn crate::progress::ProgressObserver,
) -> Result<(
StandaloneSegmentCompactionPlan,
StandaloneSegmentCompactionOutcome,
)> {
let repository = crate::store::Repository::open_with_progress(directory, observer)?;
reject_duplicate_active_segments(repository.archives())?;
certify_active_archives_with_progress(&repository, repository.archives(), observer)?;
apply_standalone_segment_cleanup_from_archives(
directory,
repository.archives(),
Some(&repository),
rule,
current_head_segment,
protected,
rewrite_policy,
expected,
observer,
#[cfg(test)]
None,
)
}
struct SweepPass {
observed_sweeps: HashMap<String, (ArchiveSweepDisposition, usize)>,
deletion_failures: Vec<DeferredFileDeletion>,
}
fn sweep_planned_archives(
directory: &Path,
archives: &[TarArchiveReader],
plan: &StandaloneSegmentCompactionPlan,
source_certificate_provider: Option<&dyn SegmentProvider>,
rewrite_policy: ArchiveRewritePolicy,
observer: &mut dyn crate::progress::ProgressObserver,
) -> Result<SweepPass> {
let provider_order: Vec<&TarArchiveReader> = archives.iter().collect();
let mut fallback_provider = None;
let mut deletion_failures = Vec::new();
let mut actually_unavailable = std::collections::HashSet::new();
let mut observed_sweeps = HashMap::new();
let planned_archives: HashMap<_, _> = plan
.archives
.iter()
.map(|planned| (planned.file_name(), planned))
.collect();
crate::progress::observe(
observer,
&crate::progress::Step::new("sweeping archives", crate::progress::WorkUnit::Archives)
.with_total(crate::progress::count(
plan.archives
.iter()
.filter(|planned| planned.changes_disk())
.count(),
)),
|observer| {
let mut swept = 0usize;
for rewrite_phase in [false, true] {
for archive in archives {
let Some(planned) = planned_archives.get(archive.file_name()) else {
continue;
};
let is_rewrite = matches!(planned, PlannedArchiveSweep::Rewrite { .. });
let is_remove = matches!(planned, PlannedArchiveSweep::Remove { .. });
if (!rewrite_phase && !is_remove) || (rewrite_phase && !is_rewrite) {
continue;
}
observer.step_advanced(crate::progress::count(swept));
let outcome = sweep_one_archive(
directory,
archive,
&plan.reclaimable,
&actually_unavailable,
&provider_order,
&mut fallback_provider,
source_certificate_provider,
rewrite_policy,
)?;
if outcome.disposition != ArchiveSweepDisposition::Unchanged {
observed_sweeps.insert(
archive.file_name().to_owned(),
(outcome.disposition, outcome.newly_unavailable.len()),
);
}
deletion_failures.extend(outcome.deletion_failures);
actually_unavailable.extend(outcome.newly_unavailable);
swept += 1;
observer.step_advanced(crate::progress::count(swept));
}
#[cfg(test)]
if !rewrite_phase {
probe_archive_sweep_phase_boundary("sweep.removals-complete-before-rewrites")?;
}
}
Ok::<(), Error>(())
},
)?;
drop(fallback_provider);
drop(provider_order);
Ok(SweepPass {
observed_sweeps,
deletion_failures,
})
}
fn account_for_swept_archives(
directory: &Path,
plan: &StandaloneSegmentCompactionPlan,
observed_sweeps: &HashMap<String, (ArchiveSweepDisposition, usize)>,
outcome: &mut StandaloneSegmentCompactionOutcome,
) -> Result<()> {
for archive in &plan.archives {
match archive {
PlannedArchiveSweep::Remove { file_name, .. }
if observed_sweeps
.get(file_name)
.is_some_and(|(disposition, _)| {
*disposition == ArchiveSweepDisposition::Removed
}) =>
{
if directory.join(file_name).try_exists()? {
if !outcome
.deletion_failures
.iter()
.any(|failure| failure.file_name == *file_name)
{
outcome.deletion_failures.push(DeferredFileDeletion {
file_name: file_name.clone(),
error: "file reappeared after the archive unlink succeeded".to_owned(),
target_was_already_absent: false,
});
}
} else {
outcome.removed_archives += 1;
outcome.removed_segments += observed_sweeps[file_name].1;
}
}
PlannedArchiveSweep::Rewrite {
file_name,
replacement_name,
..
} if observed_sweeps
.get(file_name)
.is_some_and(|(disposition, _)| {
*disposition == ArchiveSweepDisposition::Rewritten
}) =>
{
if !directory.join(replacement_name).try_exists()? {
return Err(Error::InvalidFormat {
details: format!(
"cleanup published rewrite {file_name}, but replacement \
{replacement_name} is absent"
),
});
}
outcome.rewritten_archives += 1;
outcome.removed_segments += observed_sweeps[file_name].1;
if directory.join(file_name).try_exists()?
&& !outcome
.deletion_failures
.iter()
.any(|failure| failure.file_name == *file_name)
{
outcome.deletion_failures.push(DeferredFileDeletion {
file_name: file_name.clone(),
error: "source archive remained after replacement publication".to_owned(),
target_was_already_absent: false,
});
}
}
PlannedArchiveSweep::Remove { file_name, .. }
| PlannedArchiveSweep::Rewrite { file_name, .. } => {
if observed_sweeps.contains_key(file_name) {
return Err(Error::InvalidFormat {
details: format!(
"archive sweep for {file_name} returned a disposition inconsistent with the authoritative plan"
),
});
}
}
PlannedArchiveSweep::DeferredBySavings { .. }
| PlannedArchiveSweep::DeferredAtLastGeneration { .. }
| PlannedArchiveSweep::BlockedByOccupiedGeneration { .. } => {}
}
}
Ok(())
}
#[cfg(test)]
pub(super) type StandaloneAfterPlanHook<'hook> =
dyn Fn(&StandaloneSegmentCompactionPlan) -> Result<()> + 'hook;
#[allow(
clippy::too_many_arguments,
reason = "the test-only uncertified path and production certified path share one ordered mutation engine"
)]
pub(super) fn apply_standalone_segment_cleanup_from_archives(
directory: &Path,
archives: &[TarArchiveReader],
source_certificate_provider: Option<&dyn SegmentProvider>,
rule: ReclaimRule,
current_head_segment: SegmentIdentifier,
protected: &std::collections::HashSet<SegmentIdentifier>,
rewrite_policy: ArchiveRewritePolicy,
expected: Option<&StandaloneSegmentCompactionPlan>,
observer: &mut dyn crate::progress::ProgressObserver,
#[cfg(test)] after_plan: Option<&StandaloneAfterPlanHook<'_>>,
) -> Result<(
StandaloneSegmentCompactionPlan,
StandaloneSegmentCompactionOutcome,
)> {
let plan = crate::progress::observe(
observer,
&crate::progress::Step::new(
"replanning segment reclamation",
crate::progress::WorkUnit::Archives,
)
.with_total(crate::progress::count(archives.len())),
|observer| {
analyze_standalone_segment_cleanup(
directory,
archives,
rule,
current_head_segment,
protected,
rewrite_policy,
observer,
)
},
)?;
if expected.is_some_and(|expected| expected != &plan) {
return Err(Error::InvalidFormat {
details: "the standalone segment-cleanup plan changed after confirmation; refusing \
to apply an unconfirmed archive mutation"
.to_owned(),
});
}
let archive_bytes_before = archive_file_bytes(directory)?;
#[cfg(test)]
if let Some(after_plan) = after_plan {
after_plan(&plan)?;
}
let pass = sweep_planned_archives(
directory,
archives,
&plan,
source_certificate_provider,
rewrite_policy,
observer,
)?;
sync_directory_strict(directory)?;
let mut outcome = StandaloneSegmentCompactionOutcome {
archive_bytes_before,
archive_bytes_after: archive_file_bytes(directory)?,
deletion_failures: pass.deletion_failures,
..StandaloneSegmentCompactionOutcome::default()
};
account_for_swept_archives(directory, &plan, &pass.observed_sweeps, &mut outcome)?;
outcome.deletion_failures.sort_by(|left, right| {
left.file_name
.cmp(&right.file_name)
.then_with(|| left.error.cmp(&right.error))
});
outcome.deletion_failures.dedup();
Ok((plan, outcome))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::tar_archive::archive::TarArchiveReader;
use crate::writer::store_writer::reclaim::*;
use crate::writer::store_writer::sweep_plan::*;
use crate::writer::store_writer::test_support::*;
use crate::writer::tar_writer::TarArchiveWriter;
use std::collections::HashSet;
#[test]
fn store_wide_reclaim_set_filters_cross_archive_graph_targets() {
let directory = TestDirectory::new("global-graph-filter");
let target = data_identifier(70);
let old_one = data_identifier(71);
let old_two = data_identifier(72);
let root = data_identifier(73);
let reference = generation(5, 5, false);
write_test_archive(
&directory,
"data00000a.tar",
&[TestArchiveEntry::new(target, 1, generation(0, 0, false))],
);
write_test_archive(
&directory,
"data00001a.tar",
&[
TestArchiveEntry::new(old_one, 1, generation(0, 0, false)),
TestArchiveEntry::new(old_two, 1, generation(0, 0, false)),
TestArchiveEntry::new(root, 1, reference).referencing(&[target]),
],
);
write_manifest(&directory);
let expected =
plan_cleanup_from_directory(&directory.path, reference, root, &HashSet::new())
.expect("plan");
assert_eq!(
expected.reclaimable_segments(),
&HashSet::from([target, old_one, old_two])
);
assert!(expected.archives.iter().any(|archive| matches!(
archive,
PlannedArchiveSweep::Remove { file_name, .. }
if file_name == "data00000a.tar"
)));
assert!(expected.archives.iter().any(|archive| matches!(
archive,
PlannedArchiveSweep::Rewrite {
file_name,
replacement_name,
..
} if file_name == "data00001a.tar" && replacement_name == "data00001b.tar"
)));
let (_, outcome) = apply_cleanup_from_directory(
&directory.path,
reference,
root,
&HashSet::new(),
Some(&expected),
)
.expect("apply");
assert_eq!(outcome.removed_archives, 1);
assert_eq!(outcome.rewritten_archives, 1);
assert_eq!(outcome.removed_segments, 3);
assert!(!directory.path.join("data00000a.tar").exists());
assert!(!directory.path.join("data00001a.tar").exists());
let swept = TarArchiveReader::open(&directory.path.join("data00001b.tar"))
.expect("open swept archive");
assert_eq!(swept.segment_count(), 1);
assert!(swept.contains_segment(root));
let graph = swept.segment_graph().expect("graph remains valid");
assert!(
graph
.adjacency
.iter()
.flat_map(|(_, targets)| targets)
.all(|identifier| *identifier != target),
"the target reclaimed from another tar must be filtered globally"
);
}
#[test]
fn deferred_cross_archive_target_remains_in_rewritten_graph() {
let directory = TestDirectory::new("deferred-global-graph-target");
let target = data_identifier(74);
let retained_one = data_identifier(75);
let retained_two = data_identifier(76);
let retained_three = data_identifier(77);
let old_one = data_identifier(78);
let old_two = data_identifier(79);
let root = data_identifier(80);
let reference = generation(5, 5, false);
write_test_archive(
&directory,
"data00000z.tar",
&[
TestArchiveEntry::new(target, 1, generation(0, 0, false)),
TestArchiveEntry::new(retained_one, 1, reference),
TestArchiveEntry::new(retained_two, 1, reference),
TestArchiveEntry::new(retained_three, 1, reference),
],
);
write_test_archive(
&directory,
"data00001a.tar",
&[
TestArchiveEntry::new(old_one, 1, generation(0, 0, false)),
TestArchiveEntry::new(old_two, 1, generation(0, 0, false)),
TestArchiveEntry::new(root, 1, reference).referencing(&[target]),
],
);
write_manifest(&directory);
let expected =
plan_cleanup_from_directory(&directory.path, reference, root, &HashSet::new())
.expect("plan");
assert!(expected.archives.iter().any(|archive| matches!(
archive,
PlannedArchiveSweep::DeferredAtLastGeneration { file_name, .. }
if file_name == "data00000z.tar"
)));
assert!(expected.archives.iter().any(|archive| matches!(
archive,
PlannedArchiveSweep::Rewrite { file_name, .. }
if file_name == "data00001a.tar"
)));
apply_cleanup_from_directory(
&directory.path,
reference,
root,
&HashSet::new(),
Some(&expected),
)
.expect("apply");
assert!(
directory.path.join("data00000z.tar").exists(),
"the deferred target remains physically available"
);
let swept = TarArchiveReader::open(&directory.path.join("data00001b.tar"))
.expect("open rewritten source");
let graph = swept.segment_graph().expect("graph remains valid");
assert_eq!(
graph.as_map()[&root],
[target],
"a deferred target must not be filtered by a wider global reclaim set"
);
}
#[test]
fn immediate_replan_noop_is_not_reported_as_a_completed_rewrite() {
let directory = TestDirectory::new("rewrite-replan-noop-outcome");
let old_one = data_identifier(81);
let old_two = data_identifier(82);
let root = data_identifier(83);
let reference = generation(5, 5, false);
write_test_archive(
&directory,
"data00000a.tar",
&[
TestArchiveEntry::new(old_one, 1, generation(0, 0, false)),
TestArchiveEntry::new(old_two, 1, generation(0, 0, false)),
TestArchiveEntry::new(root, 1, reference),
],
);
write_manifest(&directory);
let archives = crate::store::open_all_archives(&directory.path).expect("open archives");
let occupied = b"occupied after authoritative planning";
let after_plan = |plan: &StandaloneSegmentCompactionPlan| {
let replacement = plan
.archives
.iter()
.find_map(|archive| match archive {
PlannedArchiveSweep::Rewrite {
file_name,
replacement_name,
..
} if file_name == "data00000a.tar" => Some(replacement_name),
_ => None,
})
.expect("the authoritative outer plan must request a rewrite");
std::fs::write(directory.path.join(replacement), occupied)?;
Ok(())
};
let (plan, outcome) = apply_standalone_segment_cleanup_from_archives(
&directory.path,
&archives,
None,
standalone_rule(reference),
root,
&HashSet::new(),
ArchiveRewritePolicy::default(),
None,
&mut crate::progress::DiscardedProgress,
Some(&after_plan),
)
.expect("an occupied immediate replan is a safe no-op");
assert!(matches!(
plan.archives.as_slice(),
[PlannedArchiveSweep::Rewrite { .. }]
));
assert_eq!(outcome.rewritten_archives, 0);
assert_eq!(outcome.removed_archives, 0);
assert_eq!(outcome.removed_segments, 0);
assert!(outcome.deletion_failures.is_empty());
assert!(directory.path.join("data00000a.tar").exists());
assert_eq!(
std::fs::read(directory.path.join("data00000b.tar"))
.expect("read occupied replacement"),
occupied,
"an unrelated occupied generation must not be credited as cleanup output"
);
}
#[test]
fn sweep_preserves_survivor_brf_generation_triples_and_omits_removed_sources() {
let directory = TestDirectory::new("brf-filter-and-triples");
let root = data_identifier(80);
let removed_one = data_identifier(81);
let removed_two = data_identifier(82);
let reference = generation(6, 6, false);
let survivor_catalog_generation = generation(17, 11, true);
let removed_catalog_generation = generation(18, 12, false);
let mut writer = TarArchiveWriter::new(&directory.path, "data00000a.tar");
writer.add_binary_references(survivor_catalog_generation, root, ["live-blob".to_owned()]);
writer.add_binary_references(
removed_catalog_generation,
removed_one,
["dead-blob-one".to_owned()],
);
writer.add_binary_references(
removed_catalog_generation,
removed_two,
["dead-blob-two".to_owned()],
);
for entry in [
TestArchiveEntry::new(root, 1, reference),
TestArchiveEntry::new(removed_one, 1, generation(0, 0, false)),
TestArchiveEntry::new(removed_two, 1, generation(0, 0, false)),
] {
writer
.write_segment(entry.identifier, &entry.content, entry.generation, &[], &[])
.expect("write segment");
}
writer.close().expect("close archive");
write_manifest(&directory);
let plan = plan_cleanup_from_directory(&directory.path, reference, root, &HashSet::new())
.expect("plan");
apply_cleanup_from_directory(
&directory.path,
reference,
root,
&HashSet::new(),
Some(&plan),
)
.expect("apply");
let swept = TarArchiveReader::open(&directory.path.join("data00000b.tar"))
.expect("open swept archive");
let catalog = swept.binary_references().expect("catalog survives");
assert_eq!(catalog.generations.len(), 1);
let generation = &catalog.generations[0];
assert_eq!(
generation.generation,
survivor_catalog_generation.generation
);
assert_eq!(
generation.full_generation,
survivor_catalog_generation.full_generation
);
assert_eq!(
generation.is_compacted,
survivor_catalog_generation.is_compacted
);
assert_eq!(
generation.segments,
vec![(root, vec!["live-blob".to_owned()])]
);
assert!(catalog.generations.iter().all(|generation| {
generation
.segments
.iter()
.all(|(source, _)| *source != removed_one && *source != removed_two)
}));
}
#[test]
fn missing_brf_reconstruction_failure_leaves_original_and_no_replacement() {
let directory = TestDirectory::new("missing-brf-fail-closed");
let root = data_identifier(120);
let removed_one = data_identifier(121);
let removed_two = data_identifier(122);
let reference = generation(5, 5, false);
write_test_archive(
&directory,
"data00000a.tar",
&[
TestArchiveEntry::new(root, 64, reference),
TestArchiveEntry::new(removed_one, 64, generation(0, 0, false)),
TestArchiveEntry::new(removed_two, 64, generation(0, 0, false)),
],
);
let source_path = directory.path.join("data00000a.tar");
let mut bytes = std::fs::read(&source_path).expect("read archive");
let brf_magic = bytes
.windows(4)
.position(|window| window == [0x0A, 0x31, 0x42, 0x0A])
.expect("brf magic");
bytes[brf_magic] ^= 0x01;
std::fs::write(&source_path, &bytes).expect("corrupt only brf footer");
write_manifest(&directory);
let reader = TarArchiveReader::open(&source_path).expect("index remains valid");
assert!(reader.index().is_some());
assert!(reader.segment_graph().is_some());
assert!(reader.binary_references().is_none());
drop(reader);
let plan = plan_cleanup_from_directory(&directory.path, reference, root, &HashSet::new())
.expect("mark does not need brf");
let error = apply_cleanup_from_directory(
&directory.path,
reference,
root,
&HashSet::new(),
Some(&plan),
)
.expect_err("catalog reconstruction must fail closed on malformed data");
assert!(error.to_string().contains("magic bytes"));
assert_eq!(std::fs::read(&source_path).expect("source remains"), bytes);
assert!(!directory.path.join("data00000b.tar").exists());
}
}