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
//! Data-quality **expectations** — boolean constraints evaluated per row, the
//! Þund analogue of Dagster's blocking asset checks and SDP/DLT's
//! `@expect`/`CONSTRAINT … EXPECT` clauses.
//!
//! An expectation is a named boolean expression over a flow's output rows plus
//! a policy for what to do when a row violates it. Lowering: on the Spark
//! backend these become SDP expectation clauses (gated behind a Spark
//! capability check — they are not yet fully upstreamed into OSS Spark 4.1, so
//! `SparkBackend` reports them via the capability surface); on the native
//! backend they are a filter/abort operator on the RecordBatch stream.

use serde::{Deserialize, Serialize};

/// What happens to a row (or the whole run) when an [`Expectation`] is
/// violated.
///
/// Mirrors DLT/SDP `ON VIOLATION {DROP ROW | FAIL UPDATE}` plus the implicit
/// "warn and keep" default.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum OnViolation {
    /// Keep the violating row but record the failure (warn). The default.
    Warn,
    /// Drop the violating row from the output.
    Drop,
    /// Fail the whole run on the first violating row.
    Fail,
}

impl Default for OnViolation {
    fn default() -> Self {
        OnViolation::Warn
    }
}

/// A named data-quality constraint attached to a [`crate::ir::Flow`].
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Expectation {
    /// Constraint name (surfaced in metrics + the warehouse run record).
    pub name: String,
    /// A boolean SQL/Arrow expression over the flow's output columns; rows for
    /// which it is `false` are violations.
    pub constraint: String,
    /// What to do on violation.
    #[serde(default)]
    pub on_violation: OnViolation,
}

impl Expectation {
    /// A warn-only expectation named `name` over boolean `constraint`.
    pub fn new(name: impl Into<String>, constraint: impl Into<String>) -> Self {
        Expectation {
            name: name.into(),
            constraint: constraint.into(),
            on_violation: OnViolation::Warn,
        }
    }

    /// Set the violation policy, returning `self` for chaining.
    pub fn on(mut self, policy: OnViolation) -> Self {
        self.on_violation = policy;
        self
    }
}