Skip to main content

taquba_workflow/
keys.rs

1//! The runtime's reserved namespaces: the `workflow.` header keys, the
2//! `workflow/` prefix in the caller KV namespace and the builders of the
3//! durable keys within it, the step dedup-key prefix and the validated
4//! [`RunId`].
5
6use std::borrow::Borrow;
7use std::fmt;
8use std::ops::Deref;
9use std::str::FromStr;
10
11use serde::{Deserialize, Deserializer, Serialize};
12
13use crate::error::{Error, Result};
14
15/// Header key carrying the run identifier on every step job.
16pub const HEADER_RUN_ID: &str = "workflow.run_id";
17/// Header key carrying the zero-based step number on every step job.
18pub const HEADER_STEP: &str = "workflow.step";
19/// Reserved prefix the runtime owns on step-job headers. Submitter-supplied
20/// headers must not start with this prefix; if they do, the runtime treats
21/// them as its own and strips them before invoking the runner.
22pub const RESERVED_HEADER_PREFIX: &str = "workflow.";
23
24/// Reserved prefix the runtime owns in the caller KV namespace. Keys passed via
25/// [`RunSpec::effects`](crate::RunSpec::effects) or staged through an
26/// [`crate::EffectsHandle`] must not start with this prefix. Such a key is
27/// rejected with [`Error::ReservedKvKey`].
28pub const RESERVED_KV_PREFIX: &str = "workflow/";
29
30/// Header key marking a job as a terminal-notification job, whose
31/// payload is the run's committed outcome and whose worker is the
32/// configured [`TerminalHook`](crate::TerminalHook).
33pub const HEADER_TERMINAL: &str = "workflow.terminal";
34
35/// Header key marking a step job as a signal waiter; the value is the
36/// correlation key the waiter is registered under.
37pub const HEADER_SIGNAL_WAIT: &str = "workflow.signal_wait";
38/// Header key marking a step job whose signal was already consumed at the
39/// previous step's settlement; the payload is read from the durable
40/// delivered record.
41pub const HEADER_SIGNAL_DELIVERED: &str = "workflow.signal_delivered";
42
43/// Header key naming the group a run is a member of, set on every step
44/// job of a grouped run.
45pub const HEADER_GROUP: &str = "workflow.group";
46/// Header key naming a grouped run's member key within its group, set
47/// beside [`HEADER_GROUP`].
48pub const HEADER_GROUP_KEY: &str = "workflow.group_key";
49
50pub(crate) const DEDUP_PREFIX: &str = "run:";
51
52/// Maximum byte length of a [`RunId`], the limit Taquba applies to a
53/// caller-supplied job id.
54pub const MAX_RUN_ID_LEN: usize = 128;
55
56/// Prefix for the durable per-run record in Taquba's user KV namespace.
57pub(crate) const RUN_KV_PREFIX: &[u8] = b"workflow/runs/";
58
59/// Prefix for the durable current-step pointer: `run id -> (step, job id)`
60/// of the queue job currently representing the run, written with the
61/// step-0 enqueue, rewritten with every advance and deleted with the
62/// run's termination.
63pub(crate) const STEP_KV_PREFIX: &[u8] = b"workflow/steps/";
64
65/// Prefix for the durable waiter index: `correlation key -> job id` of the
66/// step job waiting on that key.
67pub(crate) const SIGNAL_WAIT_KV_PREFIX: &[u8] = b"workflow/signal-wait/";
68
69/// Prefix for the durable signal buffer: `correlation key -> payload` of a
70/// signal that arrived while no waiter was registered.
71pub(crate) const SIGNAL_BUF_KV_PREFIX: &[u8] = b"workflow/signal-buf/";
72
73/// Prefix for the durable delivered record: `(run id, step) -> payload` of
74/// a signal consumed on the waiter's behalf, read when that step is
75/// claimed and deleted with its settlement.
76pub(crate) const SIGNAL_DELIVERED_KV_PREFIX: &[u8] = b"workflow/signal-delivered/";
77
78/// Prefix of the terminal markers of runs, read by the memo retention
79/// sweep: entries of a [`taquba::ExpiryIndex`] with the run id as the
80/// suffix.
81pub(crate) const TERMINAL_KV_PREFIX: &[u8] = b"workflow/terminals/";
82
83/// Prefix for the durable terminal record of a run:
84/// `workflow/outcomes/{run_id}`, written in the settlement that
85/// terminates the run and removed with the run's memo entries by the
86/// memo sweep.
87pub(crate) const OUTCOME_KV_PREFIX: &[u8] = b"workflow/outcomes/";
88
89/// Prefix for the durable member records of run groups:
90/// `workflow/groups/{group_id}/{key}`, one per member, written with the
91/// member's submission and rewritten in the settlement that terminates
92/// it.
93pub(crate) const GROUP_KV_PREFIX: &[u8] = b"workflow/groups/";
94
95/// Prefix under which the member records of one group are stored.
96pub(crate) fn group_members_kv_prefix(group_id: &RunId) -> Vec<u8> {
97    prefixed(GROUP_KV_PREFIX, &format!("{group_id}/"))
98}
99
100/// Key of the member record of `key` in group `group_id`.
101pub(crate) fn group_member_kv_key(group_id: &RunId, key: &str) -> Vec<u8> {
102    prefixed(&group_members_kv_prefix(group_id), key)
103}
104
105/// Prefix of the terminal markers of run groups, read by the group
106/// retention sweep: entries of a [`taquba::ExpiryIndex`] with the group
107/// id as the suffix.
108pub(crate) const GROUP_TERMINAL_KV_PREFIX: &[u8] = b"workflow/group-terminals/";
109
110/// The SHA-256 digest of `input`.
111pub(crate) fn hash_input(input: &[u8]) -> [u8; 32] {
112    use sha2::{Digest, Sha256};
113    Sha256::digest(input).into()
114}
115
116/// The lowercase hex SHA-256 digest of `parts` concatenated.
117pub(crate) fn hex_sha256(parts: &[&[u8]]) -> String {
118    use sha2::{Digest, Sha256};
119    use std::fmt::Write;
120    let mut hasher = Sha256::new();
121    for part in parts {
122        hasher.update(part);
123    }
124    let mut hex = String::with_capacity(64);
125    for byte in hasher.finalize() {
126        let _ = write!(&mut hex, "{byte:02x}");
127    }
128    hex
129}
130
131/// A run id: 1 to [`MAX_RUN_ID_LEN`] bytes of `[A-Za-z0-9_-]`. A run id
132/// is a path segment in the memo store and a key segment in the queue's
133/// KV namespace, so it is restricted to the characters Taquba accepts
134/// in a caller-supplied job id. `RunId` is the parameter type of every
135/// key and path builder, so a key over an unvalidated id does not
136/// compile. An id is validated by [`RunId::new`], by [`str::parse`] and
137/// by deserialization. A group id is a `RunId` as well, because it is
138/// stored at the same key positions. The type dereferences to `str`
139/// and implements `PartialEq<str>`.
140#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize)]
141#[serde(transparent)]
142pub struct RunId(String);
143
144impl RunId {
145    /// Validate `id`. An empty id, an id over [`MAX_RUN_ID_LEN`] bytes
146    /// or one with a character outside `[A-Za-z0-9_-]` is
147    /// [`Error::InvalidRunId`].
148    pub fn new(id: impl Into<String>) -> Result<Self> {
149        let id = id.into();
150        // An empty id is the memo prefix itself, and the sweep then
151        // clears every run's entries.
152        let reason = if id.is_empty() {
153            "run id must not be empty"
154        } else if id.len() > MAX_RUN_ID_LEN {
155            "run id exceeds maximum length of 128 bytes"
156        } else if !id
157            .bytes()
158            .all(|b| b.is_ascii_alphanumeric() || b == b'_' || b == b'-')
159        {
160            "run id must contain only `[A-Za-z0-9_-]`"
161        } else {
162            return Ok(Self(id));
163        };
164        Err(Error::InvalidRunId { run_id: id, reason })
165    }
166
167    /// A generated id: a ULID.
168    pub(crate) fn generate() -> Self {
169        Self(ulid::Ulid::new().to_string())
170    }
171
172    /// The id that is the lowercase hex SHA-256 digest of `parts`
173    /// concatenated, valid by construction.
174    pub(crate) fn digest(parts: &[&[u8]]) -> Self {
175        Self(hex_sha256(parts))
176    }
177
178    /// The id as a string slice.
179    pub fn as_str(&self) -> &str {
180        &self.0
181    }
182
183    /// The id as an owned string.
184    pub fn into_string(self) -> String {
185        self.0
186    }
187}
188
189impl Deref for RunId {
190    type Target = str;
191
192    fn deref(&self) -> &str {
193        &self.0
194    }
195}
196
197impl AsRef<str> for RunId {
198    fn as_ref(&self) -> &str {
199        &self.0
200    }
201}
202
203impl Borrow<str> for RunId {
204    fn borrow(&self) -> &str {
205        &self.0
206    }
207}
208
209impl fmt::Display for RunId {
210    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
211        f.write_str(&self.0)
212    }
213}
214
215impl FromStr for RunId {
216    type Err = Error;
217
218    fn from_str(id: &str) -> Result<Self> {
219        Self::new(id)
220    }
221}
222
223impl PartialEq<str> for RunId {
224    fn eq(&self, other: &str) -> bool {
225        self.0 == other
226    }
227}
228
229impl PartialEq<&str> for RunId {
230    fn eq(&self, other: &&str) -> bool {
231        self.0 == *other
232    }
233}
234
235impl PartialEq<String> for RunId {
236    fn eq(&self, other: &String) -> bool {
237        &self.0 == other
238    }
239}
240
241impl From<RunId> for String {
242    fn from(id: RunId) -> Self {
243        id.0
244    }
245}
246
247impl<'de> Deserialize<'de> for RunId {
248    fn deserialize<D: Deserializer<'de>>(deserializer: D) -> std::result::Result<Self, D::Error> {
249        let id = String::deserialize(deserializer)?;
250        Self::new(id).map_err(serde::de::Error::custom)
251    }
252}
253
254/// `{prefix}{suffix}`.
255fn prefixed(prefix: &[u8], suffix: &str) -> Vec<u8> {
256    let mut k = Vec::with_capacity(prefix.len() + suffix.len());
257    k.extend_from_slice(prefix);
258    k.extend_from_slice(suffix.as_bytes());
259    k
260}
261
262pub(crate) fn run_kv_key(run_id: &RunId) -> Vec<u8> {
263    prefixed(RUN_KV_PREFIX, run_id)
264}
265
266pub(crate) fn step_kv_key(run_id: &RunId) -> Vec<u8> {
267    prefixed(STEP_KV_PREFIX, run_id)
268}
269
270pub(crate) fn outcome_kv_key(run_id: &RunId) -> Vec<u8> {
271    prefixed(OUTCOME_KV_PREFIX, run_id)
272}
273
274pub(crate) fn signal_wait_kv_key(correlation_key: &str) -> Vec<u8> {
275    prefixed(SIGNAL_WAIT_KV_PREFIX, correlation_key)
276}
277
278pub(crate) fn signal_buf_kv_key(correlation_key: &str) -> Vec<u8> {
279    prefixed(SIGNAL_BUF_KV_PREFIX, correlation_key)
280}
281
282pub(crate) fn signal_delivered_kv_key(run_id: &RunId, step_number: u32) -> Vec<u8> {
283    prefixed(
284        SIGNAL_DELIVERED_KV_PREFIX,
285        &format!("{run_id}/{step_number}"),
286    )
287}
288
289#[cfg(test)]
290mod tests {
291    use super::*;
292
293    #[test]
294    fn internal_kv_prefixes_are_under_the_reserved_prefix() {
295        for prefix in [
296            RUN_KV_PREFIX,
297            STEP_KV_PREFIX,
298            SIGNAL_WAIT_KV_PREFIX,
299            SIGNAL_BUF_KV_PREFIX,
300            SIGNAL_DELIVERED_KV_PREFIX,
301            TERMINAL_KV_PREFIX,
302            OUTCOME_KV_PREFIX,
303            GROUP_KV_PREFIX,
304            GROUP_TERMINAL_KV_PREFIX,
305        ] {
306            assert!(
307                prefix.starts_with(RESERVED_KV_PREFIX.as_bytes()),
308                "internal kv prefix `{}` is outside the reserved prefix",
309                String::from_utf8_lossy(prefix),
310            );
311        }
312    }
313
314    #[test]
315    fn run_id_rejects_empty_long_and_unsafe_ids() {
316        for bad in [
317            "",
318            "run/1",
319            "run 1",
320            "run:1",
321            &"a".repeat(MAX_RUN_ID_LEN + 1),
322        ] {
323            assert!(
324                matches!(RunId::new(bad), Err(Error::InvalidRunId { .. })),
325                "`{bad}` must be rejected",
326            );
327        }
328        assert!(RunId::new("a".repeat(MAX_RUN_ID_LEN)).is_ok());
329        assert!(rmp_serde::from_slice::<RunId>(&rmp_serde::to_vec("").unwrap()).is_err());
330        assert_eq!(
331            rmp_serde::from_slice::<RunId>(&rmp_serde::to_vec("run-1").unwrap()).unwrap(),
332            "run-1"
333        );
334    }
335}