knut-thund 0.1.7

Þund — a Rust-native, Arrow-centric streaming dataflow engine (batch + streaming) with a pluggable execution backend: native Arrow/DataFusion or lower-to-Spark-Declarative-Pipelines via Spark Connect. The 'Airflow killer' authoring+runtime for knut.
Documentation
//! [`Pipeline`] — the top-level IR: the whole dataflow graph plus its
//! identity and defaults. This is what every authoring surface produces and
//! every [`crate::backend::ExecBackend`] consumes.

use serde::{Deserialize, Serialize};
use std::collections::{BTreeMap, BTreeSet, VecDeque};

use super::dataset::Dataset;
use super::flow::Flow;
use crate::error::{Result, ThundError};

/// A complete Þund pipeline: datasets (nodes) + flows (edges) + identity.
///
/// The `(flow.reads -> flow.target)` pairs are the DAG. The graph is a
/// superset of [`knut_pipelines::DataflowGraph`]; [`Pipeline::to_sdp`] lowers
/// it for the Spark backend.
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct Pipeline {
    /// Human pipeline name.
    pub name: String,
    /// Default catalog for unqualified dataset names.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub default_catalog: Option<String>,
    /// Default database for unqualified dataset names.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub default_database: Option<String>,
    /// Storage root for checkpoints + metadata (the SDP `storage` URI).
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub storage: Option<String>,
    /// The output nodes, in declaration order.
    pub datasets: Vec<Dataset>,
    /// The flow edges, in declaration order.
    pub flows: Vec<Flow>,
}

impl Pipeline {
    /// An empty pipeline named `name`.
    pub fn new(name: impl Into<String>) -> Self {
        Pipeline {
            name: name.into(),
            ..Default::default()
        }
    }

    /// Add a dataset node, returning `self` for chaining.
    pub fn with_dataset(mut self, d: Dataset) -> Self {
        self.datasets.push(d);
        self
    }

    /// Add a flow edge, returning `self` for chaining.
    pub fn with_flow(mut self, f: Flow) -> Self {
        self.flows.push(f);
        self
    }

    /// Look up a dataset by exact name.
    pub fn dataset(&self, name: &str) -> Option<&Dataset> {
        self.datasets.iter().find(|d| d.name == name)
    }

    /// `true` if any flow is unbounded — i.e. this is a streaming pipeline
    /// that the native backend must run as a long-lived job and that the Spark
    /// backend lowers to streaming tables.
    pub fn is_streaming(&self) -> bool {
        self.flows.iter().any(|f| f.kind.is_unbounded())
    }

    /// Validate internal consistency: unique dataset/flow names, every flow
    /// target resolving to a defined dataset, and the in-graph DAG acyclic.
    ///
    /// Reads are *not* required to be in-graph datasets — a flow may read an
    /// external table (a source).
    pub fn validate(&self) -> Result<()> {
        let result = self.validate_inner();
        // Emit the IR validation gate's health into the nornir matrix.
        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(())
    }

    /// A deterministic topological order of dataset names following the flow
    /// edges (`reads` precede `target`). `None` if the in-graph DAG is cyclic.
    /// External reads (not in-graph) are ignored for ordering.
    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)
    }

    /// Lower the IR to a Spark-SDP [`knut_pipelines::DataflowGraph`].
    ///
    /// This is the heart of the Spark backend: every dataset and flow maps to
    /// its SDP form (see [`Dataset::to_sdp`] / [`Flow::to_sdp`]). Streaming
    /// knobs and expectations that SDP cannot yet express are dropped here and
    /// surfaced by the backend's capability report — the *graph shape* always
    /// lowers faithfully so the existing `knut-pipelines` client can define +
    /// run it against a real Spark Connect server.
    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());
        }
        // Emit the IR→SDP lowering surface's health: the graph shape always
        // lowers faithfully, so this records a green row carrying the lowered
        // dataset/flow counts.
        crate::functional_status(
            "knut-thund/ir",
            "lower_to_sdp",
            true,
            &format!(
                "{}: {} dataset(s), {} flow(s)",
                self.name,
                g.datasets.len(),
                g.flows.len()
            ),
        );
        g
    }
}