1use crate::KevyResult;
14
15use crate::store::Store;
16
17#[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 pub fn len(&self) -> usize {
52 self.ops.len()
53 }
54
55 pub fn is_empty(&self) -> bool {
57 self.ops.is_empty()
58 }
59
60 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 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 pub fn incr(mut self, key: &[u8]) -> Self {
77 self.ops.push(PendingOp::Incr { key: key.to_vec() });
78 self
79 }
80
81 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 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 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 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 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 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 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 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 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 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 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 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 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 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 pub fn pipeline(&self) -> Pipeline<'_> {
270 Pipeline::new(self)
271 }
272}
273
274#[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];