use std::collections::HashMap;
use taquba_workflow::RunId;
use crate::partition::{InvalidPartition, Partition};
use crate::records;
pub const HEADER_GRAPH: &str = "swale.graph";
pub const HEADER_PARTITION: &str = "swale.partition";
pub const HEADER_NODE: &str = "swale.node";
pub const HEADER_ASSET: &str = "swale.asset";
pub const HEADER_DEFINITION: &str = "swale.definition";
pub const HEADER_RERUN: &str = "swale.rerun";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TaskIdentity {
pub graph: String,
pub partition: Partition,
pub node: String,
pub asset: Option<String>,
pub definition: String,
pub rerun: u32,
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum IdentityError {
#[error("header `{0}` is absent")]
MissingHeader(&'static str),
#[error("header `{HEADER_PARTITION}`: {0}")]
Partition(#[from] InvalidPartition),
#[error("header `{HEADER_RERUN}` is not a count: `{0}`")]
Rerun(String),
}
pub fn run_id(graph: &str, partition: &Partition, node: &str, rerun: u32) -> RunId {
RunId::new(format!("{graph}-{partition}-{node}-r{rerun}"))
.expect("graph and node names are bounded so that every run id is valid")
}
impl TaskIdentity {
pub fn run_id(&self) -> RunId {
run_id(&self.graph, &self.partition, &self.node, self.rerun)
}
pub fn headers(&self) -> HashMap<String, String> {
let mut headers = HashMap::from([
(HEADER_GRAPH.to_string(), self.graph.clone()),
(HEADER_PARTITION.to_string(), self.partition.to_string()),
(HEADER_NODE.to_string(), self.node.clone()),
(HEADER_DEFINITION.to_string(), self.definition.clone()),
(HEADER_RERUN.to_string(), self.rerun.to_string()),
]);
if let Some(asset) = &self.asset {
headers.insert(HEADER_ASSET.to_string(), asset.clone());
}
headers
}
pub fn from_headers(headers: &HashMap<String, String>) -> Result<Self, IdentityError> {
let required = |name: &'static str| {
headers
.get(name)
.cloned()
.ok_or(IdentityError::MissingHeader(name))
};
let rerun = required(HEADER_RERUN)?;
Ok(TaskIdentity {
graph: required(HEADER_GRAPH)?,
partition: Partition::new(required(HEADER_PARTITION)?)?,
node: required(HEADER_NODE)?,
asset: headers.get(HEADER_ASSET).cloned(),
definition: required(HEADER_DEFINITION)?,
rerun: rerun.parse().map_err(|_| IdentityError::Rerun(rerun))?,
})
}
pub fn record_key(&self) -> Vec<u8> {
records::record_key(
&self.graph,
&self.partition,
&self.node,
self.asset.as_deref(),
)
}
}
#[cfg(test)]
mod tests {
use super::*;
fn identity(asset: Option<&str>) -> TaskIdentity {
TaskIdentity {
graph: "orders_daily".into(),
partition: Partition::new("20260915").unwrap(),
node: "extract".into(),
asset: asset.map(str::to_string),
definition: "abc".into(),
rerun: 2,
}
}
#[test]
fn run_id_has_the_documented_form() {
assert_eq!(
identity(Some("orders_raw")).run_id().as_str(),
"orders_daily-20260915-extract-r2"
);
}
#[test]
fn headers_round_trip_with_and_without_an_asset() {
for asset in [Some("orders_raw"), None] {
let identity = identity(asset);
let headers = identity.headers();
assert_eq!(headers.contains_key(HEADER_ASSET), asset.is_some());
assert!(headers.keys().all(|k| k.starts_with("swale.")));
assert_eq!(TaskIdentity::from_headers(&headers).unwrap(), identity);
}
}
#[test]
fn record_key_is_the_asset_key_or_the_task_key() {
assert_eq!(
identity(Some("orders_raw")).record_key(),
b"swale/assets/orders_raw/20260915"
);
assert_eq!(
identity(None).record_key(),
b"swale/tasks/orders_daily/20260915/extract"
);
}
#[test]
fn missing_and_malformed_headers_are_reported() {
let mut headers = identity(None).headers();
headers.insert(HEADER_RERUN.into(), "x".into());
assert_eq!(
TaskIdentity::from_headers(&headers),
Err(IdentityError::Rerun("x".into()))
);
headers.insert(HEADER_RERUN.into(), "0".into());
headers.insert(HEADER_PARTITION.into(), "2026-09-15".into());
assert!(matches!(
TaskIdentity::from_headers(&headers),
Err(IdentityError::Partition(_))
));
headers.remove(HEADER_GRAPH);
headers.insert(HEADER_PARTITION.into(), "20260915".into());
assert_eq!(
TaskIdentity::from_headers(&headers),
Err(IdentityError::MissingHeader(HEADER_GRAPH))
);
}
}