Skip to main content

taquba_workflow/
effects.rs

1use 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/// Application KV effects staged during a step, applied in the same
11/// transaction as the settlement that commits the step's outcome.
12///
13/// Obtained through [`Delivery::effects`](crate::Delivery::effects). Writes and
14/// deletes staged here are applied atomically with the acknowledgement
15/// of the [`StepOutcome`](crate::StepOutcome) the runner returns
16/// (`Continue`, `Succeed`, `Fail` and `Cancel`): either the settlement
17/// and every staged effect commit together or none of them do. Delivery
18/// is at-least-once, so a retried step stages its effects again and
19/// every staged value must be correct when applied more than once. No
20/// effects are applied when the runner returns a
21/// [`StepError`](crate::StepError) or when an external
22/// [`WorkflowRuntime::cancel`](crate::WorkflowRuntime::cancel)
23/// overrides the runner's outcome.
24///
25/// Each operation is validated when it is staged: keys must not start
26/// with the reserved `workflow/` prefix
27/// ([`RESERVED_KV_PREFIX`](crate::RESERVED_KV_PREFIX)), values are
28/// capped at [`taquba::MAX_KV_VALUE_SIZE`] each (an effects set has no
29/// aggregate cap) and a key cannot be staged for both a write and a
30/// delete within one step. The handle is sealed once `run_step`
31/// returns; staging through a clone retained past that point returns
32/// [`Error::EffectsSealed`].
33///
34/// The handle is cheap to clone and clones share one accumulator. Use
35/// [`EffectsHandle::detached`] when constructing a
36/// [`Step`](crate::Step) in tests.
37#[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    /// Seal the state and move out everything staged.
65    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/// Writes and deletes accumulated by an [`EffectsHandle`]. Stored in
91/// the step-output replay record so a replayed delivery applies the
92/// same effects.
93#[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    /// Build a handle bound to no delivery, for constructing a
101    /// [`Step`](crate::Step) in tests. A detached handle accepts and
102    /// validates effects like a delivery-bound one, is never sealed and
103    /// its staged effects are never applied.
104    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    /// Stage a write of `value` under `key` in the caller KV namespace.
115    ///
116    /// # Errors
117    ///
118    /// [`Error::ReservedKvKey`] when `key` starts with the reserved
119    /// `workflow/` prefix, [`Error::Queue`] with
120    /// [`taquba::Error::KvValueTooLarge`] when `value` exceeds
121    /// [`taquba::MAX_KV_VALUE_SIZE`], [`Error::ConflictingKvEffect`]
122    /// when `key` is already staged for a delete and
123    /// [`Error::EffectsSealed`] when the step has returned.
124    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    /// Stage a delete of `key` from the caller KV namespace.
129    ///
130    /// # Errors
131    ///
132    /// [`Error::ReservedKvKey`] when `key` starts with the reserved
133    /// `workflow/` prefix, [`Error::ConflictingKvEffect`] when `key` is
134    /// already staged for a write and [`Error::EffectsSealed`] when the
135    /// step has returned.
136    pub fn delete(&self, key: impl Into<Vec<u8>>) -> Result<()> {
137        self.inner.lock().unwrap().delete(key.into())
138    }
139
140    /// Seal the handle and move out everything staged. An effect staged
141    /// after the seal could not join the settlement, so later staging
142    /// attempts return [`Error::EffectsSealed`].
143    pub(crate) fn seal_and_take(&self) -> StagedEffects {
144        self.inner.lock().unwrap().seal_and_take()
145    }
146}
147
148/// Effects staged by a [`TerminalHook`](crate::TerminalHook) during a
149/// notification delivery, applied in the same transaction as the
150/// notification job's acknowledgement.
151///
152/// Passed to
153/// [`TerminalHook::on_termination`](crate::TerminalHook::on_termination).
154/// Beyond the KV writes and deletes of [`EffectsHandle`] (validated by
155/// the same rules), a hook stages follow-up enqueues, so work driven by
156/// a run's termination is created atomically with the notification
157/// being acknowledged. Effects are applied only when the hook returns
158/// `Ok`; a retried hook stages its effects again.
159///
160/// The handle is sealed once the hook returns; staging through a clone
161/// retained past that point returns [`Error::EffectsSealed`]. Use
162/// [`TerminalEffects::detached`] when invoking a hook directly in
163/// tests.
164#[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    /// Build a handle bound to no delivery, for invoking a
177    /// [`TerminalHook`](crate::TerminalHook) directly in tests. A
178    /// detached handle accepts and validates effects like a
179    /// delivery-bound one, is never sealed and its staged effects are
180    /// never applied.
181    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    /// Stage a follow-up enqueue, committed with the notification's
192    /// acknowledgement.
193    ///
194    /// # Errors
195    ///
196    /// [`Error::EffectsSealed`] when the hook has returned.
197    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    /// Stage a write of `value` under `key` in the caller KV namespace.
207    ///
208    /// # Errors
209    ///
210    /// As [`EffectsHandle::put`].
211    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    /// Stage a delete of `key` from the caller KV namespace.
216    ///
217    /// # Errors
218    ///
219    /// As [`EffectsHandle::delete`].
220    pub fn delete(&self, key: impl Into<Vec<u8>>) -> Result<()> {
221        self.inner.lock().unwrap().kv.delete(key.into())
222    }
223
224    /// Seal the handle and move out everything staged.
225    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}