knut-thund 0.1.1

Þ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
//! [`Flow`] — an edge in the IR: how one [`crate::ir::Dataset`] is computed
//! from its inputs. A flow is the unit that carries the batch/streaming
//! distinction, the query, the streaming semantics, the CDC spec, and the
//! data-quality expectations.

use serde::{Deserialize, Serialize};

use super::expectation::Expectation;
use super::streaming::{OutputMode, SourceSpec, Trigger, Watermark, WindowSpec};

/// Slowly-changing-dimension type for a CDC/`APPLY CHANGES` flow.
///
/// Mirrors `spark.connect.SCDType`: SCD type 1 overwrites; type 2 keeps
/// history. (OSS SDP currently only carries type 1; type 2 lowers on the
/// native backend.)
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ScdType {
    /// Overwrite the matching row (no history).
    Type1,
    /// Keep historical versions of the matching row.
    Type2,
}

/// A CDC / `APPLY CHANGES INTO` specification on a flow — apply an upstream
/// changelog (inserts/updates/deletes) to the target as an idempotent merge.
///
/// Mirrors `spark.connect.AutoCdcFlowDetails`.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CdcSpec {
    /// Primary-key columns used to match changelog rows to target rows.
    pub keys: Vec<String>,
    /// Column establishing change order (the changelog sequence).
    pub sequence_by: String,
    /// Predicate identifying delete rows in the changelog, if any.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub apply_as_deletes: Option<String>,
    /// SCD behaviour.
    pub scd_type: ScdType,
}

/// What kind of computation a flow performs — the batch/streaming switch.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum FlowKind {
    /// A bounded batch query (a `SELECT` recomputing the target).
    Batch,
    /// An unbounded streaming query with full event-time semantics.
    Streaming {
        /// Where the stream is read from.
        source: SourceSpec,
        /// Event-time watermark policy, if event-time processing is used.
        #[serde(default, skip_serializing_if = "Option::is_none")]
        watermark: Option<Watermark>,
        /// When results are emitted.
        #[serde(default)]
        trigger: Trigger,
        /// What is written each trigger.
        #[serde(default)]
        output_mode: OutputMode,
        /// Windowed-aggregation shape, if any.
        #[serde(default, skip_serializing_if = "Option::is_none")]
        window: Option<WindowSpec>,
    },
    /// A change-data-capture merge flow.
    Cdc {
        /// The CDC application spec.
        cdc: CdcSpec,
    },
}

impl FlowKind {
    /// `true` when this flow drives an unbounded (streaming or CDC) target.
    pub fn is_unbounded(&self) -> bool {
        matches!(self, FlowKind::Streaming { .. } | FlowKind::Cdc { .. })
    }

    /// A stable label for rendering.
    pub fn label(&self) -> &'static str {
        match self {
            FlowKind::Batch => "batch",
            FlowKind::Streaming { .. } => "streaming",
            FlowKind::Cdc { .. } => "cdc",
        }
    }
}

/// An edge in the IR dataflow graph: a flow writing one target dataset,
/// reading zero or more upstream datasets.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Flow {
    /// Flow name (unique within the pipeline).
    pub name: String,
    /// Dataset this flow writes to (the edge head).
    pub target: String,
    /// Upstream datasets this flow reads (edge tails). May be empty for a
    /// source flow reading from outside the graph.
    #[serde(default, skip_serializing_if = "Vec::is_empty")]
    pub reads: Vec<String>,
    /// Batch vs streaming vs CDC, with the attendant semantics.
    pub kind: FlowKind,
    /// The query that defines the flow — SQL (or a serialized relation). For
    /// the Spark backend this becomes the flow's `Relation`; required for the
    /// flow to *run* (vs register a structure-only edge).
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub query: Option<String>,
    /// Data-quality expectations on this flow's output.
    #[serde(default, skip_serializing_if = "Vec::is_empty")]
    pub expectations: Vec<Expectation>,
}

impl Flow {
    /// A batch flow named `name` writing `target`, reading `reads`.
    pub fn batch(
        name: impl Into<String>,
        target: impl Into<String>,
        reads: impl IntoIterator<Item = impl Into<String>>,
    ) -> Self {
        Flow {
            name: name.into(),
            target: target.into(),
            reads: reads.into_iter().map(Into::into).collect(),
            kind: FlowKind::Batch,
            query: None,
            expectations: Vec::new(),
        }
    }

    /// A streaming flow named `name` writing `target` from `source`.
    pub fn streaming(
        name: impl Into<String>,
        target: impl Into<String>,
        source: SourceSpec,
    ) -> Self {
        Flow {
            name: name.into(),
            target: target.into(),
            reads: Vec::new(),
            kind: FlowKind::Streaming {
                source,
                watermark: None,
                trigger: Trigger::default(),
                output_mode: OutputMode::default(),
                window: None,
            },
            query: None,
            expectations: Vec::new(),
        }
    }

    /// Attach the defining query, returning `self`.
    pub fn with_query(mut self, sql: impl Into<String>) -> Self {
        self.query = Some(sql.into());
        self
    }

    /// Attach an expectation, returning `self`.
    pub fn expect(mut self, e: Expectation) -> Self {
        self.expectations.push(e);
        self
    }

    /// Set the watermark on a streaming flow (no-op on non-streaming kinds),
    /// returning `self`.
    pub fn with_watermark(mut self, wm: Watermark) -> Self {
        if let FlowKind::Streaming { watermark, .. } = &mut self.kind {
            *watermark = Some(wm);
        }
        self
    }

    /// Lower this flow to a Spark-SDP [`knut_pipelines::Flow`].
    ///
    /// SDP can express the *shape* (target/reads/query/once) but not the
    /// streaming knobs (watermark/trigger/output-mode/window) or expectations;
    /// those are carried only on the native backend. A streaming flow with an
    /// `AvailableNow` trigger lowers with `once = true`.
    pub fn to_sdp(&self) -> knut_pipelines::Flow {
        let mut f = knut_pipelines::Flow::new(self.name.clone(), self.target.clone(), self.reads.clone());
        if let Some(q) = &self.query {
            f = f.with_query(q.clone());
        }
        if let FlowKind::Streaming {
            trigger: Trigger::AvailableNow,
            ..
        } = &self.kind
        {
            f.once = true;
        }
        f
    }
}