taquba_workflow/
effects.rs1use std::collections::{HashMap, HashSet};
2use std::sync::{Arc, Mutex};
3
4use serde::{Deserialize, Serialize};
5use taquba::EnqueueRequest;
6
7use crate::error::{Error, Result};
8use crate::keys::RESERVED_KV_PREFIX;
9
10#[derive(Debug, Clone)]
38pub struct EffectsHandle {
39 inner: Arc<Mutex<EffectsState>>,
40}
41
42#[derive(Debug, Default)]
43struct EffectsState {
44 staged: StagedEffects,
45 sealed: bool,
46}
47
48impl EffectsState {
49 fn put(&mut self, key: Vec<u8>, value: Vec<u8>) -> Result<()> {
50 self.check_key(&key)?;
51 if value.len() > taquba::MAX_KV_VALUE_SIZE {
52 return Err(Error::Queue(taquba::Error::KvValueTooLarge {
53 size: value.len(),
54 max: taquba::MAX_KV_VALUE_SIZE,
55 }));
56 }
57 if self.staged.deletes.contains(&key) {
58 return Err(Error::ConflictingKvEffect(display_key(&key)));
59 }
60 self.staged.writes.insert(key, value);
61 Ok(())
62 }
63
64 fn seal_and_take(&mut self) -> StagedEffects {
66 self.sealed = true;
67 std::mem::take(&mut self.staged)
68 }
69
70 fn delete(&mut self, key: Vec<u8>) -> Result<()> {
71 self.check_key(&key)?;
72 if self.staged.writes.contains_key(&key) {
73 return Err(Error::ConflictingKvEffect(display_key(&key)));
74 }
75 self.staged.deletes.insert(key);
76 Ok(())
77 }
78
79 fn check_key(&self, key: &[u8]) -> Result<()> {
80 if self.sealed {
81 return Err(Error::EffectsSealed);
82 }
83 if key.starts_with(RESERVED_KV_PREFIX.as_bytes()) {
84 return Err(Error::ReservedKvKey(display_key(key)));
85 }
86 Ok(())
87 }
88}
89
90#[derive(Debug, Clone, Default, Serialize, Deserialize)]
94pub(crate) struct StagedEffects {
95 pub(crate) writes: HashMap<Vec<u8>, Vec<u8>>,
96 pub(crate) deletes: HashSet<Vec<u8>>,
97}
98
99impl EffectsHandle {
100 pub fn detached() -> Self {
105 Self::for_delivery()
106 }
107
108 pub(crate) fn for_delivery() -> Self {
109 Self {
110 inner: Arc::new(Mutex::new(EffectsState::default())),
111 }
112 }
113
114 pub fn put(&self, key: impl Into<Vec<u8>>, value: impl Into<Vec<u8>>) -> Result<()> {
125 self.inner.lock().unwrap().put(key.into(), value.into())
126 }
127
128 pub fn delete(&self, key: impl Into<Vec<u8>>) -> Result<()> {
137 self.inner.lock().unwrap().delete(key.into())
138 }
139
140 pub(crate) fn seal_and_take(&self) -> StagedEffects {
144 self.inner.lock().unwrap().seal_and_take()
145 }
146}
147
148#[derive(Debug, Clone)]
165pub struct TerminalEffects {
166 inner: Arc<Mutex<TerminalState>>,
167}
168
169#[derive(Debug, Default)]
170struct TerminalState {
171 kv: EffectsState,
172 enqueues: Vec<EnqueueRequest>,
173}
174
175impl TerminalEffects {
176 pub fn detached() -> Self {
182 Self::for_delivery()
183 }
184
185 pub(crate) fn for_delivery() -> Self {
186 Self {
187 inner: Arc::new(Mutex::new(TerminalState::default())),
188 }
189 }
190
191 pub fn enqueue(&self, request: EnqueueRequest) -> Result<()> {
198 let mut state = self.inner.lock().unwrap();
199 if state.kv.sealed {
200 return Err(Error::EffectsSealed);
201 }
202 state.enqueues.push(request);
203 Ok(())
204 }
205
206 pub fn put(&self, key: impl Into<Vec<u8>>, value: impl Into<Vec<u8>>) -> Result<()> {
212 self.inner.lock().unwrap().kv.put(key.into(), value.into())
213 }
214
215 pub fn delete(&self, key: impl Into<Vec<u8>>) -> Result<()> {
221 self.inner.lock().unwrap().kv.delete(key.into())
222 }
223
224 pub(crate) fn seal_and_take(&self) -> (StagedEffects, Vec<EnqueueRequest>) {
226 let mut state = self.inner.lock().unwrap();
227 let staged = state.kv.seal_and_take();
228 (staged, std::mem::take(&mut state.enqueues))
229 }
230}
231
232fn display_key(key: &[u8]) -> String {
233 String::from_utf8_lossy(key).into_owned()
234}
235
236#[cfg(test)]
237mod tests {
238 use super::*;
239
240 #[test]
241 fn staging_validates_keys_values_and_conflicts() {
242 let handle = EffectsHandle::detached();
243 assert!(matches!(
244 handle.put("workflow/x", "v"),
245 Err(Error::ReservedKvKey(_))
246 ));
247 assert!(matches!(
248 handle.delete("workflow/x"),
249 Err(Error::ReservedKvKey(_))
250 ));
251 let oversized = vec![0u8; taquba::MAX_KV_VALUE_SIZE + 1];
252 assert!(matches!(
253 handle.put("k", oversized),
254 Err(Error::Queue(taquba::Error::KvValueTooLarge { .. }))
255 ));
256 handle.put("a", "v").unwrap();
257 assert!(matches!(
258 handle.delete("a"),
259 Err(Error::ConflictingKvEffect(_))
260 ));
261 handle.delete("b").unwrap();
262 assert!(matches!(
263 handle.put("b", "v"),
264 Err(Error::ConflictingKvEffect(_))
265 ));
266 }
267
268 #[test]
269 fn a_clone_stages_into_the_shared_accumulator_until_the_seal() {
270 let handle = EffectsHandle::for_delivery();
271 let clone = handle.clone();
272 clone.put("a", "v").unwrap();
273 clone.delete("b").unwrap();
274 let staged = handle.seal_and_take();
275 assert_eq!(staged.writes.get(b"a".as_slice()), Some(&b"v".to_vec()));
276 assert!(staged.deletes.contains(b"b".as_slice()));
277 assert!(matches!(clone.put("c", "v"), Err(Error::EffectsSealed)));
278 assert!(matches!(clone.delete("c"), Err(Error::EffectsSealed)));
279 }
280
281 #[test]
282 fn the_terminal_handle_applies_the_staging_and_seal_rules() {
283 let handle = TerminalEffects::for_delivery();
284 assert!(matches!(
285 handle.put("workflow/x", "v"),
286 Err(Error::ReservedKvKey(_))
287 ));
288 assert!(matches!(
289 handle.delete("workflow/x"),
290 Err(Error::ReservedKvKey(_))
291 ));
292 handle.put("a", "v").unwrap();
293 assert!(matches!(
294 handle.delete("a"),
295 Err(Error::ConflictingKvEffect(_))
296 ));
297 handle.delete("b").unwrap();
298 handle
299 .enqueue(taquba::EnqueueRequest {
300 queue: "side".to_string(),
301 payload: Vec::new(),
302 options: Default::default(),
303 })
304 .unwrap();
305 let (staged, enqueues) = handle.seal_and_take();
306 assert_eq!(staged.writes.len(), 1);
307 assert_eq!(staged.deletes.len(), 1);
308 assert_eq!(enqueues.len(), 1);
309 assert!(matches!(handle.put("c", "v"), Err(Error::EffectsSealed)));
310 assert!(matches!(handle.delete("c"), Err(Error::EffectsSealed)));
311 assert!(matches!(
312 handle.enqueue(taquba::EnqueueRequest {
313 queue: "side".to_string(),
314 payload: Vec::new(),
315 options: Default::default(),
316 }),
317 Err(Error::EffectsSealed)
318 ));
319 }
320}