Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 21 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,27 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Performance
- **`--io-busy-poll-us` is now deploy-safe: a per-shard contention governor
gates the spin on shared cores (O3).** The p=1 busy-poll win (GCE pinned:
ARM 0.95→1.21×, x86 1.06→1.66× vs Redis) previously inverted into a
regression whenever the shard's core was shared (OrbStack/laptop), making
the flag pinned-cores-only operator judgment. Each shard thread's 1s
chore now samples its own `nonvoluntary_ctxt_switches` from
`/proc/thread-self/status` — the kernel's direct "someone else needed
this core" signal — and flips a per-thread gate in the vendored monoio
legacy driver: one window over 25 preempts/s stops the spin immediately;
five consecutive quiet windows re-enable it (asymmetric hysteresis, no
flapping). Startup state is ungated, so pinned-core deployments see the
win from the first request and a shared-core host pays at most ~one
window of spin. Idle cost is unchanged (the existing 10ms idle-disengage
already covers quiet threads). `MOON_SPIN_ADAPTIVE=0` restores
unconditional spinning (same-binary A/B knob);
`MOON_SPIN_MAX_PREEMPTS_PER_SEC` tunes the threshold. Non-Linux has no
preemption signal — the governor is inert there (pre-O3 behavior). This
also makes `--profile standalone` (which presets busy-poll 40) safe on
non-dedicated hosts.

### Documentation
- **SPSC-wake `Notify` stays on flume — the lock-free replacement was
rejected by measurement (O2).** A fully-validated AtomicBool token +
Expand Down
2 changes: 1 addition & 1 deletion CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,7 @@ orb run -m moon-dev bash -c 'sudo apt-get update -qq && sudo apt-get install -y
- `RUST_LOG=moon=debug` — enable tracing output (uses `tracing-subscriber` with `env-filter`)
- `MOON_NO_URING=1` — force-disable io_uring everywhere (monoio runtime + tokio bridge); used in CI/containers/WSL where io_uring is unavailable. ⚠ Before 2026-07 this env was a silent NO-OP for the monoio driver (FusionDriver picked io_uring regardless); it now forces the epoll/kqueue LegacyDriver. CLI equivalent: `--io-driver epoll`. Empirically verify the driver via `ls -l /proc/<pid>/fd | grep io_uring` — monoio logs neither choice. **GCE ARM (c4a Axion): epoll beats io_uring by 2-4% at ALL pipeline depths for KV** (same-instance A/B 2026-07-03); other platforms favor io_uring — bench per platform.
- `MOON_URING=1` — opt **into** the tokio→io_uring bridge. The bridge is **default-off under the tokio runtime** (it floods errors under load and can hang the accept loop); tokio shards run plain epoll/kqueue unless this is set. No effect on the monoio runtime, which always uses io_uring unless `MOON_NO_URING` is set.
- `MOON_EPOLL_SPIN_US=<µs>` / CLI `--io-busy-poll-us <µs>` — poll-mode park (vendored-monoio patch, `vendor/monoio`, grep "moon patch"): the shard thread busy-loops zero-timeout readiness polls for the budget before blocking, deleting the per-op scheduler sleep+wake. Legacy (epoll/kqueue) driver only — the CLI flag forces it; io_uring CQEs are NOT observable from userspace without task participation (DEFER_TASKRUN posts only inside enter; TWA_SIGNAL kicks cost ~µs; SQPOLL doesn't run poll task_work on 6.17 — all measured 2026-07-03). **This is what flipped GCE p=1 c1 to a WIN vs Redis: ARM c4a 0.95→1.19/1.21, x86 c3 1.06→1.66** (same-instance A/Bs). Costs up to budget-µs CPU per idle park (~4%/core at 40µs vs 1ms timer parks). ⚠ Unpinned/shared-core environments (OrbStack default, laptop) show spin as a REGRESSION — only judge it on pinned disjoint cores.
- `MOON_EPOLL_SPIN_US=<µs>` / CLI `--io-busy-poll-us <µs>` — poll-mode park (vendored-monoio patch, `vendor/monoio`, grep "moon patch"): the shard thread busy-loops zero-timeout readiness polls for the budget before blocking, deleting the per-op scheduler sleep+wake. Legacy (epoll/kqueue) driver only — the CLI flag forces it; io_uring CQEs are NOT observable from userspace without task participation (DEFER_TASKRUN posts only inside enter; TWA_SIGNAL kicks cost ~µs; SQPOLL doesn't run poll task_work on 6.17 — all measured 2026-07-03). **This is what flipped GCE p=1 c1 to a WIN vs Redis: ARM c4a 0.95→1.19/1.21, x86 c3 1.06→1.66** (same-instance A/Bs). Costs up to budget-µs CPU per idle park (~4%/core at 40µs vs 1ms timer parks; bounded by the 10ms idle-disengage). **O3 (2026-07-19): the spin is now self-gating** — each shard thread's 1s chore samples its own `nonvoluntary_ctxt_switches` (`/proc/thread-self/status`) and gates the spin via a per-thread flag in the vendored driver while the core is contended (>25 preempts/s one window = off; 5 quiet windows = on; `src/shard/spin_governor.rs`). The old "only judge spin on pinned disjoint cores" caveat now applies only to `MOON_SPIN_ADAPTIVE=0` (unconditional-spin A/B knob); `MOON_SPIN_MAX_PREEMPTS_PER_SEC` overrides the threshold. Non-Linux: no signal, governor inert, pre-O3 behavior.
- `MOON_URING_SPIN_US`, `MOON_URING_SQPOLL[_CPU]`, `MOON_URING_PLAIN` — io_uring-side experiment gates kept as documented diagnostics; all dead ends for the p=1 path (see `tmp/KV-FULLPROOF.md` Round 2).
- `MOON_IDLE_PARK=0` — disable the adaptive idle park (#373 phase 2): pins the monoio shard loop to its fixed 1ms periodic tick instead of stretching to 10ms after 64 provably-no-op ticks. Same-binary A/B knob for idle-CPU / latency benches. All counter-based chore cadences are multiples of the 10ms stretched period, so no chore timing changes either way; the only behavioral delta with the park ON is cached-clock staleness up to 10ms (vs 1ms) for a command arriving mid-park after ≥64ms of total shard quiet.
- `MOON_XSHARD_SPIN_BUDGET` / `MOON_XSHARD_SPIN_GATE` / `MOON_XSHARD_SPIN_MAX_CONNS` — diagnostic overrides for the C2 reply-side spin (`src/shard/slice.rs`; defaults 4096 iters / gate 2 / **solo-conn 1**). Budget `0` disables the spin entirely (the same-instance A/B knob that proved the c8P1 convoy). ⚠ The solo-conn ceiling (spin only when the conn is ALONE on its shard thread) is the L1 convoy fix — raising `MAX_CONNS` re-creates the s4 c8P1 collapse (a spinning conn starves its sibling AND the shard's SPSC drain, 0.45× vs Redis; fixed = 2.75× better, see `tmp/MULTISHARD-REDESIGN.md`). Bench-only knobs: never set in production.
Expand Down
8 changes: 8 additions & 0 deletions src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -311,6 +311,14 @@ pub struct ServerConfig {
/// low-pipeline request/response workloads on dedicated cores; costs up to
/// ~N µs of CPU per idle park. Measured (GCE c1 GET p=1, 2026-07):
/// ARM c4a 0.95→1.21× vs Redis, x86 c3 1.06→1.66×. monoio runtime only.
///
/// Deploy-safe by default (O3): each shard thread watches its own
/// involuntary-preemption rate (Linux) and self-gates the spin while its
/// core is shared with other runnable threads — the shared-core
/// regression that previously made this flag pinned-cores-only judgment
/// no longer applies (one >25-preempts/s window gates the spin; 5
/// consecutive quiet windows re-enable it). `MOON_SPIN_ADAPTIVE=0`
/// restores unconditional spinning (bench A/B knob).
#[arg(long = "io-busy-poll-us", default_value_t = 0)]
pub io_busy_poll_us: u64,

Expand Down
26 changes: 26 additions & 0 deletions src/shard/event_loop.rs
Original file line number Diff line number Diff line change
Expand Up @@ -879,6 +879,14 @@ impl super::Shard {
// cadence below (all multiples of 10) still fires exactly on time.
#[cfg(feature = "runtime-monoio")]
let mut idle_park = crate::shard::idle_park::IdleParkState::new();
// O3: adaptive busy-poll contention governor. Only constructed when a
// spin budget is actually configured — otherwise there is nothing to
// gate and the 1s /proc read would be pure waste.
#[cfg(feature = "runtime-monoio")]
let mut spin_governor: Option<crate::shard::spin_governor::SpinGovernor> =
(crate::runtime::epoll_spin_configured(server_config.io_busy_poll_us)
&& crate::shard::spin_governor::adaptive_enabled())
.then(crate::shard::spin_governor::SpinGovernor::new);
// Used by tokio select! for event-driven SPSC drain; monoio drains in periodic tick.
let spsc_notify_local = spsc_notify;
#[cfg(feature = "runtime-monoio")]
Expand Down Expand Up @@ -2427,6 +2435,24 @@ impl super::Shard {
// P6 is gated here (not per-1ms tick) to avoid the read_dir
// syscall overhead of wal.stats() on the hot path.
if monoio_tick_counter % 1000 == 0 {
// O3: sample this shard thread's involuntary-preemption
// rate and gate the driver spin while the core is shared.
// The gate is thread-local in the vendored driver, so the
// flip below affects exactly this shard's parks.
if let Some(gov) = spin_governor.as_mut() {
if let Some(contended) = gov.tick() {
monoio::set_legacy_spin_contended(contended);
tracing::info!(
"Shard {}: busy-poll spin {} (involuntary-preemption governor)",
shard_id,
if contended {
"GATED — core contended"
} else {
"re-enabled — core quiet"
}
);
}
}
timers::sync_wal_v3(&mut wal_writer);
// P3+MA1+MA2: MVCC committed prune + zombie sweep + kill old snapshots
// + RECL_* + segment-stall.
Expand Down
4 changes: 4 additions & 0 deletions src/shard/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,10 @@ pub mod segment_stall;
pub mod self_msg;
pub mod shared_databases;
pub mod slice;
/// O3 adaptive busy-poll governor — the spin it gates exists only in the
/// vendored monoio legacy driver, so the module is monoio-only.
#[cfg(feature = "runtime-monoio")]
pub(crate) mod spin_governor;
pub mod spsc_handler;
/// Shared MOVE/COPY-DB two-database intercept for every `ShardMessage` SPSC
/// arm (Gap A). Split out of `spsc_handler.rs` per the repo's file-size
Expand Down
271 changes: 271 additions & 0 deletions src/shard/spin_governor.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,271 @@
//! Adaptive busy-poll contention governor (digest O3).
//!
//! `--io-busy-poll-us` is a measured p=1 win ONLY on an effectively-dedicated
//! core (GCE pinned: ARM 0.95→1.21×, x86 1.06→1.66× vs Redis); on a shared
//! core the burned spin budget steals time from whatever is contending and
//! the win inverts (documented OrbStack/laptop regression). That made the
//! flag operator-judgment-only. This governor removes the judgment call: each
//! shard thread watches its OWN involuntary-preemption rate
//! (`nonvoluntary_ctxt_switches` in `/proc/thread-self/status` — the direct
//! kernel signal for "someone else needed this core") and flips the vendored
//! driver's per-thread contention gate (`monoio::set_legacy_spin_contended`).
//!
//! Policy (hysteresis, asymmetric on purpose):
//! - ONE window over the preemption threshold → gate the spin immediately
//! (contended cores stop burning budget within a window).
//! - [`CLEAN_WINDOWS_TO_ENABLE`] consecutive clean windows → ungate (slow
//! re-enable so a briefly-quiet interloper doesn't cause flapping).
//! - Start state: NOT contended — identical to pre-O3 behavior until
//! evidence arrives, so pinned-core benchmarks see the win from request
//! one and a shared-core host pays at most ~one window of spin.
//!
//! The governor ticks from the shard's 1s chore cadence (a multiple of the
//! 10ms stretched idle tick, per the #373 divisor rule). Windows are
//! measured with real elapsed time, so idle-park stretching just lengthens
//! a window rather than skewing the rate.
//!
//! Idle cost is NOT this module's job: the driver's idle-disengage
//! (`MOON_SPIN_IDLE_DISENGAGE_US`, default 10ms) already stops spinning on
//! quiet threads.
//!
//! Knobs (diagnostics; not for production tuning):
//! - `MOON_SPIN_ADAPTIVE=0` — disable the governor entirely (always-spin,
//! exact pre-O3 behavior; the same-binary A/B knob).
//! - `MOON_SPIN_MAX_PREEMPTS_PER_SEC` — threshold, default 25/s. A pinned
//! thread alone on its core measures ~0-5/s (timer/IPI noise); sharing
//! with one busy thread measures hundreds/s under CFS.
//!
//! Non-Linux: `/proc` does not exist → [`SpinGovernor::tick`] returns `None`
//! forever and the gate is never touched (pre-O3 behavior).

use std::time::Instant;

/// Consecutive clean 1s windows required to re-enable spinning.
const CLEAN_WINDOWS_TO_ENABLE: u32 = 5;

/// Default involuntary-preemption threshold (per second).
const DEFAULT_MAX_PREEMPTS_PER_SEC: f64 = 25.0;

/// Whether the adaptive governor is enabled (`MOON_SPIN_ADAPTIVE` != "0").
pub fn adaptive_enabled() -> bool {
static ENABLED: std::sync::OnceLock<bool> = std::sync::OnceLock::new();
*ENABLED.get_or_init(|| std::env::var_os("MOON_SPIN_ADAPTIVE").is_none_or(|v| v != "0"))
}

fn max_preempts_per_sec() -> f64 {
static V: std::sync::OnceLock<f64> = std::sync::OnceLock::new();
*V.get_or_init(|| {
std::env::var("MOON_SPIN_MAX_PREEMPTS_PER_SEC")
.ok()
.and_then(|v| v.parse::<f64>().ok())
.filter(|v| *v >= 0.0)
.unwrap_or(DEFAULT_MAX_PREEMPTS_PER_SEC)
})
}

/// Read the calling thread's involuntary context-switch counter.
#[cfg(target_os = "linux")]
fn read_nonvoluntary_switches() -> Option<u64> {
let status = std::fs::read_to_string("/proc/thread-self/status").ok()?;
parse_nonvoluntary(&status)
}

#[cfg(not(target_os = "linux"))]
fn read_nonvoluntary_switches() -> Option<u64> {
None
}

/// Parse `nonvoluntary_ctxt_switches:\t<N>` out of a /proc status blob.
#[cfg_attr(not(target_os = "linux"), allow(dead_code))] // non-Linux lib target: only tests call it
fn parse_nonvoluntary(status: &str) -> Option<u64> {
status
.lines()
.find(|l| l.starts_with("nonvoluntary_ctxt_switches:"))
.and_then(|l| l.split_whitespace().nth(1))
.and_then(|v| v.parse().ok())
}

/// Pure hysteresis decision: `Some(new_contended)` when the state flips.
fn decide(
contended: bool,
clean_windows: &mut u32,
preempts_per_sec: f64,
threshold: f64,
) -> Option<bool> {
if preempts_per_sec > threshold {
*clean_windows = 0;
if !contended {
return Some(true); // fast off: one bad window gates the spin
}
return None;
}
if contended {
*clean_windows += 1;
if *clean_windows >= CLEAN_WINDOWS_TO_ENABLE {
*clean_windows = 0;
return Some(false); // slow on: sustained quiet re-enables
}
}
None
}

/// Per-shard-thread governor state. Construct once per event loop; call
/// [`tick`](Self::tick) from the 1s chore.
pub struct SpinGovernor {
last_count: Option<u64>,
last_read: Instant,
clean_windows: u32,
contended: bool,
}

impl SpinGovernor {
pub fn new() -> Self {
Self {
last_count: None,
last_read: Instant::now(),
clean_windows: 0,
contended: false,
}
}

/// Sample the preemption counter and run the hysteresis. Returns
/// `Some(contended)` when the gate should flip (caller pushes it to the
/// driver); `None` on no change or when the signal is unavailable.
pub fn tick(&mut self) -> Option<bool> {
if !adaptive_enabled() {
return None;
}
let now = Instant::now();
let count = read_nonvoluntary_switches()?;
let Some(prev) = self.last_count else {
self.last_count = Some(count); // first sample: baseline only
self.last_read = now;
return None;
};
let elapsed = now.duration_since(self.last_read).as_secs_f64();
if elapsed <= 0.05 {
// Degenerate window: keep the previous baseline so the preemptions
// in this span fold into the next window instead of being dropped.
return None;
}
self.last_count = Some(count);
self.last_read = now;
let rate = count.saturating_sub(prev) as f64 / elapsed;
let flip = decide(
self.contended,
&mut self.clean_windows,
rate,
max_preempts_per_sec(),
);
if let Some(c) = flip {
self.contended = c;
}
flip
}

/// Current gate state.
#[cfg(test)]
fn contended(&self) -> bool {
self.contended
}
}

impl Default for SpinGovernor {
fn default() -> Self {
Self::new()
}
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn parse_nonvoluntary_extracts_counter() {
let blob =
"Name:\tmoon\nvoluntary_ctxt_switches:\t123\nnonvoluntary_ctxt_switches:\t4567\n";
assert_eq!(parse_nonvoluntary(blob), Some(4567));
assert_eq!(parse_nonvoluntary("Name:\tmoon\n"), None);
assert_eq!(parse_nonvoluntary(""), None);
}

#[test]
fn decide_one_bad_window_gates_immediately() {
let mut clean = 0;
assert_eq!(decide(false, &mut clean, 100.0, 25.0), Some(true));
}

#[test]
fn decide_clean_windows_below_threshold_do_not_flip_uncontended() {
let mut clean = 0;
for _ in 0..10 {
assert_eq!(decide(false, &mut clean, 1.0, 25.0), None);
}
}

#[test]
fn decide_reenables_only_after_sustained_quiet() {
let mut clean = 0;
// 4 clean windows: still gated.
for _ in 0..CLEAN_WINDOWS_TO_ENABLE - 1 {
assert_eq!(decide(true, &mut clean, 1.0, 25.0), None);
}
// 5th clean window: ungate.
assert_eq!(decide(true, &mut clean, 1.0, 25.0), Some(false));
assert_eq!(clean, 0, "counter resets after re-enable");
}

#[test]
fn decide_bad_window_resets_clean_streak() {
let mut clean = 0;
for _ in 0..3 {
assert_eq!(decide(true, &mut clean, 1.0, 25.0), None);
}
assert_eq!(clean, 3);
// Interloper returns: streak resets, stays contended (no flip event).
assert_eq!(decide(true, &mut clean, 200.0, 25.0), None);
assert_eq!(clean, 0);
// Needs the full quiet run again.
for _ in 0..CLEAN_WINDOWS_TO_ENABLE - 1 {
assert_eq!(decide(true, &mut clean, 1.0, 25.0), None);
}
assert_eq!(decide(true, &mut clean, 1.0, 25.0), Some(false));
}

#[test]
fn decide_boundary_rate_is_clean() {
let mut clean = 0;
// Exactly at threshold: not "over" — clean.
assert_eq!(decide(false, &mut clean, 25.0, 25.0), None);
}

#[cfg(target_os = "linux")]
#[test]
fn read_nonvoluntary_switches_works_on_linux() {
assert!(read_nonvoluntary_switches().is_some());
}

#[cfg(target_os = "linux")]
#[test]
fn governor_sub_window_tick_preserves_baseline() {
let mut g = SpinGovernor::new();
assert_eq!(g.tick(), None); // baseline sample
let baseline_count = g.last_count;
let baseline_read = g.last_read;
// A second tick <50ms later is a degenerate window: it must NOT
// advance the baseline, so the skipped span's preemptions fold into
// the next full window instead of being silently discarded.
assert_eq!(g.tick(), None);
assert_eq!(g.last_count, baseline_count);
assert_eq!(g.last_read, baseline_read);
}

#[test]
fn governor_first_tick_is_baseline_only() {
let mut g = SpinGovernor::new();
// First tick can never flip (no previous sample). On non-Linux this
// is None for lack of signal; on Linux for lack of a baseline.
assert_eq!(g.tick(), None);
assert!(!g.contended());
}
}
Loading
Loading