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
//! [`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) **unless** the flow
    /// instead carries a typed [`projection`](Self::projection) /
    /// [`filter`](Self::filter) over a single input (see below).
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub query: Option<String>,
    /// **Typed column projection** — the ordered list of columns to select from
    /// this flow's single input. Unlike the opaque [`query`](Self::query) SQL
    /// string (which the engine cannot introspect), this is a *structured* node
    /// the planner reads directly: the native backend lowers it to a DataFusion
    /// scan that **prunes columns at the source** (only the listed columns are
    /// read off disk), and [`to_sdp`](Self::to_sdp) synthesizes the equivalent
    /// `SELECT` so Spark/Catalyst does its own column pruning. Empty = no
    /// projection (all columns). Combined with a [`query`](Self::query) it
    /// applies on top of the query result; on its own (no query) it reads the
    /// single [`reads`](Self::reads) input directly, giving real scan pushdown.
    #[serde(default, skip_serializing_if = "Vec::is_empty")]
    pub projection: Vec<String>,
    /// **Typed row filter** — a simple SQL predicate string (e.g.
    /// `"amount > 100"`) the planner can push into the source scan. Like
    /// [`projection`](Self::projection) this is introspectable structure, not an
    /// opaque query: the native backend lowers it to a DataFusion `filter` that
    /// (for a file source) becomes a **pushdown predicate** pruning row-groups
    /// at the source, and [`to_sdp`](Self::to_sdp) lowers it to a `WHERE` clause
    /// Spark pushes down. `None` = no filter.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub filter: 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,
            projection: Vec::new(),
            filter: 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,
            projection: Vec::new(),
            filter: 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 a **typed column projection** — the ordered columns to select —
    /// returning `self`. On its own (no [`with_query`](Self::with_query)) the
    /// native backend reads the flow's single input and prunes to exactly these
    /// columns at the scan; see [`projection`](Self::projection).
    pub fn with_projection<I, S>(mut self, cols: I) -> Self
    where
        I: IntoIterator<Item = S>,
        S: Into<String>,
    {
        self.projection = cols.into_iter().map(Into::into).collect();
        self
    }

    /// Attach a **typed row filter** predicate (e.g. `"amount > 100"`),
    /// returning `self`. The native backend pushes it into the source scan; see
    /// [`filter`](Self::filter).
    pub fn with_filter(mut self, predicate: impl Into<String>) -> Self {
        self.filter = Some(predicate.into());
        self
    }

    /// `true` when this flow carries a typed projection or filter node the
    /// planner can introspect (as opposed to only an opaque [`query`](Self::query)).
    pub fn has_pushdown(&self) -> bool {
        !self.projection.is_empty() || self.filter.is_some()
    }

    /// Synthesize the `SELECT … FROM <input> [WHERE …]` SQL equivalent of a
    /// typed [`projection`](Self::projection)/[`filter`](Self::filter) flow that
    /// carries **no** opaque [`query`](Self::query). This is how the typed
    /// pushdown node LOWERS to the Spark backend: Catalyst then does its own
    /// column pruning + filter pushdown from the generated SQL. Returns `None`
    /// when the flow has an explicit query, has no pushdown node, or does not
    /// read exactly one input (the single-input scan requirement).
    pub fn synth_query(&self) -> Option<String> {
        if self.query.is_some() || !self.has_pushdown() {
            return None;
        }
        let [input] = self.reads.as_slice() else {
            return None;
        };
        let cols = if self.projection.is_empty() {
            "*".to_string()
        } else {
            self.projection.join(", ")
        };
        let mut sql = format!("SELECT {cols} FROM {input}");
        if let Some(pred) = &self.filter {
            sql.push_str(" WHERE ");
            sql.push_str(pred);
        }
        Some(sql)
    }

    /// 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
    }

    /// Set the windowed-aggregation spec on a streaming flow (no-op on
    /// non-streaming kinds), returning `self`. Consumed by the native backend's
    /// incremental mode to close/emit/evict windows on the event-time watermark;
    /// the recompute path leaves it advisory.
    pub fn with_window(mut self, win: WindowSpec) -> Self {
        if let FlowKind::Streaming { window, .. } = &mut self.kind {
            *window = Some(win);
        }
        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`. A typed
    /// [`projection`](Self::projection)/[`filter`](Self::filter) flow with no
    /// explicit query lowers via [`synth_query`](Self::synth_query) to the
    /// equivalent `SELECT … WHERE …`, so Spark pushes the pruning/filter down.
    pub fn to_sdp(&self) -> knut_pipelines::Flow {
        let mut f =
            knut_pipelines::Flow::new(self.name.clone(), self.target.clone(), self.reads.clone());
        // Query: the explicit opaque SQL if present, else — for a typed
        // projection/filter flow with no query — the synthesized `SELECT` so
        // Spark/Catalyst does its own column pruning + filter pushdown. This is
        // the Spark half of the typed pushdown node (the native half lives in
        // the DataFusion backend's scan lowering).
        if let Some(q) = self.query.clone().or_else(|| self.synth_query()) {
            f = f.with_query(q);
        }
        if let FlowKind::Streaming {
            trigger: Trigger::AvailableNow,
            ..
        } = &self.kind
        {
            f.once = true;
        }
        f
    }
}