use std::fs;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, SystemTime};
const LOCK_SERIALIZED_STAGING_PREFIXES: &[&str] = &[
"cass-lexical-shards.",
"cass-lexical-merge.",
"cass-empty-lexical-repair-",
];
const SHARED_STAGING_PREFIXES: &[&str] = &["cass-federated-materialize-"];
const LOCK_SERIALIZED_STAGING_MIN_AGE: Duration = Duration::from_secs(60);
const SHARED_STAGING_MIN_AGE: Duration = Duration::from_secs(60 * 60);
const RECLAIM_MARKER_INFIX: &str = ".reclaim-";
static RECLAIM_SEQ: AtomicU64 = AtomicU64::new(0);
#[derive(Debug, Default)]
pub(crate) struct StagingReclaimReport {
pub reclaimed_dirs: usize,
pub reclaimed_bytes: u64,
pub skipped_young: usize,
pub errors: Vec<String>,
}
impl StagingReclaimReport {
fn merge(&mut self, other: StagingReclaimReport) {
self.reclaimed_dirs += other.reclaimed_dirs;
self.reclaimed_bytes = self.reclaimed_bytes.saturating_add(other.reclaimed_bytes);
self.skipped_young += other.skipped_young;
self.errors.extend(other.errors);
}
pub(crate) fn log(&self) {
if self.reclaimed_dirs > 0 || self.reclaimed_bytes > 0 {
tracing::info!(
reclaimed_dirs = self.reclaimed_dirs,
reclaimed_bytes = self.reclaimed_bytes,
skipped_young = self.skipped_young,
"reclaimed orphaned lexical staging directories from previous crashed runs (#324)"
);
}
for error in &self.errors {
tracing::warn!(
error = %error,
"orphaned staging dir sweep hit a non-fatal error; will retry next startup"
);
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum StagingNameClass {
NotStaging,
LockSerialized,
Shared,
}
fn classify_staging_name(name: &str) -> StagingNameClass {
for prefix in LOCK_SERIALIZED_STAGING_PREFIXES {
if name.len() > prefix.len() && name.starts_with(prefix) {
return StagingNameClass::LockSerialized;
}
}
for prefix in SHARED_STAGING_PREFIXES {
if name.len() > prefix.len() && name.starts_with(prefix) {
return StagingNameClass::Shared;
}
}
StagingNameClass::NotStaging
}
fn is_reclaim_marker(name: &str) -> bool {
name.contains(RECLAIM_MARKER_INFIX)
}
fn dir_tree_size_bytes(root: &Path) -> u64 {
let mut total = 0_u64;
for entry in walkdir::WalkDir::new(root).follow_links(false) {
let Ok(entry) = entry else { continue };
let Ok(metadata) = entry.path().symlink_metadata() else {
continue;
};
if metadata.is_file() {
total = total.saturating_add(metadata.len());
}
}
total
}
pub(crate) fn reclaim_orphaned_staging_dirs_for_data_dir(
data_dir: &Path,
now: SystemTime,
) -> StagingReclaimReport {
let staging_root = crate::search::tantivy::expected_index_dir(data_dir)
.parent()
.map_or_else(|| data_dir.join("index"), Path::to_path_buf);
let mut report = reclaim_orphaned_staging_dirs_under_index_run_lock(&staging_root, now);
if staging_root != *data_dir {
report.merge(reclaim_orphaned_staging_dirs_under_index_run_lock(
data_dir, now,
));
}
report
}
pub(crate) fn reclaim_orphaned_staging_dirs_under_index_run_lock(
staging_root: &Path,
now: SystemTime,
) -> StagingReclaimReport {
let mut report = StagingReclaimReport::default();
let entries = match fs::read_dir(staging_root) {
Ok(entries) => entries,
Err(err) if err.kind() == std::io::ErrorKind::NotFound => return report,
Err(err) => {
report
.errors
.push(format!("reading {}: {err}", staging_root.display()));
return report;
}
};
for entry in entries {
let entry = match entry {
Ok(entry) => entry,
Err(err) => {
report
.errors
.push(format!("listing {}: {err}", staging_root.display()));
continue;
}
};
let name_os = entry.file_name();
let Some(name) = name_os.to_str() else {
continue;
};
let class = classify_staging_name(name);
if class == StagingNameClass::NotStaging {
continue;
}
let path = entry.path();
let metadata = match fs::symlink_metadata(&path) {
Ok(metadata) => metadata,
Err(err) => {
if err.kind() != std::io::ErrorKind::NotFound {
report
.errors
.push(format!("stat {}: {err}", path.display()));
}
continue;
}
};
if !metadata.is_dir() {
continue;
}
let min_age = match class {
StagingNameClass::LockSerialized => LOCK_SERIALIZED_STAGING_MIN_AGE,
StagingNameClass::Shared => SHARED_STAGING_MIN_AGE,
StagingNameClass::NotStaging => unreachable!("filtered above"),
};
let old_enough = metadata
.modified()
.ok()
.and_then(|modified| now.duration_since(modified).ok())
.is_some_and(|age| age >= min_age);
if !old_enough {
report.skipped_young += 1;
continue;
}
let reclaimed_bytes = dir_tree_size_bytes(&path);
let doomed_path = if is_reclaim_marker(name) {
path.clone()
} else {
let marker_name = format!(
"{name}{RECLAIM_MARKER_INFIX}{}-{}",
std::process::id(),
RECLAIM_SEQ.fetch_add(1, Ordering::Relaxed)
);
let marker_path: PathBuf = staging_root.join(marker_name);
match fs::rename(&path, &marker_path) {
Ok(()) => marker_path,
Err(err) => {
if err.kind() != std::io::ErrorKind::NotFound {
report.errors.push(format!(
"renaming {} for reclamation: {err}",
path.display()
));
}
continue;
}
}
};
match fs::remove_dir_all(&doomed_path) {
Ok(()) => {
report.reclaimed_dirs += 1;
report.reclaimed_bytes = report.reclaimed_bytes.saturating_add(reclaimed_bytes);
}
Err(err) => {
report.errors.push(format!(
"removing orphaned staging dir {}: {err}",
doomed_path.display()
));
}
}
}
report
}
#[cfg(test)]
mod tests {
use super::*;
use std::fs;
use std::time::{Duration, SystemTime};
fn make_dir_with_payload(root: &Path, name: &str) -> PathBuf {
let dir = root.join(name);
fs::create_dir_all(dir.join("nested")).unwrap();
fs::write(dir.join("nested").join("payload.bin"), vec![0_u8; 4096]).unwrap();
dir
}
fn future_now(min_age: Duration) -> SystemTime {
SystemTime::now() + min_age + Duration::from_secs(600)
}
#[test]
fn stale_lexical_staging_dirs_are_reclaimed() {
let tmp = tempfile::tempdir().unwrap();
let shards = make_dir_with_payload(tmp.path(), "cass-lexical-shards.abc123");
let merge = make_dir_with_payload(tmp.path(), "cass-lexical-merge.xYz789");
let repair = make_dir_with_payload(tmp.path(), "cass-empty-lexical-repair-q1w2e3");
let report = reclaim_orphaned_staging_dirs_under_index_run_lock(
tmp.path(),
future_now(LOCK_SERIALIZED_STAGING_MIN_AGE),
);
assert_eq!(report.reclaimed_dirs, 3, "errors: {:?}", report.errors);
assert!(report.reclaimed_bytes >= 3 * 4096);
assert!(report.errors.is_empty(), "errors: {:?}", report.errors);
assert!(!shards.exists());
assert!(!merge.exists());
assert!(!repair.exists());
}
#[test]
fn fresh_lexical_staging_dirs_survive_the_sweep() {
let tmp = tempfile::tempdir().unwrap();
let fresh = make_dir_with_payload(tmp.path(), "cass-lexical-merge.fresh1");
let report =
reclaim_orphaned_staging_dirs_under_index_run_lock(tmp.path(), SystemTime::now());
assert_eq!(report.reclaimed_dirs, 0);
assert_eq!(report.skipped_young, 1);
assert!(fresh.exists());
assert!(fresh.join("nested").join("payload.bin").exists());
}
#[test]
fn federated_materialize_dirs_use_the_longer_threshold() {
let tmp = tempfile::tempdir().unwrap();
let materialize = make_dir_with_payload(tmp.path(), "cass-federated-materialize-r4nd0m");
let report = reclaim_orphaned_staging_dirs_under_index_run_lock(
tmp.path(),
SystemTime::now() + LOCK_SERIALIZED_STAGING_MIN_AGE + Duration::from_secs(60),
);
assert_eq!(report.reclaimed_dirs, 0);
assert_eq!(report.skipped_young, 1);
assert!(materialize.exists());
let report = reclaim_orphaned_staging_dirs_under_index_run_lock(
tmp.path(),
future_now(SHARED_STAGING_MIN_AGE),
);
assert_eq!(report.reclaimed_dirs, 1, "errors: {:?}", report.errors);
assert!(!materialize.exists());
}
#[test]
fn non_staging_entries_are_never_touched() {
let tmp = tempfile::tempdir().unwrap();
let keepers = [
"v6",
"v8",
".lexical-publish-backups",
"v7.bak.3",
"cass-lexical-shards.",
"cass-lexical-merge.",
"cass-lexical-shardsX",
"my-cass-lexical-merge.abc",
];
for name in keepers {
make_dir_with_payload(tmp.path(), name);
}
fs::write(tmp.path().join("cass-lexical-merge.imafile"), b"data").unwrap();
let report = reclaim_orphaned_staging_dirs_under_index_run_lock(
tmp.path(),
future_now(SHARED_STAGING_MIN_AGE),
);
assert_eq!(report.reclaimed_dirs, 0);
assert!(report.errors.is_empty(), "errors: {:?}", report.errors);
for name in keepers {
assert!(tmp.path().join(name).exists(), "{name} must survive");
}
assert!(tmp.path().join("cass-lexical-merge.imafile").exists());
}
#[cfg(unix)]
#[test]
fn symlinks_with_staging_names_are_skipped_and_targets_survive() {
let tmp = tempfile::tempdir().unwrap();
let victim = make_dir_with_payload(tmp.path(), "precious-user-data");
let link = tmp.path().join("cass-lexical-merge.sneaky");
std::os::unix::fs::symlink(&victim, &link).unwrap();
let report = reclaim_orphaned_staging_dirs_under_index_run_lock(
tmp.path(),
future_now(SHARED_STAGING_MIN_AGE),
);
assert_eq!(report.reclaimed_dirs, 0);
assert!(link.exists(), "symlink itself must not be removed");
assert!(victim.join("nested").join("payload.bin").exists());
}
#[test]
fn interrupted_reclaim_markers_are_finished_on_the_next_sweep() {
let tmp = tempfile::tempdir().unwrap();
let marker = make_dir_with_payload(tmp.path(), "cass-lexical-shards.abc.reclaim-42-0");
let report = reclaim_orphaned_staging_dirs_under_index_run_lock(
tmp.path(),
future_now(LOCK_SERIALIZED_STAGING_MIN_AGE),
);
assert_eq!(report.reclaimed_dirs, 1, "errors: {:?}", report.errors);
assert!(!marker.exists());
}
#[test]
fn data_dir_wrapper_sweeps_both_index_root_and_data_dir() {
let tmp = tempfile::tempdir().unwrap();
let index_root = tmp.path().join("index");
fs::create_dir_all(&index_root).unwrap();
let in_index = make_dir_with_payload(&index_root, "cass-lexical-merge.inindex");
let top_level = make_dir_with_payload(tmp.path(), "cass-lexical-shards.toplevel");
let live_index = make_dir_with_payload(&index_root, "v8");
let report = reclaim_orphaned_staging_dirs_for_data_dir(
tmp.path(),
future_now(LOCK_SERIALIZED_STAGING_MIN_AGE),
);
assert_eq!(report.reclaimed_dirs, 2, "errors: {:?}", report.errors);
assert!(!in_index.exists());
assert!(!top_level.exists());
assert!(live_index.exists(), "published index generations survive");
}
#[test]
fn missing_staging_root_is_a_quiet_no_op() {
let tmp = tempfile::tempdir().unwrap();
let report = reclaim_orphaned_staging_dirs_under_index_run_lock(
&tmp.path().join("does-not-exist"),
future_now(SHARED_STAGING_MIN_AGE),
);
assert_eq!(report.reclaimed_dirs, 0);
assert!(report.errors.is_empty());
}
}