Skip to main content

swale/
records.rs

1//! The KV records of the orchestrator and their keys. Every key has the
2//! `swale/` prefix, and every value is JSON written as an absolute value.
3//!
4//! - `swale/runs/{graph}/{partition}`: the [`GraphRunRecord`].
5//! - `swale/assets/{asset}/{partition}`: the [`NodeRecord`] of an asset node.
6//! - `swale/tasks/{graph}/{partition}/{node}`: the [`NodeRecord`] of a task
7//!   node.
8
9use serde::{Deserialize, Serialize};
10use taquba_workflow::TerminalStatus;
11
12use crate::graph::Node;
13use crate::partition::Partition;
14
15/// The prefix of every key of this crate.
16pub const KV_PREFIX: &str = "swale/";
17
18/// The prefix of every graph run record.
19pub const GRAPH_RUNS_PREFIX: &str = "swale/runs/";
20
21/// The key of the graph run record.
22pub fn graph_run_key(graph: &str, partition: &Partition) -> Vec<u8> {
23    format!("{GRAPH_RUNS_PREFIX}{graph}/{partition}").into_bytes()
24}
25
26/// The graph and partition of a graph run key, or `None` for another key.
27pub 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
35/// The key of the asset record.
36pub fn asset_key(asset: &str, partition: &Partition) -> Vec<u8> {
37    format!("{KV_PREFIX}assets/{asset}/{partition}").into_bytes()
38}
39
40/// The key of the task record.
41pub fn task_key(graph: &str, partition: &Partition, node: &str) -> Vec<u8> {
42    format!("{KV_PREFIX}tasks/{graph}/{partition}/{node}").into_bytes()
43}
44
45/// The key of the record of `node` for `partition`: the asset key of an asset
46/// node or the task key of a task node.
47pub 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/// The terminal status of a task instance.
55#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
56#[serde(rename_all = "snake_case")]
57pub enum RecordStatus {
58    /// The run succeeded.
59    Succeeded,
60    /// The run failed: the operator reported a failure, or the step was
61    /// dead-lettered.
62    Failed,
63    /// The run was cancelled.
64    Cancelled,
65}
66
67impl RecordStatus {
68    /// The lowercase name, as in the JSON form.
69    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/// The record a terminated task instance leaves: the asset record of an asset
95/// node or the task record of a task node.
96#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
97pub struct NodeRecord {
98    /// The terminal status.
99    pub status: RecordStatus,
100    /// The run id of the task instance.
101    pub run_id: String,
102    /// The definition hash the graph run records.
103    pub definition: String,
104    /// The rerun count of the task instance.
105    pub rerun: u32,
106    /// The time of the termination, in milliseconds from the Unix epoch.
107    pub terminated_at_ms: u64,
108    /// The output of a succeeded run, when the record stays within the KV
109    /// value cap.
110    #[serde(default, skip_serializing_if = "Option::is_none")]
111    pub output: Option<serde_json::Value>,
112    /// Whether the output was left out because the record exceeded the cap.
113    #[serde(default, skip_serializing_if = "std::ops::Not::not")]
114    pub output_omitted: bool,
115    /// The error of a failed or cancelled run.
116    #[serde(default, skip_serializing_if = "Option::is_none")]
117    pub error: Option<String>,
118}
119
120/// The state of a graph run.
121#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
122#[serde(rename_all = "snake_case")]
123pub enum GraphRunState {
124    /// Nodes are running or ready to run.
125    Active,
126    /// The run was cancelled, and no further node is submitted.
127    Cancelled,
128    /// Every node has a succeeded record.
129    Complete,
130    /// No node is ready or running and at least one record is failed.
131    Failed,
132}
133
134impl GraphRunState {
135    /// The lowercase name, as in the JSON form.
136    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/// The record of a graph run for one partition.
153#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
154pub struct GraphRunRecord {
155    /// The definition hash the run records.
156    pub definition: String,
157    /// The time of the request, in milliseconds from the Unix epoch.
158    pub requested_at_ms: u64,
159    /// The state.
160    pub state: GraphRunState,
161}
162
163macro_rules! json_record {
164    ($t:ty) => {
165        impl $t {
166            /// The JSON form of the record.
167            pub fn to_bytes(&self) -> Vec<u8> {
168                serde_json::to_vec(self).expect("a record serializes to JSON")
169            }
170
171            /// Parses the JSON form of the record.
172            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}