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
//! [`Dataset`] — a node in the IR: one declared *output* (what data should
//! exist), the asset-graph primitive.

use std::collections::BTreeMap;

use serde::{Deserialize, Serialize};

use super::schema::DatasetSchema;

/// The kind of output a [`Dataset`] node represents.
///
/// A strict superset of [`knut_pipelines::OutputType`]: the four SDP kinds
/// lower 1:1, and Þund adds no new *kinds* here — the batch/streaming
/// distinction lives on the [`Materialization`] and on the feeding
/// [`crate::ir::Flow`], not on the output kind, exactly as SDP models it (a
/// "streaming table" is just a `Table` fed by a streaming flow).
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum OutputType {
    /// A query result recomputed and published to the catalog (batch).
    MaterializedView,
    /// A table published to the catalog; becomes a *streaming table* when fed
    /// by a streaming flow.
    Table,
    /// An internal view, not published to the catalog.
    TemporaryView,
    /// A streaming write to an external system, not published (SDP 4.2+).
    Sink,
}

impl OutputType {
    /// A short, stable label for text/JSON rendering.
    pub fn label(self) -> &'static str {
        match self {
            OutputType::MaterializedView => "materialized_view",
            OutputType::Table => "table",
            OutputType::TemporaryView => "temporary_view",
            OutputType::Sink => "sink",
        }
    }

    /// Lower to the Spark-SDP [`knut_pipelines::OutputType`] this maps onto.
    pub fn to_sdp(self) -> knut_pipelines::OutputType {
        match self {
            OutputType::MaterializedView => knut_pipelines::OutputType::MaterializedView,
            OutputType::Table => knut_pipelines::OutputType::Table,
            OutputType::TemporaryView => knut_pipelines::OutputType::TemporaryView,
            OutputType::Sink => knut_pipelines::OutputType::Sink,
        }
    }
}

/// How a dataset's contents are maintained over time.
///
/// This is the Dagster-style reconciliation knob: a [`Materialization::Full`]
/// dataset is recomputed from scratch, [`Materialization::Incremental`] only
/// processes new/changed input (the basis for "recompute only what's stale").
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Materialization {
    /// Recompute the whole dataset (batch full-refresh).
    Full,
    /// Process only new/changed input since the last run.
    Incremental,
}

impl Default for Materialization {
    fn default() -> Self {
        Materialization::Full
    }
}

/// How a file **sink** treats an already-populated destination when a flow is
/// re-run — the declarative clear-vs-accumulate knob.
///
/// This governs the native on-disk file sink ([`Dataset::path`] /
/// [`crate::backend::native::NativeBackend::with_file_output`]). A
/// Hive-partitioned directory sink writes fresh part-files on every run, so
/// without a clear a re-run *accumulates* — the output holds two (or N)
/// generations of rows. [`WriteMode::Overwrite`] clears the destination first,
/// so a re-run holds exactly one generation.
///
/// The default is [`WriteMode::Append`] — **today's behaviour**, so existing
/// specs deserialize unchanged and keep appending (L2-additive). Opt into
/// clear-then-write with [`Dataset::overwrite`] / [`Dataset::with_write_mode`].
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum WriteMode {
    /// Add to whatever is already at the destination — the pre-existing
    /// behaviour. A partitioned directory accumulates new part-files each run;
    /// a single-file sink is truncated-and-rewritten by the file writer either
    /// way (so Append and Overwrite coincide for a single file — the distinction
    /// is observable on a partitioned/directory sink).
    Append,
    /// Clear the destination file/directory before writing, so a re-run replaces
    /// the prior contents (fresh-dir semantics) and holds exactly one
    /// generation's rows.
    Overwrite,
}

impl Default for WriteMode {
    fn default() -> Self {
        WriteMode::Append
    }
}

impl WriteMode {
    /// A short, stable label for text/JSON rendering.
    pub fn label(self) -> &'static str {
        match self {
            WriteMode::Append => "append",
            WriteMode::Overwrite => "overwrite",
        }
    }

    /// `true` when this is the default [`WriteMode::Append`] — used to keep the
    /// serialized form byte-identical for the default (L2-additive).
    pub fn is_append(&self) -> bool {
        matches!(self, WriteMode::Append)
    }
}

/// Where a graph element was declared in user source — for deep-linking a
/// viewer node/event back to the authoring surface. Mirrors
/// [`knut_pipelines::SourceLocation`].
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct SourceLocation {
    /// File the element was defined in, if known.
    pub file_name: Option<String>,
    /// 1-based line number, if known.
    pub line_number: Option<i32>,
}

impl SourceLocation {
    /// `true` when no location information is present.
    pub fn is_empty(&self) -> bool {
        self.file_name.is_none() && self.line_number.is_none()
    }
}

/// A node in the IR dataflow graph: one declared output.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Dataset {
    /// Output name, possibly qualified (`catalog.db.name`).
    pub name: String,
    /// What kind of output this is.
    pub output_type: OutputType,
    /// How it is maintained over time.
    #[serde(default)]
    pub materialization: Materialization,
    /// Arrow-typed schema; empty = infer from the feeding flow.
    #[serde(default, skip_serializing_if = "DatasetSchema::is_empty")]
    pub schema: DatasetSchema,
    /// Storage format (`delta`, `iceberg`, `parquet`), if specified.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub format: Option<String>,
    /// On-disk destination for this output, if it is a **file sink declared in
    /// the IR itself** (rather than via the imperative
    /// [`crate::backend::native::NativeBackend::with_file_output`] builder). When
    /// set, the native backend writes this dataset's materialised rows to `path`
    /// after the run — a single `parquet`/`csv` file, or a Hive-partitioned
    /// directory when [`partition_cols`](Self::partition_cols) is non-empty —
    /// using [`format`](Self::format) (defaulting to the path extension, else
    /// parquet). This is what lets a pipeline authored purely as IR/RON name
    /// where its outputs land, so a `.ron` spec round-trips a full read→write
    /// pipeline with no imperative wiring. `None` = not a declarative file sink.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub path: Option<String>,
    /// Partition columns, if any.
    #[serde(default, skip_serializing_if = "Vec::is_empty")]
    pub partition_cols: Vec<String>,
    /// How a re-run of a file sink treats an already-populated destination:
    /// [`WriteMode::Append`] (default — accumulate, today's behaviour) vs
    /// [`WriteMode::Overwrite`] (clear the target dir/file first, one generation
    /// only). Governs the native file sink ([`path`](Self::path) /
    /// `with_file_output`); serde-default = `Append` so existing configs
    /// deserialize unchanged and keep appending (L2-additive).
    #[serde(default, skip_serializing_if = "WriteMode::is_append")]
    pub write_mode: WriteMode,
    /// Optional human comment.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub comment: Option<String>,
    /// Table properties — SQL `TBLPROPERTIES ('k' = 'v', …)`.
    ///
    /// Lowered to `DefineOutput.TableDetails.table_properties` by
    /// [`Self::to_sdp`], so these genuinely reach Spark rather than being
    /// parsed and dropped.
    ///
    /// A `BTreeMap`, not a `HashMap`: these round-trip through RON specs that
    /// get diffed and warehouse-historized, and a randomized iteration order
    /// would rewrite the file on every save. Empty is skipped on serialize, so
    /// existing specs deserialize byte-identically (L2-additive).
    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
    pub properties: BTreeMap<String, String>,
    /// Where this output was declared in source.
    #[serde(default, skip_serializing_if = "SourceLocation::is_empty")]
    pub source_location: SourceLocation,
}

impl Dataset {
    /// A new dataset with just a name and kind; everything else defaulted.
    pub fn new(name: impl Into<String>, output_type: OutputType) -> Self {
        Dataset {
            name: name.into(),
            output_type,
            materialization: Materialization::default(),
            schema: DatasetSchema::new(),
            format: None,
            path: None,
            partition_cols: Vec::new(),
            write_mode: WriteMode::default(),
            comment: None,
            properties: BTreeMap::new(),
            source_location: SourceLocation::default(),
        }
    }

    /// Set table properties (`TBLPROPERTIES`), returning `self` for chaining.
    pub fn with_properties<I, K, V>(mut self, props: I) -> Self
    where
        I: IntoIterator<Item = (K, V)>,
        K: Into<String>,
        V: Into<String>,
    {
        self.properties = props
            .into_iter()
            .map(|(k, v)| (k.into(), v.into()))
            .collect();
        self
    }

    /// Set the Arrow schema, returning `self` for chaining.
    pub fn with_schema(mut self, schema: DatasetSchema) -> Self {
        self.schema = schema;
        self
    }

    /// Mark this dataset incrementally materialized, returning `self`.
    pub fn incremental(mut self) -> Self {
        self.materialization = Materialization::Incremental;
        self
    }

    /// Set the storage format, returning `self`.
    pub fn with_format(mut self, fmt: impl Into<String>) -> Self {
        self.format = Some(fmt.into());
        self
    }

    /// Declare an on-disk destination for this output, returning `self` — a
    /// **file sink named in the IR**. After the run the native backend writes
    /// this dataset's rows to `path` (single file, or a Hive-partitioned
    /// directory when [`partition_cols`](Self::partition_cols) is set), so a
    /// pipeline loaded from RON can land its results with no imperative
    /// [`with_file_output`](crate::backend::native::NativeBackend::with_file_output)
    /// call. Pair with [`with_format`](Self::with_format) to pick `parquet`/`csv`
    /// explicitly (otherwise the path extension decides, defaulting to parquet).
    pub fn with_path(mut self, path: impl Into<String>) -> Self {
        self.path = Some(path.into());
        self
    }

    /// Declare the partition columns, returning `self`. On the native backend a
    /// declared, non-empty partitioning makes a file sink write a Hive-partitioned
    /// directory (`col=val/…`) instead of a single file.
    pub fn with_partition_cols<I, S>(mut self, cols: I) -> Self
    where
        I: IntoIterator<Item = S>,
        S: Into<String>,
    {
        self.partition_cols = cols.into_iter().map(Into::into).collect();
        self
    }

    /// Set the file-sink [`WriteMode`], returning `self`. [`WriteMode::Overwrite`]
    /// makes the native backend clear the destination before writing (a re-run
    /// holds one generation); [`WriteMode::Append`] (the default) accumulates.
    pub fn with_write_mode(mut self, mode: WriteMode) -> Self {
        self.write_mode = mode;
        self
    }

    /// Shorthand for [`with_write_mode(WriteMode::Overwrite)`](Self::with_write_mode)
    /// — clear the sink destination before writing, returning `self`.
    pub fn overwrite(mut self) -> Self {
        self.write_mode = WriteMode::Overwrite;
        self
    }

    /// Lower this dataset to its Spark-SDP [`knut_pipelines::Dataset`] form.
    ///
    /// [`write_mode`](Self::write_mode) is **not** carried here: it governs the
    /// native on-disk file sink ([`path`](Self::path)) — a native-backend concept
    /// that, like `path`, is not part of the SDP declarative graph. SDP already
    /// expresses overwrite-vs-append through the *output type / materialization*
    /// (a `MaterializedView` is fully recomputed each run — overwrite semantics —
    /// while a `Table` fed by a streaming flow appends), both of which `to_sdp`
    /// already lowers. See `thund-design.md`.
    pub fn to_sdp(&self) -> knut_pipelines::Dataset {
        let mut d = knut_pipelines::Dataset::new(self.name.clone(), self.output_type.to_sdp());
        d.comment = self.comment.clone();
        d.format = self.format.clone();
        d.partition_cols = self.partition_cols.clone();
        d.schema = self.schema.to_sdp_schema_string();
        d.properties = self.properties.clone();
        d
    }
}