use serde::{Deserialize, Serialize};
use std::collections::{BTreeMap, BTreeSet, VecDeque};
use super::dataset::Dataset;
use super::flow::Flow;
use crate::error::{Result, ThundError};
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct Pipeline {
pub name: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub default_catalog: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub default_database: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub storage: Option<String>,
pub datasets: Vec<Dataset>,
pub flows: Vec<Flow>,
}
impl Pipeline {
pub fn new(name: impl Into<String>) -> Self {
Pipeline {
name: name.into(),
..Default::default()
}
}
pub fn with_dataset(mut self, d: Dataset) -> Self {
self.datasets.push(d);
self
}
pub fn with_flow(mut self, f: Flow) -> Self {
self.flows.push(f);
self
}
pub fn dataset(&self, name: &str) -> Option<&Dataset> {
self.datasets.iter().find(|d| d.name == name)
}
pub fn is_streaming(&self) -> bool {
self.flows.iter().any(|f| f.kind.is_unbounded())
}
pub fn validate(&self) -> Result<()> {
let result = self.validate_inner();
crate::functional_status(
"knut-thund/ir",
"validate",
result.is_ok(),
&match &result {
Ok(()) => self.name.clone(),
Err(e) => e.to_string(),
},
);
result
}
fn validate_inner(&self) -> Result<()> {
let mut seen_ds = BTreeSet::new();
for d in &self.datasets {
if !seen_ds.insert(d.name.as_str()) {
return Err(ThundError::DuplicateName {
kind: "dataset",
name: d.name.clone(),
});
}
}
let mut seen_flow = BTreeSet::new();
for f in &self.flows {
if !seen_flow.insert(f.name.as_str()) {
return Err(ThundError::DuplicateName {
kind: "flow",
name: f.name.clone(),
});
}
if !seen_ds.contains(f.target.as_str()) {
return Err(ThundError::DanglingFlow {
flow: f.name.clone(),
target: f.target.clone(),
});
}
}
if self.topo_order().is_none() {
return Err(ThundError::Cyclic);
}
Ok(())
}
pub fn topo_order(&self) -> Option<Vec<String>> {
let nodes: BTreeSet<&str> = self.datasets.iter().map(|d| d.name.as_str()).collect();
let mut indeg: BTreeMap<&str, usize> = nodes.iter().map(|n| (*n, 0)).collect();
let mut adj: BTreeMap<&str, BTreeSet<&str>> = BTreeMap::new();
for f in &self.flows {
if !nodes.contains(f.target.as_str()) {
continue;
}
for r in &f.reads {
if !nodes.contains(r.as_str()) {
continue;
}
if adj.entry(r.as_str()).or_default().insert(f.target.as_str()) {
*indeg.get_mut(f.target.as_str()).unwrap() += 1;
}
}
}
let mut queue: VecDeque<&str> = indeg
.iter()
.filter(|(_, d)| **d == 0)
.map(|(n, _)| *n)
.collect();
let mut out = Vec::with_capacity(nodes.len());
while let Some(n) = queue.pop_front() {
out.push(n.to_string());
if let Some(succ) = adj.get(n) {
for &m in succ {
let d = indeg.get_mut(m).unwrap();
*d -= 1;
if *d == 0 {
queue.push_back(m);
}
}
}
}
(out.len() == nodes.len()).then_some(out)
}
pub fn to_sdp(&self) -> knut_pipelines::DataflowGraph {
let mut g = knut_pipelines::DataflowGraph::new();
g.default_catalog = self.default_catalog.clone();
g.default_database = self.default_database.clone();
for d in &self.datasets {
g = g.with_dataset(d.to_sdp());
}
for f in &self.flows {
g = g.with_flow(f.to_sdp());
}
crate::functional_status(
"knut-thund/ir",
"lower_to_sdp",
true,
&format!(
"{}: {} dataset(s), {} flow(s)",
self.name,
g.datasets.len(),
g.flows.len()
),
);
g
}
}