From 1d5a0a454760c73184bd57f6afd92d01bf176fbb Mon Sep 17 00:00:00 2001 From: CMGS Date: Wed, 9 Sep 2026 23:14:22 +0900 Subject: [PATCH 01/16] perf: the auto tuner keeps the shared pipes only where they measure faster The worker rule (busy 85% and thin local batches) used to move every session on the worker to the shared pipes; a memtier sliding pipeline at 200 connections lost 16% that way, since its local path was not congested and the cross-worker hop cost more than the deeper batches returned. The rule now only starts an experiment: the worker moves its sessions, compares its command rate a second later with the second before, keeps the shared pipes on a 5% gain and otherwise takes the sessions back, waiting 30 s before the next try and doubling that up to 8 min while the answer holds; from the shared pipes it probes the other way on the same schedule. INFO reports pipe_probes, pipe_keeps, pipe_reverts and per-worker worker_shared. Nothing runs on the command path: the tick reads the worker's existing command counter. --- docs/architecture.md | 19 +-- docs/configuration.md | 2 +- docs/operations.md | 2 +- src/admin.rs | 33 +++-- src/client/tuner.rs | 278 ++++++++++++++++++++++++++++++++++-------- src/stats.rs | 4 + 6 files changed, 271 insertions(+), 67 deletions(-) diff --git a/docs/architecture.md b/docs/architecture.md index 61f5134..7702300 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -101,18 +101,21 @@ idle dispatches have run the score down, back to its worker's connections once four busy ones have restored it. On top of that each worker samples its own CPU busyness and the frames its own backend connections batch per write every 100 ms (the in-flight commands per master stand in while no -local traffic flows): while it is busy (85% and above) and the local -batches are thin (under eight frames per write) every session on it prefers the -shared pipes, since one pipe per node batches what 64 workers × 128 nodes -of per-worker connections cannot; it lets go once busyness has stayed under -60% or the depth above sixteen for three seconds — slowly, because moving -its sessions away is what lowers its own busyness — and takes the shared -pipes after 300 ms of the entry condition. Switches happen only while a session has +local traffic flows): busy (85% and above) with thin local batches (under +eight frames per write) for 300 ms only starts an experiment. The worker +moves its sessions to the shared pipes and compares the command rate a +second later against the second before the move; it keeps them there on a +5% gain, and otherwise takes them back and waits 30 seconds before trying +again, doubling that to eight minutes while the answer holds. From the +shared pipes it probes the other way on the same schedule, and lets go +unmeasured once busyness has stayed under 60% or the depth above sixteen +for three seconds — slowly, because moving its sessions away is what lowers +its own busyness. An idle worker never probes. Switches happen only while a session has nothing in flight — a session that never idles is paused for one round trip so it can drain and move — so a request-response client ends up on the shared pipe, a pipelining client on a lightly loaded worker on its own connections, and a saturated proxy in front of a wide cluster on the -shared pipes. The reply cache is +shared pipes where they measure faster. The reply cache is wired as under `yes` (owner-worker trackers, process-wide coverage, broadcast invalidation): sessions on the shared pipes fill it, sessions on the worker-local connections read it but never fill it, since those diff --git a/docs/configuration.md b/docs/configuration.md index f24bf01..6e98ac7 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -14,7 +14,7 @@ See [`mithril.conf.sample`](https://github.com/projecteru2/mithril/blob/master/m | `worker-threads` | int | CPU count | worker threads, one runtime each | | `maxclients` | int | `10000` | client connection cap, enforced at accept | | `backend-conns` | 1..512 | `1` | shared pipelined connections per node per worker, per database in use (a session on `SELECT n` opens its own pool to each node; the exclusive-connection cap for blocking and `WATCH` stays per node across databases) | -| `backend-sharding` | `yes`/`no`/`auto` | `no` | `yes`: one process-wide connection per node per database in use, owned by the worker its address hashes to, which deepens backend pipelines for unpipelined workloads; `auto`: decided per session and per worker — a session that keeps sending one command at a time moves to the shared pipes, a pipelining session stays on its worker's own connections unless that worker is saturated with thin local batches (busy ≥ 85% and fewer than eight frames per backend write), in which case all of its sessions move to the shared pipes (switches happen only while a session has nothing in flight; with the reply cache, only sessions on the shared pipes fill it — the worker-local connections carry no key tracking — while every session reads it) | +| `backend-sharding` | `yes`/`no`/`auto` | `no` | `yes`: one process-wide connection per node per database in use, owned by the worker its address hashes to, which deepens backend pipelines for unpipelined workloads; `auto`: decided per session and per worker — a session that keeps sending one command at a time moves to the shared pipes, a pipelining session stays on its worker's own connections unless that worker is saturated with thin local batches (busy ≥ 85% and fewer than eight frames per backend write), which makes it try the shared pipes for a second and keep its sessions there only where they measure faster (switches happen only while a session has nothing in flight; with the reply cache, only sessions on the shared pipes fill it — the worker-local connections carry no key tracking — while every session reads it) | | `reply-cache` | `yes`/`no` | `no` | worker-local GET/MGET reply cache; the servers track the keys the proxy caches (redirected opt-in RESP3 tracking) and invalidate them on change, so hits skip the backend round trip and a session reads its own writes | | `reply-cache-max-bytes` | bytes | `64mb` | per-worker budget for the two cache generations; the live generation holds half of it, so a worker's hot set must fit in half the budget to keep hitting (see operations.md for sizing and the RSS to expect) | | `reply-cache-max-age-secs` | 1..3600 | `10` | staleness backstop for missed invalidations; entries older than this never serve | diff --git a/docs/operations.md b/docs/operations.md index 8d8f5ab..e8a0271 100644 --- a/docs/operations.md +++ b/docs/operations.md @@ -16,7 +16,7 @@ client's current command in `cmd`. | Clients | `connected_clients` | | CPU | `used_cpu_sys`, `used_cpu_user` | | Stats | `total_connections_received`, `total_commands_processed`, `total_net_input_bytes`, `total_net_output_bytes`, `total_error_replies`, `redirections`, `redirect_waits`, session lifecycle counters (`readers_exited`, `writers_exited`, `sessions_closed`) | -| Mithril | `worker_threads`, `backend_conns_per_node`, `backend_sharding`, `slave_mode`, `reply_cache`, `cache_hits`, `cache_misses`, `cache_invalidations`, `cache_entries`, `cache_bytes`, `cache_flips`, `cache_armed_workers`, `worker_commands` (per-worker) | +| Mithril | `worker_threads`, `backend_conns_per_node`, `backend_sharding`, `slave_mode`, `reply_cache`, `cache_hits`, `cache_misses`, `cache_invalidations`, `cache_entries`, `cache_bytes`, `cache_flips`, `cache_armed_workers`, `pipe_probes`, `pipe_keeps`, `pipe_reverts` (auto-sharding experiments and how they ended), `worker_commands` (per-worker), `worker_shared` (per-worker, 1 on the shared pipes) | | Cluster | `cluster_enabled` (always 1) | | Commandstats | `cmdstat_:calls=` for every command run at least once, subcommands as `cmdstat_client\|list`; counted once accepted, summed over workers | diff --git a/src/admin.rs b/src/admin.rs index d101eca..9f3b096 100644 --- a/src/admin.rs +++ b/src/admin.rs @@ -2,7 +2,7 @@ use std::borrow::Cow; use std::fmt::Write; -use std::sync::atomic::{AtomicU16, Ordering}; +use std::sync::atomic::{AtomicU16, AtomicU64, Ordering}; use std::time::{SystemTime, UNIX_EPOCH}; use crate::acl::Acl; @@ -211,7 +211,8 @@ pub fn info(cfg: &Config, stats: &Stats, started: u64) -> Vec { backend_sharding:{}\r\nslave_mode:{}\r\nreply_cache:{}\r\n\ cache_hits:{}\r\ncache_misses:{}\r\ncache_invalidations:{}\r\n\ cache_armed_workers:{}\r\ncache_entries:{}\r\ncache_bytes:{}\r\n\ - cache_flips:{}\r\nworker_commands:{}\r\n", + cache_flips:{}\r\npipe_probes:{}\r\npipe_keeps:{}\r\npipe_reverts:{}\r\n\ + worker_commands:{}\r\nworker_shared:{}\r\n", crate::VERSION, std::process::id(), cfg.port, @@ -242,12 +243,11 @@ pub fn info(cfg: &Config, stats: &Stats, started: u64) -> Vec { stats.sum(|w| &w.cache_entries), stats.sum(|w| &w.cache_bytes), stats.sum(|w| &w.cache_flips), - stats - .workers - .iter() - .map(|w| w.commands.load(Ordering::Relaxed).to_string()) - .collect::>() - .join(","), + stats.sum(|w| &w.pipe_probes), + stats.sum(|w| &w.pipe_keeps), + stats.sum(|w| &w.pipe_reverts), + per_worker(stats, |w| &w.commands), + per_worker(stats, |w| &w.pipe_shared), ); text.push_str("\r\n# Cluster\r\ncluster_enabled:1\r\n\r\n# Commandstats\r\n"); for id in 0..command::entries() as u16 { @@ -418,6 +418,15 @@ fn split_announce(announce: &str) -> (&str, u16) { } } +fn per_worker &AtomicU64>(stats: &Stats, field: F) -> String { + stats + .workers + .iter() + .map(|w| field(w).load(Ordering::Relaxed).to_string()) + .collect::>() + .join(",") +} + fn yesno(v: bool) -> &'static str { if v { "yes" } else { "no" } } @@ -525,6 +534,9 @@ mod tests { let get = command::lookup(b"get").unwrap().id; stats.workers[0].calls.at(get).store(3, Ordering::Relaxed); stats.workers[1].calls.at(get).store(4, Ordering::Relaxed); + stats.workers[1].pipe_shared.store(1, Ordering::Relaxed); + stats.workers[0].pipe_probes.store(5, Ordering::Relaxed); + stats.workers[1].pipe_keeps.store(2, Ordering::Relaxed); let out = info(&cfg, &stats, 0); let text = String::from_utf8_lossy(&out); let field = |k: &str| { @@ -535,6 +547,11 @@ mod tests { assert_eq!(field("connected_clients").as_deref(), Some("7")); assert_eq!(field("mithril_version").as_deref(), Some(crate::VERSION)); assert_eq!(field("cache_flips").as_deref(), Some("0")); + assert_eq!(field("worker_commands").as_deref(), Some("0,0")); + assert_eq!(field("worker_shared").as_deref(), Some("0,1")); + assert_eq!(field("pipe_probes").as_deref(), Some("5")); + assert_eq!(field("pipe_keeps").as_deref(), Some("2")); + assert_eq!(field("pipe_reverts").as_deref(), Some("0")); assert_eq!(field("total_error_replies").as_deref(), Some("0")); assert_eq!(field("cluster_enabled").as_deref(), Some("1")); assert_eq!(field("cmdstat_get").as_deref(), Some("calls=7")); diff --git a/src/client/tuner.rs b/src/client/tuner.rs index 96f15ab..cd49ce2 100644 --- a/src/client/tuner.rs +++ b/src/client/tuner.rs @@ -1,17 +1,18 @@ //! Pipe selection under auto sharding: the session score and the worker-level tuner. use std::rc::Rc; +use std::sync::atomic::Ordering; use std::time::{Duration, Instant}; use super::Shared; use super::session::Session; -use crate::stats; +use crate::stats::{self, WorkerStats}; // pipelining score: a local session shares at 0, a shared one returns at PIPELINED_LOCAL pub(super) const PIPELINED_LOCAL: u8 = 4; const PIPELINED_MAX: u8 = 8; -// auto sharding moves a worker's sessions to the shared pipes while the worker is -// busy and its local backend batches would stay thin (in-flight commands per master) +// what triggers a probe: the worker is busy and its local backend batches would +// stay thin (in-flight commands per master) const TUNE_PERIOD: Duration = Duration::from_millis(100); const BUSY_ENTER: u32 = 85; const BUSY_LEAVE: u32 = 60; @@ -21,6 +22,19 @@ const DEPTH_LEAVE: u64 = 16; // lowers its own busyness, so leaving is slow and entering prompt const ENTER_TICKS: u32 = 3; const LEAVE_TICKS: u32 = 30; +// the command-rate window, ten ticks of 100 ms, so a rate is commands per second +const RATE_TICKS: usize = 10; +// ticks the probe gives the sessions to move before it measures +const SETTLE_TICKS: u32 = 10; +const PROBE_TICKS: u32 = SETTLE_TICKS + RATE_TICKS as u32; +// the shared pipes are kept only where they measure this much faster +const KEEP_GAIN_PCT: u64 = 5; +// commands per second under which the worker is idle and its rates say nothing +const MIN_RATE_PER_SEC: u64 = 100; +// ticks to the next probe, doubling while decisions confirm the current state +const PROBE_BACKOFF_TICKS: u32 = 300; +const PROBE_BACKOFF_MAX_TICKS: u32 = 4800; +const PROBE_DOUBLINGS: u32 = (PROBE_BACKOFF_MAX_TICKS / PROBE_BACKOFF_TICKS).ilog2(); impl Session { // an unpipelined session gains from the deeper batches of the shared pipe, a @@ -40,13 +54,140 @@ impl Session { } } -/// Samples this worker's CPU busyness and in-flight depth per master and publishes +#[derive(Clone, Copy, Default)] +enum Phase { + #[default] + Steady, + Probing { + baseline: u64, + ticks: u32, + }, +} + +// the worker's experiment: the busy-and-thin rule triggers it, the measured command rate decides +#[derive(Default)] +struct Probe { + prefer: bool, + probes: u64, + keeps: u64, + reverts: u64, + phase: Phase, + streak: u32, + wait: u32, + confirmed: u32, + ring: [u64; RATE_TICKS], + at: usize, + seen: usize, + last: u64, +} + +impl Probe { + fn tick(&mut self, busy_pct: u32, depth: u64, commands_now: u64) -> bool { + self.record(commands_now); + self.wait = self.wait.saturating_sub(1); + match self.phase { + Phase::Probing { baseline, ticks } => self.probing(baseline, ticks), + Phase::Steady if self.prefer => self.leave(busy_pct, depth), + Phase::Steady => self.enter(busy_pct, depth), + } + self.prefer + } + + fn publish(&self, wstats: &WorkerStats) { + wstats + .pipe_shared + .store(u64::from(self.prefer), Ordering::Relaxed); + wstats.pipe_probes.store(self.probes, Ordering::Relaxed); + wstats.pipe_keeps.store(self.keeps, Ordering::Relaxed); + wstats.pipe_reverts.store(self.reverts, Ordering::Relaxed); + } + + fn record(&mut self, commands_now: u64) { + self.ring[self.at] = commands_now.saturating_sub(self.last); + self.at = (self.at + 1) % RATE_TICKS; + self.last = commands_now; + self.seen = (self.seen + 1).min(RATE_TICKS); + } + + fn rate(&self) -> u64 { + self.ring.iter().sum() + } + + fn enter(&mut self, busy_pct: u32, depth: u64) { + if busy_pct < BUSY_ENTER || depth >= DEPTH_ENTER { + self.streak = 0; + return; + } + self.streak = (self.streak + 1).min(ENTER_TICKS); + if self.streak == ENTER_TICKS { + self.start_probe(); + } + } + + fn leave(&mut self, busy_pct: u32, depth: u64) { + if busy_pct > BUSY_LEAVE && depth < DEPTH_LEAVE { + self.streak = 0; + self.start_probe(); + return; + } + self.streak += 1; + if self.streak >= LEAVE_TICKS { + // the load went away rather than a measurement: the next busy stretch probes at once + self.streak = 0; + self.wait = 0; + self.prefer = false; + } + } + + fn start_probe(&mut self) { + let baseline = self.rate(); + if self.wait > 0 || self.seen < RATE_TICKS || baseline < MIN_RATE_PER_SEC { + return; + } + self.probes += 1; + self.streak = 0; + self.prefer = !self.prefer; + self.phase = Phase::Probing { baseline, ticks: 0 }; + } + + fn probing(&mut self, baseline: u64, ticks: u32) { + let ticks = ticks + 1; + if ticks < PROBE_TICKS { + self.phase = Phase::Probing { baseline, ticks }; + return; + } + let measured = self.rate(); + let (shared, local) = if self.prefer { + (measured, baseline) + } else { + (baseline, measured) + }; + let keep = shared * 100 >= local * (100 + KEEP_GAIN_PCT); + // a decision that flips the state starts the schedule over, one that confirms it waits longer + if keep == self.prefer { + self.confirmed = 0; + self.wait = PROBE_BACKOFF_TICKS; + } else { + self.wait = PROBE_BACKOFF_TICKS << self.confirmed; + self.confirmed = (self.confirmed + 1).min(PROBE_DOUBLINGS); + } + if keep { + self.keeps += 1; + } else { + self.reverts += 1; + } + self.prefer = keep; + self.phase = Phase::Steady; + } +} + +/// Samples this worker's CPU busyness, batch depth and command rate, and publishes /// whether its sessions should prefer the shared pipes. pub async fn auto_tuner(shared: Rc) { let mut last_ticks = stats::thread_cpu_ticks(); let mut last_at = Instant::now(); let mut busy_x16 = 0u32; - let mut streak = 0u32; + let mut probe = Probe::default(); loop { tokio::time::sleep(TUNE_PERIOD).await; let now = Instant::now(); @@ -62,9 +203,11 @@ pub async fn auto_tuner(shared: Rc) { let masters = shared.topo.load().masters.len().max(1) as u64; let (measured, writes) = shared.backends.batch_depth(); let depth = tune_depth(measured, writes, shared.inflight.get(), masters); - let cur = shared.prefer_shared.get(); - let wanted = prefer_shared(cur, (busy_x16 + 8) / 16, depth); - shared.prefer_shared.set(settle(cur, wanted, &mut streak)); + let commands = shared.wstats.commands.load(Ordering::Relaxed); + shared + .prefer_shared + .set(probe.tick((busy_x16 + 8) / 16, depth, commands)); + probe.publish(&shared.wstats); } } @@ -80,30 +223,6 @@ fn switch_pipes(sharded: bool, score: u8, worker_prefers_shared: bool) -> bool { } } -// hysteresis on both signals so a worker does not flap its sessions between pipe kinds -fn prefer_shared(current: bool, busy_pct: u32, depth: u64) -> bool { - if current { - busy_pct > BUSY_LEAVE && depth < DEPTH_LEAVE - } else { - busy_pct >= BUSY_ENTER && depth < DEPTH_ENTER - } -} - -// a change of preference must hold for ENTER_TICKS or LEAVE_TICKS samples -fn settle(current: bool, wanted: bool, streak: &mut u32) -> bool { - if wanted == current { - *streak = 0; - return current; - } - *streak += 1; - let need = if current { LEAVE_TICKS } else { ENTER_TICKS }; - if *streak < need { - return current; - } - *streak = 0; - wanted -} - // the measured local batch while local traffic flows; with none, the in-flight // commands spread over the masters — an estimate that errs toward staying on the // shared pipes, which at saturation batch at least as well as local connections @@ -153,15 +272,9 @@ mod tests { } #[test] - fn worker_preference_overrides_the_score_with_hysteresis() { + fn worker_preference_overrides_the_score() { assert!(switch_pipes(false, PIPELINED_MAX, true)); assert!(!switch_pipes(true, PIPELINED_MAX, true)); - assert!(!prefer_shared(false, 84, 2)); - assert!(!prefer_shared(false, 90, DEPTH_ENTER)); - assert!(prefer_shared(false, 90, DEPTH_ENTER - 1)); - assert!(prefer_shared(true, 61, DEPTH_LEAVE - 1)); - assert!(!prefer_shared(true, 60, 1)); - assert!(!prefer_shared(true, 99, DEPTH_LEAVE)); let mut b = 0; for _ in 0..64 { b = busy_ewma(b, BUSY_ENTER); @@ -169,16 +282,83 @@ mod tests { assert_eq!((b + 8) / 16, BUSY_ENTER); assert_eq!(tune_depth(3, 10, 900, 128), 3); assert_eq!(tune_depth(3, 0, 900, 128), 7); - let mut streak = 0; - for _ in 0..ENTER_TICKS - 1 { - assert!(!settle(false, true, &mut streak)); - } - assert!(settle(false, true, &mut streak)); - for _ in 0..LEAVE_TICKS - 1 { - assert!(settle(true, false, &mut streak)); + } + + #[test] + fn a_probe_that_gains_keeps_the_shared_pipes() { + let mut rig = Rig::default(); + assert!(!rig.run(RATE_TICKS as u32 - 1, 1_000, 90, 1)); + assert!(rig.run(1, 1_000, 90, 1)); + assert_eq!(rig.probe.probes, 1); + rig.run(SETTLE_TICKS, 1_000, 90, 1); + assert!(rig.run(RATE_TICKS as u32, 1_200, 90, 1)); + assert_eq!((rig.probe.keeps, rig.probe.reverts), (1, 0)); + assert_eq!(rig.probe.wait, PROBE_BACKOFF_TICKS); + } + + #[test] + fn a_probe_that_loses_reverts_and_waits_twice_as_long_after_the_second() { + let mut rig = Rig::default(); + rig.run(RATE_TICKS as u32, 1_000, 90, 1); + rig.run(PROBE_TICKS - 1, 800, 90, 1); + assert!(!rig.run(1, 800, 90, 1)); + assert_eq!((rig.probe.keeps, rig.probe.reverts), (0, 1)); + assert_eq!(rig.probe.wait, PROBE_BACKOFF_TICKS); + assert!(!rig.run(PROBE_BACKOFF_TICKS - 1, 1_000, 90, 1)); + assert_eq!(rig.probe.probes, 1); + assert!(rig.run(1, 1_000, 90, 1)); + assert_eq!(rig.probe.probes, 2); + rig.run(PROBE_TICKS - 1, 500, 90, 1); + assert!(!rig.run(1, 500, 90, 1)); + assert_eq!(rig.probe.reverts, 2); + assert_eq!(rig.probe.wait, PROBE_BACKOFF_TICKS * 2); + } + + #[test] + fn low_busyness_leaves_the_shared_pipes_without_a_backoff() { + let mut rig = Rig::default(); + rig.run(RATE_TICKS as u32, 1_000, 90, 1); + assert!(rig.run(PROBE_TICKS, 1_200, 90, 1)); + assert!(rig.run(LEAVE_TICKS - 1, 1_200, BUSY_LEAVE, 1)); + assert!(!rig.run(1, 1_200, BUSY_LEAVE, 1)); + assert_eq!(rig.probe.wait, 0); + assert_eq!((rig.probe.probes, rig.probe.keeps), (1, 1)); + } + + #[test] + fn a_reverse_probe_returns_to_local_when_the_shared_pipes_are_not_faster() { + let mut rig = Rig::default(); + rig.run(RATE_TICKS as u32, 1_000, 90, 1); + rig.run(PROBE_TICKS, 1_200, 90, 1); + assert!(!rig.run(PROBE_BACKOFF_TICKS, 1_200, 90, 1)); + assert_eq!(rig.probe.probes, 2); + rig.run(PROBE_TICKS - 1, 1_200, 90, 1); + assert!(!rig.run(1, 1_200, 90, 1)); + assert_eq!((rig.probe.keeps, rig.probe.reverts), (1, 1)); + assert_eq!(rig.probe.wait, PROBE_BACKOFF_TICKS); + } + + #[test] + fn an_idle_worker_never_probes() { + let mut rig = Rig::default(); + assert!(!rig.run(PROBE_BACKOFF_TICKS, MIN_RATE_PER_SEC - 10, 90, 1)); + assert_eq!(rig.probe.probes, 0); + } + + #[derive(Default)] + struct Rig { + probe: Probe, + commands: u64, + } + + impl Rig { + fn run(&mut self, ticks: u32, per_sec: u64, busy_pct: u32, depth: u64) -> bool { + let mut prefer = self.probe.prefer; + for _ in 0..ticks { + self.commands += per_sec / RATE_TICKS as u64; + prefer = self.probe.tick(busy_pct, depth, self.commands); + } + prefer } - assert!(!settle(true, false, &mut streak)); - assert!(settle(true, true, &mut streak)); - assert_eq!(streak, 0); } } diff --git a/src/stats.rs b/src/stats.rs index 39dc859..2fcf3c4 100644 --- a/src/stats.rs +++ b/src/stats.rs @@ -48,6 +48,10 @@ pub struct WorkerStats { pub cache_entries: AtomicU64, pub cache_bytes: AtomicU64, pub cache_flips: AtomicU64, + pub pipe_shared: AtomicU64, + pub pipe_probes: AtomicU64, + pub pipe_keeps: AtomicU64, + pub pipe_reverts: AtomicU64, pub calls: Calls, } From 611b6cd4cde660f4284e52a1b275668a1d797bb8 Mon Sep 17 00:00:00 2001 From: CMGS Date: Wed, 9 Sep 2026 23:27:07 +0900 Subject: [PATCH 02/16] fix: a pipelining session drains for the move back to its worker's connections The pause that lets a never-idle session drain and switch existed only for the move to the shared pipes; a session the worker had moved there stayed on them after a losing probe, since its depth never reached zero, and a reverse probe measured the shared path twice. The pause now applies whenever the session should switch. The per-worker INFO field is worker_prefers_shared: it reports the tuner's preference, not lane placement. --- docs/operations.md | 2 +- src/admin.rs | 4 ++-- src/client/tuner.rs | 10 +++++----- 3 files changed, 8 insertions(+), 8 deletions(-) diff --git a/docs/operations.md b/docs/operations.md index e8a0271..9c43e1b 100644 --- a/docs/operations.md +++ b/docs/operations.md @@ -16,7 +16,7 @@ client's current command in `cmd`. | Clients | `connected_clients` | | CPU | `used_cpu_sys`, `used_cpu_user` | | Stats | `total_connections_received`, `total_commands_processed`, `total_net_input_bytes`, `total_net_output_bytes`, `total_error_replies`, `redirections`, `redirect_waits`, session lifecycle counters (`readers_exited`, `writers_exited`, `sessions_closed`) | -| Mithril | `worker_threads`, `backend_conns_per_node`, `backend_sharding`, `slave_mode`, `reply_cache`, `cache_hits`, `cache_misses`, `cache_invalidations`, `cache_entries`, `cache_bytes`, `cache_flips`, `cache_armed_workers`, `pipe_probes`, `pipe_keeps`, `pipe_reverts` (auto-sharding experiments and how they ended), `worker_commands` (per-worker), `worker_shared` (per-worker, 1 on the shared pipes) | +| Mithril | `worker_threads`, `backend_conns_per_node`, `backend_sharding`, `slave_mode`, `reply_cache`, `cache_hits`, `cache_misses`, `cache_invalidations`, `cache_entries`, `cache_bytes`, `cache_flips`, `cache_armed_workers`, `pipe_probes`, `pipe_keeps`, `pipe_reverts` (auto-sharding experiments and how they ended), `worker_commands` (per-worker), `worker_prefers_shared` (per-worker, 1 while the auto tuner sends the worker's sessions to the shared pipes; sessions on the shared pipes by their own score, or under `backend-sharding yes`, are not counted) | | Cluster | `cluster_enabled` (always 1) | | Commandstats | `cmdstat_:calls=` for every command run at least once, subcommands as `cmdstat_client\|list`; counted once accepted, summed over workers | diff --git a/src/admin.rs b/src/admin.rs index 9f3b096..a71d6e3 100644 --- a/src/admin.rs +++ b/src/admin.rs @@ -212,7 +212,7 @@ pub fn info(cfg: &Config, stats: &Stats, started: u64) -> Vec { cache_hits:{}\r\ncache_misses:{}\r\ncache_invalidations:{}\r\n\ cache_armed_workers:{}\r\ncache_entries:{}\r\ncache_bytes:{}\r\n\ cache_flips:{}\r\npipe_probes:{}\r\npipe_keeps:{}\r\npipe_reverts:{}\r\n\ - worker_commands:{}\r\nworker_shared:{}\r\n", + worker_commands:{}\r\nworker_prefers_shared:{}\r\n", crate::VERSION, std::process::id(), cfg.port, @@ -548,7 +548,7 @@ mod tests { assert_eq!(field("mithril_version").as_deref(), Some(crate::VERSION)); assert_eq!(field("cache_flips").as_deref(), Some("0")); assert_eq!(field("worker_commands").as_deref(), Some("0,0")); - assert_eq!(field("worker_shared").as_deref(), Some("0,1")); + assert_eq!(field("worker_prefers_shared").as_deref(), Some("0,1")); assert_eq!(field("pipe_probes").as_deref(), Some("5")); assert_eq!(field("pipe_keeps").as_deref(), Some("2")); assert_eq!(field("pipe_reverts").as_deref(), Some("0")); diff --git a/src/client/tuner.rs b/src/client/tuner.rs index cd49ce2..7e42a16 100644 --- a/src/client/tuner.rs +++ b/src/client/tuner.rs @@ -38,16 +38,16 @@ const PROBE_DOUBLINGS: u32 = (PROBE_BACKOFF_MAX_TICKS / PROBE_BACKOFF_TICKS).ilo impl Session { // an unpipelined session gains from the deeper batches of the shared pipe, a - // pipelined one from its worker-local connection; switch only while nothing is in flight + // pipelined one from its worker-local connection; a switch waits for the session + // to drain, and a session that never idles is paused until it has pub(super) fn adapt_pipes(&self) { let depth = self.outstanding(); let score = self.pipelined.get(); self.pipelined.set(pipelining_score(score, depth)); let sharded = self.link.sharded.get(); - let prefer_shared = self.shared.prefer_shared.get(); - self.switch_pending - .set(prefer_shared && !sharded && depth > 0); - if depth == 0 && switch_pipes(sharded, score, prefer_shared) && !self.fanouts_pending() { + let switch = switch_pipes(sharded, score, self.shared.prefer_shared.get()); + self.switch_pending.set(switch && depth > 0); + if switch && depth == 0 && !self.fanouts_pending() { self.link.sharded.set(!sharded); self.conns.borrow_mut().by_node.clear(); } From 762682c05cef6cad495f32520e1d0eec12cf375c Mon Sep 17 00:00:00 2001 From: CMGS Date: Wed, 9 Sep 2026 23:33:29 +0900 Subject: [PATCH 03/16] fix: only a switch that stays wanted pauses a session to drain A session at score one with a reply outstanding was paused for the score-based move to the shared pipes, then the idle dispatch lowered its score again without a switch, so a lightly pipelined client paid a round trip every other command. The pause now applies to the worker's move to the shared pipes and to a pipelining session's move back, the two switches still wanted once the session has drained; an unpipelined session idles and moves on its own. --- src/client/tuner.rs | 31 ++++++++++++++++++++++++++----- 1 file changed, 26 insertions(+), 5 deletions(-) diff --git a/src/client/tuner.rs b/src/client/tuner.rs index 7e42a16..e8dbda2 100644 --- a/src/client/tuner.rs +++ b/src/client/tuner.rs @@ -38,16 +38,17 @@ const PROBE_DOUBLINGS: u32 = (PROBE_BACKOFF_MAX_TICKS / PROBE_BACKOFF_TICKS).ilo impl Session { // an unpipelined session gains from the deeper batches of the shared pipe, a - // pipelined one from its worker-local connection; a switch waits for the session - // to drain, and a session that never idles is paused until it has + // pipelined one from its worker-local connection; a switch happens while nothing + // is in flight pub(super) fn adapt_pipes(&self) { let depth = self.outstanding(); let score = self.pipelined.get(); self.pipelined.set(pipelining_score(score, depth)); let sharded = self.link.sharded.get(); - let switch = switch_pipes(sharded, score, self.shared.prefer_shared.get()); - self.switch_pending.set(switch && depth > 0); - if switch && depth == 0 && !self.fanouts_pending() { + let prefer_shared = self.shared.prefer_shared.get(); + self.switch_pending + .set(depth > 0 && must_drain(sharded, score, prefer_shared)); + if depth == 0 && switch_pipes(sharded, score, prefer_shared) && !self.fanouts_pending() { self.link.sharded.set(!sharded); self.conns.borrow_mut().by_node.clear(); } @@ -223,6 +224,15 @@ fn switch_pipes(sharded: bool, score: u8, worker_prefers_shared: bool) -> bool { } } +// a never-idle session is paused to move only for a switch still wanted once it drains +fn must_drain(sharded: bool, score: u8, worker_prefers_shared: bool) -> bool { + if worker_prefers_shared { + !sharded + } else { + sharded && score >= PIPELINED_LOCAL + } +} + // the measured local batch while local traffic flows; with none, the in-flight // commands spread over the masters — an estimate that errs toward staying on the // shared pipes, which at saturation batch at least as well as local connections @@ -271,6 +281,17 @@ mod tests { assert_eq!(pipelining_score(0, 0), 0); } + #[test] + fn only_a_stable_switch_pauses_a_session_to_drain() { + assert!(must_drain(false, 0, true)); + assert!(must_drain(false, PIPELINED_MAX, true)); + assert!(!must_drain(true, PIPELINED_MAX, true)); + assert!(must_drain(true, PIPELINED_LOCAL, false)); + assert!(!must_drain(true, PIPELINED_LOCAL - 1, false)); + assert!(!must_drain(false, 1, false)); + assert!(switch_pipes(false, 1, false)); + } + #[test] fn worker_preference_overrides_the_score() { assert!(switch_pipes(false, PIPELINED_MAX, true)); From fd9042c4bffd9223371032e72575bdbd92d9dfc4 Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 10 Sep 2026 07:47:56 +0900 Subject: [PATCH 04/16] perf: one experiment for the whole proxy, judged on its whole command rate The per-worker probe measured a worker's own rate while the cost of the shared pipes is process-wide: a worker moving its few sessions saw its own latency fall and kept, every worker concluded the same, and the proxy ended at the throughput of full sharding (memtier P16 c200: 3.61M against 4.33M without sharding). Each worker's tick now only publishes its busyness and batch depth and mirrors one process-wide preference; the lead worker runs the experiment when a majority of workers is busy with thin batches, moves every session, and compares the sum of all workers' command rates a second later with the second before. INFO reports pipes_shared, the probe counters and worker_busy. --- docs/architecture.md | 22 ++--- docs/configuration.md | 2 +- docs/operations.md | 2 +- src/admin.rs | 23 ++--- src/client/tuner.rs | 190 +++++++++++++++++++++++++++++------------- src/server.rs | 2 +- src/stats.rs | 19 +++-- 7 files changed, 172 insertions(+), 88 deletions(-) diff --git a/docs/architecture.md b/docs/architecture.md index 7702300..e815c37 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -101,16 +101,18 @@ idle dispatches have run the score down, back to its worker's connections once four busy ones have restored it. On top of that each worker samples its own CPU busyness and the frames its own backend connections batch per write every 100 ms (the in-flight commands per master stand in while no -local traffic flows): busy (85% and above) with thin local batches (under -eight frames per write) for 300 ms only starts an experiment. The worker -moves its sessions to the shared pipes and compares the command rate a -second later against the second before the move; it keeps them there on a -5% gain, and otherwise takes them back and waits 30 seconds before trying -again, doubling that to eight minutes while the answer holds. From the -shared pipes it probes the other way on the same schedule, and lets go -unmeasured once busyness has stayed under 60% or the depth above sixteen -for three seconds — slowly, because moving its sessions away is what lowers -its own busyness. An idle worker never probes. Switches happen only while a session has +local traffic flows), and one worker turns those samples into a single +experiment for the whole proxy, since the cost of the shared pipes is +process-wide and a per-worker verdict measures a free-rider gain. Half the +workers busy (85% and above) with thin local batches (under eight frames +per write) for 300 ms moves every session to the shared pipes; the proxy +then compares the commands it runs a second later against the second +before the move, keeps them on a 5% gain, and otherwise takes them back +and waits 30 seconds before trying again, doubling that to eight minutes +while the answer holds. From the shared pipes it probes the other way on +the same schedule, and lets go unmeasured once fewer than half the workers +have stayed busy for three seconds — slowly, because moving the sessions +away is what lowers their busyness. An idle proxy never probes. Switches happen only while a session has nothing in flight — a session that never idles is paused for one round trip so it can drain and move — so a request-response client ends up on the shared pipe, a pipelining client on a lightly loaded worker on its own diff --git a/docs/configuration.md b/docs/configuration.md index 6e98ac7..f197734 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -14,7 +14,7 @@ See [`mithril.conf.sample`](https://github.com/projecteru2/mithril/blob/master/m | `worker-threads` | int | CPU count | worker threads, one runtime each | | `maxclients` | int | `10000` | client connection cap, enforced at accept | | `backend-conns` | 1..512 | `1` | shared pipelined connections per node per worker, per database in use (a session on `SELECT n` opens its own pool to each node; the exclusive-connection cap for blocking and `WATCH` stays per node across databases) | -| `backend-sharding` | `yes`/`no`/`auto` | `no` | `yes`: one process-wide connection per node per database in use, owned by the worker its address hashes to, which deepens backend pipelines for unpipelined workloads; `auto`: decided per session and per worker — a session that keeps sending one command at a time moves to the shared pipes, a pipelining session stays on its worker's own connections unless that worker is saturated with thin local batches (busy ≥ 85% and fewer than eight frames per backend write), which makes it try the shared pipes for a second and keep its sessions there only where they measure faster (switches happen only while a session has nothing in flight; with the reply cache, only sessions on the shared pipes fill it — the worker-local connections carry no key tracking — while every session reads it) | +| `backend-sharding` | `yes`/`no`/`auto` | `no` | `yes`: one process-wide connection per node per database in use, owned by the worker its address hashes to, which deepens backend pipelines for unpipelined workloads; `auto`: decided per session and once for the whole proxy — a session that keeps sending one command at a time moves to the shared pipes, a pipelining session stays on its worker's own connections unless half the workers are saturated with thin local batches (busy ≥ 85% and fewer than eight frames per backend write), which makes the proxy try the shared pipes for a second and keep them only where its throughput measures higher (switches happen only while a session has nothing in flight; with the reply cache, only sessions on the shared pipes fill it — the worker-local connections carry no key tracking — while every session reads it) | | `reply-cache` | `yes`/`no` | `no` | worker-local GET/MGET reply cache; the servers track the keys the proxy caches (redirected opt-in RESP3 tracking) and invalidate them on change, so hits skip the backend round trip and a session reads its own writes | | `reply-cache-max-bytes` | bytes | `64mb` | per-worker budget for the two cache generations; the live generation holds half of it, so a worker's hot set must fit in half the budget to keep hitting (see operations.md for sizing and the RSS to expect) | | `reply-cache-max-age-secs` | 1..3600 | `10` | staleness backstop for missed invalidations; entries older than this never serve | diff --git a/docs/operations.md b/docs/operations.md index 9c43e1b..b48d848 100644 --- a/docs/operations.md +++ b/docs/operations.md @@ -16,7 +16,7 @@ client's current command in `cmd`. | Clients | `connected_clients` | | CPU | `used_cpu_sys`, `used_cpu_user` | | Stats | `total_connections_received`, `total_commands_processed`, `total_net_input_bytes`, `total_net_output_bytes`, `total_error_replies`, `redirections`, `redirect_waits`, session lifecycle counters (`readers_exited`, `writers_exited`, `sessions_closed`) | -| Mithril | `worker_threads`, `backend_conns_per_node`, `backend_sharding`, `slave_mode`, `reply_cache`, `cache_hits`, `cache_misses`, `cache_invalidations`, `cache_entries`, `cache_bytes`, `cache_flips`, `cache_armed_workers`, `pipe_probes`, `pipe_keeps`, `pipe_reverts` (auto-sharding experiments and how they ended), `worker_commands` (per-worker), `worker_prefers_shared` (per-worker, 1 while the auto tuner sends the worker's sessions to the shared pipes; sessions on the shared pipes by their own score, or under `backend-sharding yes`, are not counted) | +| Mithril | `worker_threads`, `backend_conns_per_node`, `backend_sharding`, `slave_mode`, `reply_cache`, `cache_hits`, `cache_misses`, `cache_invalidations`, `cache_entries`, `cache_bytes`, `cache_flips`, `cache_armed_workers`, `pipes_shared` (1 while the auto tuner holds every session on the shared pipes), `pipe_probes`, `pipe_keeps`, `pipe_reverts` (its experiments and how they ended), `worker_commands` (per-worker), `worker_busy` (per-worker CPU busyness, the tuner's own sample) | | Cluster | `cluster_enabled` (always 1) | | Commandstats | `cmdstat_:calls=` for every command run at least once, subcommands as `cmdstat_client\|list`; counted once accepted, summed over workers | diff --git a/src/admin.rs b/src/admin.rs index a71d6e3..c522877 100644 --- a/src/admin.rs +++ b/src/admin.rs @@ -211,8 +211,8 @@ pub fn info(cfg: &Config, stats: &Stats, started: u64) -> Vec { backend_sharding:{}\r\nslave_mode:{}\r\nreply_cache:{}\r\n\ cache_hits:{}\r\ncache_misses:{}\r\ncache_invalidations:{}\r\n\ cache_armed_workers:{}\r\ncache_entries:{}\r\ncache_bytes:{}\r\n\ - cache_flips:{}\r\npipe_probes:{}\r\npipe_keeps:{}\r\npipe_reverts:{}\r\n\ - worker_commands:{}\r\nworker_prefers_shared:{}\r\n", + cache_flips:{}\r\npipes_shared:{}\r\npipe_probes:{}\r\npipe_keeps:{}\r\n\ + pipe_reverts:{}\r\nworker_commands:{}\r\nworker_busy:{}\r\n", crate::VERSION, std::process::id(), cfg.port, @@ -243,11 +243,12 @@ pub fn info(cfg: &Config, stats: &Stats, started: u64) -> Vec { stats.sum(|w| &w.cache_entries), stats.sum(|w| &w.cache_bytes), stats.sum(|w| &w.cache_flips), - stats.sum(|w| &w.pipe_probes), - stats.sum(|w| &w.pipe_keeps), - stats.sum(|w| &w.pipe_reverts), + u8::from(stats.pipes.prefer.load(Ordering::Relaxed)), + stats.pipes.probes.load(Ordering::Relaxed), + stats.pipes.keeps.load(Ordering::Relaxed), + stats.pipes.reverts.load(Ordering::Relaxed), per_worker(stats, |w| &w.commands), - per_worker(stats, |w| &w.pipe_shared), + per_worker(stats, |w| &w.busy_pct), ); text.push_str("\r\n# Cluster\r\ncluster_enabled:1\r\n\r\n# Commandstats\r\n"); for id in 0..command::entries() as u16 { @@ -534,9 +535,10 @@ mod tests { let get = command::lookup(b"get").unwrap().id; stats.workers[0].calls.at(get).store(3, Ordering::Relaxed); stats.workers[1].calls.at(get).store(4, Ordering::Relaxed); - stats.workers[1].pipe_shared.store(1, Ordering::Relaxed); - stats.workers[0].pipe_probes.store(5, Ordering::Relaxed); - stats.workers[1].pipe_keeps.store(2, Ordering::Relaxed); + stats.workers[1].busy_pct.store(90, Ordering::Relaxed); + stats.pipes.prefer.store(true, Ordering::Relaxed); + stats.pipes.probes.store(5, Ordering::Relaxed); + stats.pipes.keeps.store(2, Ordering::Relaxed); let out = info(&cfg, &stats, 0); let text = String::from_utf8_lossy(&out); let field = |k: &str| { @@ -548,7 +550,8 @@ mod tests { assert_eq!(field("mithril_version").as_deref(), Some(crate::VERSION)); assert_eq!(field("cache_flips").as_deref(), Some("0")); assert_eq!(field("worker_commands").as_deref(), Some("0,0")); - assert_eq!(field("worker_prefers_shared").as_deref(), Some("0,1")); + assert_eq!(field("worker_busy").as_deref(), Some("0,90")); + assert_eq!(field("pipes_shared").as_deref(), Some("1")); assert_eq!(field("pipe_probes").as_deref(), Some("5")); assert_eq!(field("pipe_keeps").as_deref(), Some("2")); assert_eq!(field("pipe_reverts").as_deref(), Some("0")); diff --git a/src/client/tuner.rs b/src/client/tuner.rs index e8dbda2..0f4c468 100644 --- a/src/client/tuner.rs +++ b/src/client/tuner.rs @@ -1,4 +1,4 @@ -//! Pipe selection under auto sharding: the session score and the worker-level tuner. +//! Pipe selection under auto sharding: the session score and the proxy-wide tuner. use std::rc::Rc; use std::sync::atomic::Ordering; @@ -6,20 +6,20 @@ use std::time::{Duration, Instant}; use super::Shared; use super::session::Session; -use crate::stats::{self, WorkerStats}; +use crate::stats::{self, Pipes, Stats}; // pipelining score: a local session shares at 0, a shared one returns at PIPELINED_LOCAL pub(super) const PIPELINED_LOCAL: u8 = 4; const PIPELINED_MAX: u8 = 8; -// what triggers a probe: the worker is busy and its local backend batches would +// what triggers a probe: a worker is busy and its local backend batches would // stay thin (in-flight commands per master) const TUNE_PERIOD: Duration = Duration::from_millis(100); const BUSY_ENTER: u32 = 85; const BUSY_LEAVE: u32 = 60; const DEPTH_ENTER: u64 = 8; const DEPTH_LEAVE: u64 = 16; -// ticks the enter or leave condition must hold: moving sessions off a worker -// lowers its own busyness, so leaving is slow and entering prompt +// ticks the enter or leave condition must hold: moving sessions off the workers +// lowers their busyness, so leaving is slow and entering prompt const ENTER_TICKS: u32 = 3; const LEAVE_TICKS: u32 = 30; // the command-rate window, ten ticks of 100 ms, so a rate is commands per second @@ -29,7 +29,7 @@ const SETTLE_TICKS: u32 = 10; const PROBE_TICKS: u32 = SETTLE_TICKS + RATE_TICKS as u32; // the shared pipes are kept only where they measure this much faster const KEEP_GAIN_PCT: u64 = 5; -// commands per second under which the worker is idle and its rates say nothing +// commands per second per worker under which the proxy is idle and its rates say nothing const MIN_RATE_PER_SEC: u64 = 100; // ticks to the next probe, doubling while decisions confirm the current state const PROBE_BACKOFF_TICKS: u32 = 300; @@ -65,7 +65,7 @@ enum Phase { }, } -// the worker's experiment: the busy-and-thin rule triggers it, the measured command rate decides +// the proxy's experiment: a busy-and-thin majority triggers it, the measured command rate decides #[derive(Default)] struct Probe { prefer: bool, @@ -76,6 +76,7 @@ struct Probe { streak: u32, wait: u32, confirmed: u32, + floor: u64, ring: [u64; RATE_TICKS], at: usize, seen: usize, @@ -83,24 +84,28 @@ struct Probe { } impl Probe { - fn tick(&mut self, busy_pct: u32, depth: u64, commands_now: u64) -> bool { + fn new(workers: usize) -> Probe { + Probe { + floor: MIN_RATE_PER_SEC * workers as u64, + ..Probe::default() + } + } + + fn tick(&mut self, busy_thin: bool, still_busy: bool, commands_now: u64) -> bool { self.record(commands_now); self.wait = self.wait.saturating_sub(1); match self.phase { Phase::Probing { baseline, ticks } => self.probing(baseline, ticks), - Phase::Steady if self.prefer => self.leave(busy_pct, depth), - Phase::Steady => self.enter(busy_pct, depth), + Phase::Steady if self.prefer => self.leave(still_busy), + Phase::Steady => self.enter(busy_thin), } self.prefer } - fn publish(&self, wstats: &WorkerStats) { - wstats - .pipe_shared - .store(u64::from(self.prefer), Ordering::Relaxed); - wstats.pipe_probes.store(self.probes, Ordering::Relaxed); - wstats.pipe_keeps.store(self.keeps, Ordering::Relaxed); - wstats.pipe_reverts.store(self.reverts, Ordering::Relaxed); + fn publish(&self, pipes: &Pipes) { + pipes.probes.store(self.probes, Ordering::Relaxed); + pipes.keeps.store(self.keeps, Ordering::Relaxed); + pipes.reverts.store(self.reverts, Ordering::Relaxed); } fn record(&mut self, commands_now: u64) { @@ -114,8 +119,8 @@ impl Probe { self.ring.iter().sum() } - fn enter(&mut self, busy_pct: u32, depth: u64) { - if busy_pct < BUSY_ENTER || depth >= DEPTH_ENTER { + fn enter(&mut self, busy_thin: bool) { + if !busy_thin { self.streak = 0; return; } @@ -125,8 +130,8 @@ impl Probe { } } - fn leave(&mut self, busy_pct: u32, depth: u64) { - if busy_pct > BUSY_LEAVE && depth < DEPTH_LEAVE { + fn leave(&mut self, still_busy: bool) { + if still_busy { self.streak = 0; self.start_probe(); return; @@ -142,7 +147,7 @@ impl Probe { fn start_probe(&mut self) { let baseline = self.rate(); - if self.wait > 0 || self.seen < RATE_TICKS || baseline < MIN_RATE_PER_SEC { + if self.wait > 0 || self.seen < RATE_TICKS || baseline < self.floor { return; } self.probes += 1; @@ -182,13 +187,13 @@ impl Probe { } } -/// Samples this worker's CPU busyness, batch depth and command rate, and publishes -/// whether its sessions should prefer the shared pipes. -pub async fn auto_tuner(shared: Rc) { +/// Publishes this worker's CPU busyness and batch depth and mirrors the process-wide +/// pipe preference; the lead worker also runs the experiment that sets it. +pub async fn auto_tuner(shared: Rc, lead: bool) { let mut last_ticks = stats::thread_cpu_ticks(); let mut last_at = Instant::now(); let mut busy_x16 = 0u32; - let mut probe = Probe::default(); + let mut probe = lead.then(|| Probe::new(shared.stats.workers.len())); loop { tokio::time::sleep(TUNE_PERIOD).await; let now = Instant::now(); @@ -204,12 +209,45 @@ pub async fn auto_tuner(shared: Rc) { let masters = shared.topo.load().masters.len().max(1) as u64; let (measured, writes) = shared.backends.batch_depth(); let depth = tune_depth(measured, writes, shared.inflight.get(), masters); - let commands = shared.wstats.commands.load(Ordering::Relaxed); + shared + .wstats + .busy_pct + .store(u64::from((busy_x16 + 8) / 16), Ordering::Relaxed); + shared.wstats.batch_depth.store(depth, Ordering::Relaxed); + if let Some(probe) = probe.as_mut() { + conduct(&shared.stats, probe); + } shared .prefer_shared - .set(probe.tick((busy_x16 + 8) / 16, depth, commands)); - probe.publish(&shared.wstats); + .set(shared.stats.pipes.prefer.load(Ordering::Relaxed)); + } +} + +// the cost of the shared pipes is process-wide, so one experiment moves every session +// and the whole proxy's command rate judges it +fn conduct(stats: &Stats, probe: &mut Probe) { + let (busy_thin, still_busy, commands) = survey(stats); + let prefer = probe.tick(busy_thin, still_busy, commands); + stats.pipes.prefer.store(prefer, Ordering::Relaxed); + probe.publish(&stats.pipes); +} + +// a worker that has not ticked yet reads as idle and counts against both majorities +fn survey(stats: &Stats) -> (bool, bool, u64) { + let (mut busy_thin, mut still_busy, mut commands) = (0usize, 0usize, 0u64); + for w in &stats.workers { + let busy = w.busy_pct.load(Ordering::Relaxed); + let depth = w.batch_depth.load(Ordering::Relaxed); + busy_thin += usize::from(busy >= u64::from(BUSY_ENTER) && depth < DEPTH_ENTER); + still_busy += usize::from(busy > u64::from(BUSY_LEAVE) && depth < DEPTH_LEAVE); + commands += w.commands.load(Ordering::Relaxed); } + let workers = stats.workers.len(); + ( + busy_thin * 2 >= workers, + still_busy * 2 >= workers, + commands, + ) } // decided on an idle dispatch from the score before that dispatch counts @@ -307,77 +345,109 @@ mod tests { #[test] fn a_probe_that_gains_keeps_the_shared_pipes() { - let mut rig = Rig::default(); - assert!(!rig.run(RATE_TICKS as u32 - 1, 1_000, 90, 1)); - assert!(rig.run(1, 1_000, 90, 1)); + let mut rig = Rig::new(); + assert!(!rig.run(RATE_TICKS as u32 - 1, 1_000, true, true)); + assert!(rig.run(1, 1_000, true, true)); assert_eq!(rig.probe.probes, 1); - rig.run(SETTLE_TICKS, 1_000, 90, 1); - assert!(rig.run(RATE_TICKS as u32, 1_200, 90, 1)); + rig.run(SETTLE_TICKS, 1_000, true, true); + assert!(rig.run(RATE_TICKS as u32, 1_200, true, true)); assert_eq!((rig.probe.keeps, rig.probe.reverts), (1, 0)); assert_eq!(rig.probe.wait, PROBE_BACKOFF_TICKS); } #[test] fn a_probe_that_loses_reverts_and_waits_twice_as_long_after_the_second() { - let mut rig = Rig::default(); - rig.run(RATE_TICKS as u32, 1_000, 90, 1); - rig.run(PROBE_TICKS - 1, 800, 90, 1); - assert!(!rig.run(1, 800, 90, 1)); + let mut rig = Rig::new(); + rig.run(RATE_TICKS as u32, 1_000, true, true); + rig.run(PROBE_TICKS - 1, 800, true, true); + assert!(!rig.run(1, 800, true, true)); assert_eq!((rig.probe.keeps, rig.probe.reverts), (0, 1)); assert_eq!(rig.probe.wait, PROBE_BACKOFF_TICKS); - assert!(!rig.run(PROBE_BACKOFF_TICKS - 1, 1_000, 90, 1)); + assert!(!rig.run(PROBE_BACKOFF_TICKS - 1, 1_000, true, true)); assert_eq!(rig.probe.probes, 1); - assert!(rig.run(1, 1_000, 90, 1)); + assert!(rig.run(1, 1_000, true, true)); assert_eq!(rig.probe.probes, 2); - rig.run(PROBE_TICKS - 1, 500, 90, 1); - assert!(!rig.run(1, 500, 90, 1)); + rig.run(PROBE_TICKS - 1, 500, true, true); + assert!(!rig.run(1, 500, true, true)); assert_eq!(rig.probe.reverts, 2); assert_eq!(rig.probe.wait, PROBE_BACKOFF_TICKS * 2); } #[test] - fn low_busyness_leaves_the_shared_pipes_without_a_backoff() { - let mut rig = Rig::default(); - rig.run(RATE_TICKS as u32, 1_000, 90, 1); - assert!(rig.run(PROBE_TICKS, 1_200, 90, 1)); - assert!(rig.run(LEAVE_TICKS - 1, 1_200, BUSY_LEAVE, 1)); - assert!(!rig.run(1, 1_200, BUSY_LEAVE, 1)); + fn a_quiet_majority_leaves_the_shared_pipes_without_a_backoff() { + let mut rig = Rig::new(); + rig.run(RATE_TICKS as u32, 1_000, true, true); + assert!(rig.run(PROBE_TICKS, 1_200, true, true)); + assert!(rig.run(LEAVE_TICKS - 1, 1_200, false, false)); + assert!(!rig.run(1, 1_200, false, false)); assert_eq!(rig.probe.wait, 0); assert_eq!((rig.probe.probes, rig.probe.keeps), (1, 1)); } #[test] fn a_reverse_probe_returns_to_local_when_the_shared_pipes_are_not_faster() { - let mut rig = Rig::default(); - rig.run(RATE_TICKS as u32, 1_000, 90, 1); - rig.run(PROBE_TICKS, 1_200, 90, 1); - assert!(!rig.run(PROBE_BACKOFF_TICKS, 1_200, 90, 1)); + let mut rig = Rig::new(); + rig.run(RATE_TICKS as u32, 1_000, true, true); + rig.run(PROBE_TICKS, 1_200, true, true); + assert!(!rig.run(PROBE_BACKOFF_TICKS, 1_200, false, true)); assert_eq!(rig.probe.probes, 2); - rig.run(PROBE_TICKS - 1, 1_200, 90, 1); - assert!(!rig.run(1, 1_200, 90, 1)); + rig.run(PROBE_TICKS - 1, 1_200, false, true); + assert!(!rig.run(1, 1_200, false, true)); assert_eq!((rig.probe.keeps, rig.probe.reverts), (1, 1)); assert_eq!(rig.probe.wait, PROBE_BACKOFF_TICKS); } #[test] - fn an_idle_worker_never_probes() { - let mut rig = Rig::default(); - assert!(!rig.run(PROBE_BACKOFF_TICKS, MIN_RATE_PER_SEC - 10, 90, 1)); + fn an_idle_proxy_never_probes() { + let mut rig = Rig::new(); + let idle = MIN_RATE_PER_SEC * WORKERS as u64 - 10; + assert!(!rig.run(PROBE_BACKOFF_TICKS, idle, true, true)); assert_eq!(rig.probe.probes, 0); } - #[derive(Default)] + #[test] + fn the_experiment_starts_only_on_a_busy_and_thin_majority() { + let stats = Stats::new(3); + assert_eq!(survey(&stats), (false, false, 0)); + let set = |i: usize, busy: u64, depth: u64| { + stats.workers[i].busy_pct.store(busy, Ordering::Relaxed); + stats.workers[i].batch_depth.store(depth, Ordering::Relaxed); + stats.workers[i].commands.store(100, Ordering::Relaxed); + }; + set(0, 90, 1); + assert_eq!(survey(&stats), (false, false, 100)); + set(1, 90, 1); + assert_eq!(survey(&stats), (true, true, 200)); + set(1, 90, DEPTH_ENTER); + assert_eq!(survey(&stats), (false, true, 200)); + set(0, u64::from(BUSY_LEAVE), 1); + assert_eq!(survey(&stats), (false, false, 200)); + let even = Stats::new(4); + even.workers[0].busy_pct.store(90, Ordering::Relaxed); + even.workers[1].busy_pct.store(90, Ordering::Relaxed); + assert_eq!(survey(&even), (true, true, 0)); + } + + const WORKERS: usize = 4; + struct Rig { probe: Probe, commands: u64, } impl Rig { - fn run(&mut self, ticks: u32, per_sec: u64, busy_pct: u32, depth: u64) -> bool { + fn new() -> Rig { + Rig { + probe: Probe::new(WORKERS), + commands: 0, + } + } + + fn run(&mut self, ticks: u32, per_sec: u64, busy_thin: bool, still_busy: bool) -> bool { let mut prefer = self.probe.prefer; for _ in 0..ticks { self.commands += per_sec / RATE_TICKS as u64; - prefer = self.probe.tick(busy_pct, depth, self.commands); + prefer = self.probe.tick(busy_thin, still_busy, self.commands); } prefer } diff --git a/src/server.rs b/src/server.rs index 974c320..65bca3f 100644 --- a/src/server.rs +++ b/src/server.rs @@ -454,7 +454,7 @@ fn worker_thread( prefer_shared: Cell::new(false), }); if shared.cfg.backend_sharding == Sharding::Auto { - tokio::task::spawn_local(auto_tuner(shared.clone())); + tokio::task::spawn_local(auto_tuner(shared.clone(), worker == 0)); } let mut next_client: u64 = worker as u64; while let Some(mut admitted) = conn_rx.recv().await { diff --git a/src/stats.rs b/src/stats.rs index 2fcf3c4..d7ad187 100644 --- a/src/stats.rs +++ b/src/stats.rs @@ -2,7 +2,7 @@ use std::collections::{HashMap, VecDeque}; use std::net::SocketAddr; -use std::sync::atomic::{AtomicI64, AtomicU16, AtomicU64, AtomicUsize, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicI64, AtomicU16, AtomicU64, AtomicUsize, Ordering}; use std::sync::{Arc, Mutex, MutexGuard, PoisonError}; use std::time::{Instant, SystemTime, UNIX_EPOCH}; @@ -48,10 +48,8 @@ pub struct WorkerStats { pub cache_entries: AtomicU64, pub cache_bytes: AtomicU64, pub cache_flips: AtomicU64, - pub pipe_shared: AtomicU64, - pub pipe_probes: AtomicU64, - pub pipe_keeps: AtomicU64, - pub pipe_reverts: AtomicU64, + pub busy_pct: AtomicU64, + pub batch_depth: AtomicU64, pub calls: Calls, } @@ -143,6 +141,15 @@ impl Default for Slowlog { } } +/// The pipe preference the auto tuner sets for the whole proxy, and its experiments. +#[derive(Default)] +pub struct Pipes { + pub prefer: AtomicBool, + pub probes: AtomicU64, + pub keeps: AtomicU64, + pub reverts: AtomicU64, +} + /// Process-wide stats shared across workers. pub struct Stats { pub workers: Vec>, @@ -150,6 +157,7 @@ pub struct Stats { pub total_connections: AtomicU64, pub registry: Mutex>, pub slowlog: Slowlog, + pub pipes: Pipes, epoch: Instant, } @@ -161,6 +169,7 @@ impl Stats { total_connections: AtomicU64::new(0), registry: Mutex::new(HashMap::new()), slowlog: Slowlog::default(), + pipes: Pipes::default(), epoch: Instant::now(), }) } From 6314cfe53ee3f0ea8d23dc05a6172cb452e6d1b6 Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 10 Sep 2026 07:55:41 +0900 Subject: [PATCH 05/16] docs: the INFO field names the tuner's preference --- docs/operations.md | 2 +- src/admin.rs | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/docs/operations.md b/docs/operations.md index b48d848..e2087ef 100644 --- a/docs/operations.md +++ b/docs/operations.md @@ -16,7 +16,7 @@ client's current command in `cmd`. | Clients | `connected_clients` | | CPU | `used_cpu_sys`, `used_cpu_user` | | Stats | `total_connections_received`, `total_commands_processed`, `total_net_input_bytes`, `total_net_output_bytes`, `total_error_replies`, `redirections`, `redirect_waits`, session lifecycle counters (`readers_exited`, `writers_exited`, `sessions_closed`) | -| Mithril | `worker_threads`, `backend_conns_per_node`, `backend_sharding`, `slave_mode`, `reply_cache`, `cache_hits`, `cache_misses`, `cache_invalidations`, `cache_entries`, `cache_bytes`, `cache_flips`, `cache_armed_workers`, `pipes_shared` (1 while the auto tuner holds every session on the shared pipes), `pipe_probes`, `pipe_keeps`, `pipe_reverts` (its experiments and how they ended), `worker_commands` (per-worker), `worker_busy` (per-worker CPU busyness, the tuner's own sample) | +| Mithril | `worker_threads`, `backend_conns_per_node`, `backend_sharding`, `slave_mode`, `reply_cache`, `cache_hits`, `cache_misses`, `cache_invalidations`, `cache_entries`, `cache_bytes`, `cache_flips`, `cache_armed_workers`, `pipes_prefer_shared` (the auto tuner's process-wide preference: 1 while it sends every session to the shared pipes; a session on them by its own score is not counted), `pipe_probes`, `pipe_keeps`, `pipe_reverts` (its experiments and how they ended), `worker_commands` (per-worker), `worker_busy` (per-worker CPU busyness, the tuner's own sample) | | Cluster | `cluster_enabled` (always 1) | | Commandstats | `cmdstat_:calls=` for every command run at least once, subcommands as `cmdstat_client\|list`; counted once accepted, summed over workers | diff --git a/src/admin.rs b/src/admin.rs index c522877..cd75fed 100644 --- a/src/admin.rs +++ b/src/admin.rs @@ -211,7 +211,7 @@ pub fn info(cfg: &Config, stats: &Stats, started: u64) -> Vec { backend_sharding:{}\r\nslave_mode:{}\r\nreply_cache:{}\r\n\ cache_hits:{}\r\ncache_misses:{}\r\ncache_invalidations:{}\r\n\ cache_armed_workers:{}\r\ncache_entries:{}\r\ncache_bytes:{}\r\n\ - cache_flips:{}\r\npipes_shared:{}\r\npipe_probes:{}\r\npipe_keeps:{}\r\n\ + cache_flips:{}\r\npipes_prefer_shared:{}\r\npipe_probes:{}\r\npipe_keeps:{}\r\n\ pipe_reverts:{}\r\nworker_commands:{}\r\nworker_busy:{}\r\n", crate::VERSION, std::process::id(), @@ -551,7 +551,7 @@ mod tests { assert_eq!(field("cache_flips").as_deref(), Some("0")); assert_eq!(field("worker_commands").as_deref(), Some("0,0")); assert_eq!(field("worker_busy").as_deref(), Some("0,90")); - assert_eq!(field("pipes_shared").as_deref(), Some("1")); + assert_eq!(field("pipes_prefer_shared").as_deref(), Some("1")); assert_eq!(field("pipe_probes").as_deref(), Some("5")); assert_eq!(field("pipe_keeps").as_deref(), Some("2")); assert_eq!(field("pipe_reverts").as_deref(), Some("0")); From d49a99ec3e009fa56275ec5fb86fefdd7a8a4bec Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 10 Sep 2026 08:36:52 +0900 Subject: [PATCH 06/16] perf: a changed workload ends the tuner's wait The reverse probe ran on a fixed schedule (30 s, doubling), so after a workload change the proxy stayed on the losing answer until the next scheduled probe: the memtier P16 c200 cell, entered from the P16 c1000 cells on the shared pipes, measured 3.91M against 4.35M local because the flip came part-way through. The rate the current state was chosen on is now remembered, and a rate a quarter away from it ends the wait, so the next probe starts as soon as the trigger holds; with that the base wait is a minute, which keeps the scheduled reverse probe out of a steady workload's way. --- docs/architecture.md | 6 ++++-- src/client/tuner.rs | 26 +++++++++++++++++++++++++- 2 files changed, 29 insertions(+), 3 deletions(-) diff --git a/docs/architecture.md b/docs/architecture.md index e815c37..611220f 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -108,8 +108,10 @@ workers busy (85% and above) with thin local batches (under eight frames per write) for 300 ms moves every session to the shared pipes; the proxy then compares the commands it runs a second later against the second before the move, keeps them on a 5% gain, and otherwise takes them back -and waits 30 seconds before trying again, doubling that to eight minutes -while the answer holds. From the shared pipes it probes the other way on +and waits a minute before trying again, doubling that to eight minutes +while the answer holds; a command rate that moves a quarter away from the +one the answer was measured on ends the wait, since the workload it was +measured on is gone. From the shared pipes it probes the other way on the same schedule, and lets go unmeasured once fewer than half the workers have stayed busy for three seconds — slowly, because moving the sessions away is what lowers their busyness. An idle proxy never probes. Switches happen only while a session has diff --git a/src/client/tuner.rs b/src/client/tuner.rs index 0f4c468..21a2505 100644 --- a/src/client/tuner.rs +++ b/src/client/tuner.rs @@ -32,9 +32,11 @@ const KEEP_GAIN_PCT: u64 = 5; // commands per second per worker under which the proxy is idle and its rates say nothing const MIN_RATE_PER_SEC: u64 = 100; // ticks to the next probe, doubling while decisions confirm the current state -const PROBE_BACKOFF_TICKS: u32 = 300; +const PROBE_BACKOFF_TICKS: u32 = 600; const PROBE_BACKOFF_MAX_TICKS: u32 = 4800; const PROBE_DOUBLINGS: u32 = (PROBE_BACKOFF_MAX_TICKS / PROBE_BACKOFF_TICKS).ilog2(); +// a rate this far from the one the current state was chosen on is a changed workload +const RATE_SHIFT_PCT: u64 = 25; impl Session { // an unpipelined session gains from the deeper batches of the shared pipe, a @@ -76,6 +78,7 @@ struct Probe { streak: u32, wait: u32, confirmed: u32, + decided: u64, floor: u64, ring: [u64; RATE_TICKS], at: usize, @@ -94,6 +97,12 @@ impl Probe { fn tick(&mut self, busy_thin: bool, still_busy: bool, commands_now: u64) -> bool { self.record(commands_now); self.wait = self.wait.saturating_sub(1); + if self.decided > 0 + && self.rate().abs_diff(self.decided) * 100 > self.decided * RATE_SHIFT_PCT + { + self.decided = 0; + self.wait = 0; + } match self.phase { Phase::Probing { baseline, ticks } => self.probing(baseline, ticks), Phase::Steady if self.prefer => self.leave(still_busy), @@ -141,6 +150,7 @@ impl Probe { // the load went away rather than a measurement: the next busy stretch probes at once self.streak = 0; self.wait = 0; + self.decided = 0; self.prefer = false; } } @@ -182,6 +192,7 @@ impl Probe { } else { self.reverts += 1; } + self.decided = if keep { shared } else { local }; self.prefer = keep; self.phase = Phase::Steady; } @@ -397,6 +408,19 @@ mod tests { assert_eq!(rig.probe.wait, PROBE_BACKOFF_TICKS); } + #[test] + fn a_changed_workload_reprobes_at_once() { + let mut rig = Rig::new(); + rig.run(RATE_TICKS as u32, 1_000, true, true); + assert!(rig.run(PROBE_TICKS, 1_200, true, true)); + assert_eq!(rig.probe.wait, PROBE_BACKOFF_TICKS); + assert!(rig.run(RATE_TICKS as u32, 1_200, true, true)); + assert!(rig.run(RATE_TICKS as u32, 600, false, false)); + assert_eq!(rig.probe.wait, 0); + assert!(!rig.run(1, 600, true, true)); + assert_eq!(rig.probe.probes, 2); + } + #[test] fn an_idle_proxy_never_probes() { let mut rig = Rig::new(); From ea3c0d7b814598393bf578208017f114ec48e0fe Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 10 Sep 2026 08:39:56 +0900 Subject: [PATCH 07/16] fix: a losing trial left in the rate ring is not a changed workload After a decision the ring still holds the trial's rate for one window; a trial a quarter slower than the winner read as a workload shift on the next tick, ended the wait and started the same probe again. The shift check now waits one rate window after every decision. --- src/client/tuner.rs | 18 +++++++++++++++++- 1 file changed, 17 insertions(+), 1 deletion(-) diff --git a/src/client/tuner.rs b/src/client/tuner.rs index 21a2505..419d5b4 100644 --- a/src/client/tuner.rs +++ b/src/client/tuner.rs @@ -79,6 +79,7 @@ struct Probe { wait: u32, confirmed: u32, decided: u64, + settling: u32, floor: u64, ring: [u64; RATE_TICKS], at: usize, @@ -97,7 +98,10 @@ impl Probe { fn tick(&mut self, busy_thin: bool, still_busy: bool, commands_now: u64) -> bool { self.record(commands_now); self.wait = self.wait.saturating_sub(1); - if self.decided > 0 + // the ring still holds the losing trial for one window after a decision + self.settling = self.settling.saturating_sub(1); + if self.settling == 0 + && self.decided > 0 && self.rate().abs_diff(self.decided) * 100 > self.decided * RATE_SHIFT_PCT { self.decided = 0; @@ -193,6 +197,7 @@ impl Probe { self.reverts += 1; } self.decided = if keep { shared } else { local }; + self.settling = RATE_TICKS as u32; self.prefer = keep; self.phase = Phase::Steady; } @@ -421,6 +426,17 @@ mod tests { assert_eq!(rig.probe.probes, 2); } + #[test] + fn a_losing_trial_left_in_the_ring_is_not_a_changed_workload() { + let mut rig = Rig::new(); + rig.run(RATE_TICKS as u32, 1_000, true, true); + assert!(!rig.run(PROBE_TICKS, 600, true, true)); + assert_eq!(rig.probe.reverts, 1); + assert!(!rig.run(RATE_TICKS as u32, 1_000, true, true)); + assert_eq!(rig.probe.probes, 1); + assert_eq!(rig.probe.wait, PROBE_BACKOFF_TICKS - RATE_TICKS as u32); + } + #[test] fn an_idle_proxy_never_probes() { let mut rig = Rig::new(); From a96c935dc9855f3b39bc89b1eb111021fceb20c0 Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 10 Sep 2026 09:26:06 +0900 Subject: [PATCH 08/16] fix: a probe needs a steady baseline A probe started right after a rate shift took its baseline from the second that held the shift itself, a gap between two benchmark phases in the acceptance runs, and a reverse trial then beat a near-zero baseline: the proxy went local for a whole cell (5.5M against 8.9M). The ten ticks of the baseline must now lie within a factor of two of each other; a gap or a ramp makes the probe wait until the rate has settled. An unsteady trial window needs no guard, since a low reading keeps the previous state. --- docs/architecture.md | 3 ++- src/client/tuner.rs | 23 ++++++++++++++++++++++- 2 files changed, 24 insertions(+), 2 deletions(-) diff --git a/docs/architecture.md b/docs/architecture.md index 611220f..0e939ca 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -105,7 +105,8 @@ local traffic flows), and one worker turns those samples into a single experiment for the whole proxy, since the cost of the shared pipes is process-wide and a per-worker verdict measures a free-rider gain. Half the workers busy (85% and above) with thin local batches (under eight frames -per write) for 300 ms moves every session to the shared pipes; the proxy +per write) for 300 ms, once the command rate has been steady for a second, +moves every session to the shared pipes; the proxy then compares the commands it runs a second later against the second before the move, keeps them on a 5% gain, and otherwise takes them back and waits a minute before trying again, doubling that to eight minutes diff --git a/src/client/tuner.rs b/src/client/tuner.rs index 419d5b4..47e4283 100644 --- a/src/client/tuner.rs +++ b/src/client/tuner.rs @@ -37,6 +37,8 @@ const PROBE_BACKOFF_MAX_TICKS: u32 = 4800; const PROBE_DOUBLINGS: u32 = (PROBE_BACKOFF_MAX_TICKS / PROBE_BACKOFF_TICKS).ilog2(); // a rate this far from the one the current state was chosen on is a changed workload const RATE_SHIFT_PCT: u64 = 25; +// a baseline whose ticks spread wider than this ratio holds a gap or a ramp, not a rate +const STEADY_RATIO: u64 = 2; impl Session { // an unpipelined session gains from the deeper batches of the shared pipe, a @@ -132,6 +134,14 @@ impl Probe { self.ring.iter().sum() } + fn steady(&self) -> bool { + let (min, max) = self + .ring + .iter() + .fold((u64::MAX, 0), |(lo, hi), &v| (lo.min(v), hi.max(v))); + max <= min * STEADY_RATIO + } + fn enter(&mut self, busy_thin: bool) { if !busy_thin { self.streak = 0; @@ -161,7 +171,7 @@ impl Probe { fn start_probe(&mut self) { let baseline = self.rate(); - if self.wait > 0 || self.seen < RATE_TICKS || baseline < self.floor { + if self.wait > 0 || self.seen < RATE_TICKS || baseline < self.floor || !self.steady() { return; } self.probes += 1; @@ -437,6 +447,17 @@ mod tests { assert_eq!(rig.probe.wait, PROBE_BACKOFF_TICKS - RATE_TICKS as u32); } + #[test] + fn a_probe_waits_for_a_steady_baseline() { + let mut rig = Rig::new(); + rig.run(RATE_TICKS as u32 - 1, 1_000, true, true); + rig.run(1, 0, true, true); + assert!(!rig.run(RATE_TICKS as u32 - 1, 1_000, true, true)); + assert_eq!(rig.probe.probes, 0); + assert!(rig.run(1, 1_000, true, true)); + assert_eq!(rig.probe.probes, 1); + } + #[test] fn an_idle_proxy_never_probes() { let mut rig = Rig::new(); From 3be40087e8dbeaa1b0a68945ed2dd1c6ba726f6e Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 10 Sep 2026 09:31:51 +0900 Subject: [PATCH 09/16] fix: a probe after a rate shift waits one window, and a steady baseline spreads under a quarter A step smaller than the spread bound left both rates in the baseline second, so a trial compared against a blend and the wrong answer held for the backoff. A detected shift now starts the same settling window a decision does, and the baseline's ticks must lie within a quarter of each other, the same margin the shift check uses. --- src/client/tuner.rs | 19 ++++++++++++++----- 1 file changed, 14 insertions(+), 5 deletions(-) diff --git a/src/client/tuner.rs b/src/client/tuner.rs index 47e4283..2dd7906 100644 --- a/src/client/tuner.rs +++ b/src/client/tuner.rs @@ -37,8 +37,8 @@ const PROBE_BACKOFF_MAX_TICKS: u32 = 4800; const PROBE_DOUBLINGS: u32 = (PROBE_BACKOFF_MAX_TICKS / PROBE_BACKOFF_TICKS).ilog2(); // a rate this far from the one the current state was chosen on is a changed workload const RATE_SHIFT_PCT: u64 = 25; -// a baseline whose ticks spread wider than this ratio holds a gap or a ramp, not a rate -const STEADY_RATIO: u64 = 2; +// a baseline whose ticks spread wider than this holds a gap or a ramp, not a rate +const STEADY_SPREAD_PCT: u64 = 25; impl Session { // an unpipelined session gains from the deeper batches of the shared pipe, a @@ -100,7 +100,8 @@ impl Probe { fn tick(&mut self, busy_thin: bool, still_busy: bool, commands_now: u64) -> bool { self.record(commands_now); self.wait = self.wait.saturating_sub(1); - // the ring still holds the losing trial for one window after a decision + // the ring still holds the losing trial for one window after a decision, and + // both rates for one window after a shift self.settling = self.settling.saturating_sub(1); if self.settling == 0 && self.decided > 0 @@ -108,6 +109,7 @@ impl Probe { { self.decided = 0; self.wait = 0; + self.settling = RATE_TICKS as u32; } match self.phase { Phase::Probing { baseline, ticks } => self.probing(baseline, ticks), @@ -139,7 +141,7 @@ impl Probe { .ring .iter() .fold((u64::MAX, 0), |(lo, hi), &v| (lo.min(v), hi.max(v))); - max <= min * STEADY_RATIO + max * 100 <= min * (100 + STEADY_SPREAD_PCT) } fn enter(&mut self, busy_thin: bool) { @@ -171,7 +173,12 @@ impl Probe { fn start_probe(&mut self) { let baseline = self.rate(); - if self.wait > 0 || self.seen < RATE_TICKS || baseline < self.floor || !self.steady() { + if self.wait > 0 + || self.settling > 0 + || self.seen < RATE_TICKS + || baseline < self.floor + || !self.steady() + { return; } self.probes += 1; @@ -432,6 +439,8 @@ mod tests { assert!(rig.run(RATE_TICKS as u32, 1_200, true, true)); assert!(rig.run(RATE_TICKS as u32, 600, false, false)); assert_eq!(rig.probe.wait, 0); + assert!(rig.run(1, 600, true, true)); + assert!(rig.run(RATE_TICKS as u32, 600, false, false)); assert!(!rig.run(1, 600, true, true)); assert_eq!(rig.probe.probes, 2); } From b1db53425e740d9419aefdda3e554d8c74301f2b Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 10 Sep 2026 10:09:09 +0900 Subject: [PATCH 10/16] perf: a rate shift is judged once the rate is steady, and a pause changes nothing Every pause between two benchmark phases read as a changed workload and bought a probe on the losing arm as soon as the rate settled: the SET cell ran 4% under full sharding from one two-second local window. The shift is now only noted; once the rate is steady again it is compared with the decided one, and only a rate that settled elsewhere voids the decision. A pause, or a blip, that returns to the same rate keeps the schedule. --- docs/architecture.md | 7 +++--- src/client/tuner.rs | 54 ++++++++++++++++++++++++++++++++------------ 2 files changed, 44 insertions(+), 17 deletions(-) diff --git a/docs/architecture.md b/docs/architecture.md index 0e939ca..0f4b8b3 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -110,9 +110,10 @@ moves every session to the shared pipes; the proxy then compares the commands it runs a second later against the second before the move, keeps them on a 5% gain, and otherwise takes them back and waits a minute before trying again, doubling that to eight minutes -while the answer holds; a command rate that moves a quarter away from the -one the answer was measured on ends the wait, since the workload it was -measured on is gone. From the shared pipes it probes the other way on +while the answer holds; a command rate that settles a quarter away from +the one the answer was measured on ends the wait, since that workload is +gone, while a pause that returns to it changes nothing. From the shared +pipes it probes the other way on the same schedule, and lets go unmeasured once fewer than half the workers have stayed busy for three seconds — slowly, because moving the sessions away is what lowers their busyness. An idle proxy never probes. Switches happen only while a session has diff --git a/src/client/tuner.rs b/src/client/tuner.rs index 2dd7906..27bf6e7 100644 --- a/src/client/tuner.rs +++ b/src/client/tuner.rs @@ -82,6 +82,7 @@ struct Probe { confirmed: u32, decided: u64, settling: u32, + shifted: bool, floor: u64, ring: [u64; RATE_TICKS], at: usize, @@ -100,17 +101,7 @@ impl Probe { fn tick(&mut self, busy_thin: bool, still_busy: bool, commands_now: u64) -> bool { self.record(commands_now); self.wait = self.wait.saturating_sub(1); - // the ring still holds the losing trial for one window after a decision, and - // both rates for one window after a shift - self.settling = self.settling.saturating_sub(1); - if self.settling == 0 - && self.decided > 0 - && self.rate().abs_diff(self.decided) * 100 > self.decided * RATE_SHIFT_PCT - { - self.decided = 0; - self.wait = 0; - self.settling = RATE_TICKS as u32; - } + self.settle(); match self.phase { Phase::Probing { baseline, ticks } => self.probing(baseline, ticks), Phase::Steady if self.prefer => self.leave(still_busy), @@ -125,6 +116,28 @@ impl Probe { pipes.reverts.store(self.reverts, Ordering::Relaxed); } + // a rate that leaves the decided one is judged again once it is steady: back + // near it the workload only paused, away from it the decision is void + fn settle(&mut self) { + self.settling = self.settling.saturating_sub(1); + if self.settling > 0 || self.decided == 0 { + return; + } + let moved = self.rate().abs_diff(self.decided) * 100 > self.decided * RATE_SHIFT_PCT; + if !self.shifted { + if moved { + self.shifted = true; + self.settling = RATE_TICKS as u32; + } + } else if self.steady() { + self.shifted = false; + if moved { + self.decided = 0; + self.wait = 0; + } + } + } + fn record(&mut self, commands_now: u64) { self.ring[self.at] = commands_now.saturating_sub(self.last); self.at = (self.at + 1) % RATE_TICKS; @@ -167,6 +180,7 @@ impl Probe { self.streak = 0; self.wait = 0; self.decided = 0; + self.shifted = false; self.prefer = false; } } @@ -175,6 +189,7 @@ impl Probe { let baseline = self.rate(); if self.wait > 0 || self.settling > 0 + || self.shifted || self.seen < RATE_TICKS || baseline < self.floor || !self.steady() @@ -215,6 +230,7 @@ impl Probe { } self.decided = if keep { shared } else { local }; self.settling = RATE_TICKS as u32; + self.shifted = false; self.prefer = keep; self.phase = Phase::Steady; } @@ -437,10 +453,8 @@ mod tests { assert!(rig.run(PROBE_TICKS, 1_200, true, true)); assert_eq!(rig.probe.wait, PROBE_BACKOFF_TICKS); assert!(rig.run(RATE_TICKS as u32, 1_200, true, true)); - assert!(rig.run(RATE_TICKS as u32, 600, false, false)); + assert!(rig.run(2 * RATE_TICKS as u32, 600, false, false)); assert_eq!(rig.probe.wait, 0); - assert!(rig.run(1, 600, true, true)); - assert!(rig.run(RATE_TICKS as u32, 600, false, false)); assert!(!rig.run(1, 600, true, true)); assert_eq!(rig.probe.probes, 2); } @@ -456,6 +470,18 @@ mod tests { assert_eq!(rig.probe.wait, PROBE_BACKOFF_TICKS - RATE_TICKS as u32); } + #[test] + fn a_pause_in_the_same_workload_keeps_the_schedule() { + let mut rig = Rig::new(); + rig.run(RATE_TICKS as u32, 1_000, true, true); + assert!(rig.run(PROBE_TICKS, 1_200, true, true)); + rig.run(RATE_TICKS as u32, 1_200, true, true); + rig.run(3, 0, false, false); + assert!(rig.run(3 * RATE_TICKS as u32, 1_200, true, true)); + assert_eq!(rig.probe.probes, 1); + assert!(rig.probe.wait > 0); + } + #[test] fn a_probe_waits_for_a_steady_baseline() { let mut rig = Rig::new(); From 993286c3bc9f20f7521200a02fa230c3082dc742 Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 10 Sep 2026 10:15:18 +0900 Subject: [PATCH 11/16] fix: an idle pause is not a settled rate A pause longer than the settling window left a steady all-zero ring, which read as a rate settled elsewhere and voided the decision; the returning workload then bought a probe it did not need. A shift is now confirmed only by a steady rate above the floor; while the proxy is idle the shift stays pending. --- src/client/tuner.rs | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/src/client/tuner.rs b/src/client/tuner.rs index 27bf6e7..5b66aff 100644 --- a/src/client/tuner.rs +++ b/src/client/tuner.rs @@ -116,8 +116,9 @@ impl Probe { pipes.reverts.store(self.reverts, Ordering::Relaxed); } - // a rate that leaves the decided one is judged again once it is steady: back - // near it the workload only paused, away from it the decision is void + // a rate that leaves the decided one is judged again once it is steady and + // above the floor: back near it the workload only paused, away from it the + // decision is void; an idle proxy waits for traffic fn settle(&mut self) { self.settling = self.settling.saturating_sub(1); if self.settling > 0 || self.decided == 0 { @@ -129,7 +130,7 @@ impl Probe { self.shifted = true; self.settling = RATE_TICKS as u32; } - } else if self.steady() { + } else if self.steady() && self.rate() >= self.floor { self.shifted = false; if moved { self.decided = 0; @@ -480,6 +481,10 @@ mod tests { assert!(rig.run(3 * RATE_TICKS as u32, 1_200, true, true)); assert_eq!(rig.probe.probes, 1); assert!(rig.probe.wait > 0); + rig.run(2 * RATE_TICKS as u32, 0, false, false); + assert!(rig.run(3 * RATE_TICKS as u32, 1_200, true, true)); + assert_eq!(rig.probe.probes, 1); + assert!(rig.probe.wait > 0); } #[test] From 3e10ce1fafd37de4a710984f82893f73ab06c302 Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 10 Sep 2026 17:29:45 +0900 Subject: [PATCH 12/16] perf: only a lighter load on the shared pipes voids the decision A rate that settled above the decided one also voided it: the GET cell, entered from the seed phase at half its rate, bought a reverse probe and ran 5% under full sharding. A heavier load batches at least as well on the shared pipes and the scheduled reverse probe still checks it; from the local path the busy trigger already watches for load. Only a rate that settled a quarter under the decided one, on the shared pipes, now ends the wait. --- docs/architecture.md | 9 +++++---- src/client/tuner.rs | 30 +++++++++++++++++++++--------- 2 files changed, 26 insertions(+), 13 deletions(-) diff --git a/docs/architecture.md b/docs/architecture.md index 0f4b8b3..3a3d24b 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -110,10 +110,11 @@ moves every session to the shared pipes; the proxy then compares the commands it runs a second later against the second before the move, keeps them on a 5% gain, and otherwise takes them back and waits a minute before trying again, doubling that to eight minutes -while the answer holds; a command rate that settles a quarter away from -the one the answer was measured on ends the wait, since that workload is -gone, while a pause that returns to it changes nothing. From the shared -pipes it probes the other way on +while the answer holds; on the shared pipes a command rate that settles +a quarter under the one the answer was measured on ends the wait, since a +lighter load may be the local path's to serve, while a pause that returns +to the same rate changes nothing. From the shared pipes it probes the +other way on the same schedule, and lets go unmeasured once fewer than half the workers have stayed busy for three seconds — slowly, because moving the sessions away is what lowers their busyness. An idle proxy never probes. Switches happen only while a session has diff --git a/src/client/tuner.rs b/src/client/tuner.rs index 5b66aff..087523c 100644 --- a/src/client/tuner.rs +++ b/src/client/tuner.rs @@ -35,8 +35,8 @@ const MIN_RATE_PER_SEC: u64 = 100; const PROBE_BACKOFF_TICKS: u32 = 600; const PROBE_BACKOFF_MAX_TICKS: u32 = 4800; const PROBE_DOUBLINGS: u32 = (PROBE_BACKOFF_MAX_TICKS / PROBE_BACKOFF_TICKS).ilog2(); -// a rate this far from the one the current state was chosen on is a changed workload -const RATE_SHIFT_PCT: u64 = 25; +// a rate this far under the one the shared pipes were chosen on is a lighter workload +const RATE_FALL_PCT: u64 = 25; // a baseline whose ticks spread wider than this holds a gap or a ramp, not a rate const STEADY_SPREAD_PCT: u64 = 25; @@ -116,23 +116,25 @@ impl Probe { pipes.reverts.store(self.reverts, Ordering::Relaxed); } - // a rate that leaves the decided one is judged again once it is steady and - // above the floor: back near it the workload only paused, away from it the - // decision is void; an idle proxy waits for traffic + // on the shared pipes a rate that fell well under the decided one is judged + // again once it is steady and above the floor: back near it the workload only + // paused, still under it the local path may serve the lighter load better and + // the decision is void; a heavier load batches at least as well, and from the + // local path the busy trigger already watches fn settle(&mut self) { self.settling = self.settling.saturating_sub(1); - if self.settling > 0 || self.decided == 0 { + if self.settling > 0 || self.decided == 0 || !self.prefer { return; } - let moved = self.rate().abs_diff(self.decided) * 100 > self.decided * RATE_SHIFT_PCT; + let fell = self.rate() * 100 < self.decided * (100 - RATE_FALL_PCT); if !self.shifted { - if moved { + if fell { self.shifted = true; self.settling = RATE_TICKS as u32; } } else if self.steady() && self.rate() >= self.floor { self.shifted = false; - if moved { + if fell { self.decided = 0; self.wait = 0; } @@ -471,6 +473,16 @@ mod tests { assert_eq!(rig.probe.wait, PROBE_BACKOFF_TICKS - RATE_TICKS as u32); } + #[test] + fn a_heavier_workload_keeps_the_shared_pipes() { + let mut rig = Rig::new(); + rig.run(RATE_TICKS as u32, 1_000, true, true); + assert!(rig.run(PROBE_TICKS, 1_200, true, true)); + assert!(rig.run(4 * RATE_TICKS as u32, 2_400, true, true)); + assert_eq!(rig.probe.probes, 1); + assert!(rig.probe.wait > 0); + } + #[test] fn a_pause_in_the_same_workload_keeps_the_schedule() { let mut rig = Rig::new(); From 3df1b1cf70c1ff9d2ec310aeeb5dd5cc681bd0c0 Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 10 Sep 2026 18:46:29 +0900 Subject: [PATCH 13/16] perf: backend-sharding defaults to auto On the 128-master rig with 90-second GET and 60-second memtier cells, auto runs the P16 fan-out lane at 8.98M against 5.54M without sharding and 9.01M with full sharding, the sliding memtier lane at 4.23M against 4.31M and 3.58M, and P1 and the cache cells at parity with the best fixed mode. The 2% in the sliding lane is the experiment itself; the alternatives lose 62% or 17% somewhere. The docs name the old default where their tables were measured with it. --- README.md | 17 +++++++++-------- docs/benchmarks.md | 25 ++++++++++++++++++------- docs/configuration.md | 2 +- mithril.conf.sample | 2 +- src/config.rs | 2 +- 5 files changed, 30 insertions(+), 18 deletions(-) diff --git a/README.md b/README.md index 365783f..85dcdf9 100644 --- a/README.md +++ b/README.md @@ -13,14 +13,15 @@ thread-per-core runtime and zero-copy frame forwarding. ## Highlights - **Thread-per-core** — each worker runs a single-threaded tokio runtime with - its own backend pools (one set per node per database in use); by default - nothing crosses workers on the request path, and a central acceptor places - each connection on the least-loaded worker (configurable) so no worker - becomes the latency floor. The optional `backend-sharding` mode trades that - isolation for one process-wide pipe per node per database in use, - deepening backend batches for unpipelined workloads — and `auto` - makes that call per session, so unpipelined and pipelining clients each - get the path that is faster for them + its own backend pools (one set per node per database in use), and a + central acceptor places each connection on the least-loaded worker + (configurable) so no worker becomes the latency floor. `backend-sharding` + decides what crosses workers: `no` keeps every request on its worker's own + connections, `yes` sends each node's traffic through one process-wide + pipe, which batches deeper for unpipelined workloads, and the default + `auto` gives an unpipelined session the pipe, a pipelining one its + worker's connections, and moves everyone to the pipes only where a + measured experiment shows the whole proxy runs faster - **Zero-copy pipeline** — requests and replies travel as `bytes::Bytes` slices of the socket buffers; the RESP layer finds frame boundaries without materializing values, and replies re-order per client by sequence number so diff --git a/docs/benchmarks.md b/docs/benchmarks.md index fa1ca96..2d17525 100644 --- a/docs/benchmarks.md +++ b/docs/benchmarks.md @@ -10,9 +10,10 @@ cores (memtier_benchmark 8 threads, 20-second windows; redis-benchmark reversed on the second pass; the ranges below are the two passes. 64-byte values unless stated, 100k keys, random keys, SET:GET 1:1 unless stated. -Columns: **mithril** (default), **shard** (`backend-sharding yes`), -**cache** (`reply-cache yes`), **shard+cache** (both), mt-proxy (C++), -predixy (C++). Commit 9336a89 lineage. +Columns: **mithril** (`backend-sharding no`, the default until v0.1.6), +**shard** (`backend-sharding yes`), **cache** (`reply-cache yes`), +**shard+cache** (both), mt-proxy (C++), predixy (C++). Commit 9336a89 +lineage. ### Throughput (ops/s) @@ -70,10 +71,20 @@ rounds each, redis-benchmark unless stated: | memtier P16, 200 conns (sliding) | **4.36M / 0.76** | 3.59M / 0.92 | 3.74M / 0.86 | | memtier P1, 400 conns | 445k | 443k | 443k | -`auto` lands on the better of the two paths in every cell but one: a -sliding-pipeline client at moderate concurrency on a partly idle proxy -runs 14% slower than default, because the shared pipes cost a hop those -sessions cannot hide. If that is your only workload, set `no`. +That `auto` landed on the better of the two paths in every cell but one: +a sliding-pipeline client at moderate concurrency on a partly idle proxy +ran 14% slower than `no`, because its worker rule moved sessions on +busyness alone. Since v0.1.7 `auto` runs one experiment for the whole +proxy and keeps the shared pipes only where the proxy's command rate +measures higher (see architecture.md); on the same rig, with 90-second +GET and 60-second memtier cells over three rotated rounds: P16 1000 conns +GET `no` 5.54M, `auto` 8.98M, `yes` 9.01M; the sliding memtier cell `no` +4.31M, `auto` 4.23M, `yes` 3.58M; P1 2000 conns and the cache cells at +parity with the best fixed mode. The 2% in the sliding cell is the +experiment itself (two two-second windows in a minute, then one per +minute doubling to eight), which is why `auto` is now the default; set +`no` only when a sliding pipeline is the whole workload and that 2% +matters more than the 60% the other shapes gain. Re-measured at v0.1.6 on the same rig (four rotated rounds, default / shard / auto): P16 1000 conns GET 5.55M / 8.87M / 8.76M and SET 5.54M / diff --git a/docs/configuration.md b/docs/configuration.md index f197734..59d8212 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -14,7 +14,7 @@ See [`mithril.conf.sample`](https://github.com/projecteru2/mithril/blob/master/m | `worker-threads` | int | CPU count | worker threads, one runtime each | | `maxclients` | int | `10000` | client connection cap, enforced at accept | | `backend-conns` | 1..512 | `1` | shared pipelined connections per node per worker, per database in use (a session on `SELECT n` opens its own pool to each node; the exclusive-connection cap for blocking and `WATCH` stays per node across databases) | -| `backend-sharding` | `yes`/`no`/`auto` | `no` | `yes`: one process-wide connection per node per database in use, owned by the worker its address hashes to, which deepens backend pipelines for unpipelined workloads; `auto`: decided per session and once for the whole proxy — a session that keeps sending one command at a time moves to the shared pipes, a pipelining session stays on its worker's own connections unless half the workers are saturated with thin local batches (busy ≥ 85% and fewer than eight frames per backend write), which makes the proxy try the shared pipes for a second and keep them only where its throughput measures higher (switches happen only while a session has nothing in flight; with the reply cache, only sessions on the shared pipes fill it — the worker-local connections carry no key tracking — while every session reads it) | +| `backend-sharding` | `yes`/`no`/`auto` | `auto` | `yes`: one process-wide connection per node per database in use, owned by the worker its address hashes to, which deepens backend pipelines for unpipelined workloads; `no`: every request stays on its worker's own connections; `auto` (the default): decided per session and once for the whole proxy — a session that keeps sending one command at a time moves to the shared pipes, a pipelining session stays on its worker's own connections unless half the workers are saturated with thin local batches (busy ≥ 85% and fewer than eight frames per backend write), which makes the proxy try the shared pipes for a second and keep them only where its throughput measures higher (switches happen only while a session has nothing in flight; with the reply cache, only sessions on the shared pipes fill it — the worker-local connections carry no key tracking — while every session reads it) | | `reply-cache` | `yes`/`no` | `no` | worker-local GET/MGET reply cache; the servers track the keys the proxy caches (redirected opt-in RESP3 tracking) and invalidate them on change, so hits skip the backend round trip and a session reads its own writes | | `reply-cache-max-bytes` | bytes | `64mb` | per-worker budget for the two cache generations; the live generation holds half of it, so a worker's hot set must fit in half the budget to keep hitting (see operations.md for sizing and the RSS to expect) | | `reply-cache-max-age-secs` | 1..3600 | `10` | staleness backstop for missed invalidations; entries older than this never serve | diff --git a/mithril.conf.sample b/mithril.conf.sample index f7e368d..5c5b209 100644 --- a/mithril.conf.sample +++ b/mithril.conf.sample @@ -17,7 +17,7 @@ bootstrap 127.0.0.1:7001,127.0.0.1:7002,127.0.0.1:7003 backend-conns 1 # one process-wide pipe per node (deeper backend batches for unpipelined load): # no | yes | auto (per session: unpipelined sessions share, pipelining ones stay local) -# backend-sharding no +# backend-sharding auto # worker-local GET/MGET reply cache, kept coherent by server-side key tracking # reply-cache no diff --git a/src/config.rs b/src/config.rs index b05e1a4..7eef4b4 100644 --- a/src/config.rs +++ b/src/config.rs @@ -111,7 +111,7 @@ impl Default for Config { maxclients: 10000, bootstrap: Vec::new(), backend_conns: 1, - backend_sharding: Sharding::Off, + backend_sharding: Sharding::Auto, reply_cache: false, reply_cache_max_bytes: 64 << 20, reply_cache_max_age_secs: 10, From d65e4260e205c3dc309f2b00804402fb3b924e24 Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 10 Sep 2026 18:55:18 +0900 Subject: [PATCH 14/16] fix: one worker has nothing to share, and the docs qualify the shared-nothing claim With a single worker the local pool already holds the one connection per node the shared pipes would add, so auto only bought the fabric and mutex-backed reply queues; it now resolves to no. The threading model, the reply queue's module line and the sample name the setting the shared-nothing guarantee holds under. --- docs/architecture.md | 10 ++++++---- mithril.conf.sample | 3 ++- src/client/queue.rs | 2 +- src/config.rs | 16 ++++++++++++++++ 4 files changed, 25 insertions(+), 6 deletions(-) diff --git a/docs/architecture.md b/docs/architecture.md index 3a3d24b..0bd9945 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -3,10 +3,12 @@ ## Threading model Mithril is thread-per-core. Each worker thread runs a single-threaded tokio -runtime with its own backend connection pools; nothing is shared between -workers on the request path, so there are no locks and no atomic -reference-count traffic per request (`Rc`, not `Arc`, everywhere inside a -worker). +runtime with its own backend connection pools; under `backend-sharding no` +nothing is shared between workers on the request path, so there are no +locks and no atomic reference-count traffic per request (`Rc`, not `Arc`, +everywhere inside a worker). The shared pipes of `yes` and `auto` (the +default) cross workers through the fabric and a mutex-backed reply queue, +described under backend sharding below. One acceptor thread owns the listen socket and hands accepted connections to workers over bounded channels, placed least-loaded by default (configurable). Kernel `SO_REUSEPORT` hashing was diff --git a/mithril.conf.sample b/mithril.conf.sample index 5c5b209..3cd7582 100644 --- a/mithril.conf.sample +++ b/mithril.conf.sample @@ -16,7 +16,8 @@ bootstrap 127.0.0.1:7001,127.0.0.1:7002,127.0.0.1:7003 # sticky shared connections per node per worker backend-conns 1 # one process-wide pipe per node (deeper backend batches for unpipelined load): -# no | yes | auto (per session: unpipelined sessions share, pipelining ones stay local) +# no | yes | auto (unpipelined sessions share; the whole proxy moves to the pipes +# only where a measured experiment shows it runs faster) # backend-sharding auto # worker-local GET/MGET reply cache, kept coherent by server-side key tracking diff --git a/src/client/queue.rs b/src/client/queue.rs index 0b3182c..3633514 100644 --- a/src/client/queue.rs +++ b/src/client/queue.rs @@ -1,4 +1,4 @@ -//! Per-session reply queue: worker-local by default, mutex-backed when owner workers deliver. +//! Per-session reply queue: worker-local under `backend-sharding no`, mutex-backed when owner workers deliver. use std::cell::{Cell, RefCell}; use std::collections::VecDeque; diff --git a/src/config.rs b/src/config.rs index 7eef4b4..ec74daf 100644 --- a/src/config.rs +++ b/src/config.rs @@ -220,6 +220,10 @@ impl Config { if self.backend_sharding == Sharding::On && self.backend_conns > 1 { crate::log_warn!("backend-conns is ignored under backend-sharding"); } + // one worker's own connections already carry every request: nothing to share + if self.workers == 1 && self.backend_sharding == Sharding::Auto { + self.backend_sharding = Sharding::Off; + } Ok(self) } } @@ -308,4 +312,16 @@ mod tests { assert!(cfg.set("no-such-key", "1").is_err()); assert!(cfg.set("backend-conns", "0").is_err()); } + + #[test] + fn one_worker_has_nothing_to_share() { + let mut cfg = Config::default(); + cfg.set("bootstrap", "127.0.0.1:7001").unwrap(); + cfg.set("worker-threads", "1").unwrap(); + assert_eq!(cfg.finish().unwrap().backend_sharding, Sharding::Off); + let mut cfg = Config::default(); + cfg.set("bootstrap", "127.0.0.1:7001").unwrap(); + cfg.set("worker-threads", "2").unwrap(); + assert_eq!(cfg.finish().unwrap().backend_sharding, Sharding::Auto); + } } From b78dff4bb3e5bd2318f0a2ad24b349cc1c5463cc Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 10 Sep 2026 18:59:06 +0900 Subject: [PATCH 15/16] fix: one worker with several connections per node still has a lane to choose The one-worker normalization also rewrote auto where backend-conns spreads a worker's sessions over several connections per node, which the shared pipe consolidates; it now applies only with one connection per node. --- src/config.rs | 10 ++++++++-- 1 file changed, 8 insertions(+), 2 deletions(-) diff --git a/src/config.rs b/src/config.rs index ec74daf..32b06d0 100644 --- a/src/config.rs +++ b/src/config.rs @@ -220,8 +220,9 @@ impl Config { if self.backend_sharding == Sharding::On && self.backend_conns > 1 { crate::log_warn!("backend-conns is ignored under backend-sharding"); } - // one worker's own connections already carry every request: nothing to share - if self.workers == 1 && self.backend_sharding == Sharding::Auto { + // one worker with one connection per node already carries every request on + // the pipe the shared lane would add: nothing to share + if self.workers == 1 && self.backend_conns == 1 && self.backend_sharding == Sharding::Auto { self.backend_sharding = Sharding::Off; } Ok(self) @@ -323,5 +324,10 @@ mod tests { cfg.set("bootstrap", "127.0.0.1:7001").unwrap(); cfg.set("worker-threads", "2").unwrap(); assert_eq!(cfg.finish().unwrap().backend_sharding, Sharding::Auto); + let mut cfg = Config::default(); + cfg.set("bootstrap", "127.0.0.1:7001").unwrap(); + cfg.set("worker-threads", "1").unwrap(); + cfg.set("backend-conns", "2").unwrap(); + assert_eq!(cfg.finish().unwrap().backend_sharding, Sharding::Auto); } } From ab0028c14d4b67391a39b357f38993271eb5bcdb Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 10 Sep 2026 19:10:01 +0900 Subject: [PATCH 16/16] review: the probe's phase is an option, its settled second one predicate, its comments one line each The four cleanup lenses over the branch: must_drain is a filter over switch_pipes (equivalent on every input, one threshold instead of two), the ring's first fill is the same settling window a decision starts, the four rate checks at a probe's start and the two in the shift confirmation are one settled() predicate, the probing phase is an Option instead of an enum, the backoff is stored and doubled instead of derived from a shift, and the module's comments are back to one line each (25 lines from 34). benchmarks.md no longer describes the retired worker rule above the table measured with it, and the configuration row defers the mechanism to architecture.md. --- docs/benchmarks.md | 8 ++-- docs/configuration.md | 2 +- src/client/tuner.rs | 95 +++++++++++++++---------------------------- src/config.rs | 3 +- 4 files changed, 38 insertions(+), 70 deletions(-) diff --git a/docs/benchmarks.md b/docs/benchmarks.md index 2d17525..535fd5f 100644 --- a/docs/benchmarks.md +++ b/docs/benchmarks.md @@ -54,11 +54,9 @@ unpipelined saturation cells — which `backend-sharding` recovers. ### backend-sharding auto (same rig, commit e6e1592 lineage) -`auto` decides per worker and per session: a worker that is busy with -thin local batches moves every session it hosts to the shared pipes, a -lightly loaded worker keeps its own connections, and an unpipelined -session prefers the shared pipe on its own (see architecture.md). Two -rounds each, redis-benchmark unless stated: +The table was measured with the first `auto`, whose worker rule moved +sessions on busyness alone; the paragraph after it has the current +mechanism. Two rounds each, redis-benchmark unless stated: | cell | default | shard | auto | |---|---|---|---| diff --git a/docs/configuration.md b/docs/configuration.md index 59d8212..f2d635d 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -14,7 +14,7 @@ See [`mithril.conf.sample`](https://github.com/projecteru2/mithril/blob/master/m | `worker-threads` | int | CPU count | worker threads, one runtime each | | `maxclients` | int | `10000` | client connection cap, enforced at accept | | `backend-conns` | 1..512 | `1` | shared pipelined connections per node per worker, per database in use (a session on `SELECT n` opens its own pool to each node; the exclusive-connection cap for blocking and `WATCH` stays per node across databases) | -| `backend-sharding` | `yes`/`no`/`auto` | `auto` | `yes`: one process-wide connection per node per database in use, owned by the worker its address hashes to, which deepens backend pipelines for unpipelined workloads; `no`: every request stays on its worker's own connections; `auto` (the default): decided per session and once for the whole proxy — a session that keeps sending one command at a time moves to the shared pipes, a pipelining session stays on its worker's own connections unless half the workers are saturated with thin local batches (busy ≥ 85% and fewer than eight frames per backend write), which makes the proxy try the shared pipes for a second and keep them only where its throughput measures higher (switches happen only while a session has nothing in flight; with the reply cache, only sessions on the shared pipes fill it — the worker-local connections carry no key tracking — while every session reads it) | +| `backend-sharding` | `yes`/`no`/`auto` | `auto` | `yes`: one process-wide connection per node per database in use, owned by the worker its address hashes to, which deepens backend pipelines for unpipelined workloads; `no`: every request stays on its worker's own connections; `auto` (the default): an unpipelined session uses the shared pipe, a pipelining one its worker's connections, and the whole proxy moves to the pipes where a measured experiment shows it runs faster (the rule is in architecture.md; a session switches only while nothing is in flight; with the reply cache, only sessions on the shared pipes fill it while every session reads it) | | `reply-cache` | `yes`/`no` | `no` | worker-local GET/MGET reply cache; the servers track the keys the proxy caches (redirected opt-in RESP3 tracking) and invalidate them on change, so hits skip the backend round trip and a session reads its own writes | | `reply-cache-max-bytes` | bytes | `64mb` | per-worker budget for the two cache generations; the live generation holds half of it, so a worker's hot set must fit in half the budget to keep hitting (see operations.md for sizing and the RSS to expect) | | `reply-cache-max-age-secs` | 1..3600 | `10` | staleness backstop for missed invalidations; entries older than this never serve | diff --git a/src/client/tuner.rs b/src/client/tuner.rs index 087523c..1b268ca 100644 --- a/src/client/tuner.rs +++ b/src/client/tuner.rs @@ -11,15 +11,13 @@ use crate::stats::{self, Pipes, Stats}; // pipelining score: a local session shares at 0, a shared one returns at PIPELINED_LOCAL pub(super) const PIPELINED_LOCAL: u8 = 4; const PIPELINED_MAX: u8 = 8; -// what triggers a probe: a worker is busy and its local backend batches would -// stay thin (in-flight commands per master) +// a worker busy with thin local batches (in-flight commands per master) is what triggers a probe const TUNE_PERIOD: Duration = Duration::from_millis(100); const BUSY_ENTER: u32 = 85; const BUSY_LEAVE: u32 = 60; const DEPTH_ENTER: u64 = 8; const DEPTH_LEAVE: u64 = 16; -// ticks the enter or leave condition must hold: moving sessions off the workers -// lowers their busyness, so leaving is slow and entering prompt +// leaving is slow, since moving the sessions away is what lowers the workers' busyness const ENTER_TICKS: u32 = 3; const LEAVE_TICKS: u32 = 30; // the command-rate window, ten ticks of 100 ms, so a rate is commands per second @@ -34,16 +32,14 @@ const MIN_RATE_PER_SEC: u64 = 100; // ticks to the next probe, doubling while decisions confirm the current state const PROBE_BACKOFF_TICKS: u32 = 600; const PROBE_BACKOFF_MAX_TICKS: u32 = 4800; -const PROBE_DOUBLINGS: u32 = (PROBE_BACKOFF_MAX_TICKS / PROBE_BACKOFF_TICKS).ilog2(); // a rate this far under the one the shared pipes were chosen on is a lighter workload const RATE_FALL_PCT: u64 = 25; // a baseline whose ticks spread wider than this holds a gap or a ramp, not a rate const STEADY_SPREAD_PCT: u64 = 25; impl Session { - // an unpipelined session gains from the deeper batches of the shared pipe, a - // pipelined one from its worker-local connection; a switch happens while nothing - // is in flight + // an unpipelined session gains from the shared pipe's batching, a pipelined one from + // its worker's connection; a switch happens while nothing is in flight pub(super) fn adapt_pipes(&self) { let depth = self.outstanding(); let score = self.pipelined.get(); @@ -59,16 +55,6 @@ impl Session { } } -#[derive(Clone, Copy, Default)] -enum Phase { - #[default] - Steady, - Probing { - baseline: u64, - ticks: u32, - }, -} - // the proxy's experiment: a busy-and-thin majority triggers it, the measured command rate decides #[derive(Default)] struct Probe { @@ -76,23 +62,24 @@ struct Probe { probes: u64, keeps: u64, reverts: u64, - phase: Phase, + probing: Option<(u64, u32)>, streak: u32, wait: u32, - confirmed: u32, + backoff: u32, decided: u64, settling: u32, shifted: bool, floor: u64, ring: [u64; RATE_TICKS], at: usize, - seen: usize, last: u64, } impl Probe { fn new(workers: usize) -> Probe { Probe { + backoff: PROBE_BACKOFF_TICKS, + settling: RATE_TICKS as u32, floor: MIN_RATE_PER_SEC * workers as u64, ..Probe::default() } @@ -102,10 +89,10 @@ impl Probe { self.record(commands_now); self.wait = self.wait.saturating_sub(1); self.settle(); - match self.phase { - Phase::Probing { baseline, ticks } => self.probing(baseline, ticks), - Phase::Steady if self.prefer => self.leave(still_busy), - Phase::Steady => self.enter(busy_thin), + match self.probing { + Some((baseline, ticks)) => self.advance(baseline, ticks), + None if self.prefer => self.leave(still_busy), + None => self.enter(busy_thin), } self.prefer } @@ -116,11 +103,8 @@ impl Probe { pipes.reverts.store(self.reverts, Ordering::Relaxed); } - // on the shared pipes a rate that fell well under the decided one is judged - // again once it is steady and above the floor: back near it the workload only - // paused, still under it the local path may serve the lighter load better and - // the decision is void; a heavier load batches at least as well, and from the - // local path the busy trigger already watches + // on the shared pipes a rate that fell a quarter under the decided one voids the + // decision once it has settled there; back near it, the workload only paused fn settle(&mut self) { self.settling = self.settling.saturating_sub(1); if self.settling > 0 || self.decided == 0 || !self.prefer { @@ -132,7 +116,7 @@ impl Probe { self.shifted = true; self.settling = RATE_TICKS as u32; } - } else if self.steady() && self.rate() >= self.floor { + } else if self.settled() { self.shifted = false; if fell { self.decided = 0; @@ -145,19 +129,21 @@ impl Probe { self.ring[self.at] = commands_now.saturating_sub(self.last); self.at = (self.at + 1) % RATE_TICKS; self.last = commands_now; - self.seen = (self.seen + 1).min(RATE_TICKS); } fn rate(&self) -> u64 { self.ring.iter().sum() } - fn steady(&self) -> bool { + // the ring holds one steady second of real traffic + fn settled(&self) -> bool { let (min, max) = self .ring .iter() .fold((u64::MAX, 0), |(lo, hi), &v| (lo.min(v), hi.max(v))); - max * 100 <= min * (100 + STEADY_SPREAD_PCT) + self.settling == 0 + && self.rate() >= self.floor + && max * 100 <= min * (100 + STEADY_SPREAD_PCT) } fn enter(&mut self, busy_thin: bool) { @@ -189,26 +175,19 @@ impl Probe { } fn start_probe(&mut self) { - let baseline = self.rate(); - if self.wait > 0 - || self.settling > 0 - || self.shifted - || self.seen < RATE_TICKS - || baseline < self.floor - || !self.steady() - { + if self.wait > 0 || self.shifted || !self.settled() { return; } self.probes += 1; self.streak = 0; self.prefer = !self.prefer; - self.phase = Phase::Probing { baseline, ticks: 0 }; + self.probing = Some((self.rate(), 0)); } - fn probing(&mut self, baseline: u64, ticks: u32) { + fn advance(&mut self, baseline: u64, ticks: u32) { let ticks = ticks + 1; if ticks < PROBE_TICKS { - self.phase = Phase::Probing { baseline, ticks }; + self.probing = Some((baseline, ticks)); return; } let measured = self.rate(); @@ -218,13 +197,13 @@ impl Probe { (baseline, measured) }; let keep = shared * 100 >= local * (100 + KEEP_GAIN_PCT); - // a decision that flips the state starts the schedule over, one that confirms it waits longer + // a flip restarts the schedule, a confirmation waits longer if keep == self.prefer { - self.confirmed = 0; + self.backoff = PROBE_BACKOFF_TICKS; self.wait = PROBE_BACKOFF_TICKS; } else { - self.wait = PROBE_BACKOFF_TICKS << self.confirmed; - self.confirmed = (self.confirmed + 1).min(PROBE_DOUBLINGS); + self.wait = self.backoff; + self.backoff = (self.backoff * 2).min(PROBE_BACKOFF_MAX_TICKS); } if keep { self.keeps += 1; @@ -235,12 +214,11 @@ impl Probe { self.settling = RATE_TICKS as u32; self.shifted = false; self.prefer = keep; - self.phase = Phase::Steady; + self.probing = None; } } -/// Publishes this worker's CPU busyness and batch depth and mirrors the process-wide -/// pipe preference; the lead worker also runs the experiment that sets it. +/// Publishes the worker's load samples and mirrors the pipe preference; the lead worker runs the experiment. pub async fn auto_tuner(shared: Rc, lead: bool) { let mut last_ticks = stats::thread_cpu_ticks(); let mut last_at = Instant::now(); @@ -275,8 +253,7 @@ pub async fn auto_tuner(shared: Rc, lead: bool) { } } -// the cost of the shared pipes is process-wide, so one experiment moves every session -// and the whole proxy's command rate judges it +// the shared pipes cost the whole proxy: one experiment moves every session, the proxy's rate judges fn conduct(stats: &Stats, probe: &mut Probe) { let (busy_thin, still_busy, commands) = survey(stats); let prefer = probe.tick(busy_thin, still_busy, commands); @@ -316,16 +293,10 @@ fn switch_pipes(sharded: bool, score: u8, worker_prefers_shared: bool) -> bool { // a never-idle session is paused to move only for a switch still wanted once it drains fn must_drain(sharded: bool, score: u8, worker_prefers_shared: bool) -> bool { - if worker_prefers_shared { - !sharded - } else { - sharded && score >= PIPELINED_LOCAL - } + (worker_prefers_shared || sharded) && switch_pipes(sharded, score, worker_prefers_shared) } -// the measured local batch while local traffic flows; with none, the in-flight -// commands spread over the masters — an estimate that errs toward staying on the -// shared pipes, which at saturation batch at least as well as local connections +// with no local writes, in-flight commands per master stand in, erring toward the shared pipes fn tune_depth(measured: u32, writes: u32, inflight: u64, masters: u64) -> u64 { if writes > 0 { u64::from(measured) diff --git a/src/config.rs b/src/config.rs index 32b06d0..87aa25e 100644 --- a/src/config.rs +++ b/src/config.rs @@ -220,8 +220,7 @@ impl Config { if self.backend_sharding == Sharding::On && self.backend_conns > 1 { crate::log_warn!("backend-conns is ignored under backend-sharding"); } - // one worker with one connection per node already carries every request on - // the pipe the shared lane would add: nothing to share + // one worker with one connection per node has nothing to share if self.workers == 1 && self.backend_conns == 1 && self.backend_sharding == Sharding::Auto { self.backend_sharding = Sharding::Off; }