use std::collections::{BTreeMap, BTreeSet};
use std::time::Duration;
use crate::operator::{OperatorError, OperatorSet};
use crate::template::Template;
pub const MAX_GRAPH_AND_NODE_NAME_LEN: usize = 98;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum Partitioning {
Daily,
Hourly,
#[default]
Unpartitioned,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum TriggerRule {
#[default]
AllSucceeded,
AllDone,
OneFailed,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum NodeKind {
Asset {
produces: String,
consumes: Vec<String>,
},
Task {
after: Vec<String>,
trigger_rule: TriggerRule,
},
}
#[derive(Debug, Clone, PartialEq)]
pub struct NodeSpec {
pub name: String,
pub kind: NodeKind,
pub operator: String,
pub pool: String,
pub retries: u32,
pub params: toml::Table,
}
#[derive(Debug, Clone, PartialEq)]
pub struct GraphSpec {
pub name: String,
pub schedule: Option<String>,
pub catchup: Option<Duration>,
pub partitioning: Partitioning,
pub nodes: Vec<NodeSpec>,
}
#[derive(Debug, Clone, PartialEq)]
pub struct Node {
spec: NodeSpec,
upstreams: Vec<String>,
downstreams: Vec<String>,
}
impl Node {
pub fn name(&self) -> &str {
&self.spec.name
}
pub fn kind(&self) -> &NodeKind {
&self.spec.kind
}
pub fn asset(&self) -> Option<&str> {
match &self.spec.kind {
NodeKind::Asset { produces, .. } => Some(produces),
NodeKind::Task { .. } => None,
}
}
pub fn trigger_rule(&self) -> TriggerRule {
match &self.spec.kind {
NodeKind::Asset { .. } => TriggerRule::AllSucceeded,
NodeKind::Task { trigger_rule, .. } => *trigger_rule,
}
}
pub fn operator(&self) -> &str {
&self.spec.operator
}
pub fn pool(&self) -> &str {
&self.spec.pool
}
pub fn retries(&self) -> u32 {
self.spec.retries
}
pub fn params(&self) -> &toml::Table {
&self.spec.params
}
pub fn upstreams(&self) -> &[String] {
&self.upstreams
}
pub fn downstreams(&self) -> &[String] {
&self.downstreams
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct Graph {
name: String,
schedule: Option<String>,
catchup: Option<Duration>,
partitioning: Partitioning,
nodes: Vec<Node>,
index: BTreeMap<String, usize>,
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum Problem {
#[error("graph name `{0}` is not `[a-z0-9_]+`")]
InvalidGraphName(String),
#[error("schedule `{expression}` is not a cron expression: {message}")]
InvalidSchedule {
expression: String,
message: String,
},
#[error("catchup: {0}")]
InvalidCatchup(String),
#[error("the graph is empty")]
NoNodes,
#[error(
"node `{0}`: an asset node declares its upstreams with `consumes` and does not declare `after`"
)]
AssetNodeAfter(String),
#[error(
"node `{0}`: an asset node does not declare `trigger_rule`, because its upstreams must all succeed"
)]
AssetNodeTriggerRule(String),
#[error(
"node `{0}`: a task node declares its upstreams with `after` and does not declare `consumes`"
)]
TaskNodeConsumes(String),
#[error("node `{0}` is defined more than once")]
DuplicateNode(String),
#[error("node `{0}`: the name is not `[a-z0-9_]+`")]
InvalidNodeName(String),
#[error(
"node `{0}`: the graph and node names exceed {MAX_GRAPH_AND_NODE_NAME_LEN} bytes together"
)]
NamesTooLong(String),
#[error("node `{node}`: asset name `{asset}` is not `[a-z0-9_]+`")]
InvalidAssetName {
node: String,
asset: String,
},
#[error("asset `{asset}` is produced by more than one node: {}", .nodes.join(", "))]
DuplicateAsset {
asset: String,
nodes: Vec<String>,
},
#[error("node `{node}`: consumes `{asset}`, which no node produces")]
UnknownAsset {
node: String,
asset: String,
},
#[error("node `{node}`: `after` refers to unknown node `{after}`")]
UnknownNode {
node: String,
after: String,
},
#[error("node `{node}`: unknown operator `{operator}`")]
UnknownOperator {
node: String,
operator: String,
},
#[error("node `{node}`: the parameters of operator `{operator}` are invalid: {message}")]
InvalidParams {
node: String,
operator: String,
message: String,
},
#[error("node `{node}`: parameter `{param}` is not a template: {message}")]
InvalidTemplate {
node: String,
param: String,
message: String,
},
#[error(
"node `{node}`: parameter `{param}` refers to `{upstream}`, which is not an upstream of the node"
)]
UnknownTemplateUpstream {
node: String,
param: String,
upstream: String,
},
#[error("the graph has a cycle: {}", .0.join(" -> "))]
Cycle(Vec<String>),
}
impl Graph {
pub fn build(spec: GraphSpec, operators: &OperatorSet) -> Result<Graph, Vec<Problem>> {
let mut problems = Vec::new();
if !is_name(&spec.name) {
problems.push(Problem::InvalidGraphName(spec.name.clone()));
}
if let Some(expression) = &spec.schedule
&& let Err(e) = croner::Cron::new(expression).parse()
{
problems.push(Problem::InvalidSchedule {
expression: expression.clone(),
message: e.to_string(),
});
}
if spec.nodes.is_empty() {
problems.push(Problem::NoNodes);
}
let mut names = BTreeSet::new();
for node in &spec.nodes {
if !is_name(&node.name) {
problems.push(Problem::InvalidNodeName(node.name.clone()));
} else if spec.name.len() + node.name.len() > MAX_GRAPH_AND_NODE_NAME_LEN {
problems.push(Problem::NamesTooLong(node.name.clone()));
}
if !names.insert(node.name.as_str()) {
problems.push(Problem::DuplicateNode(node.name.clone()));
}
}
let mut producers: BTreeMap<&str, Vec<&str>> = BTreeMap::new();
for node in &spec.nodes {
if let NodeKind::Asset { produces, consumes } = &node.kind {
for asset in std::iter::once(produces).chain(consumes) {
if !is_name(asset) {
problems.push(Problem::InvalidAssetName {
node: node.name.clone(),
asset: asset.clone(),
});
}
}
producers.entry(produces).or_default().push(&node.name);
}
}
for (asset, nodes) in &producers {
if nodes.len() > 1 {
problems.push(Problem::DuplicateAsset {
asset: asset.to_string(),
nodes: nodes.iter().map(|n| n.to_string()).collect(),
});
}
}
let mut upstreams: Vec<BTreeSet<String>> = Vec::with_capacity(spec.nodes.len());
for node in &spec.nodes {
let mut set = BTreeSet::new();
match &node.kind {
NodeKind::Asset { consumes, .. } => {
for asset in consumes {
match producers.get(asset.as_str()) {
Some(nodes) => set.extend(nodes.iter().map(|n| n.to_string())),
None => problems.push(Problem::UnknownAsset {
node: node.name.clone(),
asset: asset.clone(),
}),
}
}
}
NodeKind::Task { after, .. } => {
for name in after {
if names.contains(name.as_str()) {
set.insert(name.clone());
} else {
problems.push(Problem::UnknownNode {
node: node.name.clone(),
after: name.clone(),
});
}
}
}
}
upstreams.push(set);
}
for (node, upstream) in spec.nodes.iter().zip(&upstreams) {
match operators.check(&node.operator, &node.params) {
Ok(()) => {}
Err(OperatorError::Unknown) => problems.push(Problem::UnknownOperator {
node: node.name.clone(),
operator: node.operator.clone(),
}),
Err(OperatorError::InvalidParams(message)) => {
problems.push(Problem::InvalidParams {
node: node.name.clone(),
operator: node.operator.clone(),
message,
})
}
}
check_templates(
&node.name,
"",
&toml::Value::Table(node.params.clone()),
upstream,
&mut problems,
);
}
if !problems.is_empty() {
return Err(problems);
}
if let Some(cycle) = find_cycle(&spec.nodes, &upstreams) {
return Err(vec![Problem::Cycle(cycle)]);
}
let mut downstreams: BTreeMap<&str, BTreeSet<String>> = BTreeMap::new();
for (node, upstream) in spec.nodes.iter().zip(&upstreams) {
for up in upstream {
downstreams
.entry(up.as_str())
.or_default()
.insert(node.name.clone());
}
}
let mut nodes = Vec::with_capacity(spec.nodes.len());
let mut index = BTreeMap::new();
for (i, node_spec) in spec.nodes.into_iter().enumerate() {
let downstream = downstreams
.remove(node_spec.name.as_str())
.unwrap_or_default();
index.insert(node_spec.name.clone(), i);
nodes.push(Node {
upstreams: upstreams[i].iter().cloned().collect(),
downstreams: downstream.into_iter().collect(),
spec: node_spec,
});
}
Ok(Graph {
name: spec.name,
schedule: spec.schedule,
catchup: spec.catchup,
partitioning: spec.partitioning,
nodes,
index,
})
}
pub fn name(&self) -> &str {
&self.name
}
pub fn schedule(&self) -> Option<&str> {
self.schedule.as_deref()
}
pub fn catchup(&self) -> Option<Duration> {
self.catchup
}
pub fn partitioning(&self) -> Partitioning {
self.partitioning
}
pub fn nodes(&self) -> &[Node] {
&self.nodes
}
pub fn node(&self, name: &str) -> Option<&Node> {
self.index.get(name).map(|&i| &self.nodes[i])
}
pub fn roots(&self) -> impl Iterator<Item = &Node> {
self.nodes.iter().filter(|n| n.upstreams.is_empty())
}
pub fn leaves(&self) -> impl Iterator<Item = &Node> {
self.nodes.iter().filter(|n| n.downstreams.is_empty())
}
pub fn edge_count(&self) -> usize {
self.nodes.iter().map(|n| n.upstreams.len()).sum()
}
}
pub fn is_name(text: &str) -> bool {
!text.is_empty()
&& text
.bytes()
.all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b == b'_')
}
fn check_templates(
node: &str,
path: &str,
value: &toml::Value,
upstreams: &BTreeSet<String>,
problems: &mut Vec<Problem>,
) {
let child = |key: &str| {
if path.is_empty() {
key.to_string()
} else {
format!("{path}.{key}")
}
};
match value {
toml::Value::String(text) => match Template::parse(text) {
Ok(template) => {
let mut seen = BTreeSet::new();
for upstream in template.upstream_nodes() {
if !upstreams.contains(upstream) && seen.insert(upstream) {
problems.push(Problem::UnknownTemplateUpstream {
node: node.to_string(),
param: path.to_string(),
upstream: upstream.to_string(),
});
}
}
}
Err(e) => problems.push(Problem::InvalidTemplate {
node: node.to_string(),
param: path.to_string(),
message: e.to_string(),
}),
},
toml::Value::Array(items) => {
for (i, item) in items.iter().enumerate() {
check_templates(node, &child(&i.to_string()), item, upstreams, problems);
}
}
toml::Value::Table(table) => {
for (key, item) in table {
check_templates(node, &child(key), item, upstreams, problems);
}
}
_ => {}
}
}
fn find_cycle(nodes: &[NodeSpec], upstreams: &[BTreeSet<String>]) -> Option<Vec<String>> {
#[derive(Clone, Copy, PartialEq, Eq)]
enum Mark {
Unvisited,
Active,
Done,
}
fn visit(
i: usize,
nodes: &[NodeSpec],
upstreams: &[BTreeSet<String>],
index: &BTreeMap<&str, usize>,
marks: &mut [Mark],
path: &mut Vec<usize>,
) -> Option<Vec<String>> {
marks[i] = Mark::Active;
path.push(i);
for up in &upstreams[i] {
let j = index[up.as_str()];
match marks[j] {
Mark::Active => {
let start = path
.iter()
.position(|&p| p == j)
.expect("an active node is on the path");
let mut cycle: Vec<String> = path[start..]
.iter()
.map(|&p| nodes[p].name.clone())
.collect();
cycle.push(nodes[j].name.clone());
return Some(cycle);
}
Mark::Unvisited => {
if let Some(cycle) = visit(j, nodes, upstreams, index, marks, path) {
return Some(cycle);
}
}
Mark::Done => {}
}
}
path.pop();
marks[i] = Mark::Done;
None
}
let index: BTreeMap<&str, usize> = nodes
.iter()
.enumerate()
.map(|(i, n)| (n.name.as_str(), i))
.collect();
let mut marks = vec![Mark::Unvisited; nodes.len()];
let mut path = Vec::new();
for i in 0..nodes.len() {
if marks[i] == Mark::Unvisited
&& let Some(cycle) = visit(i, nodes, upstreams, &index, &mut marks, &mut path)
{
return Some(cycle);
}
}
None
}
#[cfg(test)]
mod tests {
use super::*;
fn asset(name: &str, produces: &str, consumes: &[&str]) -> NodeSpec {
NodeSpec {
name: name.into(),
kind: NodeKind::Asset {
produces: produces.into(),
consumes: consumes.iter().map(|s| s.to_string()).collect(),
},
operator: "subprocess".into(),
pool: "default".into(),
retries: 0,
params: r#"argv = ["true"]"#.parse().unwrap(),
}
}
fn task(name: &str, after: &[&str]) -> NodeSpec {
NodeSpec {
name: name.into(),
kind: NodeKind::Task {
after: after.iter().map(|s| s.to_string()).collect(),
trigger_rule: TriggerRule::AllDone,
},
..asset(name, "", &[])
}
}
fn spec(nodes: Vec<NodeSpec>) -> GraphSpec {
GraphSpec {
name: "g".into(),
schedule: Some("0 2 * * *".into()),
catchup: None,
partitioning: Partitioning::Daily,
nodes,
}
}
fn build(nodes: Vec<NodeSpec>) -> Result<Graph, Vec<Problem>> {
Graph::build(spec(nodes), &OperatorSet::builtin())
}
#[test]
fn edges_follow_from_assets_and_after_and_nodes_keep_definition_order() {
let graph = build(vec![
asset("load", "c", &["b"]),
asset("transform", "b", &["a"]),
task("notify", &["load", "transform"]),
asset("extract", "a", &[]),
])
.unwrap();
let names: Vec<&str> = graph.nodes().iter().map(Node::name).collect();
assert_eq!(names, ["load", "transform", "notify", "extract"]);
assert_eq!(graph.node("load").unwrap().upstreams(), ["transform"]);
assert_eq!(
graph.node("transform").unwrap().downstreams(),
["load", "notify"]
);
assert_eq!(
graph.node("notify").unwrap().upstreams(),
["load", "transform"]
);
assert_eq!(
graph.roots().map(Node::name).collect::<Vec<_>>(),
["extract"]
);
assert_eq!(
graph.leaves().map(Node::name).collect::<Vec<_>>(),
["notify"]
);
assert_eq!(graph.edge_count(), 4);
assert_eq!(graph.node("load").unwrap().asset(), Some("c"));
assert_eq!(graph.node("notify").unwrap().asset(), None);
assert_eq!(
graph.node("load").unwrap().trigger_rule(),
TriggerRule::AllSucceeded
);
assert_eq!(
graph.node("notify").unwrap().trigger_rule(),
TriggerRule::AllDone
);
}
#[test]
fn every_fault_is_reported_at_once() {
let mut bad = spec(vec![
asset("Extract", "Raw", &["missing"]),
asset("x", "dup", &[]),
asset("x", "dup", &[]),
NodeSpec {
operator: "sql".into(),
..asset("y", "y", &[])
},
NodeSpec {
params: "".parse().unwrap(),
..asset("z", "z", &[])
},
task("t", &["nowhere"]),
]);
bad.name = "Bad-Graph".into();
bad.schedule = Some("every day".into());
let problems = Graph::build(bad, &OperatorSet::builtin()).unwrap_err();
assert!(problems.contains(&Problem::InvalidGraphName("Bad-Graph".into())));
assert!(
matches!(&problems[1], Problem::InvalidSchedule { expression, .. } if expression == "every day")
);
assert!(problems.contains(&Problem::InvalidNodeName("Extract".into())));
assert!(problems.contains(&Problem::InvalidAssetName {
node: "Extract".into(),
asset: "Raw".into(),
}));
assert!(problems.contains(&Problem::UnknownAsset {
node: "Extract".into(),
asset: "missing".into(),
}));
assert!(problems.contains(&Problem::DuplicateNode("x".into())));
assert!(problems.contains(&Problem::DuplicateAsset {
asset: "dup".into(),
nodes: vec!["x".into(), "x".into()],
}));
assert!(problems.contains(&Problem::UnknownOperator {
node: "y".into(),
operator: "sql".into(),
}));
assert!(matches!(
problems.iter().find(|p| matches!(p, Problem::InvalidParams { .. })),
Some(Problem::InvalidParams { node, operator, .. }) if node == "z" && operator == "subprocess"
));
assert!(problems.contains(&Problem::UnknownNode {
node: "t".into(),
after: "nowhere".into(),
}));
}
#[test]
fn empty_graph_is_rejected() {
assert_eq!(build(vec![]), Err(vec![Problem::NoNodes]));
}
#[test]
fn names_are_bounded_so_every_run_id_fits() {
let long = "n".repeat(MAX_GRAPH_AND_NODE_NAME_LEN - 1);
assert!(build(vec![asset(&long, "a", &[])]).is_ok());
let too_long = "n".repeat(MAX_GRAPH_AND_NODE_NAME_LEN);
assert_eq!(
build(vec![asset(&too_long, "a", &[])]),
Err(vec![Problem::NamesTooLong(too_long)])
);
}
#[test]
fn templates_are_checked_at_every_depth_against_declared_upstreams() {
let mut load = asset("load", "c", &["a"]);
load.params = r#"
argv = ["{{ upstream.extract.rows }} {{ upstream.extract.rows }}"]
[nested]
list = ["ok {{ partition }}", "{{ upstream.other.x }}"]
broken = "{{ partition"
"#
.parse()
.unwrap();
let problems = build(vec![
asset("extract", "a", &[]),
asset("other", "o", &[]),
load,
])
.unwrap_err();
assert_eq!(
problems,
vec![
Problem::InvalidParams {
node: "load".into(),
operator: "subprocess".into(),
message: problems
.iter()
.find_map(|p| match p {
Problem::InvalidParams { message, .. } => Some(message.clone()),
_ => None,
})
.unwrap(),
},
Problem::InvalidTemplate {
node: "load".into(),
param: "nested.broken".into(),
message: "`{{` at byte 0 is not closed".into(),
},
Problem::UnknownTemplateUpstream {
node: "load".into(),
param: "nested.list.1".into(),
upstream: "other".into(),
},
]
);
}
#[test]
fn cycle_is_reported_as_its_path() {
let problems = build(vec![
asset("root", "r", &[]),
asset("a", "a", &["b", "r"]),
asset("b", "b", &["a"]),
asset("after_cycle", "d", &["b"]),
])
.unwrap_err();
assert_eq!(
problems,
vec![Problem::Cycle(vec!["a".into(), "b".into(), "a".into()])]
);
assert_eq!(
problems[0].to_string(),
"the graph has a cycle: a -> b -> a"
);
}
#[test]
fn is_name_accepts_lowercase_digits_and_underscore_only() {
assert!(is_name("orders_daily_2"));
for bad in ["", "Orders", "orders-daily", "orders daily", "ordérs"] {
assert!(!is_name(bad), "{bad}");
}
}
}