1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
//! The public entry point: configure and run the thread-per-core server.
use crate::Commands;
use kevy_persist::Fsync;
use std::path::PathBuf;
/// Default slots in each per-core-pair SPSC ring. A full ring spills
/// to a local backlog (see [`Shard`]), so this only bounds the
/// lock-free fast path, not capacity. Overridable via the
/// `[advanced] ring_capacity` config field threaded through
/// [`Runtime::with_advanced`].
const DEFAULT_RING_CAPACITY: usize = 1024;
/// The public entry point: configure and run the thread-per-core server.
#[derive(Debug)]
pub struct Runtime<C: Commands> {
pub(crate) ip: [u8; 4],
pub(crate) port: u16,
pub(crate) nshards: usize,
pub(crate) commands: C,
/// Directory for per-shard snapshot files (`dump-<id>.rdb`) and AOF logs.
pub(crate) data_dir: PathBuf,
/// Whether the append-only log is enabled.
pub(crate) enable_aof: bool,
/// fsync policy for the AOF. Default `EverySec` matches Redis.
pub(crate) appendfsync: Fsync,
/// auto-trigger BGREWRITEAOF when AOF grew this many % above the size
/// at the previous rewrite. `0` disables. Default `100` (matches Redis).
pub(crate) auto_aof_rewrite_pct: u32,
pub(crate) auto_aof_rewrite_bytes: u64,
pub(crate) auto_aof_rewrite_interval_secs: u64,
pub(crate) replay_resync: bool,
/// Floor below which auto-rewrite is skipped. Default `64 MiB`.
pub(crate) auto_aof_rewrite_min_size: u64,
/// Reactor SPSC ring slot count. See [`DEFAULT_RING_CAPACITY`].
pub(crate) ring_capacity: usize,
/// Reactor busy-poll iter limit before parking. Stored as `u32`
/// for the per-shard counter; the [`Shard`] field carries it
/// forward into the loop.
pub(crate) spin_limit: u32,
/// `Some(N)` = only shards `0..N` arm accept SQE. `None`
/// = every shard accepts (the default; byte-identical to the
/// pre-flag behaviour).
pub(crate) accept_shards: Option<usize>,
/// Total cap on active client conns. `0` = unlimited.
pub(crate) max_clients: usize,
/// Reactor blocking-wait timeout in ms when parked.
pub(crate) park_timeout_ms: u32,
/// Wall-clock-read throttle for the tick check (TTL reaper / live
/// config refresh / auto-AOF-rewrite).
pub(crate) tick_check_every: u32,
/// `[slowlog].slower_than_micros`. Default: `-1` (OFF — zero
/// hot-path cost: every command would otherwise pay an
/// `Instant::now()` pair around dispatch). Set to `10_000` to match
/// Redis's default 10 ms threshold; see [`Self::with_slowlog`] /
/// `CONFIG SET slowlog-log-slower-than 10000`.
pub(crate) slowlog_slower_than_micros: i64,
/// `[slowlog].max_len`. Per-shard cap.
pub(crate) slowlog_max_len: u32,
/// Single-node cluster mode: slot-based key routing (CRC16 `{hashtag}`
/// → contiguous ranges) + one deterministic extra listener per shard at
/// `cluster_port_base + id`. `None` = off (default, zero change).
pub(crate) cluster_port_base: Option<u16>,
/// Replication: when `true`, each shard runs a
/// `ReplicationSource` with `replication_buffer_size` byte budget;
/// every applied mutation is pushed to the backlog. This wires
/// only the producer side — the TCP listener + streaming loop are
/// gated separately by `replication_port_base`. Default `false`.
pub(crate) enable_replication: bool,
/// FEED.* consumer surface. When set, every shard keeps a
/// backlog (even with no replicas) and persists the (generation,
/// offset) cursor via the feed sidecars.
pub(crate) feed_enabled: bool,
/// Per-shard backlog byte budget for the feed (`[feed]
/// feed_buffer_size`); the effective budget is
/// `max(replication_buffer_size, feed_buffer_size)` when both
/// features are on (one backlog, two readers).
pub(crate) feed_buffer_size: u64,
/// Per-shard backlog byte budget when `enable_replication` is set.
/// Fed from `[replication] replication_buffer_size`. Default
/// `256 MiB` (matches the kevy-config default).
pub(crate) replication_buffer_size: u64,
/// v3-cluster replication listener: shard `i` binds at
/// `replication_port_base + i` (mirrors cluster listener pattern;
/// per Issue Ledger I2). `None` = no listener (producer side runs
/// without a network surface, backlog accumulates and evicts —
/// useful for benchmarks). Default `None`.
pub(crate) replication_port_base: Option<u16>,
/// Per-shard SlotTable reconnect-window in ms. After a
/// streaming replica disconnects, its `(replica_id, sent_offset)`
/// is recorded in the shard's `slots` map; slots past this age
/// are reaped on the next shard tick. Default `60_000` (60 s)
/// matches the kevy-config default.
pub(crate) replication_reconnect_window_ms: u32,
/// Per-shard replica inboxes installed by
/// [`Self::with_replica_inboxes`]. Each entry is consumed
/// (via `Option::take`) when its shard is constructed, so the
/// receiver flows from this Vec to the matching `Shard.replica_inbox`.
/// Empty when no replica mode is configured.
pub(crate) replica_inboxes: Vec<Option<crate::replica_inbox::ReplicaInboxReceiver>>,
/// UDS: when `Some(path)`, ALSO bind a Unix-domain stream
/// listener at `path` on shard 0 (single global socket, like valkey's
/// `unixsocket` config). Lets benches/local clients skip TCP loopback
/// overhead. TCP listener stays bound regardless.
#[allow(dead_code)] // consumed during run() via take-into-Shard
pub(crate) unix_socket_path: Option<PathBuf>,
/// Transparent-tiering RAM budget for the WHOLE process, in
/// resolved bytes (the server resolves auto/percent forms before
/// building the runtime). Split evenly across shards at
/// construction. `None` = check the minimal `KEVY_TIER_BUDGET`
/// plain-bytes env knob (back-compat), else tiering off.
pub(crate) tier_budget: Option<u64>,
/// Cold-tier spill dir override (`[tiering] spill_dir`). `None` =
/// `<data_dir>/tier/`.
pub(crate) tier_dir: Option<PathBuf>,
}
impl<C: Commands> Runtime<C> {
/// Start configuring a runtime for `commands`. `Runtime` is its own
/// builder: chain [`Self::bind`] / [`Self::shards`] / the `with_*`
/// setters, then call `run`. Defaults: bind `127.0.0.1:6004`, one
/// shard, AOF on (`EverySec`), data dir `"."`.
#[must_use]
pub fn builder(commands: C) -> Self {
Runtime {
ip: [127, 0, 0, 1],
port: 6004,
nshards: 1,
commands,
data_dir: PathBuf::from("."),
enable_aof: true,
appendfsync: Fsync::EverySec,
auto_aof_rewrite_pct: 100,
auto_aof_rewrite_bytes: 0,
auto_aof_rewrite_interval_secs: 0,
replay_resync: false,
auto_aof_rewrite_min_size: 64 * 1024 * 1024,
ring_capacity: DEFAULT_RING_CAPACITY,
spin_limit: 256,
accept_shards: None,
max_clients: 10_000,
park_timeout_ms: 50,
tick_check_every: 256,
slowlog_slower_than_micros: -1,
slowlog_max_len: 128,
cluster_port_base: None,
enable_replication: false,
feed_enabled: false,
feed_buffer_size: 64 * 1024 * 1024,
replica_inboxes: Vec::new(),
replication_buffer_size: 256 * 1024 * 1024,
replication_port_base: None,
replication_reconnect_window_ms: 60_000,
unix_socket_path: None,
tier_budget: None,
tier_dir: None,
}
}
/// The process-level tiering budget in bytes (`None` = env-knob
/// fallback / off) and per-shard slice of it. The vlog dir root is
/// [`Self::tier_root`].
pub(crate) fn resolved_tier_budget(&self) -> Option<u64> {
self.tier_budget
.or_else(|| std::env::var("KEVY_TIER_BUDGET").ok().and_then(|v| v.parse::<u64>().ok()))
}
/// One shard's slice of the process tiering budget (even split,
/// floored at 1 byte so a tiny budget still tiers rather than
/// silently disabling).
pub(crate) fn per_shard_tier_budget(total: u64, nshards: usize) -> u64 {
(total / nshards.max(1) as u64).max(1)
}
/// The cold-tier root dir: the `[tiering] spill_dir` override, or
/// `<data_dir>/tier/`.
pub(crate) fn tier_root(&self) -> PathBuf {
self.tier_dir.clone().unwrap_or_else(|| self.data_dir.join("tier"))
}
/// Listen address for the client TCP listener (every shard binds it
/// via SO_REUSEPORT). Default `127.0.0.1:6004`.
#[must_use]
pub fn bind(mut self, ip: [u8; 4], port: u16) -> Self {
self.ip = ip;
self.port = port;
self
}
/// Shard (reactor thread) count. Clamped to at least 1. Default 1.
#[must_use]
pub fn shards(mut self, n: usize) -> Self {
self.nshards = n.max(1);
self
}
/// Spawn one thread per shard and run until `stop` is set.
/// UDS: also bind a Unix-domain stream listener at `path`. Lets
/// local clients (and benchmarks) skip the TCP loopback round-trip.
/// Bound on shard 0 only (no SO_REUSEPORT for AF_UNIX, single global
/// socket like valkey's `unixsocket` config). TCP listener stays
/// bound at the configured `port` regardless.
#[must_use]
pub fn with_unix_socket(mut self, path: PathBuf) -> Self {
self.unix_socket_path = Some(path);
self
}
// `run` (and its per-stage build helpers) lives in
// [`crate::runtime_run`] — same `impl<C: Commands> Runtime<C>`,
// split out so this file stays under the 500-LOC house rule.
}