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/architecture.md b/docs/architecture.md index 61f5134..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 @@ -101,18 +103,28 @@ 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), 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, 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 +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 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/benchmarks.md b/docs/benchmarks.md index fa1ca96..535fd5f 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) @@ -53,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 | |---|---|---|---| @@ -70,10 +69,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 f24bf01..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` | `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` | `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/docs/operations.md b/docs/operations.md index 8d8f5ab..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`, `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`, `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/mithril.conf.sample b/mithril.conf.sample index f7e368d..3cd7582 100644 --- a/mithril.conf.sample +++ b/mithril.conf.sample @@ -16,8 +16,9 @@ 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) -# backend-sharding no +# 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 # reply-cache no diff --git a/src/admin.rs b/src/admin.rs index d101eca..cd75fed 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\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(), cfg.port, @@ -242,12 +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 - .workers - .iter() - .map(|w| w.commands.load(Ordering::Relaxed).to_string()) - .collect::>() - .join(","), + 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.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 { @@ -418,6 +419,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 +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].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| { @@ -535,6 +549,12 @@ 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_busy").as_deref(), Some("0,90")); + 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")); 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/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/client/tuner.rs b/src/client/tuner.rs index 96f15ab..1b268ca 100644 --- a/src/client/tuner.rs +++ b/src/client/tuner.rs @@ -1,30 +1,45 @@ -//! 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; use std::time::{Duration, Instant}; use super::Shared; use super::session::Session; -use crate::stats; +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; -// 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) +// 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 a worker -// lowers its own 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 +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 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 = 600; +const PROBE_BACKOFF_MAX_TICKS: u32 = 4800; +// 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; switch only 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(); @@ -32,7 +47,7 @@ impl Session { let sharded = self.link.sharded.get(); let prefer_shared = self.shared.prefer_shared.get(); self.switch_pending - .set(prefer_shared && !sharded && depth > 0); + .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(); @@ -40,13 +55,175 @@ impl Session { } } -/// Samples this worker's CPU busyness and in-flight depth per master and publishes -/// whether its sessions should prefer the shared pipes. -pub async fn auto_tuner(shared: Rc) { +// the proxy's experiment: a busy-and-thin majority triggers it, the measured command rate decides +#[derive(Default)] +struct Probe { + prefer: bool, + probes: u64, + keeps: u64, + reverts: u64, + probing: Option<(u64, u32)>, + streak: u32, + wait: u32, + backoff: u32, + decided: u64, + settling: u32, + shifted: bool, + floor: u64, + ring: [u64; RATE_TICKS], + at: 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() + } + } + + 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); + self.settle(); + 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 + } + + 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); + } + + // 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 { + return; + } + let fell = self.rate() * 100 < self.decided * (100 - RATE_FALL_PCT); + if !self.shifted { + if fell { + self.shifted = true; + self.settling = RATE_TICKS as u32; + } + } else if self.settled() { + self.shifted = false; + if fell { + 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; + self.last = commands_now; + } + + fn rate(&self) -> u64 { + self.ring.iter().sum() + } + + // 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))); + self.settling == 0 + && self.rate() >= self.floor + && max * 100 <= min * (100 + STEADY_SPREAD_PCT) + } + + fn enter(&mut self, busy_thin: bool) { + if !busy_thin { + self.streak = 0; + return; + } + self.streak = (self.streak + 1).min(ENTER_TICKS); + if self.streak == ENTER_TICKS { + self.start_probe(); + } + } + + fn leave(&mut self, still_busy: bool) { + if still_busy { + 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.decided = 0; + self.shifted = false; + self.prefer = false; + } + } + + fn start_probe(&mut self) { + if self.wait > 0 || self.shifted || !self.settled() { + return; + } + self.probes += 1; + self.streak = 0; + self.prefer = !self.prefer; + self.probing = Some((self.rate(), 0)); + } + + fn advance(&mut self, baseline: u64, ticks: u32) { + let ticks = ticks + 1; + if ticks < PROBE_TICKS { + self.probing = Some((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 flip restarts the schedule, a confirmation waits longer + if keep == self.prefer { + self.backoff = PROBE_BACKOFF_TICKS; + self.wait = PROBE_BACKOFF_TICKS; + } else { + self.wait = self.backoff; + self.backoff = (self.backoff * 2).min(PROBE_BACKOFF_MAX_TICKS); + } + if keep { + self.keeps += 1; + } else { + self.reverts += 1; + } + self.decided = if keep { shared } else { local }; + self.settling = RATE_TICKS as u32; + self.shifted = false; + self.prefer = keep; + self.probing = None; + } +} + +/// 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(); let mut busy_x16 = 0u32; - let mut streak = 0u32; + let mut probe = lead.then(|| Probe::new(shared.stats.workers.len())); loop { tokio::time::sleep(TUNE_PERIOD).await; let now = Instant::now(); @@ -62,12 +239,46 @@ 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)); + 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(shared.stats.pipes.prefer.load(Ordering::Relaxed)); } } +// 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); + 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 fn switch_pipes(sharded: bool, score: u8, worker_prefers_shared: bool) -> bool { if worker_prefers_shared { @@ -80,33 +291,12 @@ 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 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 { + (worker_prefers_shared || sharded) && switch_pipes(sharded, score, worker_prefers_shared) } -// 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 +// 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) @@ -153,15 +343,20 @@ mod tests { } #[test] - fn worker_preference_overrides_the_score_with_hysteresis() { + 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)); 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 +364,176 @@ 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)); + } + + #[test] + fn a_probe_that_gains_keeps_the_shared_pipes() { + 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, 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::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, true, true)); + assert_eq!(rig.probe.probes, 1); + assert!(rig.run(1, 1_000, true, true)); + assert_eq!(rig.probe.probes, 2); + 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 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::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, 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 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(2 * 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 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 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(); + 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); + 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] + 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(); + 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); + } + + #[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 new() -> Rig { + Rig { + probe: Probe::new(WORKERS), + commands: 0, + } } - assert!(settle(false, true, &mut streak)); - for _ in 0..LEAVE_TICKS - 1 { - assert!(settle(true, false, &mut streak)); + + 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_thin, still_busy, self.commands); + } + prefer } - assert!(!settle(true, false, &mut streak)); - assert!(settle(true, true, &mut streak)); - assert_eq!(streak, 0); } } diff --git a/src/config.rs b/src/config.rs index b05e1a4..87aa25e 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, @@ -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 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; + } Ok(self) } } @@ -308,4 +312,21 @@ 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); + 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); + } } 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 39dc859..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,6 +48,8 @@ pub struct WorkerStats { pub cache_entries: AtomicU64, pub cache_bytes: AtomicU64, pub cache_flips: AtomicU64, + pub busy_pct: AtomicU64, + pub batch_depth: AtomicU64, pub calls: Calls, } @@ -139,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>, @@ -146,6 +157,7 @@ pub struct Stats { pub total_connections: AtomicU64, pub registry: Mutex>, pub slowlog: Slowlog, + pub pipes: Pipes, epoch: Instant, } @@ -157,6 +169,7 @@ impl Stats { total_connections: AtomicU64::new(0), registry: Mutex::new(HashMap::new()), slowlog: Slowlog::default(), + pipes: Pipes::default(), epoch: Instant::now(), }) }