use crate::error::{Error, Result};
pub const HEADER_RUN_ID: &str = "workflow.run_id";
pub const HEADER_STEP: &str = "workflow.step";
pub const RESERVED_HEADER_PREFIX: &str = "workflow.";
pub const RESERVED_KV_PREFIX: &str = "workflow/";
pub const HEADER_TERMINAL: &str = "workflow.terminal";
pub const HEADER_SIGNAL_WAIT: &str = "workflow.signal_wait";
pub const HEADER_SIGNAL_DELIVERED: &str = "workflow.signal_delivered";
pub(crate) const DEDUP_PREFIX: &str = "run:";
pub const MAX_RUN_ID_LEN: usize = 128;
pub(crate) const RUN_KV_PREFIX: &[u8] = b"workflow/runs/";
pub(crate) const STEP_KV_PREFIX: &[u8] = b"workflow/steps/";
pub(crate) const SIGNAL_WAIT_KV_PREFIX: &[u8] = b"workflow/signal-wait/";
pub(crate) const SIGNAL_BUF_KV_PREFIX: &[u8] = b"workflow/signal-buf/";
pub(crate) const SIGNAL_DELIVERED_KV_PREFIX: &[u8] = b"workflow/signal-delivered/";
pub(crate) const TERMINAL_KV_PREFIX: &[u8] = b"workflow/terminals/";
pub(crate) const BULK_KV_PREFIX: &[u8] = b"workflow/bulk/batches/";
pub(crate) fn bulk_items_kv_prefix(batch_id: &str) -> Vec<u8> {
let mut k = Vec::from(BULK_KV_PREFIX);
k.extend_from_slice(batch_id.as_bytes());
k.extend_from_slice(b"/items/");
k
}
pub(crate) fn bulk_item_kv_key(batch_id: &str, key: &str) -> Vec<u8> {
let mut k = bulk_items_kv_prefix(batch_id);
k.extend_from_slice(key.as_bytes());
k
}
pub(crate) fn terminal_kv_key(run_id: &str, terminal_at_ms: u64) -> Vec<u8> {
timestamped_kv_key(TERMINAL_KV_PREFIX, run_id, terminal_at_ms)
}
pub(crate) const BULK_TERMINAL_KV_PREFIX: &[u8] = b"workflow/bulk/terminals/";
pub(crate) fn bulk_terminal_kv_key(batch_id: &str, terminal_at_ms: u64) -> Vec<u8> {
timestamped_kv_key(BULK_TERMINAL_KV_PREFIX, batch_id, terminal_at_ms)
}
pub(crate) fn timestamped_kv_key(prefix: &[u8], id: &str, ts_ms: u64) -> Vec<u8> {
let mut k = Vec::from(prefix);
k.extend_from_slice(format!("{ts_ms:020}/").as_bytes());
k.extend_from_slice(id.as_bytes());
k
}
pub(crate) fn parse_timestamped_kv_key(prefix: &[u8], key: &[u8]) -> Option<(String, u64)> {
let suffix = key.strip_prefix(prefix)?;
let text = std::str::from_utf8(suffix).ok()?;
let (ts, id) = text.split_once('/')?;
Some((id.to_string(), ts.parse().ok()?))
}
pub(crate) fn hash_input(input: &[u8]) -> [u8; 32] {
use sha2::{Digest, Sha256};
Sha256::digest(input).into()
}
pub(crate) fn hex_sha256(parts: &[&[u8]]) -> String {
use sha2::{Digest, Sha256};
use std::fmt::Write;
let mut hasher = Sha256::new();
for part in parts {
hasher.update(part);
}
let mut hex = String::with_capacity(64);
for byte in hasher.finalize() {
let _ = write!(&mut hex, "{byte:02x}");
}
hex
}
pub(crate) fn validate_run_id(run_id: &str) -> Result<()> {
let reason = if run_id.is_empty() {
"run id must not be empty"
} else if run_id.len() > MAX_RUN_ID_LEN {
"run id exceeds maximum length of 128 bytes"
} else if !run_id
.bytes()
.all(|b| b.is_ascii_alphanumeric() || b == b'_' || b == b'-')
{
"run id must contain only `[A-Za-z0-9_-]`"
} else {
return Ok(());
};
Err(Error::InvalidRunId {
run_id: run_id.to_string(),
reason,
})
}
pub(crate) fn run_kv_key(run_id: &str) -> Vec<u8> {
let mut k = Vec::with_capacity(RUN_KV_PREFIX.len() + run_id.len());
k.extend_from_slice(RUN_KV_PREFIX);
k.extend_from_slice(run_id.as_bytes());
k
}
pub(crate) fn step_kv_key(run_id: &str) -> Vec<u8> {
let mut k = Vec::with_capacity(STEP_KV_PREFIX.len() + run_id.len());
k.extend_from_slice(STEP_KV_PREFIX);
k.extend_from_slice(run_id.as_bytes());
k
}
pub(crate) fn signal_wait_kv_key(correlation_key: &str) -> Vec<u8> {
let mut k = Vec::with_capacity(SIGNAL_WAIT_KV_PREFIX.len() + correlation_key.len());
k.extend_from_slice(SIGNAL_WAIT_KV_PREFIX);
k.extend_from_slice(correlation_key.as_bytes());
k
}
pub(crate) fn signal_buf_kv_key(correlation_key: &str) -> Vec<u8> {
let mut k = Vec::with_capacity(SIGNAL_BUF_KV_PREFIX.len() + correlation_key.len());
k.extend_from_slice(SIGNAL_BUF_KV_PREFIX);
k.extend_from_slice(correlation_key.as_bytes());
k
}
pub(crate) fn signal_delivered_kv_key(run_id: &str, step_number: u32) -> Vec<u8> {
let mut k = Vec::from(SIGNAL_DELIVERED_KV_PREFIX);
k.extend_from_slice(run_id.as_bytes());
k.push(b'/');
k.extend_from_slice(step_number.to_string().as_bytes());
k
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn internal_kv_prefixes_are_under_the_reserved_prefix() {
for prefix in [
RUN_KV_PREFIX,
SIGNAL_WAIT_KV_PREFIX,
SIGNAL_BUF_KV_PREFIX,
SIGNAL_DELIVERED_KV_PREFIX,
TERMINAL_KV_PREFIX,
BULK_KV_PREFIX,
BULK_TERMINAL_KV_PREFIX,
] {
assert!(
prefix.starts_with(RESERVED_KV_PREFIX.as_bytes()),
"internal kv prefix `{}` is outside the reserved prefix",
String::from_utf8_lossy(prefix),
);
}
}
}