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.
//! Integration proof for the native backend's **on-disk bounded file source** —
//! the laptop-consumer entry point: point thund at a local `orders.csv` /
//! `orders.parquet`, author a plain **batch** flow whose SQL reads it, run on
//! the REAL DataFusion engine, and get an aggregated table back. No live infra,
//! no Kafka/Spark, no in-memory pre-load.
//!
//! This closes a coverage gap: before `NativeBackend::with_file_input`, every
//! native-exec test seeded data *in memory* (`with_input` / `with_stream_input`)
//! — the file-registration path (`register_named_file`, CSV/parquet) had zero
//! coverage and no runnable example. Here a real file on disk is the source.
//!
//! Gated on `--features native` (the DataFusion tree); a no-op otherwise.
#![cfg(feature = "native")]

use knut_thund::backend::{ExecBackend, native::NativeBackend};
use knut_thund::ir::{Dataset, Flow, OutputType, Pipeline};

use datafusion::arrow::array::{Array, Int64Array, RecordBatch, StringArray, StringViewArray};
use std::io::Write;

/// Read a `Utf8`/`Utf8View` string column by name into owned `String`s.
fn str_col(b: &RecordBatch, name: &str) -> Vec<String> {
    let col = b
        .column_by_name(name)
        .unwrap_or_else(|| panic!("column `{name}` present"))
        .as_any();
    if let Some(a) = col.downcast_ref::<StringArray>() {
        (0..a.len()).map(|i| a.value(i).to_string()).collect()
    } else if let Some(a) = col.downcast_ref::<StringViewArray>() {
        (0..a.len()).map(|i| a.value(i).to_string()).collect()
    } else {
        panic!("column `{name}` is neither Utf8 nor Utf8View");
    }
}

/// The pipeline under test: one materialised view fed by one batch flow that
/// reads the registered `orders` table and rolls it up per customer.
fn rollup_pipeline() -> Pipeline {
    Pipeline::new("orders_rollup")
        .with_dataset(Dataset::new("by_customer", OutputType::MaterializedView))
        .with_flow(
            Flow::batch("agg", "by_customer", ["orders"])
                .with_query("SELECT customer, SUM(amount) AS total FROM orders GROUP BY customer"),
        )
}

/// Collect the `by_customer` output as `(customer, total)` pairs, sorted by
/// customer so the assertion is deterministic (GROUP BY order is unspecified).
fn sorted_totals(batches: &[RecordBatch]) -> Vec<(String, i64)> {
    let mut out: Vec<(String, i64)> = Vec::new();
    for b in batches {
        let col = b.column(0).as_any();
        // CSV yields `Utf8` (StringArray); DataFusion reads parquet strings as
        // `Utf8View` (StringViewArray) by default — accept either.
        let cust: Vec<String> = if let Some(a) = col.downcast_ref::<StringArray>() {
            (0..a.len()).map(|i| a.value(i).to_string()).collect()
        } else if let Some(a) = col.downcast_ref::<StringViewArray>() {
            (0..a.len()).map(|i| a.value(i).to_string()).collect()
        } else {
            panic!("customer column is neither Utf8 nor Utf8View");
        };
        let total = b
            .column(1)
            .as_any()
            .downcast_ref::<Int64Array>()
            .expect("total is Int64");
        for (i, c) in cust.into_iter().enumerate() {
            out.push((c, total.value(i)));
        }
    }
    out.sort();
    out
}

/// A real CSV on disk → registered `orders` table → batch SQL → materialised
/// rollup. Asserts the exact aggregated rows (RED-when-broken).
#[test]
fn csv_file_source_batch_rollup() {
    let dir = tempfile::tempdir().expect("tempdir");
    let path = dir.path().join("orders.csv");
    {
        let mut f = std::fs::File::create(&path).expect("create csv");
        // customer,amount — two customers, five rows.
        writeln!(f, "customer,amount").unwrap();
        writeln!(f, "alice,10").unwrap();
        writeln!(f, "bob,5").unwrap();
        writeln!(f, "alice,20").unwrap();
        writeln!(f, "bob,7").unwrap();
        writeln!(f, "alice,3").unwrap();
    }

    let run = NativeBackend::new()
        .with_file_input("orders", path.to_str().unwrap(), "csv")
        .run(&rollup_pipeline())
        .expect("native run over the on-disk CSV completes");

    let totals = sorted_totals(run.output("by_customer").expect("materialised output"));

    // alice: 10+20+3 = 33, bob: 5+7 = 12 — computed from the FILE, not memory.
    assert_eq!(
        totals,
        vec![("alice".to_string(), 33), ("bob".to_string(), 12)],
        "the batch flow scanned the local CSV and rolled it up per customer"
    );

    // Emit the file-source path's health into the nornir matrix (gated →
    // release strips it). GREEN only because a real file produced real rows.
    knut_thund::functional_status(
        "knut-thund/native_file_source",
        "csv_batch",
        totals == vec![("alice".to_string(), 33), ("bob".to_string(), 12)],
        "orders.csv -> registered `orders` -> batch rollup",
    );
}

/// The same rollup, but the source is a real **parquet** file on disk — proves
/// the shared `register_named_file` helper's parquet arm, not just CSV.
#[test]
fn parquet_file_source_batch_rollup() {
    use datafusion::arrow::array::{Int64Array as I64, StringArray as Utf8};
    use datafusion::arrow::datatypes::{DataType, Field, Schema};
    use datafusion::parquet::arrow::ArrowWriter;
    use std::sync::Arc;

    let dir = tempfile::tempdir().expect("tempdir");
    let path = dir.path().join("orders.parquet");

    // Write orders.parquet directly with the arrow parquet writer (datafusion
    // re-export → same arrow tree the engine reads with).
    let schema = Arc::new(Schema::new(vec![
        Field::new("customer", DataType::Utf8, false),
        Field::new("amount", DataType::Int64, false),
    ]));
    let batch = RecordBatch::try_new(
        schema.clone(),
        vec![
            Arc::new(Utf8::from(vec!["alice", "bob", "alice", "bob", "alice"])),
            Arc::new(I64::from(vec![10, 5, 20, 7, 3])),
        ],
    )
    .expect("build orders batch");
    {
        let file = std::fs::File::create(&path).expect("create parquet");
        let mut w = ArrowWriter::try_new(file, schema, None).expect("arrow writer");
        w.write(&batch).expect("write parquet batch");
        w.close().expect("close parquet");
    }

    let run = NativeBackend::new()
        .with_file_input("orders", path.to_str().unwrap(), "parquet")
        .run(&rollup_pipeline())
        .expect("native run over the on-disk parquet completes");

    let totals = sorted_totals(run.output("by_customer").expect("materialised output"));
    assert_eq!(
        totals,
        vec![("alice".to_string(), 33), ("bob".to_string(), 12)],
        "the batch flow scanned the local parquet and rolled it up per customer"
    );

    knut_thund::functional_status(
        "knut-thund/native_file_source",
        "parquet_batch",
        totals == vec![("alice".to_string(), 33), ("bob".to_string(), 12)],
        "orders.parquet -> registered `orders` -> batch rollup",
    );
}

// ---------------------------------------------------------------------------
// Partitioned file SOURCE — reading a Hive-partitioned directory back, with the
// partition column recovered from the `col=val/…` path segments. The read-side
// sibling of the partitioned file sink: this closes the partitioned round trip
// through thund's OWN source path (not a hand-rolled directory walk).
// ---------------------------------------------------------------------------

/// Collect a partitioned read-back as `(region, customer, total)`, sorted.
fn sorted_regional(batches: &[RecordBatch]) -> Vec<(String, String, i64)> {
    let mut out: Vec<(String, String, i64)> = Vec::new();
    for b in batches {
        let region = str_col(b, "region");
        let customer = str_col(b, "customer");
        let total = b
            .column_by_name("total")
            .expect("total column present")
            .as_any()
            .downcast_ref::<Int64Array>()
            .expect("total is Int64");
        for i in 0..b.num_rows() {
            out.push((region[i].clone(), customer[i].clone(), total.value(i)));
        }
    }
    out.sort();
    out
}

/// Write a regional `orders.csv` fixture (`region,customer,amount`).
fn write_regional_orders_csv(path: &std::path::Path) {
    let mut f = std::fs::File::create(path).expect("create orders.csv");
    writeln!(f, "region,customer,amount").unwrap();
    writeln!(f, "US,alice,10").unwrap();
    writeln!(f, "EU,bob,5").unwrap();
    writeln!(f, "US,alice,20").unwrap();
    writeln!(f, "EU,bob,7").unwrap();
    writeln!(f, "US,carol,3").unwrap();
}

/// Full partitioned round trip: write a Hive-partitioned parquet directory with
/// the partitioned SINK (`region=US/…`, `region=EU/…`), then read it back with
/// `with_partitioned_file_input("...", root, "parquet", ["region"])` and assert
/// the full `(region, customer, total)` rows — the `region` column reconstructed
/// from the directory names, through thund's own source path.
///
/// RED-when-broken: if `partition_cols` were NOT wired into the source
/// (`table_partition_cols`), the read-back schema would carry no `region`
/// column, so the `SELECT region, …` passthrough would fail to plan and the run
/// would error — the `.expect(...)` below would panic and the test fail.
#[test]
fn partitioned_parquet_source_round_trip() {
    let dir = tempfile::tempdir().expect("tempdir");
    let src = dir.path().join("orders.csv");
    write_regional_orders_csv(&src);

    // (1) Write the partitioned directory with the partitioned sink.
    let root = dir.path().join("warehouse").join("by_region");
    let write = NativeBackend::new()
        .with_file_input("orders", src.to_str().unwrap(), "csv")
        .with_file_output("by_customer", root.to_str().unwrap(), "parquet")
        .run(
            &Pipeline::new("orders_by_region")
                .with_dataset(
                    Dataset::new("by_customer", OutputType::Sink).with_partition_cols(["region"]),
                )
                .with_flow(Flow::batch("agg", "by_customer", ["orders"]).with_query(
                    "SELECT region, customer, SUM(amount) AS total \
                     FROM orders GROUP BY region, customer",
                )),
        )
        .expect("partitioned sink write completes");
    assert_eq!(
        write
            .sink("by_customer")
            .expect("sink report")
            .partition_cols,
        vec!["region".to_string()],
        "the sink partitioned by region"
    );
    assert!(root.is_dir(), "the sink produced a Hive directory root");

    // (2) Read the partitioned directory back through thund's OWN source path,
    //     recovering `region` from the `region=<val>` subdir names.
    let back = NativeBackend::new()
        .with_partitioned_file_input("t", root.to_str().unwrap(), "parquet", ["region"])
        .run(
            &Pipeline::new("read_partitioned")
                .with_dataset(Dataset::new("t_out", OutputType::MaterializedView))
                .with_flow(
                    Flow::batch("scan", "t_out", ["t"])
                        .with_query("SELECT region, customer, total FROM t"),
                ),
        )
        .expect("partitioned read-back completes (region recovered from the path)");

    let got = sorted_regional(back.output("t_out").expect("read-back output"));
    let expected = vec![
        ("EU".to_string(), "bob".to_string(), 12),
        ("US".to_string(), "alice".to_string(), 30),
        ("US".to_string(), "carol".to_string(), 3),
    ];
    assert_eq!(
        got, expected,
        "the partition column `region` is recovered from the path and the rows survive the round trip"
    );

    knut_thund::functional_status(
        "knut-thund/native_file_source",
        "partitioned_parquet",
        got == expected,
        "region=US/ , region=EU/ (Hive parquet dir) -> partitioned source -> region recovered from path",
    );
}

/// The partitioned source works for **CSV** Hive directories too (same
/// `table_partition_cols` seam, exercised via the csv writer). Proves the csv
/// arm and that a partitioned CSV round-trips just like parquet.
#[test]
fn partitioned_csv_source_round_trip() {
    let dir = tempfile::tempdir().expect("tempdir");
    let src = dir.path().join("orders.csv");
    write_regional_orders_csv(&src);

    let root = dir.path().join("csv_warehouse");
    NativeBackend::new()
        .with_file_input("orders", src.to_str().unwrap(), "csv")
        .with_file_output("by_customer", root.to_str().unwrap(), "csv")
        .run(
            &Pipeline::new("orders_by_region_csv")
                .with_dataset(
                    Dataset::new("by_customer", OutputType::Sink).with_partition_cols(["region"]),
                )
                .with_flow(Flow::batch("agg", "by_customer", ["orders"]).with_query(
                    "SELECT region, customer, SUM(amount) AS total \
                     FROM orders GROUP BY region, customer",
                )),
        )
        .expect("partitioned csv sink write completes");

    let back = NativeBackend::new()
        .with_partitioned_file_input("t", root.to_str().unwrap(), "csv", ["region"])
        .run(
            &Pipeline::new("read_partitioned_csv")
                .with_dataset(Dataset::new("t_out", OutputType::MaterializedView))
                .with_flow(
                    Flow::batch("scan", "t_out", ["t"])
                        .with_query("SELECT region, customer, total FROM t"),
                ),
        )
        .expect("partitioned csv read-back completes");

    let got = sorted_regional(back.output("t_out").expect("read-back output"));
    let expected = vec![
        ("EU".to_string(), "bob".to_string(), 12),
        ("US".to_string(), "alice".to_string(), 30),
        ("US".to_string(), "carol".to_string(), 3),
    ];
    assert_eq!(
        got, expected,
        "partitioned CSV round-trips with region from the path"
    );

    knut_thund::functional_status(
        "knut-thund/native_file_source",
        "partitioned_csv",
        got == expected,
        "region=US/ , region=EU/ (Hive csv dir) -> partitioned source -> region recovered from path",
    );
}