From 4a712b0ff4956255328cc9356f1fb834a6a98aab Mon Sep 17 00:00:00 2001 From: Tin Dang Date: Fri, 17 Jul 2026 12:34:42 +0700 Subject: [PATCH 1/2] =?UTF-8?q?fix(persistence):=20latch=20a=20degraded=20?= =?UTF-8?q?state=20when=20the=20data=20dir=20vanishes=20=E2=80=94=20no=20p?= =?UTF-8?q?er-tick=20ENOENT=20error=20loop=20(#366)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Deleting --dir under a running server (operator rm -rf, tmp-cleaner sweep, mount loss) used to leave the per-shard persistence tick retrying its file operations on ENOENT forever: WAL flush at the 1ms cadence, the WAL-overflow directory scan warning once per second per shard, checkpoints failing every tick. Observed in the wild as an idle 4-shard server at ~667% CPU (sys:user ~50:1) with a wedged PING. Design-for-failure, mirroring the MOONERR diskfull guard: - DiskMonitor's existing 5s statvfs poll now distinguishes ENOENT/ENOTDIR on the monitored path (deterministic "the dir is gone") from transient errors (which keep retaining previous state, as before), and latches a new `dir_lost` flag with ONE loud ERROR log. The latch feeds the existing write-stall gate: writes are refused with a distinct `MOONERR dirmissing` message at both dispatch sites; reads and PING keep serving. - The persistence tick consults `is_dir_lost()` (one relaxed atomic load) and skips WAL flush, checkpoint tick, WAL-overflow scan, and auto-save entirely while latched — zero per-tick syscalls, zero per-tick log lines. Manual BGSAVE still fails one attempt loudly through the existing finalize-error path (fan-in counter intact). - Self-heal: the next poll that sees the path again clears the latch and resumes writes, warning explicitly that data accepted since the deletion is not guaranteed durable until restart (WAL fds point at unlinked inodes). - The latch only arms when persistence is configured (an in-memory server losing its CWD has nothing to protect) and works with --disk-free-min-pct 0: poll_global now runs whenever either the free-space guard or the dir guard is active. Measured on the incident configuration (4 shards, appendonly yes): SET → MOONERR dirmissing within one poll, GET/PING unaffected, CPU flat at the idle baseline, stderr frozen after the single ERROR line (pre-fix: +4 warn lines/s forever). Red/green: tests/dir_deleted_degraded.rs pins the four contract points (refusal ≤ one poll, responsiveness while latched, quiesced log+CPU, self-heal) and failed at the refusal assertion pre-fix; two unit tests pin the latch transitions and the disarmed-guard default. Gates: fmt, clippy (default + tokio,jemalloc), 4369 lib tests, integration test green under monoio (macOS), tokio (macOS), and Linux VM. Closes #366 author: Tin Dang --- CHANGELOG.md | 17 ++ src/main.rs | 10 +- src/server/conn/handler_monoio/dispatch.rs | 7 +- src/server/conn/handler_sharded/dispatch.rs | 8 +- src/shard/disk_monitor.rs | 224 +++++++++++++++++--- src/shard/persistence_tick.rs | 29 +++ tests/dir_deleted_degraded.rs | 209 ++++++++++++++++++ 7 files changed, 468 insertions(+), 36 deletions(-) create mode 100644 tests/dir_deleted_degraded.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index fba32b323..6be269804 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,23 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Fixed +- **Data-dir deleted under a running server now latches a degraded state + instead of error-looping the persistence tick (#366).** When `--dir` + vanishes (operator `rm -rf`, tmp-cleaner sweep, mount loss), the per-shard + 1ms WAL/checkpoint tick used to retry its file operations on `ENOENT` + forever — observed in the wild at ~667% CPU on an idle 4-shard server, + with a wedged `PING`. The disk monitor's 5s poll now distinguishes + `statvfs` `ENOENT`/`ENOTDIR` on the monitored path and latches `dir_lost` + (one loud `ERROR` log): write commands are refused with `MOONERR + dirmissing` (reads and `PING` keep serving, mirroring the diskfull guard), + and the persistence tick skips WAL flush / checkpoint / WAL-overflow scan + / auto-save entirely — zero per-tick syscalls and zero per-tick log lines + while latched. The latch self-heals: the next poll that sees the directory + again resumes writes (with an explicit warning that data accepted since + the deletion is not guaranteed durable until restart). Only engages when + persistence is configured, and works with `--disk-free-min-pct 0`. + ### Performance - **Disk-offload: shard event loop no longer pays manifest-commit fsyncs (task #59 levers 1+2).** Under spill load, `apply_completion_vec` ran one diff --git a/src/main.rs b/src/main.rs index 491a97aef..1352bdd3a 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1025,10 +1025,16 @@ fn main() -> anyhow::Result<()> { // MA12: Initialise disk free-space monitor. // Monitors the WAL/persistence volume. When disk_free_min_pct == 0, the - // monitor is inactive (poll_global is a no-op, is_write_paused always false). + // free-space guard is inactive, but a persistence-enabled server still + // polls so the dir-lost latch (#366) can refuse writes and quiesce the + // persistence tick if the data directory is deleted out from under it. { let monitor_path = persistence_dir.as_deref().unwrap_or(&config.dir); - moon::shard::disk_monitor::init_global(config.disk_free_min_pct, monitor_path); + moon::shard::disk_monitor::init_global( + config.disk_free_min_pct, + monitor_path, + persistence_dir.is_some(), + ); } // Wave 3: Initialise proactive RSS memory watchdog ("mem-full guard"). diff --git a/src/server/conn/handler_monoio/dispatch.rs b/src/server/conn/handler_monoio/dispatch.rs index 70ec47a89..d17312cd3 100644 --- a/src/server/conn/handler_monoio/dispatch.rs +++ b/src/server/conn/handler_monoio/dispatch.rs @@ -814,7 +814,12 @@ pub(super) fn try_enforce_readonly( #[inline] pub(super) fn try_enforce_disk_full(cmd: &[u8], responses: &mut Vec) -> bool { if metadata::is_write(cmd) && crate::shard::segment_stall::is_any_write_stall_active() { - let msg: &'static [u8] = if crate::shard::disk_monitor::is_write_paused() { + // dir-lost first: it also sets `is_write_paused`, and "diskfull" + // would send the operator hunting free space instead of the missing + // data directory (#366). + let msg: &'static [u8] = if crate::shard::disk_monitor::is_dir_lost() { + b"MOONERR dirmissing: data directory was removed; writes refused until it is restored" + } else if crate::shard::disk_monitor::is_write_paused() { b"MOONERR diskfull: writes paused until free space recovers" } else if crate::shard::mem_monitor::is_write_paused() { b"MOONERR memfull: writes paused until memory pressure recovers" diff --git a/src/server/conn/handler_sharded/dispatch.rs b/src/server/conn/handler_sharded/dispatch.rs index 9ba0ec0d1..e4d9e4b55 100644 --- a/src/server/conn/handler_sharded/dispatch.rs +++ b/src/server/conn/handler_sharded/dispatch.rs @@ -779,8 +779,12 @@ pub(super) fn try_enforce_readonly( #[inline] pub(super) fn try_enforce_disk_full(cmd: &[u8], responses: &mut Vec) -> bool { if metadata::is_write(cmd) && crate::shard::segment_stall::is_any_write_stall_active() { - // Distinguish the stall source for operator clarity. - let msg: &'static [u8] = if crate::shard::disk_monitor::is_write_paused() { + // Distinguish the stall source for operator clarity. dir-lost first: + // it also sets `is_write_paused`, and "diskfull" would send the + // operator hunting free space instead of the missing data dir (#366). + let msg: &'static [u8] = if crate::shard::disk_monitor::is_dir_lost() { + b"MOONERR dirmissing: data directory was removed; writes refused until it is restored" + } else if crate::shard::disk_monitor::is_write_paused() { b"MOONERR diskfull: writes paused until free space recovers" } else if crate::shard::mem_monitor::is_write_paused() { b"MOONERR memfull: writes paused until memory pressure recovers" diff --git a/src/shard/disk_monitor.rs b/src/shard/disk_monitor.rs index 8e7753423..50998e609 100644 --- a/src/shard/disk_monitor.rs +++ b/src/shard/disk_monitor.rs @@ -17,6 +17,29 @@ //! rises above `min_pct + 5` percentage points. This prevents flapping around the //! threshold when compaction reclaims space incrementally. //! +//! ## Dir-lost latch (issue #366) +//! +//! When the monitored path itself vanishes (`statvfs` → `ENOENT`/`ENOTDIR`: +//! operator `rm -rf`, tmp-cleaner sweep, mount disappearing), a server with +//! persistence configured must NOT keep retrying its per-tick WAL/checkpoint +//! file operations — that is a 1ms syscall error loop that burned ~100% CPU +//! per shard thread in the wild. Instead `poll` latches `dir_lost`: +//! +//! - `paused()` (and thus the `MOONERR` write gate) reports `true`, with a +//! distinct `dirmissing` message at the dispatch sites. +//! - The persistence tick consults [`is_dir_lost`] and skips WAL flush / +//! checkpoint / overflow-scan / auto-save work entirely — zero per-tick +//! syscalls and zero per-tick log lines while latched. +//! - The latch self-heals: the next 5s poll that sees the path again clears +//! it and writes resume. Data written between deletion and restart is NOT +//! guaranteed durable (the WAL fd points at unlinked inodes) — the heal +//! log line says so. +//! +//! The latch only engages when the server was started WITH persistence +//! (`dir_guard`, see [`init_global`]); an in-memory server losing its CWD +//! has nothing to protect. Transient `statvfs` errors other than +//! path-missing keep the previous state, exactly as before. +//! //! ## Platform support //! //! - **Linux** and **macOS**: both use `libc::statvfs` (POSIX; available on both). @@ -49,8 +72,15 @@ pub struct DiskMonitor { free_bytes: AtomicU64, /// True when writes should be refused (free < min_pct). paused: AtomicBool, + /// True when the monitored path itself has vanished (issue #366). + /// Latched by `poll`, cleared when the path reappears. + dir_lost: AtomicBool, /// Minimum free-space percentage before writes are paused (inclusive). min_pct: u8, + /// Whether the dir-lost latch is armed. Only true when the server was + /// started with persistence configured — an in-memory server losing the + /// directory it happened to be monitoring has nothing to protect. + dir_guard: bool, /// Path whose volume is monitored (typically the WAL / data directory). monitored_path: Box, } @@ -62,16 +92,28 @@ impl DiskMonitor { /// * `path` — any path on the volume to monitor (typically `persistence_dir`). /// /// Starts unpaused with `free_bytes = u64::MAX` (optimistic: no pause before - /// the first `poll` fires). + /// the first `poll` fires). The dir-lost latch is DISARMED; arm it with + /// [`Self::with_dir_guard`] when persistence is configured. pub fn new(min_pct: u8, path: impl AsRef) -> Self { Self { free_bytes: AtomicU64::new(u64::MAX), paused: AtomicBool::new(false), + dir_lost: AtomicBool::new(false), min_pct, + dir_guard: false, monitored_path: path.as_ref().into(), } } + /// Arm (or disarm) the dir-lost latch. Builder-style so the single + /// production call site ([`init_global`]) stays one expression and the + /// many test constructors keep the disarmed default. + #[must_use] + pub fn with_dir_guard(mut self, armed: bool) -> Self { + self.dir_guard = armed; + self + } + /// Latest free bytes on the monitored volume. /// /// Returns `u64::MAX` until the first `poll` completes successfully. @@ -80,12 +122,22 @@ impl DiskMonitor { self.free_bytes.load(Ordering::Relaxed) } - /// Returns `true` when writes should be refused. + /// Returns `true` when writes should be refused (disk nearly full OR the + /// data directory has vanished). /// - /// Extremely cheap: single `AtomicBool::load(Relaxed)`. + /// Extremely cheap: two `AtomicBool::load(Relaxed)`. #[inline] pub fn paused(&self) -> bool { - self.paused.load(Ordering::Relaxed) + self.paused.load(Ordering::Relaxed) || self.dir_lost.load(Ordering::Relaxed) + } + + /// Returns `true` while the monitored data directory is missing. + /// + /// Extremely cheap: single `AtomicBool::load(Relaxed)`. The persistence + /// tick consults this to skip all file operations while latched. + #[inline] + pub fn dir_lost(&self) -> bool { + self.dir_lost.load(Ordering::Relaxed) } /// Current configured minimum free-space percentage. @@ -100,24 +152,44 @@ impl DiskMonitor { /// but must not be called on every write (it is not free). pub fn poll(&self) { match query_free_bytes(&self.monitored_path) { - Some((free, total)) => { + Ok((free, total)) => { + if self.dir_guard && self.dir_lost.swap(false, Ordering::Release) { + tracing::warn!( + path = %self.monitored_path.display(), + "disk_monitor: data directory reappeared — writes resume; \ + data accepted since the deletion is NOT guaranteed durable \ + until the server is restarted (WAL fds point at unlinked \ + inodes)", + ); + } self.free_bytes.store(free, Ordering::Relaxed); self.update_paused(free, total); // Wire P10 INFO metrics (MA12 → RECL_*). crate::command::info_reclamation::RECL_DISK_FREE_BYTES .store(free, Ordering::Relaxed); - crate::command::info_reclamation::RECL_WRITE_STALL_ACTIVE.store( - if self.paused.load(Ordering::Relaxed) { - 1 - } else { - 0 - }, - Ordering::Relaxed, - ); + crate::command::info_reclamation::RECL_WRITE_STALL_ACTIVE + .store(if self.paused() { 1 } else { 0 }, Ordering::Relaxed); + } + Err(QueryError::PathMissing) if self.dir_guard => { + // Issue #366: the data dir itself is gone. Latch loud ONCE, + // refuse writes, and let the persistence tick skip its file + // operations — never a per-tick retry/error loop. + if !self.dir_lost.swap(true, Ordering::Release) { + tracing::error!( + path = %self.monitored_path.display(), + "disk_monitor: data directory MISSING (deleted or \ + unmounted?) — persistence suspended and writes refused \ + (MOONERR dirmissing); restore the directory to resume", + ); + } + crate::command::info_reclamation::RECL_WRITE_STALL_ACTIVE + .store(1, Ordering::Relaxed); } - None => { + Err(_) => { // Cannot read filesystem info — leave previous state unchanged - // to avoid spurious pauses due to transient errors. + // to avoid spurious pauses due to transient errors. (This arm + // also swallows PathMissing when the dir-guard is disarmed: + // an in-memory server keeps running unchanged.) tracing::warn!( path = %self.monitored_path.display(), "disk_monitor: statvfs failed; retaining previous pause state", @@ -184,10 +256,14 @@ static GLOBAL_DISK_MONITOR: OnceLock> = OnceLock::new(); /// Must be called once at server startup (before any shard event loops start). /// Calling it more than once is a no-op — the first call wins. /// -/// * `min_pct` — from `ServerConfig::disk_free_min_pct`. -/// * `path` — the WAL / persistence directory (or `"."` if none configured). -pub fn init_global(min_pct: u8, path: impl AsRef) { - let monitor = Arc::new(DiskMonitor::new(min_pct, path)); +/// * `min_pct` — from `ServerConfig::disk_free_min_pct`. +/// * `path` — the WAL / persistence directory (or `"."` if none configured). +/// * `guard_dir` — arm the dir-lost latch (issue #366). Pass `true` exactly +/// when persistence is configured: the latch protects the WAL/AOF data +/// directory, and works even when `min_pct == 0` relaxes the disk-FULL +/// guard (the standard crash-test configuration). +pub fn init_global(min_pct: u8, path: impl AsRef, guard_dir: bool) { + let monitor = Arc::new(DiskMonitor::new(min_pct, path).with_dir_guard(guard_dir)); // Ignore the error: if already set, the existing instance remains. let _ = GLOBAL_DISK_MONITOR.set(monitor); } @@ -203,16 +279,33 @@ pub fn global() -> Option> { /// Poll the global monitor (shard 0 calls this every 5 seconds). /// -/// No-op if `init_global` has not been called or if `min_pct == 0` -/// (monitoring disabled). +/// No-op if `init_global` has not been called, or if BOTH the disk-full +/// guard (`min_pct == 0`) and the dir-lost latch (`guard_dir == false`) are +/// disabled — a persistence-enabled server polls even with `min_pct == 0` +/// so the dir-lost latch stays armed. pub fn poll_global() { if let Some(m) = GLOBAL_DISK_MONITOR.get() { - if m.min_pct > 0 { + if m.min_pct > 0 || m.dir_guard { m.poll(); } } } +/// Returns `true` while the persistence data directory is missing (#366). +/// +/// **Hot/tick-path function.** Single `AtomicBool::load(Relaxed)`; `false` +/// when monitoring is uninitialised or the dir-guard is disarmed. The +/// persistence tick uses this to skip WAL flush / checkpoint / overflow-scan +/// / auto-save file operations while the directory is gone, instead of +/// error-looping on ENOENT every tick. +#[inline] +pub fn is_dir_lost() -> bool { + match GLOBAL_DISK_MONITOR.get() { + Some(m) => m.dir_lost(), + None => false, + } +} + /// Returns `true` if writes should be refused due to low disk space. /// /// **Hot-path function.** Single `AtomicBool::load(Relaxed)` when monitoring @@ -239,16 +332,27 @@ pub fn global_free_bytes() -> u64 { // ── Platform-specific statvfs wrapper ────────────────────────────────────── -/// Returns `Some((free_bytes, total_bytes))` for the volume containing `path`. -/// Returns `None` on any error. +/// Why `query_free_bytes` failed. `PathMissing` is the only variant the +/// caller distinguishes: it deterministically means the monitored path no +/// longer exists (the issue #366 trigger), unlike transient I/O errors. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum QueryError { + /// `ENOENT` / `ENOTDIR` — the monitored path is gone. + PathMissing, + /// Any other failure (permissions, interior NULs, transient I/O). + Other, +} + +/// Returns `Ok((free_bytes, total_bytes))` for the volume containing `path`. #[cfg(unix)] -fn query_free_bytes(path: &Path) -> Option<(u64, u64)> { +fn query_free_bytes(path: &Path) -> Result<(u64, u64), QueryError> { use std::ffi::CString; use std::mem::MaybeUninit; // Convert the path to a C string. Paths with interior NULs are pathological // — treat them as unreadable. - let c_path = CString::new(path.as_os_str().as_encoded_bytes()).ok()?; + let c_path = + CString::new(path.as_os_str().as_encoded_bytes()).map_err(|_| QueryError::Other)?; // `statvfs` is a POSIX C function. We pass a properly terminated C string // and a properly sized output buffer. `MaybeUninit` avoids reading @@ -258,7 +362,14 @@ fn query_free_bytes(path: &Path) -> Option<(u64, u64)> { // SAFETY: POSIX statvfs with valid C-string and MaybeUninit out-pointer. let rc = unsafe { libc::statvfs(c_path.as_ptr(), stat.as_mut_ptr()) }; if rc != 0 { - return None; + // Distinguish "the path is gone" (deterministic — the #366 trigger) + // from transient failures. errno is valid immediately after a failed + // libc call on this thread. + let errno = std::io::Error::last_os_error().raw_os_error(); + return Err(match errno { + Some(libc::ENOENT) | Some(libc::ENOTDIR) => QueryError::PathMissing, + _ => QueryError::Other, + }); } // SAFETY: statvfs returned 0 (success), so the struct is fully initialised. let stat = unsafe { stat.assume_init() }; @@ -276,13 +387,14 @@ fn query_free_bytes(path: &Path) -> Option<(u64, u64)> { let free = stat.f_bavail as u64 * bsize; let total = stat.f_blocks as u64 * bsize; - Some((free, total)) + Ok((free, total)) } -/// Non-Unix stub: monitoring is a no-op, writes are never paused by disk. +/// Non-Unix stub: monitoring is a no-op, writes are never paused by disk and +/// the dir-lost latch never engages. #[cfg(not(unix))] -fn query_free_bytes(_path: &Path) -> Option<(u64, u64)> { - Some((u64::MAX, u64::MAX)) +fn query_free_bytes(_path: &Path) -> Result<(u64, u64), QueryError> { + Ok((u64::MAX, u64::MAX)) } // ── Tests ─────────────────────────────────────────────────────────────────── @@ -436,6 +548,56 @@ mod tests { assert_eq!(m.free_bytes(), u64::MAX, "initial free_bytes must be MAX"); } + // ── Dir-lost latch (issue #366) ───────────────────────────────────────── + + /// Fresh scratch directory unique to this test process + tag. + fn scratch_dir(tag: &str) -> std::path::PathBuf { + let d = std::env::temp_dir().join(format!("moon-diskmon-{}-{tag}", std::process::id())); + let _ = std::fs::remove_dir_all(&d); + std::fs::create_dir_all(&d).expect("create scratch dir"); + d + } + + #[test] + fn test_dir_lost_latches_pauses_and_self_heals() { + let d = scratch_dir("latch"); + // min_pct 0: the disk-FULL guard is relaxed (crash-test config), but + // the dir-lost latch must still work. + let m = DiskMonitor::new(0, &d).with_dir_guard(true); + + m.poll(); + assert!(!m.dir_lost(), "existing dir must not latch"); + assert!(!m.paused(), "existing dir must not pause writes"); + + std::fs::remove_dir_all(&d).expect("delete scratch dir"); + m.poll(); + assert!(m.dir_lost(), "deleted dir must latch dir_lost"); + assert!(m.paused(), "dir_lost must refuse writes via paused()"); + // Idempotent: repeated polls keep the latch without re-logging loud. + m.poll(); + assert!(m.dir_lost() && m.paused()); + + std::fs::create_dir_all(&d).expect("recreate scratch dir"); + m.poll(); + assert!(!m.dir_lost(), "reappeared dir must clear the latch"); + assert!(!m.paused(), "writes must resume after the dir reappears"); + + let _ = std::fs::remove_dir_all(&d); + } + + #[test] + fn test_dir_lost_disarmed_without_guard() { + let d = scratch_dir("noguard"); + // Default constructor: guard disarmed (in-memory server) — a missing + // path must keep the previous state, exactly the pre-#366 behavior. + let m = DiskMonitor::new(0, &d); + + std::fs::remove_dir_all(&d).expect("delete scratch dir"); + m.poll(); + assert!(!m.dir_lost(), "disarmed guard must never latch"); + assert!(!m.paused(), "disarmed guard must never pause writes"); + } + /// Smoke test: poll() on the real "/" path must not panic and must update /// free_bytes from the u64::MAX sentinel to some smaller value. #[test] diff --git a/src/shard/persistence_tick.rs b/src/shard/persistence_tick.rs index d6acb6a93..1a57bcc79 100644 --- a/src/shard/persistence_tick.rs +++ b/src/shard/persistence_tick.rs @@ -116,6 +116,17 @@ pub(crate) fn check_auto_save_trigger( let new_epoch = snapshot_trigger_rx.borrow(); if new_epoch > *last_snapshot_epoch && snapshot_state.is_none() { *last_snapshot_epoch = new_epoch; + // #366: consume the epoch but skip the save while the data directory + // is missing — the snapshot temp file can't be created, and retrying + // per save-point trigger just spams one doomed attempt per epoch. + if crate::shard::disk_monitor::is_dir_lost() { + tracing::warn!( + "Shard {}: auto-save epoch {} skipped — data directory missing", + shard_id, + new_epoch + ); + return; + } if let Some(dir) = persistence_dir { // When disk-offload is enabled, write snapshot to the offload shard directory // so v3 recovery can find it alongside WAL v3 segments and manifest. @@ -220,6 +231,13 @@ pub(crate) fn finalize_snapshot_error( pub(crate) fn flush_wal_v3_if_needed( wal_v3: &mut Option, ) { + // #366: while the data directory is missing, every flush that needs a + // path operation (segment rotation) fails — at the 1ms tick cadence that + // is a hot syscall+log error loop. The dir-lost latch already logged + // loudly and refused new writes; skip until the directory heals. + if crate::shard::disk_monitor::is_dir_lost() { + return; + } if let Some(wal) = wal_v3 { if let Err(e) = wal.flush_if_needed() { tracing::error!("WAL v3 flush failed: {}", e); @@ -1107,6 +1125,11 @@ pub(crate) fn maybe_force_checkpoint_on_wal_overflow( max_checkpoint_lag_ms: u64, graph_save: &mut dyn FnMut(u64) -> bool, ) -> bool { + // #366: no overflow scan while the data directory is missing — the + // per-second directory read would warn-loop forever. + if crate::shard::disk_monitor::is_dir_lost() { + return false; + } // Condition 1: total on-disk WAL exceeds the configured ceiling. let total_wal = match wal.stats() { Ok(s) => { @@ -1230,6 +1253,12 @@ pub(crate) fn handle_checkpoint_tick( tombstone_retain_secs: u64, graph_save: &mut dyn FnMut(u64) -> bool, ) -> bool { + // #366: checkpoints are pure file work (page pwrite, manifest rename, + // WAL recycle) — all doomed while the data directory is missing. Skip; + // the CheckpointManager simply resumes when the latch clears. + if crate::shard::disk_monitor::is_dir_lost() { + return false; + } match checkpoint_mgr.advance_tick() { CheckpointAction::Nothing => false, CheckpointAction::FlushPages(count) => { diff --git a/tests/dir_deleted_degraded.rs b/tests/dir_deleted_degraded.rs new file mode 100644 index 000000000..6bf5d91f1 --- /dev/null +++ b/tests/dir_deleted_degraded.rs @@ -0,0 +1,209 @@ +//! Issue #366: data-dir deleted under a running server must latch into a +//! degraded state — not spin the persistence tick in a syscall error loop. +//! +//! Observed in the wild (macOS, runtime-tokio, 4 shards, appendonly yes): an +//! orphaned server whose `--dir` was cleaned climbed to ~667% CPU while idle +//! (per-thread sys-time dominating user ~50:1) and wedged PING. The per-tick +//! WAL/checkpoint file operations retried on ENOENT with no failure latch, no +//! backoff, and no escalation. +//! +//! Contract pinned here (design-for-failure, mirrors the diskfull guard): +//! +//! 1. Within one disk-monitor poll (5s) of the dir vanishing, write commands +//! are refused with `MOONERR dirmissing` — like `MOONERR diskfull`. +//! 2. Reads and PING keep answering; the server must not wedge. +//! 3. The per-tick error loop stops: stderr must go quiet (no 1/s WAL-scan +//! warnings) and CPU must stay near idle baseline while latched. +//! 4. If the directory reappears, writes resume within one poll (self-heal). +//! +//! Run with: +//! cargo test --release --test dir_deleted_degraded + +#![cfg(unix)] +#![allow(clippy::unwrap_used)] + +mod common; + +use std::io::{Read, Write}; +use std::path::PathBuf; +use std::process::Child; +use std::time::{Duration, Instant}; + +/// Kill-on-drop guard so a mid-test panic can't orphan a server whose data +/// dir this test deletes — exactly the leak class that produced issue #366. +struct MoonGuard(Option); + +impl Drop for MoonGuard { + fn drop(&mut self) { + if let Some(mut child) = self.0.take() { + common::sigkill(&mut child); + } + } +} + +struct Conn(std::net::TcpStream); + +impl Conn { + fn open(port: u16) -> Self { + let deadline = Instant::now() + Duration::from_secs(30); + loop { + if let Ok(s) = std::net::TcpStream::connect(("127.0.0.1", port)) { + s.set_read_timeout(Some(Duration::from_secs(5))).unwrap(); + return Self(s); + } + assert!(Instant::now() < deadline, "server never accepted"); + std::thread::sleep(Duration::from_millis(50)); + } + } + + /// Send an inline command, return the raw reply up to one read() chunk. + fn cmd(&mut self, line: &str) -> String { + self.0 + .write_all(format!("{line}\r\n").as_bytes()) + .expect("write"); + let mut buf = [0u8; 64 * 1024]; + let n = self.0.read(&mut buf).expect("read"); + String::from_utf8_lossy(&buf[..n]).into_owned() + } +} + +/// Total process CPU time (user + system) in milliseconds, via `ps` so the +/// same probe works on macOS and Linux. `cputime` renders as `[DD-]HH:MM:SS` +/// or `MM:SS.cc` depending on platform — parse colon fields generically. +fn process_cpu_ms(pid: u32) -> u64 { + let out = std::process::Command::new("ps") + .args(["-o", "cputime=", "-p", &pid.to_string()]) + .output() + .expect("ps"); + let s = String::from_utf8_lossy(&out.stdout); + let s = s.trim().replace('-', ":"); // DD-HH:MM:SS → DD:HH:MM:SS + let mut secs = 0.0f64; + for field in s.split(':') { + secs = secs * 60.0 + field.parse::().unwrap_or(0.0); + } + (secs * 1000.0) as u64 +} + +#[test] +fn dir_deleted_latches_write_refusal_stays_responsive_and_self_heals() { + let data_dir: PathBuf = std::env::temp_dir().join(format!( + "moon-dir-deleted-{}-{}", + std::process::id(), + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_nanos() + )); + std::fs::create_dir_all(&data_dir).unwrap(); + let log_path = data_dir.with_extension("stderr.log"); + + let bin = common::find_moon_binary(); + let (child, port) = common::spawn_listening(|port| { + let log = std::fs::File::create(&log_path).expect("create stderr log"); + std::process::Command::new(&bin) + .args([ + "--port", + &port.to_string(), + "--dir", + &data_dir.to_string_lossy(), + "--shards", + "4", + "--appendonly", + "yes", + // The dir-lost latch must work even with the disk-FULL guard + // relaxed (the standard crash-test configuration). + "--disk-free-min-pct", + "0", + ]) + .stdout(std::process::Stdio::null()) + .stderr(std::process::Stdio::from(log)) + .spawn() + .expect("spawn moon") + }); + let pid = child.id(); + let _guard = MoonGuard(Some(child)); + + let mut c = Conn::open(port); + assert!( + c.cmd("SET before-loss v1").starts_with("+OK"), + "healthy server must accept writes" + ); + + // ── The incident trigger: data dir vanishes under the running server ── + std::fs::remove_dir_all(&data_dir).unwrap(); + + // (1) Writes must latch to MOONERR dirmissing within ~one 5s poll. + let refusal = wait_for(Duration::from_secs(20), || { + let r = c.cmd("SET after-loss v2"); + r.contains("MOONERR dirmissing").then_some(r) + }); + assert!( + refusal.is_some(), + "writes were still accepted 20s after the data dir was deleted — \ + the dir-lost latch never engaged" + ); + + // (2) Reads and PING keep answering while latched. + assert!( + c.cmd("PING").starts_with("+PONG"), + "PING must stay responsive while persistence is degraded" + ); + assert!( + c.cmd("GET before-loss").contains("v1"), + "reads must keep serving in-memory data while latched" + ); + + // (3a) The per-tick error loop must stop: stderr goes quiet once latched. + // Pre-fix every shard warns once per second (WAL overflow scan), so + // 5s accumulates ≥20 lines; allow a small tail for in-flight lines. + let len_before = std::fs::metadata(&log_path).map(|m| m.len()).unwrap_or(0); + std::thread::sleep(Duration::from_secs(5)); + let len_after = std::fs::metadata(&log_path).map(|m| m.len()).unwrap_or(0); + assert!( + len_after - len_before < 400, + "stderr grew {} bytes in 5s while dir-lost latched — the persistence \ + tick is still error-looping", + len_after - len_before + ); + + // (3b) CPU stays near idle baseline (the 667%-spin regression pin). The + // threshold is deliberately generous — 40% of ONE core across the + // whole 4-shard process — so a loaded CI box cannot flake it, while + // the incident state (≈6+ cores) fails it by an order of magnitude. + let cpu_before = process_cpu_ms(pid); + std::thread::sleep(Duration::from_secs(5)); + let cpu_delta = process_cpu_ms(pid).saturating_sub(cpu_before); + assert!( + cpu_delta < 2_000, + "process burned {cpu_delta}ms CPU in 5s wall while idle+latched — \ + persistence tick is spinning" + ); + + // (4) Self-heal: restore the directory, writes resume within one poll. + std::fs::create_dir_all(&data_dir).unwrap(); + let resumed = wait_for(Duration::from_secs(20), || { + c.cmd("SET after-heal v3").starts_with("+OK").then_some(()) + }); + assert!( + resumed.is_some(), + "writes never resumed within 20s of the data dir reappearing" + ); + assert!(c.cmd("GET after-heal").contains("v3")); + + let _ = std::fs::remove_dir_all(&data_dir); + let _ = std::fs::remove_file(&log_path); +} + +/// Poll `probe` every 500ms until it returns Some or `timeout` elapses. +fn wait_for(timeout: Duration, mut probe: impl FnMut() -> Option) -> Option { + let deadline = Instant::now() + timeout; + loop { + if let Some(v) = probe() { + return Some(v); + } + if Instant::now() >= deadline { + return None; + } + std::thread::sleep(Duration::from_millis(500)); + } +} From 5b01285e26fc1a815ccdd8a29a0c9f6cb134da57 Mon Sep 17 00:00:00 2001 From: Tin Dang Date: Fri, 17 Jul 2026 13:15:46 +0700 Subject: [PATCH 2/2] =?UTF-8?q?fix(persistence):=20close=20the=20shallow-h?= =?UTF-8?q?eal=20hole=20+=20flush=20backoff=20=E2=80=94=20#369=20review=20?= =?UTF-8?q?round=20(B1,=20S1-class,=20S3,=20N1,=20N2)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adversarial review of the dir-lost latch found one blocker and three gaps; all fixed here, two skipped with reasons (PR thread). B1 (blocker) — shallow heal re-armed the original bug UNLATCHED: the heal check is statvfs on the top-level --dir only, but WalWriterV3 creates /shard-N/wal-v3/ exactly once at startup and open_new_segment() only `create(true)`s the LEAF file. An operator answering "MOONERR dirmissing" with `mkdir -p ` cleared the latch while the nested WAL dir stayed missing — the first post-heal segment rotation then failed ENOENT on every 1ms tick, reproducing issue #366 with the latch already released. Two independent layers: - open_new_segment() now runs fs::create_dir_all(&wal_dir) first — rotation self-repairs a vanished parent (rotations are rare; one mkdir syscall is free). Unit-pinned by test_rotation_self_repairs_deleted_wal_dir (tiny 256-byte segment, delete the wal dir under an open writer, force rotation, assert it recreates the dir and lands the new segment). - flush_wal_v3_if_needed() arms a 1s backoff on ANY flush failure (flush_backing_off/note_flush_failure/clear_flush_backoff on WalWriterV3) — no future failure class (EACCES, quota, …) can loop at tick cadence. Healthy-path cost: one Option discriminant check. Unit-pinned by test_flush_backoff_arms_and_clears. S3 — disk-offload without AOF/save never armed the guard: spill and the vector warm tier write to the data dir even with appendonly=no. init_global now arms on persistence OR disk_offload_enabled(), and check_warm_transitions (10s cadence segment writer) gates on is_dir_lost() like the other tick functions. N1 — module doc no longer claims clean-unmount coverage (mountpoint dir survives an unmount; catching it needs an st_dev comparison — follow-up material). N2 — QueryError variants are only constructed under cfg(unix); added cfg_attr(not(unix), allow(dead_code)) so the main-push-only Windows Check job can't fail on the dead-variant lint. Skipped with reason: - S2 (explicit BGSAVE not fast-failed while latched): the existing finalize-error path already fails the one attempt cleanly and keeps the BGSAVE fan-in counter correct; a fast-fail reply would need its own counter accounting and risks the stuck-in-progress bug class. - S1 residual (expiry/eviction WAL appends growing the buffer during an outage): append-time flush errors predate this PR (same class as diskfull) and are now bounded by the 1s backoff; full write-path gating of internal ops is follow-up scope. author: Tin Dang --- CHANGELOG.md | 7 ++- src/main.rs | 5 +- src/persistence/wal_v3/segment.rs | 95 +++++++++++++++++++++++++++++++ src/shard/disk_monitor.rs | 7 ++- src/shard/persistence_tick.rs | 16 +++++- 5 files changed, 126 insertions(+), 4 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 6be269804..3135fe7b0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -21,7 +21,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 while latched. The latch self-heals: the next poll that sees the directory again resumes writes (with an explicit warning that data accepted since the deletion is not guaranteed durable until restart). Only engages when - persistence is configured, and works with `--disk-free-min-pct 0`. + persistence or disk-offload is configured, and works with + `--disk-free-min-pct 0`. Two hardening layers close the shallow-heal hole + (adversarial review): WAL segment rotation re-creates a missing parent + directory (`mkdir -p ` after an incident no longer leaves the nested + `wal-v3/` dir dead), and any tick-flush failure arms a 1s retry backoff so + no flush error class can ever loop at the 1ms tick cadence again. ### Performance - **Disk-offload: shard event loop no longer pays manifest-commit fsyncs diff --git a/src/main.rs b/src/main.rs index 1352bdd3a..a0746849d 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1030,10 +1030,13 @@ fn main() -> anyhow::Result<()> { // persistence tick if the data directory is deleted out from under it. { let monitor_path = persistence_dir.as_deref().unwrap_or(&config.dir); + // Arm the dir-lost latch for any config that writes to the data dir: + // AOF/save persistence, OR disk-offload (spill + vector warm tier + // write there even with appendonly=no). moon::shard::disk_monitor::init_global( config.disk_free_min_pct, monitor_path, - persistence_dir.is_some(), + persistence_dir.is_some() || config.disk_offload_enabled(), ); } diff --git a/src/persistence/wal_v3/segment.rs b/src/persistence/wal_v3/segment.rs index 2795fc7f0..2ab726587 100644 --- a/src/persistence/wal_v3/segment.rs +++ b/src/persistence/wal_v3/segment.rs @@ -353,6 +353,10 @@ pub struct WalWriterV3 { /// Set when agent spawn failed once — never retried (inline fsync /// fallback, logged once). sync_agent_unavailable: bool, + /// Backoff deadline after a tick-flush failure (#366): the 1ms tick must + /// never retry a failing flush at tick cadence (syscall+log error loop). + /// `None` while healthy. + flush_backoff_until: Option, } /// Bound on every blocking durability wait (checkpoint ordering gates, @@ -393,6 +397,7 @@ impl WalWriterV3 { max_wal_bytes: DEFAULT_MAX_WAL_BYTES, sync_agent: None, sync_agent_unavailable: false, + flush_backoff_until: None, }; writer.open_new_segment()?; @@ -555,6 +560,31 @@ impl WalWriterV3 { } } + /// True while the tick-flush is backing off after a failure (#366). + /// + /// Healthy path cost: one `Option` discriminant check — `Instant::now()` + /// only runs while a backoff is actually armed. + #[inline] + pub fn flush_backing_off(&self) -> bool { + self.flush_backoff_until + .is_some_and(|until| std::time::Instant::now() < until) + } + + /// Arm a 1s tick-flush backoff after a flush failure (#366): the 1ms + /// tick must never retry a failing flush at tick cadence. + pub fn note_flush_failure(&mut self) { + self.flush_backoff_until = + Some(std::time::Instant::now() + std::time::Duration::from_secs(1)); + } + + /// Clear the backoff after a successful flush. + #[inline] + pub fn clear_flush_backoff(&mut self) { + if self.flush_backoff_until.is_some() { + self.flush_backoff_until = None; + } + } + /// Fsync without writing — call after `flush_write` / `flush_if_needed` /// to make all previously written records durable. pub fn sync_data(&mut self) -> std::io::Result<()> { @@ -808,6 +838,11 @@ impl WalWriterV3 { /// Open a new segment file and write its 64-byte header. fn open_new_segment(&mut self) -> std::io::Result<()> { + // Self-repair a vanished parent (#366 heal path): `create(true)` + // only creates the LEAF file, so a shallow `mkdir -p ` after a + // dir-lost incident would leave this rotation failing ENOENT on + // every retry. Rotations are rare — one extra mkdir syscall is free. + fs::create_dir_all(&self.wal_dir)?; let path = WalSegment::segment_path(&self.wal_dir, self.current_sequence); let mut file = OpenOptions::new() .create(true) @@ -1084,6 +1119,66 @@ mod tests { assert_eq!(WalSegment::segment_name(0), "000000000000.wal"); } + /// #366 review blocker B1: after a dir-lost incident an operator does a + /// shallow `mkdir -p ` — the nested `wal-v3/` directory is still + /// gone. Segment rotation must self-repair the parent instead of failing + /// ENOENT on every retry (which would re-create the tick error loop the + /// dir-lost latch exists to kill, this time UNLATCHED because the + /// top-level path exists again). + #[test] + fn test_rotation_self_repairs_deleted_wal_dir() { + let tmp = tempfile::tempdir().unwrap(); + let wal_dir = tmp.path().join("wal"); + // Tiny segment so a handful of appends forces rotation. + let mut writer = WalWriterV3::new(0, &wal_dir, 256).unwrap(); + + // The shallow-heal state: the whole wal dir vanishes while the + // writer's current-segment fd stays open (unlinked inode). + fs::remove_dir_all(&wal_dir).unwrap(); + + // Push enough bytes through append+flush to cross segment_size and + // force rotate_segment() -> open_new_segment(). + for i in 0..16u32 { + writer.append( + WalRecordType::Command, + format!("SET key{i} {}", "v".repeat(64)).as_bytes(), + ); + } + writer + .flush_write() + .expect("rotation must recreate the missing wal dir, not fail ENOENT"); + + assert!(wal_dir.exists(), "wal dir must have been re-created"); + assert!( + writer.current_segment_sequence() > 1, + "rotation must actually have happened for this test to prove anything" + ); + // And the new segment file is a real on-disk file. + let seg = WalSegment::segment_path(&wal_dir, writer.current_segment_sequence()); + assert!(seg.exists(), "post-rotation segment must exist on disk"); + } + + /// #366: the tick flush must never retry a failing flush at 1ms tick + /// cadence — a failure arms a 1s backoff, success clears it. + #[test] + fn test_flush_backoff_arms_and_clears() { + let tmp = tempfile::tempdir().unwrap(); + let wal_dir = tmp.path().join("wal"); + let mut writer = WalWriterV3::new(0, &wal_dir, DEFAULT_SEGMENT_SIZE).unwrap(); + + assert!( + !writer.flush_backing_off(), + "healthy writer never backs off" + ); + writer.note_flush_failure(); + assert!(writer.flush_backing_off(), "failure must arm the backoff"); + writer.clear_flush_backoff(); + assert!( + !writer.flush_backing_off(), + "success must clear the backoff" + ); + } + #[test] fn test_writer_creates_segment() { let tmp = tempfile::tempdir().unwrap(); diff --git a/src/shard/disk_monitor.rs b/src/shard/disk_monitor.rs index 50998e609..ef1085d33 100644 --- a/src/shard/disk_monitor.rs +++ b/src/shard/disk_monitor.rs @@ -20,7 +20,9 @@ //! ## Dir-lost latch (issue #366) //! //! When the monitored path itself vanishes (`statvfs` → `ENOENT`/`ENOTDIR`: -//! operator `rm -rf`, tmp-cleaner sweep, mount disappearing), a server with +//! operator `rm -rf`, tmp-cleaner sweep; NOT a clean unmount — the mountpoint +//! directory usually survives that, and catching it would need an `st_dev` +//! comparison), a server with //! persistence configured must NOT keep retrying its per-tick WAL/checkpoint //! file operations — that is a 1ms syscall error loop that burned ~100% CPU //! per shard thread in the wild. Instead `poll` latches `dir_lost`: @@ -336,6 +338,9 @@ pub fn global_free_bytes() -> u64 { /// caller distinguishes: it deterministically means the monitored path no /// longer exists (the issue #366 trigger), unlike transient I/O errors. #[derive(Debug, Clone, Copy, PartialEq, Eq)] +// Non-unix: the stub `query_free_bytes` never fails, so the variants are +// never constructed there — silence the dead-variant lint for that cfg only. +#[cfg_attr(not(unix), allow(dead_code))] enum QueryError { /// `ENOENT` / `ENOTDIR` — the monitored path is gone. PathMissing, diff --git a/src/shard/persistence_tick.rs b/src/shard/persistence_tick.rs index 1a57bcc79..b2ae3bfa2 100644 --- a/src/shard/persistence_tick.rs +++ b/src/shard/persistence_tick.rs @@ -239,8 +239,17 @@ pub(crate) fn flush_wal_v3_if_needed( return; } if let Some(wal) = wal_v3 { + // Belt-and-suspenders for failure classes the latch can't see + // (EACCES, a shallow heal missing nested dirs, …): a failed flush + // arms a 1s backoff so no flush error can ever loop at tick cadence. + if wal.flush_backing_off() { + return; + } if let Err(e) = wal.flush_if_needed() { - tracing::error!("WAL v3 flush failed: {}", e); + wal.note_flush_failure(); + tracing::error!("WAL v3 flush failed (retrying in 1s): {}", e); + } else { + wal.clear_flush_backoff(); } } } @@ -268,6 +277,11 @@ pub(crate) fn check_warm_transitions( shard_id: usize, wal: &mut Option, ) { + // #366: warm-tier transitions write segment files under the offload dir + // — doomed while the data directory is missing; skip until it heals. + if crate::shard::disk_monitor::is_dir_lost() { + return; + } let count = vector_store.try_warm_transitions_all_idle( shard_dir, manifest,