diff --git a/crates/engine/src/log/file.rs b/crates/engine/src/log/file.rs index cde2fcb..d238f82 100644 --- a/crates/engine/src/log/file.rs +++ b/crates/engine/src/log/file.rs @@ -250,6 +250,57 @@ pub fn reclaim_tmp_filename(file_id: u32) -> String { format!("data-{:010}.log.tmp", file_id) } +pub fn replaces_filename(file_id: u32) -> String { + format!("data-{:010}.log.replaces", file_id) +} + +/// File ids a compact said it superseded. Open skips those even if the +/// unlinked `data-*.log` reappears (crash between rename and unlink). +pub fn load_replaced_ids(dir: &Path) -> std::collections::HashSet { + let mut skip = std::collections::HashSet::new(); + let Ok(rd) = std::fs::read_dir(dir) else { + return skip; + }; + for entry in rd.flatten() { + let name = entry.file_name(); + let name = name.to_string_lossy(); + if !name.ends_with(".log.replaces") { + continue; + } + let Ok(body) = std::fs::read_to_string(entry.path()) else { + continue; + }; + for line in body.lines() { + if let Ok(id) = line.parse::() { + skip.insert(id); + } + } + } + skip +} + +pub fn write_replaces(dir: &Path, new_id: u32, replaced: &[u32]) -> std::io::Result<()> { + let mut ids: Vec = replaced.to_vec(); + for &id in replaced { + let p = dir.join(replaces_filename(id)); + if let Ok(body) = std::fs::read_to_string(p) { + for line in body.lines() { + if let Ok(x) = line.parse::() { + ids.push(x); + } + } + } + } + ids.sort_unstable(); + ids.dedup(); + let body = ids + .iter() + .map(u32::to_string) + .collect::>() + .join("\n"); + std::fs::write(dir.join(replaces_filename(new_id)), body) +} + /// An open log file. Used for both active (writable) and sealed (read-only) /// files; the only difference is whether `append` is called. /// diff --git a/crates/engine/src/log/mod.rs b/crates/engine/src/log/mod.rs index 4a73c8d..da77c8b 100644 --- a/crates/engine/src/log/mod.rs +++ b/crates/engine/src/log/mod.rs @@ -3012,6 +3012,68 @@ mod enospc_recovery_tests { } } +#[cfg(test)] +mod reclaim_leak_tests { + use super::*; + use crate::log::config::LogConfig; + use crate::log::file::data_filename; + use bytes::Bytes; + use tempfile::TempDir; + + fn run(f: F) -> F::Output { + monoio::RuntimeBuilder::::new() + .enable_timer() + .build() + .expect("monoio runtime") + .block_on(f) + } + + /// Reclaim unlinks old sealed files best-effort. If one leaks back onto + /// disk, apply_footer_entries insert-only resurrects deletes. + #[test] + fn leaked_old_sealed_file_must_not_resurrect_deleted_key() { + run(async { + let dir = TempDir::new().unwrap(); + let path = dir.path().to_path_buf(); + let cfg = LogConfig { + rotate_threshold: 1 << 40, + fanout: 2, + value_sep_threshold: 1 << 20, + }; + { + let log = NamespaceLog::open(path.clone(), cfg).await.unwrap(); + log.put_full(Bytes::from_static(b"keep"), b"one", &[], None) + .await + .unwrap(); + log.put_full(Bytes::from_static(b"deleted"), b"gone", &[], None) + .await + .unwrap(); + log.seal_active_for_shutdown().await.unwrap(); + } + let leaked = path.join("leaked-data-0.log"); + std::fs::copy(path.join(data_filename(0)), &leaked).unwrap(); + + { + let log = NamespaceLog::open(path.clone(), cfg).await.unwrap(); + log.tombstone(b"deleted").await.unwrap(); + log.reclaim().await.unwrap(); + } + + std::fs::copy(&leaked, path.join(data_filename(0))).unwrap(); + + let log = NamespaceLog::open(path, cfg).await.unwrap(); + assert!( + log.index.borrow().get(b"keep").is_some(), + "live key must survive" + ); + assert!( + log.index.borrow().get(b"deleted").is_none(), + "tombstoned key must not come back from a leaked pre-delete sealed file" + ); + }); + } +} + #[cfg(test)] mod fd_footprint_tests { use super::*; diff --git a/crates/engine/src/log/reclaim.rs b/crates/engine/src/log/reclaim.rs index e80c254..00250f1 100644 --- a/crates/engine/src/log/reclaim.rs +++ b/crates/engine/src/log/reclaim.rs @@ -129,6 +129,15 @@ pub async fn reclaim_namespace( // already unlinked below — losing the compacted data. crate::log::file::sync_dir(&dir).await; + // Durable "these ids are dead" before unlink. Open ignores them even + // if unlink fails or a crash leaves the old data-*.log behind. + let replaced: Vec = sealed_files.iter().map(|f| f.file_id).collect(); + if let Err(e) = crate::log::file::write_replaces(&dir, next_file_id, &replaced) { + warn!(error = %e, "failed to write replaces sidecar; leaked inputs may resurrect deletes"); + } else { + crate::log::file::sync_dir(&dir).await; + } + let live_keys = new_entries.len() as u64; // Unlink all old sealed files concurrently via io_uring. diff --git a/crates/engine/src/log/recover.rs b/crates/engine/src/log/recover.rs index 11a4e75..e59093b 100644 --- a/crates/engine/src/log/recover.rs +++ b/crates/engine/src/log/recover.rs @@ -31,6 +31,8 @@ pub struct OpenedFiles { pub async fn open_namespace(dir: PathBuf) -> Result { std::fs::create_dir_all(&dir)?; let mut data_files = list_data_files(&dir)?; + let replaced = crate::log::file::load_replaced_ids(&dir); + data_files.retain(|(id, _)| !replaced.contains(id)); let mut index = NsIndex::new(); let mut sealed: Vec = Vec::new();