knut-thund 0.2.0

Þ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.
//! The SQL authoring surface — parse, lower, and round-trip.
//!
//! GATED: this whole file is `#![cfg(feature = "sql")]`. Under a default
//! `cargo test -p knut-thund` it compiles to ZERO tests and reports a green
//! "ok" — the same silent-dark trap the `native`-gated files carry. CI
//! therefore runs an explicit `--features sql,dsl` lane; see
//! `.github/workflows/ci.yml`.

#![cfg(feature = "sql")]

use knut_thund::authoring::sql;
use knut_thund::ir::{FlowKind, OnViolation, OutputType};

/// A medallion slice exercising every clause the grammar covers.
const SCRIPT: &str = r#"
CREATE MATERIALIZED VIEW bronze_events (
  event_id   BIGINT NOT NULL,
  event_ts   TIMESTAMP,
  amount     DOUBLE,
  CONSTRAINT amount_non_negative EXPECT (amount >= 0) ON VIOLATION DROP ROW
)
PARTITIONED BY (months(event_ts))
COMMENT 'raw events, month-partitioned'
TBLPROPERTIES ('owner' = 'data-eng')
AS SELECT event_id, event_ts, amount FROM landing_events;

CREATE MATERIALIZED VIEW silver_events AS
SELECT event_id, event_ts, amount FROM bronze_events WHERE amount > 0;

CREATE STREAMING TABLE gold_live AS
SELECT event_id, amount FROM STREAM(silver_events);

CREATE FLOW gold_backfill AS
INSERT INTO gold_live BY NAME
SELECT event_id, amount FROM STREAM(bronze_events);
"#;

#[test]
fn parses_every_clause_into_the_ir() {
    let p = sql::from_sql("medallion", SCRIPT).unwrap();

    assert_eq!(p.datasets.len(), 3);
    assert_eq!(p.flows.len(), 4);

    let bronze = p.dataset("bronze_events").unwrap();
    assert_eq!(bronze.output_type, OutputType::MaterializedView);
    assert_eq!(
        bronze.comment.as_deref(),
        Some("raw events, month-partitioned")
    );
    assert_eq!(bronze.partition_cols, vec!["months(event_ts)".to_string()]);
    assert_eq!(bronze.schema.fields.len(), 3);
    assert_eq!(bronze.schema.fields[0].name, "event_id");
    assert_eq!(bronze.schema.fields[0].arrow_type, "Int64");
    assert!(!bronze.schema.fields[0].nullable, "BIGINT NOT NULL");
    assert!(bronze.schema.fields[1].nullable, "TIMESTAMP is nullable");
    assert_eq!(
        bronze.schema.fields[1].arrow_type,
        "Timestamp(Microsecond, None)"
    );

    // The expectation rides on the flow that computes the dataset.
    let bf = p
        .flows
        .iter()
        .find(|f| f.target == "bronze_events")
        .unwrap();
    assert_eq!(bf.expectations.len(), 1);
    assert_eq!(bf.expectations[0].name, "amount_non_negative");
    assert_eq!(bf.expectations[0].on_violation, OnViolation::Drop);

    assert_eq!(
        p.dataset("gold_live").unwrap().output_type,
        OutputType::Table
    );
}

#[test]
fn stream_marks_the_flow_not_the_query() {
    let p = sql::from_sql("medallion", SCRIPT).unwrap();
    let live = p.flows.iter().find(|f| f.target == "gold_live").unwrap();

    assert!(matches!(live.kind, FlowKind::Streaming { .. }));
    assert_eq!(live.reads, vec!["silver_events".to_string()]);
    // STREAM(...) is erased from the SQL — it was never a property of the
    // SELECT, so the query handed to a backend is plain, portable SQL.
    let q = live.query.as_ref().unwrap();
    assert!(
        !q.to_uppercase().contains("STREAM"),
        "query still says STREAM: {q}"
    );
    assert!(q.contains("silver_events"));
}

#[test]
fn dag_is_inferred_from_reads() {
    let p = sql::from_sql("medallion", SCRIPT).unwrap();
    let order = p.topo_order().expect("acyclic");
    let pos = |n: &str| order.iter().position(|x| x == n).unwrap();
    assert!(pos("bronze_events") < pos("silver_events"));
    assert!(pos("silver_events") < pos("gold_live"));
}

#[test]
fn standalone_create_flow_appends_to_an_existing_target() {
    let p = sql::from_sql("medallion", SCRIPT).unwrap();
    let backfill = p.flows.iter().find(|f| f.name == "gold_backfill").unwrap();
    assert_eq!(backfill.target, "gold_live");
    assert_eq!(backfill.reads, vec!["bronze_events".to_string()]);
    // Two flows, one target — the multi-flow streaming-table shape.
    assert_eq!(
        p.flows.iter().filter(|f| f.target == "gold_live").count(),
        2
    );
}

#[test]
fn sql_round_trips_through_itself() {
    let p1 = sql::from_sql("medallion", SCRIPT).unwrap();
    let rendered = sql::to_sql(&p1).unwrap();
    let p2 = sql::from_sql("medallion", &rendered)
        .unwrap_or_else(|e| panic!("re-parse failed: {e}\n--- rendered ---\n{rendered}"));

    assert_eq!(p1.datasets.len(), p2.datasets.len());
    for d1 in &p1.datasets {
        let d2 = p2
            .dataset(&d1.name)
            .expect("dataset survives the round trip");
        assert_eq!(d1.output_type, d2.output_type, "{}", d1.name);
        assert_eq!(d1.partition_cols, d2.partition_cols, "{}", d1.name);
        assert_eq!(d1.comment, d2.comment, "{}", d1.name);
        assert_eq!(d1.schema, d2.schema, "{}", d1.name);
        assert_eq!(d1.properties, d2.properties, "{}", d1.name);
    }
}

/// The claim the whole "one IR, many surfaces" design rests on: SQL and RON are
/// two renderings of the *same* value, so a pipeline can cross between them
/// without drift.
#[cfg(feature = "dsl")]
#[test]
fn sql_to_ron_to_sql_is_the_same_pipeline() {
    use knut_thund::authoring::dsl;

    let p1 = sql::from_sql("medallion", SCRIPT).unwrap();
    let ron = dsl::to_ron(&p1).unwrap();
    let p2 = dsl::from_ron(&ron).unwrap();

    assert_eq!(p1, p2, "RON round-trip must be lossless");
    assert_eq!(sql::to_sql(&p1).unwrap(), sql::to_sql(&p2).unwrap());
}

#[test]
fn partition_transforms_are_preserved_for_the_backend_to_judge() {
    // Spark's SDP `PartitionHelper` accepts only identity transforms, so the
    // Spark lowering must screen these. The IR's job is to carry them intact.
    let p = sql::from_sql(
        "p",
        "CREATE MATERIALIZED VIEW t
         PARTITIONED BY (months(event_ts), bucket(16, user_id), truncate(10, name))
         AS SELECT * FROM src;",
    )
    .unwrap();
    assert_eq!(
        p.dataset("t").unwrap().partition_cols,
        vec![
            "months(event_ts)".to_string(),
            "bucket(16, user_id)".to_string(),
            "truncate(10, name)".to_string(),
        ]
    );
}

#[test]
fn tblproperties_are_captured_and_reach_sdp() {
    let p = sql::from_sql("medallion", SCRIPT).unwrap();
    let bronze = p.dataset("bronze_events").unwrap();
    assert_eq!(
        bronze.properties.get("owner").map(String::as_str),
        Some("data-eng")
    );

    // The whole point of the field: it must survive the lowering, not just the
    // parse. `to_sdp` feeds `DefineOutput.TableDetails.table_properties`.
    let sdp = p.to_sdp();
    let lowered = sdp
        .datasets
        .iter()
        .find(|d| d.name == "bronze_events")
        .unwrap();
    assert_eq!(
        lowered.properties.get("owner").map(String::as_str),
        Some("data-eng")
    );
}

#[test]
fn tblproperties_accept_dotted_bare_and_empty_forms() {
    let p = sql::from_sql(
        "p",
        "CREATE MATERIALIZED VIEW a
           TBLPROPERTIES ('delta.appendOnly' = 'true', bare = 42, 'q' = 'it''s')
           AS SELECT 1 AS x;
         CREATE MATERIALIZED VIEW b TBLPROPERTIES () AS SELECT 1 AS x;",
    )
    .unwrap();

    let a = p.dataset("a").unwrap();
    // A dotted key must not be truncated at the first `.`.
    assert_eq!(
        a.properties.get("delta.appendOnly").map(String::as_str),
        Some("true")
    );
    assert_eq!(a.properties.get("bare").map(String::as_str), Some("42"));
    assert_eq!(a.properties.get("q").map(String::as_str), Some("it's"));
    assert!(p.dataset("b").unwrap().properties.is_empty());
}

#[test]
fn tblproperties_round_trip_through_sql() {
    let p1 = sql::from_sql("medallion", SCRIPT).unwrap();
    let p2 = sql::from_sql("medallion", &sql::to_sql(&p1).unwrap()).unwrap();
    assert_eq!(
        p1.dataset("bronze_events").unwrap().properties,
        p2.dataset("bronze_events").unwrap().properties
    );
}

#[test]
fn a_cyclic_script_is_rejected() {
    let err = sql::from_sql(
        "p",
        "CREATE MATERIALIZED VIEW a AS SELECT * FROM b;
         CREATE MATERIALIZED VIEW b AS SELECT * FROM a;",
    )
    .unwrap_err();
    assert!(
        matches!(err, knut_thund::ThundError::Cyclic),
        "expected a cycle error, got {err}"
    );
}

#[test]
fn non_definition_sql_is_rejected_with_a_useful_message() {
    let err = sql::from_sql("p", "SELECT 1;").unwrap_err();
    let msg = err.to_string();
    assert!(msg.contains("CREATE"), "unhelpful message: {msg}");
}

#[test]
fn expect_syntax_spark_rejects_is_accepted_here() {
    // Verified against Spark 4.1.3: `CONSTRAINT ... EXPECT` is a Databricks-DLT
    // extension and fails there with PARSE_SYNTAX_ERROR near 'EXPECT'. The IR
    // has always modelled it, so the native backend can honour it.
    let p = sql::from_sql(
        "p",
        "CREATE MATERIALIZED VIEW t (
           x INT,
           CONSTRAINT x_pos EXPECT (x > 0) ON VIOLATION FAIL UPDATE
         ) AS SELECT x FROM src;",
    )
    .unwrap();
    let f = p.flows.iter().find(|f| f.target == "t").unwrap();
    assert_eq!(f.expectations[0].on_violation, OnViolation::Fail);
}