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