From e9936ed2da7d3c780e7bda45a02fbc399d6a105a Mon Sep 17 00:00:00 2001 From: Tin Dang Date: Tue, 7 Jul 2026 10:34:02 +0700 Subject: [PATCH 1/2] fix(persistence): TopLevel-monoio AOF writer honors EverySec bound when idle MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The TopLevel monoio AOF writer blocked on an UNTIMED rx.recv(), with its EverySec deadline check living only inside the batch-commit path. A batch written under appendfsync everysec gets no per-batch fsync, so if the client stopped writing right after a burst, the buffered bytes only became durable when the NEXT message happened to arrive — the 1s fsync bound was deferred indefinitely while idle. Exposure is host-crash-only: the per-batch flush() already reaches the kernel page cache, so a plain process kill (and therefore any kill-9 crash test) loses nothing — which is exactly why this can't be red/green tested from userspace and is instead validated by mirroring the audited PerShard pattern plus full AOF regression reruns. The loop now matches the PerShard writers line-for-line: - bounded rx.recv_timeout() on the wave-5 IdleWait ladder (50ms -> 250ms -> 1s; reset on message, escalate on empty timeout) - mark_pending() after an everysec batch is buffered without its fsync, pinning the poll at the 50ms floor until the fsync lands - the EverySec deadline check moves from inside the batch path to the END of the loop, so it runs on message AND timeout iterations — timeout wake-ups are what restore the ~1s bound while idle - clear_pending() when the proactive fsync succeeds, letting the idle cadence escalate again Also updates the IdleWait struct docs, which claimed "TopLevel monoio blocks on an untimed rx.recv() and needs none of this" — that was the documented follow-up from the wave-5 review, now closed. Validation: clippy clean on both feature sets; AOF regression suites rerun green under default (monoio) features — persistence lib tests, wal_group_commit, coordinator_local_leg_durability (incl. ignored), crash_matrix_per_shard_aof, aof_toplevel_multishard_refusal, aof_fsync_err_subscribe_ordering. Follow-up from the RSS/CPU wave-5 review (PR #232). author: Tin Dang --- CHANGELOG.md | 29 ++++ src/persistence/aof/writer_task.rs | 253 ++++++++++++++++------------- 2 files changed, 172 insertions(+), 110 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index e6c6e2327..0841f74c1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,35 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Fixed — TopLevel-monoio AOF writer: EverySec fsync deferred indefinitely when idle (PR #TBD) + +- **`src/persistence/aof/writer_task.rs`**: the TopLevel monoio AOF writer + blocked on an **untimed** `rx.recv()`, with its EverySec deadline check + living only inside the batch-commit path. A batch written under + `appendfsync everysec` gets no per-batch fsync, so if the client stopped + writing right after a burst, the buffered bytes only became durable when + the NEXT message happened to arrive — the 1s fsync bound was deferred + indefinitely while idle. Exposure is host-crash-only (the per-batch + `flush()` already reaches the kernel page cache, so a plain process kill + loses nothing), which is exactly the window EverySec exists to bound. + The loop now mirrors the audited PerShard writers: bounded + `recv_timeout` on the wave-5 `IdleWait` ladder (50ms → 250ms → 1s, + pinned at the floor while a batch awaits its fsync) plus an end-of-loop + proactive fsync that runs on message AND timeout iterations. This + supersedes wave 5's "TopLevel monoio needs none of this" note — it was + the one writer loop left without the ≤ ~1s idle durability bound. +- **`tests/crash_matrix_per_shard_aof.rs` harness hardening** (found while + validating the above): (1) `redis_set` asserted only redis-cli's exit + status, which is 0 even for server ERROR replies — a tripped diskfull + guard (host <5% free) turned every SET into a silent no-op and surfaced + as a bogus "200 keys missing after recovery"; the helper now pins the + reply to `+OK`. (2) The three server-spawning tests split-brain when run + in parallel: `unique_port()` hands out OS-sequential ephemeral ports, the + other tests offset +1/+2, and SO_REUSEPORT lets two tests bind the SAME + port without an error — one test's redis-cli traffic lands on another + test's server (rotating total-loss false alarms). A shared mutex now + serializes them. + ### Fixed — RSS/CPU remediation wave 5 (PR #TBD) - **Item A — mmap the exact-rerank f16 sidecar on segment reload** diff --git a/src/persistence/aof/writer_task.rs b/src/persistence/aof/writer_task.rs index 8d9176d0f..cc2452257 100644 --- a/src/persistence/aof/writer_task.rs +++ b/src/persistence/aof/writer_task.rs @@ -25,9 +25,8 @@ use super::group_commit::{GroupCommitSink, commit_group_commit_batch}; /// Idle-adaptive wake cadence for a background AOF writer's channel poll /// (RSS/CPU wave 5, item B). /// -/// The steady-state writer loops (PerShard monoio/tokio, TopLevel tokio — -/// TopLevel monoio blocks on an untimed `rx.recv()` and needs none of this) -/// poll their channel with a bounded timeout so the EverySec proactive-fsync +/// All four steady-state writer loops (PerShard monoio/tokio, TopLevel +/// monoio/tokio) poll their channel with a bounded timeout so the EverySec proactive-fsync /// deadline check that follows every wake still fires when no new Appends /// ever arrive. A FIXED cadence forever (previously 50ms monoio / 200ms /// tokio) means an idle server's AOF writer thread wakes 5-20 times a @@ -356,13 +355,34 @@ pub async fn aof_writer_task( // Read once at task startup; zero cost in production (var absent). let fail_fsync_for_test = std::env::var("MOON_TEST_AOF_FSYNC_FAIL").as_deref() == Ok("1"); + // Idle-adaptive channel-poll wake cadence (RSS/CPU wave 5, item B) — + // see `IdleWait` docs near the top of this file. This loop used to + // block on an UNTIMED `rx.recv()`, with the EverySec deadline check + // only inside the batch path: a batch buffered without a per-batch + // fsync only became durable when the NEXT message happened to + // arrive, so idle-after-a-burst deferred the fsync indefinitely + // (host-crash exposure only — the per-batch flush already reaches + // the kernel page cache, so a plain process kill loses nothing). + // The bounded recv + end-of-loop proactive fsync below restore the + // ~1s EverySec bound exactly like the PerShard writers. + let mut idle_wait = IdleWait::new(); + loop { - // Group commit: block for one message, then opportunistically drain - // whatever else is already queued into a bounded batch so a single - // fsync makes the whole batch durable (TopLevel = plain RESP bytes). - let first = match rx.recv() { - Ok(m) => m, - Err(_) => { + // Group commit: wait (bounded) for one message, then + // opportunistically drain whatever else is already queued into a + // bounded batch so a single fsync makes the whole batch durable + // (TopLevel = plain RESP bytes). On timeout, fall through (None) + // to the EverySec proactive fsync at the end of the loop. + let first = match rx.recv_timeout(idle_wait.current()) { + Ok(m) => { + idle_wait.on_message(); + Some(m) + } + Err(flume::RecvTimeoutError::Timeout) => { + idle_wait.on_timeout(); + None + } + Err(flume::RecvTimeoutError::Disconnected) => { // Channel disconnected — final sync + shut down. if !write_error { if let Err(e) = file.flush().and_then(|_| file.sync_data()) { @@ -373,123 +393,136 @@ pub async fn aof_writer_task( break; } }; - let mut batch = collect_group_commit_batch( - first, - || rx.try_recv().ok(), - AOF_GROUP_COMMIT_MAX_BATCH, - AOF_GROUP_COMMIT_MAX_BYTES, - ); - // -- commit the data batch (one fsync under Always; deadline under everysec) -- - if !batch.data.is_empty() { - if write_error { - // Persistent I/O failure latched: drop appends and fail every - // AppendSync waiter — never a false durability claim. - let _ = group_commit::ack_batch(&mut batch, BatchAck::WriteFailed); - } else { - let do_fsync = matches!(fsync, FsyncPolicy::Always); - let mut sink = FileGroupSink { - file: &mut file, - fail_sync: fail_fsync_for_test, - }; - let outcome = commit_group_commit_batch(&mut sink, &mut batch, do_fsync); - if outcome.write_failed { - // A torn write may leave a partial record — latch so no - // further bytes are appended after the tear. - error!( - "AOF batch write failed (seq {}). Persistence degraded.", - manifest.seq - ); - write_error = true; - } - // EverySec: the batch was written but not per-batch-fsynced - // (do_fsync=false; there are no AppendSync waiters under - // everysec). Honor the 1s deadline exactly as the old - // per-Append path did. - if fsync == FsyncPolicy::EverySec - && !write_error - && last_fsync.elapsed() >= std::time::Duration::from_secs(1) - { - let t = Instant::now(); - if let Err(e) = file.flush().and_then(|_| file.sync_data()) { - error!("AOF sync failed (seq {}, everysec): {}", manifest.seq, e); - // Non-fatal for everysec: retry next interval - } else { - crate::admin::metrics_setup::record_aof_fsync( - t.elapsed().as_micros() as u64 + if let Some(first) = first { + let mut batch = collect_group_commit_batch( + first, + || rx.try_recv().ok(), + AOF_GROUP_COMMIT_MAX_BATCH, + AOF_GROUP_COMMIT_MAX_BYTES, + ); + + // -- commit the data batch (one fsync under Always; deadline under everysec) -- + if !batch.data.is_empty() { + if write_error { + // Persistent I/O failure latched: drop appends and fail every + // AppendSync waiter — never a false durability claim. + let _ = group_commit::ack_batch(&mut batch, BatchAck::WriteFailed); + } else { + let do_fsync = matches!(fsync, FsyncPolicy::Always); + let mut sink = FileGroupSink { + file: &mut file, + fail_sync: fail_fsync_for_test, + }; + let outcome = commit_group_commit_batch(&mut sink, &mut batch, do_fsync); + if outcome.write_failed { + // A torn write may leave a partial record — latch so no + // further bytes are appended after the tear. + error!( + "AOF batch write failed (seq {}). Persistence degraded.", + manifest.seq ); - last_fsync = Instant::now(); + write_error = true; + } + // EverySec: the batch was written but not per-batch-fsynced + // (do_fsync=false; there are no AppendSync waiters under + // everysec). The end-of-loop proactive fsync makes it + // durable within the 1s bound — pin the idle wait at its + // fast floor until that fsync clears it. + if fsync == FsyncPolicy::EverySec && !write_error { + idle_wait.mark_pending(); } } } - } - // -- handle the control message that ended the drain (if any) -- - // A control message is NEVER absorbed into the batch: the batch above - // is already committed before the control message is handled - // (batch_straddles_control is structurally impossible). - match batch.deferred_control { - None => {} - Some(AofMessage::Shutdown) => { - if !write_error { - if let Err(e) = file.flush().and_then(|_| file.sync_data()) { - error!("AOF final sync failed (seq {}): {}", manifest.seq, e); + // -- handle the control message that ended the drain (if any) -- + // A control message is NEVER absorbed into the batch: the batch above + // is already committed before the control message is handled + // (batch_straddles_control is structurally impossible). + match batch.deferred_control { + None => {} + Some(AofMessage::Shutdown) => { + if !write_error { + if let Err(e) = file.flush().and_then(|_| file.sync_data()) { + error!("AOF final sync failed (seq {}): {}", manifest.seq, e); + } } + info!("AOF writer shutting down (monoio, seq {})", manifest.seq); + break; } - info!("AOF writer shutting down (monoio, seq {})", manifest.seq); - break; - } - Some(AofMessage::Rewrite(db)) => { - if !write_error { - if let Err(e) = file.flush().and_then(|_| file.sync_data()) { - error!("AOF pre-rewrite sync failed (seq {}): {}", manifest.seq, e); + Some(AofMessage::Rewrite(db)) => { + if !write_error { + if let Err(e) = file.flush().and_then(|_| file.sync_data()) { + error!("AOF pre-rewrite sync failed (seq {}): {}", manifest.seq, e); + } } - } - match do_rewrite_single(&db, &mut manifest, &mut file, &rx) { - Ok(()) => { - write_error = false; // Reset on successful rewrite + match do_rewrite_single(&db, &mut manifest, &mut file, &rx) { + Ok(()) => { + write_error = false; // Reset on successful rewrite + } + Err(e) => error!("AOF rewrite failed (seq {}): {}", manifest.seq, e), } - Err(e) => error!("AOF rewrite failed (seq {}): {}", manifest.seq, e), + crate::command::persistence::AOF_REWRITE_IN_PROGRESS + .store(false, std::sync::atomic::Ordering::SeqCst); } - crate::command::persistence::AOF_REWRITE_IN_PROGRESS - .store(false, std::sync::atomic::Ordering::SeqCst); - } - Some(AofMessage::RewriteSharded(shard_dbs)) => { - if !write_error { - if let Err(e) = file.flush().and_then(|_| file.sync_data()) { - error!("AOF pre-rewrite sync failed (seq {}): {}", manifest.seq, e); + Some(AofMessage::RewriteSharded(shard_dbs)) => { + if !write_error { + if let Err(e) = file.flush().and_then(|_| file.sync_data()) { + error!("AOF pre-rewrite sync failed (seq {}): {}", manifest.seq, e); + } } - } - // C4 TopLevel cooperative fold: pass the wired fold channels - // (producer + notifier for shard 0) so do_rewrite_sharded can - // use the AofFold SPSC protocol instead of the deleted RwLock - // path. `fold_channels` is `None` only if main.rs failed to - // wire them at startup (Arc::get_mut race — logged at boot). - match do_rewrite_sharded( - &shard_dbs, - &mut manifest, - &mut file, - &rx, - fold_channels.as_ref(), - ) { - Ok(()) => { - write_error = false; + // C4 TopLevel cooperative fold: pass the wired fold channels + // (producer + notifier for shard 0) so do_rewrite_sharded can + // use the AofFold SPSC protocol instead of the deleted RwLock + // path. `fold_channels` is `None` only if main.rs failed to + // wire them at startup (Arc::get_mut race — logged at boot). + match do_rewrite_sharded( + &shard_dbs, + &mut manifest, + &mut file, + &rx, + fold_channels.as_ref(), + ) { + Ok(()) => { + write_error = false; + } + Err(e) => error!("AOF rewrite failed (seq {}): {}", manifest.seq, e), } - Err(e) => error!("AOF rewrite failed (seq {}): {}", manifest.seq, e), + crate::command::persistence::AOF_REWRITE_IN_PROGRESS + .store(false, std::sync::atomic::Ordering::SeqCst); } - crate::command::persistence::AOF_REWRITE_IN_PROGRESS - .store(false, std::sync::atomic::Ordering::SeqCst); + // [F6] A TopLevel writer never owns per-shard files; receiving + // RewritePerShard means a routing bug. Self-abort so the + // coordinator's countdown completes and the flag clears. + Some(AofMessage::RewritePerShard { coord, .. }) => { + warn!( + "AOF TopLevel writer received RewritePerShard — routing bug; aborting" + ); + coord.mark_failed(); + coord.shard_done(); + } + // collect_group_commit_batch only ever defers a control message. + Some(_) => {} } - // [F6] A TopLevel writer never owns per-shard files; receiving - // RewritePerShard means a routing bug. Self-abort so the - // coordinator's countdown completes and the flag clears. - Some(AofMessage::RewritePerShard { coord, .. }) => { - warn!("AOF TopLevel writer received RewritePerShard — routing bug; aborting"); - coord.mark_failed(); - coord.shard_done(); + } + + // EverySec proactive fsync — runs after every loop iteration + // (message processed OR timeout); the only path that guarantees + // the ~1s durability bound when no further messages arrive after + // a buffered batch (idle-after-a-burst). + if fsync == FsyncPolicy::EverySec + && !write_error + && last_fsync.elapsed() >= std::time::Duration::from_secs(1) + { + let t = Instant::now(); + if let Err(e) = file.flush().and_then(|_| file.sync_data()) { + error!("AOF sync failed (seq {}, everysec): {}", manifest.seq, e); + // Non-fatal for everysec: retry next interval + } else { + crate::admin::metrics_setup::record_aof_fsync(t.elapsed().as_micros() as u64); + last_fsync = Instant::now(); + idle_wait.clear_pending(); } - // collect_group_commit_batch only ever defers a control message. - Some(_) => {} } } return; From fa0488447c56163049556315737062fba0c4357e Mon Sep 17 00:00:00 2001 From: Tin Dang Date: Tue, 7 Jul 2026 10:34:02 +0700 Subject: [PATCH 2/2] =?UTF-8?q?test(persistence):=20crash-matrix=20harness?= =?UTF-8?q?=20=E2=80=94=20pin=20SET=20reply,=20serialize=20server=20tests?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two pre-existing harness defects in crash_matrix_per_shard_aof.rs surfaced while validating the TopLevel-monoio idle-fsync fix; both convert environmental conditions into bogus total-data-loss failures: 1. redis_set asserted only redis-cli's EXIT STATUS, which is 0 even for server ERROR replies. With the host disk under Moon's 5% diskfull guard, every SET replied "MOONERR diskfull: writes paused" yet the test sailed on and later reported "200 missing, 0 mismatched" — pointing at the recovery path instead of the guard. The helper now pins the reply to +OK and names the likely culprit (and the MOON_DISK_FREE_MIN_PCT=0 escape hatch) in the assert message. 2. The three server-spawning tests SPLIT-BRAIN when run in parallel (the libtest default): unique_port() hands out OS-sequential ephemeral ports, the other two tests offset by +1/+2, and moon's per-shard SO_REUSEPORT listeners bind an already-taken port WITHOUT an error — so one test's redis-cli traffic silently lands on another test's server, observed as rotating "200 missing" / "LLEN 0 expected 20" false alarms (which test fails depends on OS port allocation order). A shared static Mutex now serializes the three tests; costs ~seconds, kills the whole interference class. Validated: suite green under default parallel threads AND --test-threads=1, with MOON_DISK_FREE_MIN_PCT=0 exported (host disk currently 96% full — below the guard). author: Tin Dang --- tests/crash_matrix_per_shard_aof.rs | 33 +++++++++++++++++++++++++++++ 1 file changed, 33 insertions(+) diff --git a/tests/crash_matrix_per_shard_aof.rs b/tests/crash_matrix_per_shard_aof.rs index f806a9091..fcffbc0ce 100644 --- a/tests/crash_matrix_per_shard_aof.rs +++ b/tests/crash_matrix_per_shard_aof.rs @@ -38,6 +38,22 @@ use std::time::Duration; const KEY_COUNT: usize = 200; +/// Serializes the server-spawning tests in this binary. Run in parallel they +/// intermittently SPLIT-BRAIN: `unique_port()` hands out OS-sequential +/// ephemeral ports and the other tests offset by +1/+2, so two concurrently +/// starting tests can land on the SAME port — and moon's per-shard +/// SO_REUSEPORT listeners bind it without an error, silently splitting one +/// test's redis-cli traffic across another test's server (observed as "200 +/// missing" total-loss false alarms). A shared lock costs ~seconds and kills +/// the whole class. +static SERVER_TEST_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(()); + +fn serialize_server_test() -> std::sync::MutexGuard<'static, ()> { + SERVER_TEST_LOCK + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) +} + fn unique_port() -> u16 { // Ask the OS to assign an available ephemeral port by binding to 0. // The socket is immediately dropped after reading the port — there is a @@ -119,6 +135,20 @@ fn redis_set(port: u16, key: &str, value: &str) { value, String::from_utf8_lossy(&out.stderr) ); + // redis-cli exits 0 even for server ERROR replies (e.g. "MOONERR + // diskfull: writes paused"), so exit status alone silently converts a + // rejected write into a "key missing after recovery" false alarm at the + // end of the test. Pin the actual reply. + let reply = String::from_utf8_lossy(&out.stdout); + assert_eq!( + reply.trim(), + "OK", + "redis-cli SET {} {} did not reply +OK (diskfull guard tripping? \ + set MOON_DISK_FREE_MIN_PCT=0): {:?}", + key, + value, + reply + ); } fn redis_get(port: u16, key: &str) -> Option { @@ -207,6 +237,7 @@ fn sigkill(child: &mut Child) { #[test] #[ignore] // Requires built release binary + redis-cli; run explicitly. fn crash_01_lite_per_shard_aof_recovers_after_sigkill() { + let _serial = serialize_server_test(); let port = unique_port(); let dir = unique_dir("crash01"); std::fs::create_dir_all(&dir).expect("create test dir"); @@ -288,6 +319,7 @@ fn crash_01_lite_per_shard_aof_recovers_after_sigkill() { #[test] #[ignore] // Requires built release binary + redis-cli; run explicitly. fn crash_01_lite_always_per_shard_aof_recovers_after_sigkill() { + let _serial = serialize_server_test(); // Offset port so this test never collides with the everysec test // when both run on the same dev host. let port = unique_port().saturating_add(1); @@ -371,6 +403,7 @@ fn crash_01_lite_always_per_shard_aof_recovers_after_sigkill() { #[test] #[ignore] // Requires built release binary + redis-cli; run explicitly. fn pipeline_batch_no_double_write_after_crash_recovery() { + let _serial = serialize_server_test(); const N: usize = 20; let port = unique_port().saturating_add(2);