use super::apply::apply_prepared;
#[cfg(unix)]
use super::apply_identity::{
current_apply_credentials, metadata_source_apply_identity_issue, possible_created_group_ids,
};
use super::apply_identity::{
validate_apply_environment, validate_apply_identity, validate_plan_apply_identity,
};
use super::options::{CompactionOptions, MaintenanceTask};
use super::plan::{CompactionOutcome, CompactionPlan};
use super::planning::{
ManifestUpgradeOnFirstInstall, attach_completed_repairs, build_plan,
canonical_repository_directory, validate_options, validate_repository_shape,
};
#[cfg(unix)]
use crate::error::Error;
use crate::error::Result;
use crate::progress::{DiscardedProgress, ProgressObserver};
use crate::writer::repository_lock::RepositoryLock;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::SystemTime;
pub struct PreparedCompaction {
pub(super) directory: PathBuf,
pub(super) options: CompactionOptions,
pub(super) plan: CompactionPlan,
pub(super) repaired: Vec<crate::writer::store_writer::RepairedArchive>,
pub(super) repository_lock: Arc<RepositoryLock>,
}
impl PreparedCompaction {
pub fn prepare(directory: &Path, options: CompactionOptions) -> Result<Self> {
Self::prepare_with_progress(directory, options, &mut DiscardedProgress)
}
pub fn prepare_with_progress(
directory: &Path,
options: CompactionOptions,
observer: &mut dyn ProgressObserver,
) -> Result<Self> {
validate_options(&options)?;
let directory = canonical_repository_directory(directory)?;
validate_repository_shape(&directory)?;
crate::store::check_manifest(&directory, crate::store::ArchivePresence::Present)?;
validate_apply_environment(&directory)?;
validate_apply_identity(&directory)?;
let repository_lock = Arc::new(RepositoryLock::acquire(&directory)?);
validate_repository_shape(&directory)?;
crate::store::check_manifest(&directory, crate::store::ArchivePresence::Present)?;
validate_apply_environment(&directory)?;
validate_apply_identity(&directory)?;
repository_lock.validate_path_identity(&directory)?;
let repaired = Self::repair_before_planning(&directory, &options, observer)?;
let now = SystemTime::now();
let plan = build_plan(&directory, &options, now, observer).map_err(|error| {
attach_completed_repairs(error, &repaired)
})?;
validate_plan_apply_identity(&directory, &plan)
.map_err(|error| attach_completed_repairs(error, &repaired))?;
Ok(Self {
directory,
options,
plan,
repaired,
repository_lock,
})
}
fn repair_before_planning(
directory: &Path,
options: &CompactionOptions,
observer: &mut dyn ProgressObserver,
) -> Result<Vec<crate::writer::store_writer::RepairedArchive>> {
if !options.contains(MaintenanceTask::RepairArchives) {
return Ok(Vec::new());
}
let survey = crate::writer::store_writer::survey_indexless_archive_numbers(directory)?;
if !survey.unrepairable.is_empty() {
return Err(crate::writer::store_writer::unrepairable_archives_refusal(
&survey.unrepairable,
));
}
if survey.repairable == 0 {
return Ok(Vec::new());
}
crate::writer::store_writer::reject_duplicate_archive_generations(directory)?;
Self::validate_repair_target_identity(directory)?;
let manifest_upgrade = &mut ManifestUpgradeOnFirstInstall::new(directory);
crate::writer::store_writer::repair_indexless_archive_numbers(
directory,
observer,
manifest_upgrade,
)
}
fn validate_repair_target_identity(directory: &Path) -> Result<()> {
#[cfg(unix)]
{
use std::os::unix::fs::{MetadataExt as _, PermissionsExt as _};
let credentials = current_apply_credentials()?;
let directory_metadata = std::fs::symlink_metadata(directory)?;
let possible_created_gids = possible_created_group_ids(
directory_metadata.gid(),
directory_metadata.permissions().mode(),
&credentials,
);
for name in crate::writer::store_writer::repair_target_names(directory)? {
let path = directory.join(&name);
let metadata = std::fs::symlink_metadata(&path)?;
if let Some(issue) = metadata_source_apply_identity_issue(
&path,
metadata.uid(),
metadata.gid(),
metadata.permissions().mode(),
&possible_created_gids,
&credentials,
) {
return Err(Error::InvalidFormat { details: issue });
}
}
}
#[cfg(not(unix))]
let _ = directory;
Ok(())
}
#[must_use]
pub fn plan(&self) -> &CompactionPlan {
&self.plan
}
#[must_use]
pub fn repaired_archives(&self) -> usize {
self.repaired.len()
}
pub fn apply(self) -> Result<CompactionOutcome> {
apply_prepared(self, &mut DiscardedProgress)
}
pub fn apply_with_progress(
self,
observer: &mut dyn ProgressObserver,
) -> Result<CompactionOutcome> {
apply_prepared(self, observer)
}
}
pub fn plan_compaction(directory: &Path, options: &CompactionOptions) -> Result<CompactionPlan> {
plan_compaction_with_progress(directory, options, &mut DiscardedProgress)
}
pub fn plan_compaction_with_progress(
directory: &Path,
options: &CompactionOptions,
observer: &mut dyn ProgressObserver,
) -> Result<CompactionPlan> {
validate_options(options)?;
let directory = canonical_repository_directory(directory)?;
validate_repository_shape(&directory)?;
build_plan(&directory, options, SystemTime::now(), observer)
}
pub fn compact(directory: &Path, options: CompactionOptions) -> Result<CompactionOutcome> {
compact_with_progress(directory, options, &mut DiscardedProgress)
}
pub fn compact_with_progress(
directory: &Path,
options: CompactionOptions,
observer: &mut dyn ProgressObserver,
) -> Result<CompactionOutcome> {
PreparedCompaction::prepare_with_progress(directory, options, observer)?
.apply_with_progress(observer)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::store::Repository;
use crate::tar_archive::file_name::ArchiveFileName;
use crate::writer::commit::create_checkpoint;
use crate::writer::maintenance::options::*;
use crate::writer::maintenance::test_support::*;
use crate::writer::record_writer::ChildNodesToWrite;
use crate::writer::record_writer::PropertyToWrite;
use crate::writer::record_writer::PropertyValuesToWrite;
use crate::writer::segment_builder::GarbageCollectionGeneration;
use crate::writer::store_writer::WritableRepository;
use std::num::NonZeroUsize;
#[test]
fn an_unrepairable_archive_refuses_before_anything_is_rewritten() {
let directory = TestDirectory::new("repair-unrepairable");
{
let store = WritableRepository::open(&directory.path).expect("bootstrap");
let generation = store.writing_generation().expect("generation");
let mut writer = store.record_writer(generation);
writer
.write_string("forces a second archive")
.expect("string");
writer.finish().expect("finish");
store.close().expect("close");
}
break_index_magic(&directory.path.join("data00000a.tar"));
std::fs::write(directory.path.join("data00500a.tar"), vec![0x5au8; 4096])
.expect("unrecoverable residue");
let before = file_bytes(&directory.path);
let options = CompactionOptions::default().with_task(MaintenanceTask::RepairArchives);
let preview = plan_compaction(&directory.path, &options)
.expect_err("the preview must refuse an unrepairable archive");
assert!(
preview.to_string().contains("data00500a.tar"),
"the preview names the archive that dooms the run: {preview}"
);
match PreparedCompaction::prepare(&directory.path, options) {
Ok(_) => panic!("prepare must refuse too"),
Err(error) => assert!(
error.to_string().contains("data00500a.tar"),
"prepare names it as well: {error}"
),
}
assert_eq!(
file_bytes(&directory.path),
before,
"and neither path rewrote a single archive"
);
}
#[cfg(unix)]
#[test]
fn repository_root_symlink_is_resolved_to_the_canonical_target() {
use std::os::unix::fs::symlink;
let directory = TestDirectory::repository("root-symlink-target");
let link = directory.path.with_extension("repository-link");
let _ = std::fs::remove_file(&link);
symlink(&directory.path, &link).expect("create repository root symlink");
let expected = std::fs::canonicalize(&directory.path).expect("canonical target");
let plan =
plan_compaction(&link, &CompactionOptions::default()).expect("plan through alias");
assert_eq!(plan.directory(), expected);
let prepared = PreparedCompaction::prepare(
&link,
CompactionOptions::default().with_tasks(std::iter::empty()),
)
.expect("prepare through alias");
assert_eq!(prepared.plan().directory(), expected);
drop(prepared);
std::fs::remove_file(link).expect("remove repository link");
}
#[test]
fn journal_cleanup_preserves_mixed_physical_terminators_exactly() {
let directory = TestDirectory::repository("mixed-journal-terminators");
let head = Repository::open(&directory.path)
.expect("repository")
.head_record_identifier();
let retained = format!("{head} tag-lf 1\n{head} tag-crlf 2\r\n{head} tag-cr 3\r");
let source = format!("{retained}parser-skipped\n");
std::fs::write(directory.path.join("journal.log"), source.as_bytes())
.expect("write mixed journal fixture");
let options = CompactionOptions::default().with_tasks([MaintenanceTask::Journal]);
let plan = plan_compaction(&directory.path, &options).expect("plan mixed cleanup");
assert_eq!(plan.journal_line_removals().len(), 1);
assert_eq!(plan.journal_line_removals()[0].line_number(), 4);
compact(&directory.path, options).expect("apply mixed cleanup");
assert_eq!(
std::fs::read(directory.path.join("journal.log")).expect("read rewritten journal"),
retained.as_bytes(),
"LF, CRLF, and bare-CR terminators must remain byte-exact"
);
Repository::open(&directory.path).expect("mixed-terminator repository remains healthy");
}
#[test]
fn prepared_cleanup_is_excluded_by_an_existing_writer_lock() {
let directory = TestDirectory::repository("lock-exclusion");
let writer = WritableRepository::open(&directory.path).expect("hold writer lock");
plan_compaction(&directory.path, &CompactionOptions::default())
.expect("lock-free preview remains read-only");
assert!(
PreparedCompaction::prepare(&directory.path, CompactionOptions::default()).is_err()
);
writer.close().expect("close writer");
Repository::open(&directory.path).expect("repository healthy");
}
#[test]
fn redundant_journal_staging_is_removed_and_second_run_is_a_true_noop() {
let directory = TestDirectory::repository("temporary-idempotence");
let staging = directory.path.join("journal.log.compacting");
std::fs::copy(directory.path.join("journal.log"), &staging).expect("copy staging");
let forensic_staging = directory.path.join("journal.log.recovered");
std::fs::write(&forensic_staging, b"unterminated recovery evidence")
.expect("write ambiguous staging journal");
std::fs::write(directory.path.join("gc.log"), b"operator gc state\n").expect("seed gc log");
compact(&directory.path, CompactionOptions::default()).expect("first cleanup");
assert!(!staging.exists());
assert!(forensic_staging.exists());
assert_eq!(
std::fs::read(directory.path.join("gc.log")).expect("gc log"),
b"operator gc state\n"
);
let before = file_bytes(&directory.path);
let second =
plan_compaction(&directory.path, &CompactionOptions::default()).expect("second plan");
assert!(second.is_empty(), "second plan: {:?}", second.actions());
let outcome = PreparedCompaction::prepare(&directory.path, CompactionOptions::default())
.expect("prepare no-op")
.apply()
.expect("apply no-op");
assert_eq!(outcome.head_before, outcome.head_after);
assert_eq!(file_bytes(&directory.path), before);
}
#[test]
fn checkpoint_cleanup_allocates_after_zero_byte_next_archive_residue() {
let directory = TestDirectory::repository("checkpoint-zero-byte-next-archive");
let store = WritableRepository::open(&directory.path).expect("open writer");
create_checkpoint(&store, 1, &[]).expect("checkpoint");
store.close().expect("close writer");
std::thread::sleep(std::time::Duration::from_millis(10));
let repository = Repository::open(&directory.path).expect("open checkpoint repository");
let active_maximum = repository
.archives()
.iter()
.filter_map(|archive| ArchiveFileName::parse(archive.file_name()))
.map(|name| name.archive_number)
.max()
.expect("fixture has an active archive");
drop(repository);
let occupied_number = active_maximum.checked_add(1).expect("fixture namespace");
let certified_number = occupied_number.checked_add(1).expect("fixture namespace");
let occupied_name = format!("data{occupied_number:05}a.tar");
std::fs::write(directory.path.join(&occupied_name), b"")
.expect("install zero-byte otherwise-next residue");
let options =
CompactionOptions::default().with_tasks([MaintenanceTask::ExpiredCheckpoints]);
let plan = plan_compaction(&directory.path, &options).expect("plan checkpoint cleanup");
assert_eq!(plan.checkpoint_archive_number, Some(certified_number));
let outcome = compact(&directory.path, options).expect("apply checkpoint cleanup");
assert_eq!(outcome.removed_checkpoints, 1);
assert_eq!(
std::fs::read(directory.path.join(&occupied_name)).expect("read zero-byte residue"),
b"",
"cleanup must neither truncate nor reuse physical residue"
);
assert!(
directory
.path
.join(format!("data{certified_number:05}a.tar"))
.exists()
);
let repository = Repository::open(&directory.path).expect("healthy repository");
assert!(repository.checkpoints().expect("checkpoints").is_empty());
}
#[test]
fn the_history_price_counts_the_bulk_segments_held_behind_the_data_ones() {
let directory = TestDirectory::repository("history-price-bulk");
{
let store = WritableRepository::open(&directory.path).expect("open binary writer");
let mut writer = store.record_writer(GarbageCollectionGeneration {
generation: 0,
full_generation: 0,
is_compacted: false,
});
let content: Vec<u8> = (0..1024 * 1024).map(|index| (index % 251) as u8).collect();
let binary = writer.write_binary_content(&content).expect("binary");
let file = writer
.write_node(
Some("nt:file"),
&[],
&ChildNodesToWrite::Zero,
&[PropertyToWrite {
name: "data".to_owned(),
property_type: crate::content::property::PropertyType::Binary,
values: PropertyValuesToWrite::Single(binary),
}],
)
.expect("file node");
let head = writer
.write_node(
None,
&[],
&ChildNodesToWrite::One {
name: "root".to_owned(),
node: file,
},
&[],
)
.expect("binary super root");
writer.finish().expect("finish binary segments");
assert!(store.compare_and_set_head(store.head(), head));
store.close().expect("close binary writer");
}
{
let store = WritableRepository::open(&directory.path).expect("open new head writer");
let mut writer = store.record_writer(GarbageCollectionGeneration {
generation: 2,
full_generation: 2,
is_compacted: false,
});
let root = writer
.write_node(None, &[], &ChildNodesToWrite::Zero, &[])
.expect("new content root");
let head = writer
.write_node(
None,
&[],
&ChildNodesToWrite::One {
name: "root".to_owned(),
node: root,
},
&[],
)
.expect("new super root");
writer.finish().expect("finish new head");
assert!(store.compare_and_set_head(store.head(), head));
store.close().expect("close new head writer");
}
let priced = plan_compaction(
&directory.path,
&CompactionOptions::default().with_tasks([MaintenanceTask::Segments]),
)
.expect("priced plan");
let (_priced_segments, priced_bytes) = priced.history_protected_reclaimable();
let outcome = compact(
&directory.path,
CompactionOptions::default()
.with_tasks([MaintenanceTask::Segments, MaintenanceTask::Journal])
.with_journal_revision_retention(NonZeroUsize::new(1).expect("one revision")),
)
.expect("bounded cleanup");
let freed = outcome.archive_bytes_before - outcome.archive_bytes_after;
assert!(
freed > 1024 * 1024,
"retiring the history must free the binary content: {freed}"
);
assert_eq!(
priced_bytes, freed,
"the quoted price must be what retiring the history delivers"
);
}
#[test]
fn an_unselected_task_reports_no_step_of_its_own() {
let (directory, _old_head, _new_head) = history_veto_fixture("unselected-task-step");
let mut observer = StepNameObserver { names: Vec::new() };
compact_with_progress(
&directory.path,
CompactionOptions::default()
.with_tasks([
MaintenanceTask::Segments,
MaintenanceTask::Journal,
MaintenanceTask::StaleTemporaries,
])
.with_journal_revision_retention(NonZeroUsize::new(1).expect("one revision")),
&mut observer,
)
.expect("bounded cleanup");
for unselected in ["removing old recovery backups", "removing stale archives"] {
assert!(
!observer.names.iter().any(|name| name == unselected),
"unselected task reported {unselected:?}: {:?}",
observer.names
);
}
assert!(
observer
.names
.iter()
.any(|name| name == "removing stale temporary files"),
"a selected task still reports its step: {:?}",
observer.names
);
}
#[test]
fn dry_run_is_byte_exact_and_never_creates_the_lock_file() {
let directory = TestDirectory::repository("dry-run");
std::fs::remove_file(directory.path.join("repo.lock")).expect("remove old lock inode");
let before = file_bytes(&directory.path);
let mtimes_before = file_mtimes(&directory.path);
let plan = plan_compaction(&directory.path, &CompactionOptions::default()).expect("plan");
assert!(plan.is_empty());
assert_eq!(file_bytes(&directory.path), before);
assert_eq!(file_mtimes(&directory.path), mtimes_before);
assert!(!directory.path.join("repo.lock").exists());
}
#[cfg(unix)]
#[test]
fn prepared_cleanup_is_bound_to_an_ancestor_symlinks_resolved_target() {
use std::os::unix::fs::symlink;
let first_parent = TestDirectory::new("ancestor-alias-first");
let second_parent = TestDirectory::new("ancestor-alias-second");
let first_repository = first_parent.path.join("segmentstore");
let second_repository = second_parent.path.join("segmentstore");
for repository in [&first_repository, &second_repository] {
std::fs::create_dir(repository).expect("create repository directory");
WritableRepository::open(repository)
.expect("bootstrap repository")
.close()
.expect("close repository");
std::fs::copy(
repository.join("journal.log"),
repository.join("journal.log.compacting"),
)
.expect("create removable staging file");
}
let alias = first_parent.path.with_extension("ancestor-alias");
let _ = std::fs::remove_file(&alias);
symlink(&first_parent.path, &alias).expect("create ancestor alias");
let aliased_repository = alias.join("segmentstore");
let options = CompactionOptions::default().with_tasks([MaintenanceTask::StaleTemporaries]);
let prepared = PreparedCompaction::prepare(&aliased_repository, options)
.expect("prepare first target");
assert_eq!(
prepared.plan().directory(),
std::fs::canonicalize(&first_repository).expect("canonical first repository")
);
std::fs::remove_file(&alias).expect("remove first alias");
symlink(&second_parent.path, &alias).expect("retarget ancestor alias");
prepared.apply().expect("apply captured first target");
assert!(!first_repository.join("journal.log.compacting").exists());
assert!(second_repository.join("journal.log.compacting").exists());
Repository::open(&first_repository).expect("first repository remains healthy");
Repository::open(&second_repository).expect("second repository remains healthy");
std::fs::remove_file(alias).expect("remove ancestor alias");
}
#[cfg(unix)]
#[test]
fn relative_repository_path_is_stored_as_an_absolute_canonical_target() {
let directory = TestDirectory::repository("relative-canonical-target");
let current = std::fs::canonicalize(std::env::current_dir().expect("current directory"))
.expect("canonical current directory");
let target = std::fs::canonicalize(&directory.path).expect("canonical repository");
let relative = relative_path_from(¤t, &target);
assert!(!relative.is_absolute());
let plan =
plan_compaction(&relative, &CompactionOptions::default()).expect("relative plan");
assert_eq!(plan.directory(), target);
assert!(plan.directory().is_absolute());
}
#[test]
fn prepared_cleanup_refuses_a_replaced_lock_inode() {
let directory = TestDirectory::repository("replaced-lock");
let staging = directory.path.join("journal.log.compacting");
std::fs::copy(directory.path.join("journal.log"), &staging)
.expect("create removable staging file");
let options = CompactionOptions::default().with_tasks([MaintenanceTask::StaleTemporaries]);
let prepared = PreparedCompaction::prepare(&directory.path, options).expect("prepare");
let lock_path = directory.path.join("repo.lock");
std::fs::remove_file(&lock_path).expect("unlink held lock pathname");
std::fs::write(&lock_path, b"replacement inode").expect("replace lock inode");
let error = prepared
.apply()
.expect_err("replacement lock must abort apply");
assert!(error.to_string().contains("lock inode"));
assert!(staging.exists());
Repository::open(&directory.path).expect("repository remains healthy");
}
}