use serde::{Deserialize, Serialize};
use taquba_workflow::TerminalStatus;
use crate::graph::Node;
use crate::partition::Partition;
pub const KV_PREFIX: &str = "swale/";
pub const GRAPH_RUNS_PREFIX: &str = "swale/runs/";
pub fn graph_run_key(graph: &str, partition: &Partition) -> Vec<u8> {
format!("{GRAPH_RUNS_PREFIX}{graph}/{partition}").into_bytes()
}
pub fn parse_graph_run_key(key: &[u8]) -> Option<(String, Partition)> {
let rest = std::str::from_utf8(key)
.ok()?
.strip_prefix(GRAPH_RUNS_PREFIX)?;
let (graph, partition) = rest.split_once('/')?;
Some((graph.to_string(), Partition::new(partition).ok()?))
}
pub fn asset_key(asset: &str, partition: &Partition) -> Vec<u8> {
format!("{KV_PREFIX}assets/{asset}/{partition}").into_bytes()
}
pub fn task_key(graph: &str, partition: &Partition, node: &str) -> Vec<u8> {
format!("{KV_PREFIX}tasks/{graph}/{partition}/{node}").into_bytes()
}
pub fn node_record_key(graph: &str, partition: &Partition, node: &Node) -> Vec<u8> {
match node.asset() {
Some(asset) => asset_key(asset, partition),
None => task_key(graph, partition, node.name()),
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RecordStatus {
Succeeded,
Failed,
Cancelled,
}
impl RecordStatus {
pub fn as_str(&self) -> &'static str {
match self {
RecordStatus::Succeeded => "succeeded",
RecordStatus::Failed => "failed",
RecordStatus::Cancelled => "cancelled",
}
}
}
impl std::fmt::Display for RecordStatus {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.as_str())
}
}
impl From<TerminalStatus> for RecordStatus {
fn from(status: TerminalStatus) -> Self {
match status {
TerminalStatus::Succeeded => RecordStatus::Succeeded,
TerminalStatus::Failed => RecordStatus::Failed,
TerminalStatus::Cancelled => RecordStatus::Cancelled,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct NodeRecord {
pub status: RecordStatus,
pub run_id: String,
pub definition: String,
pub rerun: u32,
pub terminated_at_ms: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output: Option<serde_json::Value>,
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub output_omitted: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum GraphRunState {
Active,
Cancelled,
Complete,
Failed,
}
impl GraphRunState {
pub fn as_str(&self) -> &'static str {
match self {
GraphRunState::Active => "active",
GraphRunState::Cancelled => "cancelled",
GraphRunState::Complete => "complete",
GraphRunState::Failed => "failed",
}
}
}
impl std::fmt::Display for GraphRunState {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct GraphRunRecord {
pub definition: String,
pub requested_at_ms: u64,
pub state: GraphRunState,
}
macro_rules! json_record {
($t:ty) => {
impl $t {
pub fn to_bytes(&self) -> Vec<u8> {
serde_json::to_vec(self).expect("a record serializes to JSON")
}
pub fn from_bytes(bytes: &[u8]) -> Result<Self, serde_json::Error> {
serde_json::from_slice(bytes)
}
}
};
}
json_record!(NodeRecord);
json_record!(GraphRunRecord);
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn keys_have_the_documented_layout() {
let partition = Partition::new("20260915").unwrap();
assert_eq!(
graph_run_key("orders_daily", &partition),
b"swale/runs/orders_daily/20260915"
);
assert_eq!(
asset_key("orders_raw", &partition),
b"swale/assets/orders_raw/20260915"
);
assert_eq!(
task_key("orders_daily", &partition, "notify"),
b"swale/tasks/orders_daily/20260915/notify"
);
assert_eq!(
parse_graph_run_key(b"swale/runs/orders_daily/20260915"),
Some(("orders_daily".to_string(), partition))
);
assert_eq!(
parse_graph_run_key(b"swale/assets/orders_raw/20260915"),
None
);
assert_eq!(
parse_graph_run_key(b"swale/runs/orders_daily/2026-09"),
None
);
}
#[test]
fn records_round_trip_through_json_with_optional_fields_left_out() {
let record = NodeRecord {
status: RecordStatus::Succeeded,
run_id: "g-p-n-r0".into(),
definition: "abc".into(),
rerun: 0,
terminated_at_ms: 5,
output: None,
output_omitted: false,
error: None,
};
let json = String::from_utf8(record.to_bytes()).unwrap();
assert!(!json.contains("output"), "{json}");
assert!(!json.contains("error"), "{json}");
assert_eq!(NodeRecord::from_bytes(json.as_bytes()).unwrap(), record);
let run = GraphRunRecord {
definition: "abc".into(),
requested_at_ms: 7,
state: GraphRunState::Active,
};
assert_eq!(GraphRunRecord::from_bytes(&run.to_bytes()).unwrap(), run);
assert_eq!(
RecordStatus::from(TerminalStatus::Cancelled),
RecordStatus::Cancelled
);
assert_eq!(RecordStatus::Failed.to_string(), "failed");
assert_eq!(GraphRunState::Complete.to_string(), "complete");
}
}