Skip to main content

taquba_workflow/
kv.rs

1use std::sync::Arc;
2
3use bytes::Bytes;
4use taquba::Queue;
5
6use crate::error::Result;
7
8/// Read access to Taquba's caller KV namespace during a step.
9///
10/// Obtained through [`Delivery::kv`](crate::Delivery::kv). [`get`](Self::get)
11/// reads from committed state: a value written by an earlier settlement (a
12/// previous step's [`EffectsHandle`](crate::EffectsHandle) writes, a
13/// [`RunSpec::effects`](crate::RunSpec::effects) write, a direct
14/// [`taquba::Queue::kv_put`]) is visible, and an effect staged by the current
15/// step becomes visible only once this step's settlement commits it. The read
16/// is also not transactional with that settlement: a value read here can change
17/// before the step's outcome commits. Delivery is at-least-once, so a read that
18/// misses (for example a marker a crashed settlement never committed)
19/// re-executes work that must be idempotent downstream.
20///
21/// The handle is cheap to clone and exposes no write or settlement
22/// operation. Use [`KvReadHandle::detached`] when constructing a
23/// [`Step`](crate::Step) in tests.
24#[derive(Clone)]
25pub struct KvReadHandle {
26    queue: Option<Arc<Queue>>,
27}
28
29impl std::fmt::Debug for KvReadHandle {
30    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
31        match &self.queue {
32            Some(_) => f.debug_struct("KvReadHandle").finish_non_exhaustive(),
33            None => f.write_str("KvReadHandle::detached"),
34        }
35    }
36}
37
38impl KvReadHandle {
39    /// Build a handle bound to no queue, for constructing a
40    /// [`Step`](crate::Step) in tests. [`get`](Self::get) on a detached
41    /// handle returns `Ok(None)` for every key.
42    pub fn detached() -> Self {
43        Self { queue: None }
44    }
45
46    pub(crate) fn for_delivery(queue: Arc<Queue>) -> Self {
47        Self { queue: Some(queue) }
48    }
49
50    /// Read the committed value under `key` from the caller KV
51    /// namespace, `None` when no value exists.
52    ///
53    /// # Errors
54    ///
55    /// [`Error::Queue`](crate::Error::Queue) when the underlying read
56    /// fails.
57    pub async fn get(&self, key: &[u8]) -> Result<Option<Bytes>> {
58        match &self.queue {
59            Some(queue) => Ok(queue.view().kv_get(key).await?),
60            None => Ok(None),
61        }
62    }
63}
64
65#[cfg(test)]
66mod tests {
67    use super::*;
68
69    #[tokio::test]
70    async fn a_detached_handle_reads_no_value() {
71        let handle = KvReadHandle::detached();
72        assert_eq!(handle.get(b"any").await.unwrap(), None);
73    }
74}