Skip to main content

kevy_embedded/
ops_pipeline.rs

1//! Non-atomic batched-fsync write queue: `Store::pipeline`.
2//!
3//! Builder-style queue: enqueue any number of writes via fluent
4//! methods, then `commit()` applies them in queue order. Per-shard
5//! AOF appends are batched into one fsync per shard at commit
6//! time, cutting fsync cost from `N` to `min(N, shard_count)`.
7//!
8//! `Pipeline` is not atomic — each op acquires its own per-shard
9//! write lock as it is applied, so other writers see intermediate
10//! states. For transactional semantics use
11//! [`Store::atomic`](crate::Store::atomic).
12
13use crate::KevyResult;
14
15use crate::store::Store;
16
17/// Builder-style write queue. Returned by [`Store::pipeline`]; call
18/// fluent methods to enqueue + `commit()` to apply with batched
19/// AOF fsync.
20pub struct Pipeline<'a> {
21    store: &'a Store,
22    ops: Vec<PendingOp>,
23}
24
25enum PendingOp {
26    Set { key: Vec<u8>, value: Vec<u8> },
27    Del { keys: Vec<Vec<u8>> },
28    Incr { key: Vec<u8> },
29    IncrBy { key: Vec<u8>, delta: i64 },
30    HSet { key: Vec<u8>, pairs: Vec<(Vec<u8>, Vec<u8>)> },
31    HDel { key: Vec<u8>, fields: Vec<Vec<u8>> },
32    HIncrBy { key: Vec<u8>, field: Vec<u8>, delta: i64 },
33    ZAdd { key: Vec<u8>, pairs: Vec<(f64, Vec<u8>)> },
34    ZAddFlags { key: Vec<u8>, pairs: Vec<(f64, Vec<u8>)>, flags: kevy_store::ZaddFlags },
35    ZRem { key: Vec<u8>, members: Vec<Vec<u8>> },
36    ZIncrBy { key: Vec<u8>, delta: f64, member: Vec<u8> },
37    SAdd { key: Vec<u8>, members: Vec<Vec<u8>> },
38    SRem { key: Vec<u8>, members: Vec<Vec<u8>> },
39    LPush { key: Vec<u8>, values: Vec<Vec<u8>> },
40    RPush { key: Vec<u8>, values: Vec<Vec<u8>> },
41}
42
43impl<'a> Pipeline<'a> {
44    pub(crate) fn new(store: &'a Store) -> Self {
45        Self { store, ops: Vec::new() }
46    }
47
48    /// Number of ops queued so far.
49    pub fn len(&self) -> usize {
50        self.ops.len()
51    }
52
53    /// `true` when no ops are queued.
54    pub fn is_empty(&self) -> bool {
55        self.ops.is_empty()
56    }
57
58    // ---- fluent enqueue --------------------------------------------
59
60    /// Queue `SET key value`. Nothing is applied (or AOF-logged)
61    /// until [`commit`](Self::commit) — true of every enqueue method.
62    pub fn set(mut self, key: &[u8], value: &[u8]) -> Self {
63        self.ops.push(PendingOp::Set { key: key.to_vec(), value: value.to_vec() });
64        self
65    }
66
67    /// Queue `DEL key [key ...]`.
68    pub fn del(mut self, keys: &[&[u8]]) -> Self {
69        self.ops.push(PendingOp::Del { keys: keys.iter().map(|k| k.to_vec()).collect() });
70        self
71    }
72
73    /// Queue `INCR key` (commit errors if the value is not an integer).
74    pub fn incr(mut self, key: &[u8]) -> Self {
75        self.ops.push(PendingOp::Incr { key: key.to_vec() });
76        self
77    }
78
79    /// Queue `INCRBY key delta`; negative `delta` does DECR-style work.
80    pub fn incr_by(mut self, key: &[u8], delta: i64) -> Self {
81        self.ops.push(PendingOp::IncrBy { key: key.to_vec(), delta });
82        self
83    }
84
85    /// Queue `HSET key field value [field value ...]`.
86    pub fn hset(mut self, key: &[u8], pairs: &[(&[u8], &[u8])]) -> Self {
87        self.ops.push(PendingOp::HSet {
88            key: key.to_vec(),
89            pairs: pairs.iter().map(|(f, v)| (f.to_vec(), v.to_vec())).collect(),
90        });
91        self
92    }
93
94    /// Queue `HDEL key field [field ...]`.
95    pub fn hdel(mut self, key: &[u8], fields: &[&[u8]]) -> Self {
96        self.ops.push(PendingOp::HDel {
97            key: key.to_vec(),
98            fields: fields.iter().map(|f| f.to_vec()).collect(),
99        });
100        self
101    }
102
103    /// Queue `HINCRBY key field delta`.
104    pub fn hincrby(mut self, key: &[u8], field: &[u8], delta: i64) -> Self {
105        self.ops.push(PendingOp::HIncrBy { key: key.to_vec(), field: field.to_vec(), delta });
106        self
107    }
108
109    /// Queue `ZADD key score member [score member ...]`.
110    pub fn zadd(mut self, key: &[u8], pairs: &[(f64, &[u8])]) -> Self {
111        self.ops.push(PendingOp::ZAdd {
112            key: key.to_vec(),
113            pairs: pairs.iter().map(|(s, m)| (*s, m.to_vec())).collect(),
114        });
115        self
116    }
117
118    /// Flags-aware `ZADD` — e.g. the `GT` monotonic-heal form.
119    pub fn zadd_flags(
120        mut self,
121        key: &[u8],
122        pairs: &[(f64, &[u8])],
123        flags: kevy_store::ZaddFlags,
124    ) -> Self {
125        self.ops.push(PendingOp::ZAddFlags {
126            key: key.to_vec(),
127            pairs: pairs.iter().map(|(s, m)| (*s, m.to_vec())).collect(),
128            flags,
129        });
130        self
131    }
132
133    /// Queue `ZREM key member [member ...]`.
134    pub fn zrem(mut self, key: &[u8], members: &[&[u8]]) -> Self {
135        self.ops.push(PendingOp::ZRem {
136            key: key.to_vec(),
137            members: members.iter().map(|m| m.to_vec()).collect(),
138        });
139        self
140    }
141
142    /// Queue `ZINCRBY key delta member`.
143    pub fn zincrby(mut self, key: &[u8], delta: f64, member: &[u8]) -> Self {
144        self.ops.push(PendingOp::ZIncrBy { key: key.to_vec(), delta, member: member.to_vec() });
145        self
146    }
147
148    /// Queue `SADD key member [member ...]`.
149    pub fn sadd(mut self, key: &[u8], members: &[&[u8]]) -> Self {
150        self.ops.push(PendingOp::SAdd {
151            key: key.to_vec(),
152            members: members.iter().map(|m| m.to_vec()).collect(),
153        });
154        self
155    }
156
157    /// Queue `SREM key member [member ...]`.
158    pub fn srem(mut self, key: &[u8], members: &[&[u8]]) -> Self {
159        self.ops.push(PendingOp::SRem {
160            key: key.to_vec(),
161            members: members.iter().map(|m| m.to_vec()).collect(),
162        });
163        self
164    }
165
166    /// Queue `LPUSH key value [value ...]`.
167    pub fn lpush(mut self, key: &[u8], values: &[&[u8]]) -> Self {
168        self.ops.push(PendingOp::LPush {
169            key: key.to_vec(),
170            values: values.iter().map(|v| v.to_vec()).collect(),
171        });
172        self
173    }
174
175    /// Queue `RPUSH key value [value ...]`.
176    pub fn rpush(mut self, key: &[u8], values: &[&[u8]]) -> Self {
177        self.ops.push(PendingOp::RPush {
178            key: key.to_vec(),
179            values: values.iter().map(|v| v.to_vec()).collect(),
180        });
181        self
182    }
183
184    /// Apply every queued op in order. Each op acquires its own
185    /// per-shard write lock — other writers see intermediate states
186    /// between ops; for transactional semantics use [`Store::atomic`]
187    /// instead. AOF appends batch into one fsync per shard.
188    pub fn commit(mut self) -> KevyResult<()> {
189        let ops = std::mem::take(&mut self.ops);
190        for op in ops {
191            self.apply_one(op)?;
192        }
193        Ok(())
194    }
195
196    // fn-length exemption: pure data-driven op match table — one flat
197    // arm per PendingOp variant, only arg plumbing + one store call.
198    // LOC-WAIVER: data-driven op dispatch table — one store-call arm per PendingOp variant.
199    fn apply_one(&self, op: PendingOp) -> KevyResult<()> {
200        match op {
201            PendingOp::Set { key, value } => {
202                self.store.set(&key, &value)?;
203            }
204            PendingOp::Del { keys } => {
205                let refs: Vec<&[u8]> = keys.iter().map(|k| k.as_slice()).collect();
206                self.store.del(&refs)?;
207            }
208            PendingOp::Incr { key } => {
209                self.store.incr(&key)?;
210            }
211            PendingOp::IncrBy { key, delta } => {
212                self.store.incr_by(&key, delta)?;
213            }
214            PendingOp::HSet { key, pairs } => {
215                let refs: Vec<(&[u8], &[u8])> =
216                    pairs.iter().map(|(f, v)| (f.as_slice(), v.as_slice())).collect();
217                self.store.hset(&key, &refs)?;
218            }
219            PendingOp::HDel { key, fields } => {
220                let refs: Vec<&[u8]> = fields.iter().map(|f| f.as_slice()).collect();
221                self.store.hdel(&key, &refs)?;
222            }
223            PendingOp::HIncrBy { key, field, delta } => {
224                self.store.hincrby(&key, &field, delta)?;
225            }
226            PendingOp::ZAdd { key, pairs } => {
227                let refs: Vec<(f64, &[u8])> =
228                    pairs.iter().map(|(s, m)| (*s, m.as_slice())).collect();
229                self.store.zadd(&key, &refs)?;
230            }
231            PendingOp::ZAddFlags { key, pairs, flags } => {
232                let refs: Vec<(f64, &[u8])> =
233                    pairs.iter().map(|(s, m)| (*s, m.as_slice())).collect();
234                self.store.zadd_flags(&key, &refs, flags)?;
235            }
236            PendingOp::ZRem { key, members } => {
237                let refs: Vec<&[u8]> = members.iter().map(|m| m.as_slice()).collect();
238                self.store.zrem(&key, &refs)?;
239            }
240            PendingOp::ZIncrBy { key, delta, member } => {
241                self.store.zincrby(&key, delta, &member)?;
242            }
243            PendingOp::SAdd { key, members } => {
244                let refs: Vec<&[u8]> = members.iter().map(|m| m.as_slice()).collect();
245                self.store.sadd(&key, &refs)?;
246            }
247            PendingOp::SRem { key, members } => {
248                let refs: Vec<&[u8]> = members.iter().map(|m| m.as_slice()).collect();
249                self.store.srem(&key, &refs)?;
250            }
251            PendingOp::LPush { key, values } => {
252                let refs: Vec<&[u8]> = values.iter().map(|v| v.as_slice()).collect();
253                self.store.lpush(&key, &refs)?;
254            }
255            PendingOp::RPush { key, values } => {
256                let refs: Vec<&[u8]> = values.iter().map(|v| v.as_slice()).collect();
257                self.store.rpush(&key, &refs)?;
258            }
259        }
260        Ok(())
261    }
262}
263
264impl Store {
265    /// Begin a [`Pipeline`] — fluent write queue. Add ops via
266    /// `.set(...).hset(...).zadd(...)` then call `.commit()`.
267    pub fn pipeline(&self) -> Pipeline<'_> {
268        Pipeline::new(self)
269    }
270}
271
272/// Parity manifest: command names `Pipeline` implements.
273/// Cross-checked against `kevy_resp::ops_table` in
274/// `store_tests_op_table.rs` — update BOTH when adding an op.
275#[cfg_attr(not(test), allow(dead_code))]
276pub(crate) const PIPELINE_OPS: &[&str] = &[
277    "SET", "DEL", "INCR", "INCRBY", "HSET", "HDEL", "HINCRBY", "ZADD", "ZREM", "ZINCRBY", "SADD",
278    "SREM", "LPUSH", "RPUSH",
279];