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
25/// passed via [`RunSpec::kv_writes`](crate::RunSpec::kv_writes) or staged through an
26/// [`crate::EffectsHandle`] must not start with this prefix; they are
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 for the durable terminal marker: the time-ordered index of
79/// runs that have reached a terminal state, read by the memo-retention
80/// sweep. Written only when [`WorkflowRuntimeBuilder::memo_retention`]
81/// is set, in the same transaction that settles the run.
82pub(crate) const TERMINAL_KV_PREFIX: &[u8] = b"workflow/terminals/";
83
84/// Prefix for the durable terminal record of a run:
85/// `workflow/outcomes/{run_id}`, written in the settlement that
86/// terminates the run and removed with the run's memo entries by the
87/// memo sweep.
88pub(crate) const OUTCOME_KV_PREFIX: &[u8] = b"workflow/outcomes/";
89
90/// Prefix for the durable member records of run groups:
91/// `workflow/groups/{group_id}/{key}`, one per member, written with the
92/// member's submission and rewritten in the settlement that terminates
93/// it.
94pub(crate) const GROUP_KV_PREFIX: &[u8] = b"workflow/groups/";
95
96/// Prefix under which the member records of one group are stored.
97pub(crate) fn group_members_kv_prefix(group_id: &RunId) -> Vec<u8> {
98    prefixed(GROUP_KV_PREFIX, &format!("{group_id}/"))
99}
100
101/// Key of the member record of `key` in group `group_id`.
102pub(crate) fn group_member_kv_key(group_id: &RunId, key: &str) -> Vec<u8> {
103    prefixed(&group_members_kv_prefix(group_id), key)
104}
105
106/// Key of the terminal marker for `run_id`, terminated at
107/// `terminal_at_ms`. The zero-padded timestamp leads the suffix, so a
108/// prefix scan returns markers oldest first and the sweep's expired set
109/// is the front of the range. The value is empty: both fields are in
110/// the key.
111pub(crate) fn terminal_kv_key(run_id: &RunId, terminal_at_ms: u64) -> Vec<u8> {
112    timestamped_kv_key(TERMINAL_KV_PREFIX, run_id, terminal_at_ms)
113}
114
115/// Prefix for the durable terminal markers of run groups, read by the
116/// group retention sweep: `workflow/group-terminals/{ts:020}/{group_id}`.
117pub(crate) const GROUP_TERMINAL_KV_PREFIX: &[u8] = b"workflow/group-terminals/";
118
119/// Key of the terminal marker of group `group_id`, whose members all
120/// terminated by `terminal_at_ms`.
121pub(crate) fn group_terminal_kv_key(group_id: &RunId, terminal_at_ms: u64) -> Vec<u8> {
122    timestamped_kv_key(GROUP_TERMINAL_KV_PREFIX, group_id, terminal_at_ms)
123}
124
125/// `{prefix}{ts:020}/{id}`: a marker whose zero-padded timestamp leads the
126/// suffix, so a prefix scan returns markers oldest first.
127pub(crate) fn timestamped_kv_key(prefix: &[u8], id: &RunId, ts_ms: u64) -> Vec<u8> {
128    prefixed(prefix, &format!("{ts_ms:020}/{id}"))
129}
130
131/// The `(id, ts_ms)` of a key built by [`timestamped_kv_key`], `None`
132/// for a key outside `prefix`, with a malformed timestamp or with an
133/// id that is not a valid run id.
134pub(crate) fn parse_timestamped_kv_key(prefix: &[u8], key: &[u8]) -> Option<(RunId, u64)> {
135    let suffix = key.strip_prefix(prefix)?;
136    let text = std::str::from_utf8(suffix).ok()?;
137    let (ts, id) = text.split_once('/')?;
138    Some((RunId::new(id).ok()?, ts.parse().ok()?))
139}
140
141/// The SHA-256 digest of `input`.
142pub(crate) fn hash_input(input: &[u8]) -> [u8; 32] {
143    use sha2::{Digest, Sha256};
144    Sha256::digest(input).into()
145}
146
147/// The lowercase hex SHA-256 digest of `parts` concatenated.
148pub(crate) fn hex_sha256(parts: &[&[u8]]) -> String {
149    use sha2::{Digest, Sha256};
150    use std::fmt::Write;
151    let mut hasher = Sha256::new();
152    for part in parts {
153        hasher.update(part);
154    }
155    let mut hex = String::with_capacity(64);
156    for byte in hasher.finalize() {
157        let _ = write!(&mut hex, "{byte:02x}");
158    }
159    hex
160}
161
162/// A run id: 1 to [`MAX_RUN_ID_LEN`] bytes of `[A-Za-z0-9_-]`. A run id
163/// is a path segment in the memo store and a key segment in the queue's
164/// KV namespace, so it is restricted to the characters Taquba accepts
165/// in a caller-supplied job id. `RunId` is the parameter type of every
166/// key and path builder, so a key over an unvalidated id does not
167/// compile. An id is validated by [`RunId::new`], by [`str::parse`] and
168/// by deserialization. A group id is a `RunId` as well, because it is
169/// stored at the same key positions. The type dereferences to `str`
170/// and implements `PartialEq<str>`.
171#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize)]
172#[serde(transparent)]
173pub struct RunId(String);
174
175impl RunId {
176    /// Validate `id`. An empty id, an id over [`MAX_RUN_ID_LEN`] bytes
177    /// or one with a character outside `[A-Za-z0-9_-]` is
178    /// [`Error::InvalidRunId`].
179    pub fn new(id: impl Into<String>) -> Result<Self> {
180        let id = id.into();
181        // An empty id is the memo prefix itself, and the sweep then
182        // clears every run's entries.
183        let reason = if id.is_empty() {
184            "run id must not be empty"
185        } else if id.len() > MAX_RUN_ID_LEN {
186            "run id exceeds maximum length of 128 bytes"
187        } else if !id
188            .bytes()
189            .all(|b| b.is_ascii_alphanumeric() || b == b'_' || b == b'-')
190        {
191            "run id must contain only `[A-Za-z0-9_-]`"
192        } else {
193            return Ok(Self(id));
194        };
195        Err(Error::InvalidRunId { run_id: id, reason })
196    }
197
198    /// A generated id: a ULID.
199    pub(crate) fn generate() -> Self {
200        Self(ulid::Ulid::new().to_string())
201    }
202
203    /// The id that is the lowercase hex SHA-256 digest of `parts`
204    /// concatenated, valid by construction.
205    pub(crate) fn digest(parts: &[&[u8]]) -> Self {
206        Self(hex_sha256(parts))
207    }
208
209    /// The id as a string slice.
210    pub fn as_str(&self) -> &str {
211        &self.0
212    }
213
214    /// The id as an owned string.
215    pub fn into_string(self) -> String {
216        self.0
217    }
218}
219
220impl Deref for RunId {
221    type Target = str;
222
223    fn deref(&self) -> &str {
224        &self.0
225    }
226}
227
228impl AsRef<str> for RunId {
229    fn as_ref(&self) -> &str {
230        &self.0
231    }
232}
233
234impl Borrow<str> for RunId {
235    fn borrow(&self) -> &str {
236        &self.0
237    }
238}
239
240impl fmt::Display for RunId {
241    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
242        f.write_str(&self.0)
243    }
244}
245
246impl FromStr for RunId {
247    type Err = Error;
248
249    fn from_str(id: &str) -> Result<Self> {
250        Self::new(id)
251    }
252}
253
254impl PartialEq<str> for RunId {
255    fn eq(&self, other: &str) -> bool {
256        self.0 == other
257    }
258}
259
260impl PartialEq<&str> for RunId {
261    fn eq(&self, other: &&str) -> bool {
262        self.0 == *other
263    }
264}
265
266impl PartialEq<String> for RunId {
267    fn eq(&self, other: &String) -> bool {
268        &self.0 == other
269    }
270}
271
272impl From<RunId> for String {
273    fn from(id: RunId) -> Self {
274        id.0
275    }
276}
277
278impl<'de> Deserialize<'de> for RunId {
279    fn deserialize<D: Deserializer<'de>>(deserializer: D) -> std::result::Result<Self, D::Error> {
280        let id = String::deserialize(deserializer)?;
281        Self::new(id).map_err(serde::de::Error::custom)
282    }
283}
284
285/// `{prefix}{suffix}`.
286fn prefixed(prefix: &[u8], suffix: &str) -> Vec<u8> {
287    let mut k = Vec::with_capacity(prefix.len() + suffix.len());
288    k.extend_from_slice(prefix);
289    k.extend_from_slice(suffix.as_bytes());
290    k
291}
292
293pub(crate) fn run_kv_key(run_id: &RunId) -> Vec<u8> {
294    prefixed(RUN_KV_PREFIX, run_id)
295}
296
297pub(crate) fn step_kv_key(run_id: &RunId) -> Vec<u8> {
298    prefixed(STEP_KV_PREFIX, run_id)
299}
300
301pub(crate) fn outcome_kv_key(run_id: &RunId) -> Vec<u8> {
302    prefixed(OUTCOME_KV_PREFIX, run_id)
303}
304
305pub(crate) fn signal_wait_kv_key(correlation_key: &str) -> Vec<u8> {
306    prefixed(SIGNAL_WAIT_KV_PREFIX, correlation_key)
307}
308
309pub(crate) fn signal_buf_kv_key(correlation_key: &str) -> Vec<u8> {
310    prefixed(SIGNAL_BUF_KV_PREFIX, correlation_key)
311}
312
313pub(crate) fn signal_delivered_kv_key(run_id: &RunId, step_number: u32) -> Vec<u8> {
314    prefixed(
315        SIGNAL_DELIVERED_KV_PREFIX,
316        &format!("{run_id}/{step_number}"),
317    )
318}
319
320#[cfg(test)]
321mod tests {
322    use super::*;
323
324    #[test]
325    fn internal_kv_prefixes_are_under_the_reserved_prefix() {
326        for prefix in [
327            RUN_KV_PREFIX,
328            STEP_KV_PREFIX,
329            SIGNAL_WAIT_KV_PREFIX,
330            SIGNAL_BUF_KV_PREFIX,
331            SIGNAL_DELIVERED_KV_PREFIX,
332            TERMINAL_KV_PREFIX,
333            OUTCOME_KV_PREFIX,
334            GROUP_KV_PREFIX,
335            GROUP_TERMINAL_KV_PREFIX,
336        ] {
337            assert!(
338                prefix.starts_with(RESERVED_KV_PREFIX.as_bytes()),
339                "internal kv prefix `{}` is outside the reserved prefix",
340                String::from_utf8_lossy(prefix),
341            );
342        }
343    }
344
345    #[test]
346    fn run_id_rejects_empty_long_and_unsafe_ids() {
347        for bad in [
348            "",
349            "run/1",
350            "run 1",
351            "run:1",
352            &"a".repeat(MAX_RUN_ID_LEN + 1),
353        ] {
354            assert!(
355                matches!(RunId::new(bad), Err(Error::InvalidRunId { .. })),
356                "`{bad}` must be rejected",
357            );
358        }
359        assert!(RunId::new("a".repeat(MAX_RUN_ID_LEN)).is_ok());
360        assert!(rmp_serde::from_slice::<RunId>(&rmp_serde::to_vec("").unwrap()).is_err());
361        assert_eq!(
362            rmp_serde::from_slice::<RunId>(&rmp_serde::to_vec("run-1").unwrap()).unwrap(),
363            "run-1"
364        );
365    }
366
367    #[test]
368    fn terminal_marker_keys_sort_oldest_first_and_round_trip() {
369        let old = terminal_kv_key(&RunId::new("run-b").unwrap(), 1_000);
370        let young = terminal_kv_key(&RunId::new("run-a").unwrap(), 2_000);
371        assert!(
372            old < young,
373            "ordering must follow the timestamp ahead of the id"
374        );
375        assert_eq!(
376            parse_timestamped_kv_key(TERMINAL_KV_PREFIX, &young),
377            Some((RunId::new("run-a").unwrap(), 2_000)),
378        );
379        assert_eq!(
380            parse_timestamped_kv_key(
381                TERMINAL_KV_PREFIX,
382                b"workflow/terminals/00000000000000002000/"
383            ),
384            None,
385            "a marker with an empty id is malformed"
386        );
387        assert_eq!(
388            parse_timestamped_kv_key(TERMINAL_KV_PREFIX, b"workflow/runs/run-a"),
389            None
390        );
391    }
392}