use std::path::Path;
use std::time::Instant;
use crate::db::lock::{DataDirLock, LOCK_FILE, LockError};
use crate::tree::Hash;
#[path = "shard_scan.rs"]
mod shard_scan;
use shard_scan::{ShardScan, scan_shard};
use super::consult::{MetadataPass, consult_metadata};
use super::error::{MarkSourceId, VacuumError};
use super::inventory::{StoreInventory, inventory_top_level};
use super::manifest::read_manifest;
use super::mark::{ProbeOutcome, StoreMarks, commit_marks, probe_store, walk_declared};
use super::report::{
MarkTallies, MetadataReport, NodeTotals, ShardPresence, ShardReport, SweepBlocker, VacuumMode,
VacuumReport,
};
use super::{CANONICAL_REFS_DIR, CANONICAL_SNAPSHOT_REGISTRY, VacuumOptions};
const TRUST_BOUNDARY: &str = "This run held the A4 data-dir writer lock, which excludes every \
writer routed through a live Database. It does NOT exclude direct \
DiskStore/tree/WAL/branch-layer/sync access opened on paths under this data_dir by other \
code: directories under a database data_dir are owned by that database, and direct access \
forfeits vacuum safety (STORAGE-VACUUM.md §6). If any such tool touches this directory, \
stop before sweeping.";
pub fn run(options: &VacuumOptions) -> Result<VacuumReport, VacuumError> {
let started = Instant::now();
let data_dir = options.data_dir.as_path();
read_config_checked(data_dir)?;
let lock = DataDirLock::acquire(data_dir).map_err(map_lock_error)?;
let anchor_created = lock.created_anchor();
let _lock = lock;
let config = read_config_checked(data_dir)?;
let shard_count = config.shard_count;
let manifest = read_manifest(data_dir)?;
let metadata = consult_metadata(data_dir, options, manifest.as_ref())?;
let top = inventory_top_level(
data_dir,
shard_count,
&[
crate::db::CONFIG_FILE,
LOCK_FILE,
super::manifest::MANIFEST_FILE,
CANONICAL_REFS_DIR,
CANONICAL_SNAPSHOT_REGISTRY,
],
&metadata.listed_paths,
)?;
let mut tallies = MarkTallies::default();
let branch_roots = group_branch_roots(&metadata, shard_count, &mut tallies)?;
let mut snapshots = group_snapshot_roots(&metadata, shard_count, &mut tallies)?;
let mut verified = 0_usize;
let mut totals = NodeTotals::default();
let mut shard_reports = Vec::new();
let mut never_materialised = 0_usize;
let mut blockers = metadata.blockers;
for shard_id in 0..shard_count {
let scan = scan_shard(data_dir, shard_id)?;
let shard_branch_roots: &[(Hash, MarkSourceId)] =
branch_roots.get(&shard_id).map_or(&[], Vec::as_slice);
if scan.presence == ShardPresence::NeverMaterialised && shard_branch_roots.is_empty() {
never_materialised += 1;
continue;
}
let report = process_shard(
scan,
shard_branch_roots,
&mut snapshots,
&mut tallies,
&mut verified,
)?;
totals.accumulate(report.nodes);
if !report.malformed.is_empty() {
blockers.push(SweepBlocker::MalformedStoreEntries {
shard_id,
count: report.malformed.len(),
});
}
shard_reports.push(report);
}
for snapshot in &snapshots {
if !snapshot.resolved {
return Err(VacuumError::UnresolvableRoot {
source: snapshot.source.clone(),
root: snapshot.root,
missing: snapshot.last_missing,
});
}
}
for id in top.shards_beyond_count {
blockers.push(SweepBlocker::ShardBeyondCount { id, shard_count });
}
#[cfg(not(unix))]
blockers.push(SweepBlocker::UnsupportedDurability);
Ok(VacuumReport {
mode: VacuumMode::Stats,
data_dir: options.data_dir.clone(),
shard_count,
duration_ms: u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX),
lock_anchor_created: anchor_created,
never_materialised_shards: never_materialised,
shards: shard_reports,
totals,
metadata: MetadataReport {
manifest: metadata.manifest_report,
attestation: options.attest_metadata_complete,
sources: metadata.sources,
supplied_missing: metadata.supplied_missing,
},
marks: tallies,
verified_nodes: verified,
uninventoried: top.uninventoried,
sweep_blockers: blockers,
trust_boundary: TRUST_BOUNDARY.to_owned(),
})
}
fn map_lock_error(error: LockError) -> VacuumError {
match error {
LockError::AlreadyLocked {
lock_path,
created_anchor,
} => VacuumError::Locked {
lock_path,
anchor_created: created_anchor,
},
LockError::AnchorNotARegularFile { lock_path } => VacuumError::LockIo {
error: std::io::Error::other(
"writer.lock is not a regular file — a symlinked anchor would make the \
create-capable open mint an inode OUTSIDE the data dir; remove or replace \
it with a plain empty file, then re-run",
),
lock_path,
anchor_created: false,
},
LockError::Io {
lock_path,
error,
created_anchor,
} => VacuumError::LockIo {
lock_path,
error,
anchor_created: created_anchor,
},
}
}
fn read_config_checked(data_dir: &Path) -> Result<crate::db::DatabaseConfig, VacuumError> {
let config_unreadable = |reason: String| VacuumError::ConfigUnreadable {
path: data_dir.join(crate::db::CONFIG_FILE),
reason,
};
let config = crate::db::config::read_config(data_dir)
.map_err(|error| config_unreadable(error.to_string()))?;
if config.shard_count == 0 {
return Err(config_unreadable(
"shard_count is 0; the engine requires at least 1 shard, so this config \
cannot be the engine's own work"
.to_owned(),
));
}
Ok(config)
}
fn group_branch_roots(
metadata: &MetadataPass,
shard_count: usize,
tallies: &mut MarkTallies,
) -> Result<std::collections::HashMap<usize, Vec<(Hash, MarkSourceId)>>, VacuumError> {
let mut per_shard: std::collections::HashMap<usize, Vec<(Hash, MarkSourceId)>> =
std::collections::HashMap::new();
for scan in &metadata.refs_scans {
for record in &scan.records {
for shard_ref in &record.shards {
if shard_ref.shard_id >= shard_count {
return Err(VacuumError::ShardBeyondCount {
id: shard_ref.shard_id,
shard_count,
});
}
let roots = per_shard.entry(shard_ref.shard_id).or_default();
roots.push((
shard_ref.fork_anchor,
MarkSourceId::BranchAnchor {
branch: record.name.clone(),
shard_id: shard_ref.shard_id,
},
));
roots.push((
shard_ref.head,
MarkSourceId::BranchHead {
branch: record.name.clone(),
shard_id: shard_ref.shard_id,
},
));
tallies.branch_anchors += 1;
tallies.branch_heads += 1;
}
}
}
Ok(per_shard)
}
#[derive(Clone)]
enum SnapshotCandidates {
AllShards,
Declared(std::sync::Arc<[usize]>),
}
impl SnapshotCandidates {
fn applies_to(&self, shard_id: usize) -> bool {
match self {
Self::AllShards => true,
Self::Declared(ids) => ids.contains(&shard_id),
}
}
}
struct SnapshotWork {
root: Hash,
source: MarkSourceId,
candidates: SnapshotCandidates,
resolved: bool,
last_missing: Hash,
}
fn group_snapshot_roots(
metadata: &MetadataPass,
shard_count: usize,
tallies: &mut MarkTallies,
) -> Result<Vec<SnapshotWork>, VacuumError> {
let mut work = Vec::new();
for registry in &metadata.registries {
let candidates = match ®istry.declared_shards {
Some(declared) => {
for &shard_id in declared {
if shard_id >= shard_count {
return Err(VacuumError::ShardBeyondCount {
id: shard_id,
shard_count,
});
}
}
SnapshotCandidates::Declared(std::sync::Arc::from(declared.as_slice()))
}
None => SnapshotCandidates::AllShards,
};
for entry in ®istry.entries {
work.push(SnapshotWork {
root: entry.root_hash,
source: MarkSourceId::Snapshot {
name: entry.name.clone(),
registry: registry.path.clone(),
},
candidates: candidates.clone(),
resolved: false,
last_missing: entry.root_hash,
});
tallies.snapshot_roots += 1;
}
}
Ok(work)
}
fn process_shard(
scan: ShardScan,
branch_roots: &[(Hash, MarkSourceId)],
snapshots: &mut [SnapshotWork],
tallies: &mut MarkTallies,
verified: &mut usize,
) -> Result<ShardReport, VacuumError> {
let ShardScan {
shard_id,
presence,
committed_root,
store_dir,
mut inventory,
mut shard_malformed,
} = scan;
let mut marks = StoreMarks::new();
if let Some(root) = committed_root {
walk_declared(
&inventory,
&store_dir,
root,
&MarkSourceId::WalRoot { shard_id },
&mut marks,
verified,
)?;
tallies.wal_roots += 1;
}
for (root, source) in branch_roots {
walk_declared(&inventory, &store_dir, *root, source, &mut marks, verified)?;
}
for snapshot in snapshots.iter_mut() {
if !snapshot.candidates.applies_to(shard_id) {
continue;
}
match probe_store(&inventory, &store_dir, snapshot.root, &marks, verified)? {
ProbeOutcome::Resolved(resolved) => {
commit_marks(&inventory, &resolved, &mut marks);
snapshot.resolved = true;
}
ProbeOutcome::Missing(missing) => snapshot.last_missing = missing,
}
}
let nodes = shard_totals(&inventory, &marks);
shard_malformed.append(&mut inventory.malformed);
Ok(ShardReport {
shard_id,
presence,
wal_committed_root: committed_root.is_some(),
nodes,
temp_debris_files: inventory.temp_debris_files,
temp_debris_bytes: inventory.temp_debris_bytes,
malformed: shard_malformed,
})
}
fn shard_totals(inventory: &StoreInventory, marks: &StoreMarks) -> NodeTotals {
let enumerated_nodes = inventory.nodes.len();
let enumerated_bytes: u64 = inventory
.nodes
.values()
.map(|file| file.compressed_len)
.sum();
let marked_nodes = marks.len();
let marked_bytes: u64 = marks.values().sum();
NodeTotals {
enumerated_nodes,
enumerated_bytes,
marked_nodes,
marked_bytes,
unmarked_nodes: enumerated_nodes.saturating_sub(marked_nodes),
unmarked_bytes: enumerated_bytes.saturating_sub(marked_bytes),
}
}
#[cfg(test)]
#[path = "stats_candidate_tests.rs"]
mod candidate_shape_tests;