1use serde::{Deserialize, Serialize};
10use taquba_workflow::TerminalStatus;
11
12use crate::graph::Node;
13use crate::partition::Partition;
14
15pub const KV_PREFIX: &str = "swale/";
17
18pub const GRAPH_RUNS_PREFIX: &str = "swale/runs/";
20
21pub fn graph_run_key(graph: &str, partition: &Partition) -> Vec<u8> {
23 format!("{GRAPH_RUNS_PREFIX}{graph}/{partition}").into_bytes()
24}
25
26pub fn parse_graph_run_key(key: &[u8]) -> Option<(String, Partition)> {
28 let rest = std::str::from_utf8(key)
29 .ok()?
30 .strip_prefix(GRAPH_RUNS_PREFIX)?;
31 let (graph, partition) = rest.split_once('/')?;
32 Some((graph.to_string(), Partition::new(partition).ok()?))
33}
34
35pub fn asset_key(asset: &str, partition: &Partition) -> Vec<u8> {
37 format!("{KV_PREFIX}assets/{asset}/{partition}").into_bytes()
38}
39
40pub fn task_key(graph: &str, partition: &Partition, node: &str) -> Vec<u8> {
42 format!("{KV_PREFIX}tasks/{graph}/{partition}/{node}").into_bytes()
43}
44
45pub fn node_record_key(graph: &str, partition: &Partition, node: &Node) -> Vec<u8> {
48 match node.asset() {
49 Some(asset) => asset_key(asset, partition),
50 None => task_key(graph, partition, node.name()),
51 }
52}
53
54#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
56#[serde(rename_all = "snake_case")]
57pub enum RecordStatus {
58 Succeeded,
60 Failed,
63 Cancelled,
65}
66
67impl RecordStatus {
68 pub fn as_str(&self) -> &'static str {
70 match self {
71 RecordStatus::Succeeded => "succeeded",
72 RecordStatus::Failed => "failed",
73 RecordStatus::Cancelled => "cancelled",
74 }
75 }
76}
77
78impl std::fmt::Display for RecordStatus {
79 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
80 f.write_str(self.as_str())
81 }
82}
83
84impl From<TerminalStatus> for RecordStatus {
85 fn from(status: TerminalStatus) -> Self {
86 match status {
87 TerminalStatus::Succeeded => RecordStatus::Succeeded,
88 TerminalStatus::Failed => RecordStatus::Failed,
89 TerminalStatus::Cancelled => RecordStatus::Cancelled,
90 }
91 }
92}
93
94#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
97pub struct NodeRecord {
98 pub status: RecordStatus,
100 pub run_id: String,
102 pub definition: String,
104 pub rerun: u32,
106 pub terminated_at_ms: u64,
108 #[serde(default, skip_serializing_if = "Option::is_none")]
111 pub output: Option<serde_json::Value>,
112 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
114 pub output_omitted: bool,
115 #[serde(default, skip_serializing_if = "Option::is_none")]
117 pub error: Option<String>,
118}
119
120#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
122#[serde(rename_all = "snake_case")]
123pub enum GraphRunState {
124 Active,
126 Cancelled,
128 Complete,
130 Failed,
132}
133
134impl GraphRunState {
135 pub fn as_str(&self) -> &'static str {
137 match self {
138 GraphRunState::Active => "active",
139 GraphRunState::Cancelled => "cancelled",
140 GraphRunState::Complete => "complete",
141 GraphRunState::Failed => "failed",
142 }
143 }
144}
145
146impl std::fmt::Display for GraphRunState {
147 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
148 f.write_str(self.as_str())
149 }
150}
151
152#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
154pub struct GraphRunRecord {
155 pub definition: String,
157 pub requested_at_ms: u64,
159 pub state: GraphRunState,
161}
162
163macro_rules! json_record {
164 ($t:ty) => {
165 impl $t {
166 pub fn to_bytes(&self) -> Vec<u8> {
168 serde_json::to_vec(self).expect("a record serializes to JSON")
169 }
170
171 pub fn from_bytes(bytes: &[u8]) -> Result<Self, serde_json::Error> {
173 serde_json::from_slice(bytes)
174 }
175 }
176 };
177}
178
179json_record!(NodeRecord);
180json_record!(GraphRunRecord);
181
182#[cfg(test)]
183mod tests {
184 use super::*;
185
186 #[test]
187 fn keys_have_the_documented_layout() {
188 let partition = Partition::new("20260915").unwrap();
189 assert_eq!(
190 graph_run_key("orders_daily", &partition),
191 b"swale/runs/orders_daily/20260915"
192 );
193 assert_eq!(
194 asset_key("orders_raw", &partition),
195 b"swale/assets/orders_raw/20260915"
196 );
197 assert_eq!(
198 task_key("orders_daily", &partition, "notify"),
199 b"swale/tasks/orders_daily/20260915/notify"
200 );
201 assert_eq!(
202 parse_graph_run_key(b"swale/runs/orders_daily/20260915"),
203 Some(("orders_daily".to_string(), partition))
204 );
205 assert_eq!(
206 parse_graph_run_key(b"swale/assets/orders_raw/20260915"),
207 None
208 );
209 assert_eq!(
210 parse_graph_run_key(b"swale/runs/orders_daily/2026-09"),
211 None
212 );
213 }
214
215 #[test]
216 fn records_round_trip_through_json_with_optional_fields_left_out() {
217 let record = NodeRecord {
218 status: RecordStatus::Succeeded,
219 run_id: "g-p-n-r0".into(),
220 definition: "abc".into(),
221 rerun: 0,
222 terminated_at_ms: 5,
223 output: None,
224 output_omitted: false,
225 error: None,
226 };
227 let json = String::from_utf8(record.to_bytes()).unwrap();
228 assert!(!json.contains("output"), "{json}");
229 assert!(!json.contains("error"), "{json}");
230 assert_eq!(NodeRecord::from_bytes(json.as_bytes()).unwrap(), record);
231
232 let run = GraphRunRecord {
233 definition: "abc".into(),
234 requested_at_ms: 7,
235 state: GraphRunState::Active,
236 };
237 assert_eq!(GraphRunRecord::from_bytes(&run.to_bytes()).unwrap(), run);
238 assert_eq!(
239 RecordStatus::from(TerminalStatus::Cancelled),
240 RecordStatus::Cancelled
241 );
242 assert_eq!(RecordStatus::Failed.to_string(), "failed");
243 assert_eq!(GraphRunState::Complete.to_string(), "complete");
244 }
245}