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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 9 additions & 8 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
36 changes: 24 additions & 12 deletions docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
33 changes: 21 additions & 12 deletions docs/benchmarks.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down Expand Up @@ -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 |
|---|---|---|---|
Expand All @@ -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 /
Expand Down
2 changes: 1 addition & 1 deletion docs/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down
2 changes: 1 addition & 1 deletion docs/operations.md
Original file line number Diff line number Diff line change
Expand Up @@ -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_<command>:calls=<n>` for every command run at least once, subcommands as `cmdstat_client\|list`; counted once accepted, summed over workers |

Expand Down
5 changes: 3 additions & 2 deletions mithril.conf.sample
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
36 changes: 28 additions & 8 deletions src/admin.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -211,7 +211,8 @@ pub fn info(cfg: &Config, stats: &Stats, started: u64) -> Vec<u8> {
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,
Expand Down Expand Up @@ -242,12 +243,12 @@ pub fn info(cfg: &Config, stats: &Stats, started: u64) -> Vec<u8> {
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::<Vec<_>>()
.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 {
Expand Down Expand Up @@ -418,6 +419,15 @@ fn split_announce(announce: &str) -> (&str, u16) {
}
}

fn per_worker<F: Fn(&stats::WorkerStats) -> &AtomicU64>(stats: &Stats, field: F) -> String {
stats
.workers
.iter()
.map(|w| field(w).load(Ordering::Relaxed).to_string())
.collect::<Vec<_>>()
.join(",")
}

fn yesno(v: bool) -> &'static str {
if v { "yes" } else { "no" }
}
Expand Down Expand Up @@ -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| {
Expand All @@ -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"));
Expand Down
2 changes: 1 addition & 1 deletion src/client/queue.rs
Original file line number Diff line number Diff line change
@@ -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;
Expand Down
Loading