Skip to main content

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