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

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

/// 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>,
    /// Partition columns, if any.
    #[serde(default, skip_serializing_if = "Vec::is_empty")]
    pub partition_cols: Vec<String>,
    /// Optional human comment.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub comment: Option<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,
            partition_cols: Vec::new(),
            comment: None,
            source_location: SourceLocation::default(),
        }
    }

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

    /// Lower this dataset to its Spark-SDP [`knut_pipelines::Dataset`] form.
    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
    }
}