1use 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
15pub const HEADER_RUN_ID: &str = "workflow.run_id";
17pub const HEADER_STEP: &str = "workflow.step";
19pub const RESERVED_HEADER_PREFIX: &str = "workflow.";
23
24pub const RESERVED_KV_PREFIX: &str = "workflow/";
29
30pub const HEADER_TERMINAL: &str = "workflow.terminal";
34
35pub const HEADER_SIGNAL_WAIT: &str = "workflow.signal_wait";
38pub const HEADER_SIGNAL_DELIVERED: &str = "workflow.signal_delivered";
42
43pub const HEADER_GROUP: &str = "workflow.group";
46pub const HEADER_GROUP_KEY: &str = "workflow.group_key";
49
50pub(crate) const DEDUP_PREFIX: &str = "run:";
51
52pub const MAX_RUN_ID_LEN: usize = 128;
55
56pub(crate) const RUN_KV_PREFIX: &[u8] = b"workflow/runs/";
58
59pub(crate) const STEP_KV_PREFIX: &[u8] = b"workflow/steps/";
64
65pub(crate) const SIGNAL_WAIT_KV_PREFIX: &[u8] = b"workflow/signal-wait/";
68
69pub(crate) const SIGNAL_BUF_KV_PREFIX: &[u8] = b"workflow/signal-buf/";
72
73pub(crate) const SIGNAL_DELIVERED_KV_PREFIX: &[u8] = b"workflow/signal-delivered/";
77
78pub(crate) const TERMINAL_KV_PREFIX: &[u8] = b"workflow/terminals/";
83
84pub(crate) const OUTCOME_KV_PREFIX: &[u8] = b"workflow/outcomes/";
89
90pub(crate) const GROUP_KV_PREFIX: &[u8] = b"workflow/groups/";
95
96pub(crate) fn group_members_kv_prefix(group_id: &RunId) -> Vec<u8> {
98 prefixed(GROUP_KV_PREFIX, &format!("{group_id}/"))
99}
100
101pub(crate) fn group_member_kv_key(group_id: &RunId, key: &str) -> Vec<u8> {
103 prefixed(&group_members_kv_prefix(group_id), key)
104}
105
106pub(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
115pub(crate) const GROUP_TERMINAL_KV_PREFIX: &[u8] = b"workflow/group-terminals/";
118
119pub(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
125pub(crate) fn timestamped_kv_key(prefix: &[u8], id: &RunId, ts_ms: u64) -> Vec<u8> {
128 prefixed(prefix, &format!("{ts_ms:020}/{id}"))
129}
130
131pub(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
141pub(crate) fn hash_input(input: &[u8]) -> [u8; 32] {
143 use sha2::{Digest, Sha256};
144 Sha256::digest(input).into()
145}
146
147pub(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#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize)]
172#[serde(transparent)]
173pub struct RunId(String);
174
175impl RunId {
176 pub fn new(id: impl Into<String>) -> Result<Self> {
180 let id = id.into();
181 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 pub(crate) fn generate() -> Self {
200 Self(ulid::Ulid::new().to_string())
201 }
202
203 pub(crate) fn digest(parts: &[&[u8]]) -> Self {
206 Self(hex_sha256(parts))
207 }
208
209 pub fn as_str(&self) -> &str {
211 &self.0
212 }
213
214 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
285fn 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}