kevy_config/cluster.rs
1//! `[cluster]` section schema — single-node cluster mode plus the
2//! quorum-election peer list.
3//!
4//! The peer list is a flat comma-separated string
5//! (`peers = "id@host:port,..."`) rather than TOML's
6//! `[[array_of_tables]]`: the hand-rolled 0-dep parser does not
7//! support arrays of tables, and the structural future-need is
8//! bounded (per-peer TLS, auth, and region are explicitly out of
9//! charter for the election subsystem).
10
11
12/// `[cluster]` section — single-node cluster mode: keys route by
13/// Redis-cluster slot (CRC16 `{hashtag}` & 16383) and every shard `i`
14/// gets a second, deterministic listener at `port_base + i` that answers
15/// wrong-shard keys with `-MOVED`, so stock cluster-aware clients
16/// (`redis-benchmark --cluster`, `redis-cli -c`) can address shards
17/// directly. The main SO_REUSEPORT port keeps full forward-anywhere
18/// behaviour for non-cluster clients. Not hot-settable: the routing
19/// scheme is a startup property of the data dir (`shards.meta`).
20///
21/// The struct is `Clone` but not `Copy` (since `peers` and `scopes`
22/// hold owned vectors). Most call sites just clone the per-tick
23/// `Config` snapshot via `Arc<Config>`, so this is invisible in the
24/// hot path.
25#[derive(Debug, Clone, PartialEq, Eq, Default)]
26pub struct ClusterSection {
27 /// Enable cluster mode. Default `false` (zero change).
28 pub enabled: bool,
29 /// First cluster port (shard `i` listens at `port_base + i`).
30 /// `0` (default) = `server.port + 1`.
31 pub port_base: u16,
32 /// This node's stable id for the quorum election (≤ 32 B
33 /// ASCII; unique across the cluster). Default empty —
34 /// `kevy-elect` is dormant unless both `node_id` and `peers`
35 /// are set, so existing configs need no edit.
36 pub node_id: String,
37 /// First election-control listener port; shard `i` binds at
38 /// `elect_port_base + i`. Default `0` → `server.port + 200`
39 /// (locked by the `resolved_elect_port_base` unit test).
40 pub elect_port_base: u16,
41 /// Operator-declared peer list for `kevy-elect`. Empty when
42 /// failover is not configured. Each entry is one cluster node
43 /// (including potentially *this* node — kevy-elect filters
44 /// self by matching `node_id`).
45 pub peers: Vec<PeerEntry>,
46 /// Scope declarations: each entry pins a key prefix to a writer
47 /// node (and optional fallback). Empty when scope-based
48 /// multi-writer is off. Same flat-string TOML shape as `peers`
49 /// — `scopes = "prefix=writer[|fallback],..."`.
50 pub scopes: Vec<ScopeEntry>,
51}
52
53/// One scope declaration parsed from the TOML
54/// `scopes = "prefix=writer[|fallback],..."` shape. Mirrors the
55/// `kevy_scope::Scope` data; kept duplicated here so kevy-config
56/// stays leaf-level and doesn't depend on kevy-scope (the dependency
57/// direction is kevy-scope ← kevy-config consumer, not the other
58/// way).
59#[derive(Debug, Clone, PartialEq, Eq)]
60pub struct ScopeEntry {
61 /// Key-prefix bytes the scope owns. Bytes (not String) because
62 /// kevy keys are arbitrary; common keys (`app:billing:`) are
63 /// UTF-8 but the type signature stays honest.
64 pub prefix: Vec<u8>,
65 /// Declared writer's node id.
66 pub writer: String,
67 /// Optional fallback node id (F4).
68 pub fallback: Option<String>,
69}
70
71impl ScopeEntry {
72 /// Render back to the `prefix=writer[|fallback]` token shape —
73 /// exact inverse of [`Self::parse_one`] for entries it produced.
74 /// The prefix is parsed from TOML text, so it is UTF-8 whenever
75 /// this is used on a config round-trip (`CONFIG REWRITE`).
76 pub fn to_token(&self) -> String {
77 let prefix = String::from_utf8_lossy(&self.prefix);
78 match &self.fallback {
79 Some(fb) => format!("{prefix}={}|{fb}", self.writer),
80 None => format!("{prefix}={}", self.writer),
81 }
82 }
83
84 /// Parse one `prefix=writer[|fallback]` token. The first `=`
85 /// splits prefix from owner spec; the writer may carry an
86 /// optional `|fallback` suffix. Returns `None` on any shape
87 /// problem (missing `=`, empty fields, prefix containing `,`).
88 pub fn parse_one(token: &str) -> Option<Self> {
89 // Reject commas inside the token — `parse_list` already
90 // split on commas, so a comma here means the operator typed
91 // `prefix=a,b` (ambiguous owner list); we treat that as a
92 // parse error rather than silently take only `a`.
93 if token.contains(',') {
94 return None;
95 }
96 let (prefix, owners) = token.split_once('=')?;
97 if prefix.is_empty() || owners.is_empty() {
98 return None;
99 }
100 let (writer, fallback) = match owners.split_once('|') {
101 Some((w, f)) if !w.is_empty() && !f.is_empty() => (w, Some(f.to_string())),
102 Some(_) => return None, // `|` present but one side empty
103 None => (owners, None),
104 };
105 Some(ScopeEntry {
106 prefix: prefix.as_bytes().to_vec(),
107 writer: writer.to_string(),
108 fallback,
109 })
110 }
111
112 /// Parse a `scopes = "..."` value — comma-separated list of
113 /// `prefix=writer[|fallback]` tokens. Empty + whitespace-only
114 /// tokens are dropped; trailing comma tolerated. Same
115 /// error-on-first-bad-token contract as
116 /// [`PeerEntry::parse_list`].
117 pub fn parse_list(s: &str) -> Result<Vec<ScopeEntry>, String> {
118 let mut out = Vec::new();
119 for raw in s.split(',') {
120 let token = raw.trim();
121 if token.is_empty() {
122 continue;
123 }
124 match Self::parse_one(token) {
125 Some(p) => out.push(p),
126 None => return Err(token.to_string()),
127 }
128 }
129 Ok(out)
130 }
131}
132
133/// One peer in the `kevy-elect` quorum, parsed from the TOML
134/// shape `peers = "id@host:port,id@host:port,..."` (per
135/// T1.5.4.5 decision (b) — a parser-extension-free representation
136/// that works with kevy-config's flat KV-only TOML).
137///
138/// **v1.55** extends the syntax with an optional second port for
139/// **client-facing** address (used by `-MISDIRECTED writer is`
140/// replies): `id@host:elect_port:client_port`. When the extended
141/// form is used, kevy-elect still binds the elect_port, while
142/// kevy-scope's MISDIRECTED encoder reports `host:client_port` to
143/// the client so the client can actually reconnect to the writer.
144/// Without the extended form, MISDIRECTED reports `host:elect_port`
145/// (the v1.45 documented behaviour, retained for compat).
146#[derive(Debug, Clone, PartialEq, Eq)]
147pub struct PeerEntry {
148 /// Peer's stable node id.
149 pub node_id: String,
150 /// Peer's host (IPv4 dotted literal or DNS-resolvable name).
151 pub host: String,
152 /// Peer's election-control port (= peer's
153 /// `cluster.elect_port_base + 0`, the shard 0 listener).
154 pub port: u16,
155 /// **v1.55** — Peer's client-facing TCP port (the port other
156 /// kevy nodes / `redis-cli` connect to for normal operations).
157 /// `None` = unset (legacy syntax `id@host:port`); MISDIRECTED
158 /// replies fall back to `port` in that case. Set via extended
159 /// syntax `id@host:elect_port:client_port`.
160 pub client_port: Option<u16>,
161}
162
163impl PeerEntry {
164 /// Render back to the `id@host:port[:client_port]` token shape —
165 /// exact inverse of [`Self::parse_one`] for entries it produced.
166 pub fn to_token(&self) -> String {
167 match self.client_port {
168 Some(cp) => format!("{}@{}:{}:{cp}", self.node_id, self.host, self.port),
169 None => format!("{}@{}:{}", self.node_id, self.host, self.port),
170 }
171 }
172
173 /// Parse one peer token. Accepts two shapes:
174 /// - **Legacy**: `id@host:port` (v1.x — `port` = elect port).
175 /// - **v1.55+**: `id@host:elect_port:client_port` (sets
176 /// `client_port` so MISDIRECTED reports a port the client
177 /// can actually connect to).
178 ///
179 /// Returns `None` on any shape problem (empty fields, non-numeric
180 /// ports, port overflow).
181 pub fn parse_one(token: &str) -> Option<Self> {
182 let (node_id, rest) = token.split_once('@')?;
183 if node_id.is_empty() {
184 return None;
185 }
186 // Find the last colon (== client_port if extended, else elect_port).
187 let last_colon = rest.rfind(':')?;
188 let after_last: u16 = rest[last_colon + 1..].parse().ok()?;
189 let before_last = &rest[..last_colon];
190 // Try the extended form: split `before_last` on its own last colon.
191 if let Some(second_last) = before_last.rfind(':') {
192 // before_last = `host:elect_port`; after_last = `client_port`.
193 let host = &before_last[..second_last];
194 if host.is_empty() {
195 return None;
196 }
197 if let Ok(elect) = before_last[second_last + 1..].parse::<u16>() {
198 return Some(PeerEntry {
199 node_id: node_id.to_string(),
200 host: host.to_string(),
201 port: elect,
202 client_port: Some(after_last),
203 });
204 }
205 }
206 // Legacy form: `host:port` (port = elect).
207 let host = before_last;
208 if host.is_empty() {
209 return None;
210 }
211 Some(PeerEntry {
212 node_id: node_id.to_string(),
213 host: host.to_string(),
214 port: after_last,
215 client_port: None,
216 })
217 }
218
219 /// Parse the `peers = "..."` value — a comma-separated list of
220 /// `id@host:port` tokens. Empty + all-whitespace tokens are
221 /// dropped silently (a trailing comma after the last entry is
222 /// tolerated). Returns `Err(token)` on the first unparseable
223 /// token, with the offending token in the error for diagnostic.
224 pub fn parse_list(s: &str) -> Result<Vec<PeerEntry>, String> {
225 let mut out = Vec::new();
226 for raw in s.split(',') {
227 let token = raw.trim();
228 if token.is_empty() {
229 continue;
230 }
231 match Self::parse_one(token) {
232 Some(p) => out.push(p),
233 None => return Err(token.to_string()),
234 }
235 }
236 Ok(out)
237 }
238}
239
240#[cfg(test)]
241mod peer_entry_tests {
242 use super::*;
243
244 #[test]
245 fn parse_one_basic() {
246 let p = PeerEntry::parse_one("node-1@10.0.0.1:6004").unwrap();
247 assert_eq!(p.node_id, "node-1");
248 assert_eq!(p.host, "10.0.0.1");
249 assert_eq!(p.port, 6004);
250 assert_eq!(p.client_port, None);
251 }
252
253 #[test]
254 fn parse_one_v1_55_extended_form_sets_client_port() {
255 // v1.55: `id@host:elect_port:client_port` syntax — addresses
256 // v1.45.x finding (MISDIRECTED reply uses elect_port instead
257 // of main client port).
258 let p = PeerEntry::parse_one("node-1@10.0.0.1:6011:6004").unwrap();
259 assert_eq!(p.node_id, "node-1");
260 assert_eq!(p.host, "10.0.0.1");
261 assert_eq!(p.port, 6011);
262 assert_eq!(p.client_port, Some(6004));
263 }
264
265 #[test]
266 fn parse_one_extended_form_dns_host() {
267 let p = PeerEntry::parse_one("primary@db-east.local:6011:6004").unwrap();
268 assert_eq!(p.host, "db-east.local");
269 assert_eq!(p.port, 6011);
270 assert_eq!(p.client_port, Some(6004));
271 }
272
273 #[test]
274 fn parse_one_dns_host() {
275 let p = PeerEntry::parse_one("primary@db-east.local:6105").unwrap();
276 assert_eq!(p.host, "db-east.local");
277 assert_eq!(p.port, 6105);
278 }
279
280 #[test]
281 fn parse_one_rejects_empty_id_host_or_bad_port() {
282 assert!(PeerEntry::parse_one("@host:6004").is_none());
283 assert!(PeerEntry::parse_one("id@:6004").is_none());
284 assert!(PeerEntry::parse_one("id@host:NaN").is_none());
285 assert!(PeerEntry::parse_one("id@host:99999").is_none()); // u16 overflow
286 assert!(PeerEntry::parse_one("no-at-or-colon").is_none());
287 }
288
289 #[test]
290 fn parse_list_three_peers_trim_tolerated() {
291 let s = "a@1.1.1.1:6004, b@1.1.1.2:6004 ,c@1.1.1.3:6004";
292 let peers = PeerEntry::parse_list(s).unwrap();
293 assert_eq!(peers.len(), 3);
294 assert_eq!(peers[1].node_id, "b");
295 }
296
297 #[test]
298 fn parse_list_trailing_comma_ok() {
299 let peers = PeerEntry::parse_list("a@h:1,b@h:2,").unwrap();
300 assert_eq!(peers.len(), 2);
301 }
302
303 #[test]
304 fn parse_list_first_bad_token_errs() {
305 let err = PeerEntry::parse_list("a@h:1,bad-token,c@h:3").unwrap_err();
306 assert_eq!(err, "bad-token");
307 }
308
309 #[test]
310 fn parse_list_empty_is_empty() {
311 assert_eq!(PeerEntry::parse_list("").unwrap(), Vec::<PeerEntry>::new());
312 assert_eq!(PeerEntry::parse_list(" ").unwrap(), Vec::<PeerEntry>::new());
313 }
314
315 #[test]
316 fn to_token_round_trips() {
317 for tok in ["node-1@10.0.0.1:6004", "node-1@10.0.0.1:6011:6004", "p@db-east.local:6105"] {
318 let p = PeerEntry::parse_one(tok).unwrap();
319 assert_eq!(p.to_token(), tok);
320 assert_eq!(PeerEntry::parse_one(&p.to_token()), Some(p));
321 }
322 }
323}
324
325#[cfg(test)]
326mod scope_entry_tests {
327 use super::*;
328
329 #[test]
330 fn parse_one_writer_only() {
331 let s = ScopeEntry::parse_one("app:billing:=embed-billing-1").unwrap();
332 assert_eq!(s.prefix, b"app:billing:");
333 assert_eq!(s.writer, "embed-billing-1");
334 assert_eq!(s.fallback, None);
335 }
336
337 #[test]
338 fn parse_one_writer_and_fallback() {
339 let s = ScopeEntry::parse_one("app:billing:=embed-1|fb-server-eu").unwrap();
340 assert_eq!(s.writer, "embed-1");
341 assert_eq!(s.fallback.as_deref(), Some("fb-server-eu"));
342 }
343
344 #[test]
345 fn parse_one_prefix_with_colons() {
346 // Colon-heavy prefixes are the common case; only `=` and `,`
347 // are reserved.
348 let s = ScopeEntry::parse_one("ns:tenant:42:=w").unwrap();
349 assert_eq!(s.prefix, b"ns:tenant:42:");
350 }
351
352 #[test]
353 fn parse_one_rejects_empty_prefix_or_writer() {
354 assert!(ScopeEntry::parse_one("=writer").is_none());
355 assert!(ScopeEntry::parse_one("prefix=").is_none());
356 assert!(ScopeEntry::parse_one("no-equals").is_none());
357 }
358
359 #[test]
360 fn parse_one_rejects_empty_fallback_side() {
361 assert!(ScopeEntry::parse_one("p=writer|").is_none());
362 assert!(ScopeEntry::parse_one("p=|fb").is_none());
363 }
364
365 #[test]
366 fn parse_one_rejects_embedded_comma() {
367 // The split-on-comma in `parse_list` makes commas inside a
368 // token a parse error — operator probably typo'd
369 // `prefix=writer,fallback` instead of `prefix=writer|fallback`.
370 assert!(ScopeEntry::parse_one("p=writer,other").is_none());
371 }
372
373 #[test]
374 fn parse_list_two_scopes() {
375 let v = ScopeEntry::parse_list("app:billing:=w-bill|fb, app:auth:=w-auth").unwrap();
376 assert_eq!(v.len(), 2);
377 assert_eq!(v[0].writer, "w-bill");
378 assert_eq!(v[0].fallback.as_deref(), Some("fb"));
379 assert_eq!(v[1].writer, "w-auth");
380 assert!(v[1].fallback.is_none());
381 }
382
383 #[test]
384 fn parse_list_first_bad_token_errs() {
385 let err = ScopeEntry::parse_list("p1=w1,no-eq,p3=w3").unwrap_err();
386 assert_eq!(err, "no-eq");
387 }
388
389 #[test]
390 fn to_token_round_trips() {
391 for tok in ["app:billing:=embed-billing-1", "app:billing:=embed-1|fb-server-eu"] {
392 let s = ScopeEntry::parse_one(tok).unwrap();
393 assert_eq!(s.to_token(), tok);
394 assert_eq!(ScopeEntry::parse_one(&s.to_token()), Some(s));
395 }
396 }
397}
398