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