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
//! # knut-thund (Þund) — a Rust-native, Arrow-centric streaming dataflow engine
//!
//! Þund is knut's **engine**: it authors, ingests, and runs **batch +
//! streaming** data pipelines, with a **pluggable execution backend** —
//! either a native Arrow/DataFusion engine (Arroyo-style, RecordBatch streams
//! end to end) or a lowering to **Spark Declarative Pipelines (SDP)** over
//! Spark Connect (Spark 4.1+/4.2). It is the "Airflow killer" half of the
//! pipelines story; `knut-pipelines` remains the *viewer*, which Þund reuses
//! (and which provides the Spark Connect client this engine lowers onto).
//!
//! ## The three layers
//!
//! | Layer | What | Module |
//! |---|---|---|
//! | **Authoring** | typed Rust builder + thin RON DSL + facett graph editor | [`authoring`] |
//! | **IR** | one Arrow-typed dataflow graph: datasets/flows/expectations, batch + streaming | [`ir`] |
//! | **Backends** | pluggable `ExecBackend`: native Arrow/DataFusion or Spark-Connect-SDP | [`backend`] |
//!
//! All three authoring surfaces produce the same [`ir::Pipeline`]; every
//! [`backend::ExecBackend`] consumes it. The IR is the contract; front-ends
//! and back-ends are interchangeable around it. See the design doc
//! `knut/.nornir/thund.md`.
//!
//! ## Why it kills Airflow
//!
//! Airflow's DAG encodes *task order* and is blind to the data it touches, so
//! it can only re-execute, never *reconcile*; data flows through the metadata
//! DB (XCom); the scheduler re-parses every DAG every loop; it is Python-only.
//! Þund takes Dagster's **asset model** (declare the artifact, infer the DAG
//! from declared inputs/outputs → recompute only what's stale), is **Arrow
//! end-to-end** like Dremio (one columnar format storage→exec→wire, Flight
//! transport, no XCom), runs **streaming first-class** (unlike any
//! orchestrator), and offers **CLI/Rust/DSL/visual parity** over one IR.
//!
//! ## Example
//!
//! ```
//! use knut_thund::ir::{Pipeline, Dataset, Flow, OutputType};
//!
//! let p = Pipeline::new("daily_sales")
//!     .with_dataset(Dataset::new("clean_orders", OutputType::MaterializedView))
//!     .with_flow(
//!         Flow::batch("f_clean", "clean_orders", ["raw_orders"])
//!             .with_query("SELECT * FROM raw_orders WHERE amount > 0"),
//!     );
//! p.validate().unwrap();
//! assert_eq!(p.topo_order().unwrap(), vec!["clean_orders".to_string()]);
//!
//! // Lower to Spark Declarative Pipelines (works without a live Spark):
//! let sdp = p.to_sdp();
//! assert_eq!(sdp.datasets.len(), 1);
//! ```

#![deny(rustdoc::broken_intra_doc_links)]

pub mod authoring;
pub mod backend;
pub mod error;
pub mod ir;

pub use backend::{Capabilities, ExecBackend, RunEvent, RunHandle, RunPhase};
pub use error::{Result, ThundError};
pub use ir::Pipeline;

/// **Introspection / emit marker** — record one functional-status row for the
/// nornir test matrix. Wraps `nornir_testmatrix::functional_status` behind the
/// `testmatrix` feature (a compiled-out `#[inline]` no-op otherwise, with no
/// nornir dep).
///
/// Every surface calls this with honest health: the IR [`Pipeline::validate`]
/// gate, the IR→SDP lowering [`Pipeline::to_sdp`], and the RON DSL round-trip
/// ([`authoring::dsl`], feature `dsl`). The **execution** backends report their
/// real outcome too (see [`backend::native`] / [`backend::spark`]): with the
/// `native` feature the DataFusion backend's `execute` row is green when the
/// batch dataflow actually ran (red on a genuine execution/expectation
/// failure); with the `spark` feature the Spark-Connect backend's `live_run`
/// row is green when the run streamed to completion (red on a transport/run
/// error), and its IR→SDP `lower_to_sdp` sub-step, which always works, is
/// green. Built *without* those features, each backend's `run` records a red
/// row noting the feature (and hence the engine) is not compiled in.
#[inline]
pub fn functional_status(component: &str, check: &str, ok: bool, detail: &str) {
    #[cfg(feature = "testmatrix")]
    nornir_testmatrix::functional_status(component, check, ok, detail);
    #[cfg(not(feature = "testmatrix"))]
    {
        let _ = (component, check, ok, detail);
    }
}

#[cfg(test)]
mod tests {
    use super::backend::{native::NativeBackend, spark::SparkBackend, ExecBackend};
    use super::ir::*;
    use super::ThundError;

    /// Build a small streaming pipeline: a Kafka source flow feeding a
    /// streaming table, with a watermark and a data-quality expectation.
    fn sample_streaming() -> Pipeline {
        Pipeline::new("clicks")
            .with_dataset(
                Dataset::new("click_counts", OutputType::Table)
                    .incremental()
                    .with_schema(
                        DatasetSchema::new()
                            .field("window_end", "Timestamp(Microsecond, None)", false)
                            .field("n", "Int64", false),
                    ),
            )
            .with_flow(
                Flow::streaming(
                    "f_clicks",
                    "click_counts",
                    SourceSpec::Kafka {
                        bootstrap: "localhost:9092".into(),
                        topic: "clicks".into(),
                        format: "json".into(),
                    },
                )
                .with_watermark(Watermark {
                    event_time_column: "ts".into(),
                    allowed_lateness_ms: 5_000,
                    idle_timeout_ms: 30_000,
                })
                .with_query("SELECT window_end, count(*) n FROM clicks GROUP BY TUMBLE(ts, '10s')")
                .expect(Expectation::new("n_positive", "n > 0").on(OnViolation::Drop)),
            )
    }

    #[test]
    fn validate_and_topo_order() {
        let p = Pipeline::new("p")
            .with_dataset(Dataset::new("a", OutputType::Table))
            .with_dataset(Dataset::new("b", OutputType::MaterializedView))
            .with_flow(Flow::batch("fa", "a", Vec::<String>::new()))
            .with_flow(Flow::batch("fb", "b", ["a"]));
        p.validate().expect("valid");
        // a precedes b because fb reads a.
        assert_eq!(p.topo_order().unwrap(), vec!["a".to_string(), "b".to_string()]);
    }

    #[test]
    fn dangling_flow_is_rejected() {
        let p = Pipeline::new("p").with_flow(Flow::batch("f", "ghost", Vec::<String>::new()));
        let err = p.validate().unwrap_err();
        assert!(matches!(err, ThundError::DanglingFlow { .. }), "got {err:?}");
    }

    #[test]
    fn cycle_is_rejected() {
        let p = Pipeline::new("p")
            .with_dataset(Dataset::new("a", OutputType::Table))
            .with_dataset(Dataset::new("b", OutputType::Table))
            .with_flow(Flow::batch("fa", "a", ["b"]))
            .with_flow(Flow::batch("fb", "b", ["a"]));
        assert!(matches!(p.validate().unwrap_err(), ThundError::Cyclic));
        assert!(p.topo_order().is_none());
    }

    #[test]
    fn streaming_pipeline_detected_and_lowers_to_sdp() {
        let p = sample_streaming();
        p.validate().expect("valid");
        assert!(p.is_streaming(), "kafka flow makes it streaming");

        // The graph SHAPE lowers faithfully to SDP even though SDP can't
        // express the watermark/expectation knobs.
        let sdp = p.to_sdp();
        sdp.validate().expect("lowered SDP is valid");
        assert_eq!(sdp.datasets.len(), 1);
        assert_eq!(sdp.datasets[0].output_type, knut_pipelines::OutputType::Table);
        // Schema lowered to an SDP schema string.
        assert_eq!(
            sdp.datasets[0].schema.as_deref(),
            Some("window_end Timestamp(Microsecond, None), n Int64")
        );
        assert_eq!(sdp.flows.len(), 1);
        assert_eq!(sdp.flows[0].query.as_deref(), Some(
            "SELECT window_end, count(*) n FROM clicks GROUP BY TUMBLE(ts, '10s')"
        ));
    }

    #[test]
    fn spark_backend_rejects_expectations_via_capabilities() {
        // SparkBackend reports expectations:false, so `check` must refuse a
        // pipeline that carries one — BEFORE any network call. This is the
        // SanityCheckPlan-style up-front rejection.
        let p = sample_streaming(); // has an expectation
        let spark = SparkBackend::new("sc://localhost:15002");
        let err = spark.check(&p).unwrap_err();
        assert!(
            matches!(err, ThundError::Unsupported { ref what, .. } if what.contains("expectations")),
            "got {err:?}"
        );
    }

    #[test]
    fn native_backend_accepts_everything_in_the_ir() {
        // The native backend's capabilities cover the whole IR, so `check`
        // passes the same pipeline Spark refused.
        let p = sample_streaming();
        let native = NativeBackend::new();
        native.check(&p).expect("native accepts streaming + expectations + watermark");
        let caps = native.capabilities();
        assert!(caps.streaming && caps.event_time && caps.cdc && caps.expectations);
    }

    #[test]
    fn capabilities_differ_between_backends() {
        let n = NativeBackend::new().capabilities();
        let s = SparkBackend::new("sc://x:15002").capabilities();
        assert!(n.event_time && !s.event_time, "native has event-time, spark does not");
        assert_eq!(s.output_modes, vec!["append".to_string()]);
        assert!(n.output_modes.contains(&"complete".to_string()));
    }

    #[cfg(feature = "dsl")]
    #[test]
    fn dsl_roundtrips_through_ir() {
        let p = sample_streaming();
        let ron = crate::authoring::dsl::to_ron(&p).expect("serialize");
        let back = crate::authoring::dsl::from_ron(&ron).expect("parse");
        assert_eq!(p, back, "RON DSL is a lossless view of the IR");
    }
}