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