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.
//! **rel2graph, end to end, from a real invocation** — not the bench, not a test.
//!
//! One reviewed [`MappingSpec`] → the canonical [`GraphPlan`] → BOTH projections
//! (the PySpark job and the thund [`Pipeline`]), asserts they agree, then RUNS the
//! thund pipeline on the native Arrow/DataFusion engine over seeded source rows so
//! the graph-sink `UNWIND … MERGE` construction actually executes. With no live
//! FalkorDB target the mutation is constructed + recorded (honest
//! construction-only, `written = false`) — call
//! [`NativeBackend::with_falkordb_sink`] to load it for real.
//!
//! This is the reachable-from-a-real-invocation proof for the graph-sink +
//! two-projections-agree work that otherwise only the bench/tests touch.
//!
//! Run: `cargo run -p knut-thund --example rel2graph_native --features native,rel2graph`

use std::sync::Arc;

use knut_bifrost::codegen::CodegenOptions;
use knut_bifrost::{GraphPlan, MappingSpec};
use knut_thund::backend::{ExecBackend, RunHandle, native::NativeBackend};
use knut_thund::rel2graph::{to_pipeline_ron, to_thund_pipeline};

use datafusion::arrow::array::{ArrayRef, Int64Array, RecordBatch, StringArray};
use datafusion::arrow::datatypes::{DataType, Field, Schema};

/// A small but real HR mapping: employees + departments, with a normal FK edge
/// (`WORKS_IN`) and the self-referential `REPORTS_TO` (reads `manager_id`, MERGEs
/// on `employee_id` — the a7208dc case that must NOT collapse to a self-loop).
const SPEC: &str = r#"
version: 1
tables:
  - table: employees
    pipeline:
      - vertex: { label: Employee, id_column: employee_id, properties: [name] }
      - fk_edges:
          type: WORKS_IN
          src: { label: Employee,   column: employee_id }
          dst: { label: Department, column: dept_id }
      - fk_edges:
          type: REPORTS_TO
          src: { label: Employee, column: employee_id }
          dst: { label: Employee, column: manager_id }
  - table: departments
    pipeline:
      - vertex: { label: Department, id_column: dept_id, properties: [name] }
"#;

fn employees() -> RecordBatch {
    let schema = Arc::new(Schema::new(vec![
        Field::new("employee_id", DataType::Int64, false),
        Field::new("name", DataType::Utf8, false),
        Field::new("manager_id", DataType::Int64, true),
        Field::new("dept_id", DataType::Int64, false),
    ]));
    let ids: ArrayRef = Arc::new(Int64Array::from(vec![1, 2, 3, 4]));
    let names: ArrayRef = Arc::new(StringArray::from(vec!["Ada", "Bo", "Cy", "Di"]));
    // Employee 1 is the org root (null manager); 2,3 report to 1; 4 reports to 2.
    let mgr: ArrayRef = Arc::new(Int64Array::from(vec![None, Some(1), Some(1), Some(2)]));
    let dept: ArrayRef = Arc::new(Int64Array::from(vec![10, 10, 20, 20]));
    RecordBatch::try_new(schema, vec![ids, names, mgr, dept]).unwrap()
}

fn departments() -> RecordBatch {
    let schema = Arc::new(Schema::new(vec![
        Field::new("dept_id", DataType::Int64, false),
        Field::new("name", DataType::Utf8, false),
    ]));
    let ids: ArrayRef = Arc::new(Int64Array::from(vec![10, 20]));
    let names: ArrayRef = Arc::new(StringArray::from(vec!["Eng", "Ops"]));
    RecordBatch::try_new(schema, vec![ids, names]).unwrap()
}

fn main() {
    // 1. One reviewed mapping → the canonical GraphPlan.
    let spec = MappingSpec::from_yaml(SPEC).expect("mapping spec parses");
    let plan = GraphPlan::from_mapping(&spec, &CodegenOptions::default()).expect("plan");
    println!("── canonical plan contract ──\n{}\n", plan.contract());

    // 2. BOTH projections, from the SAME plan.
    let pyspark = plan.to_pyspark();
    let pipeline = to_thund_pipeline(&plan).expect("thund projection");
    let ron = to_pipeline_ron(&plan).expect("pipeline.ron");

    // 3. Assert the two projections agree (the invariant that otherwise only
    //    `knut_thund::rel2graph::tests::two_projections_agree` pins).
    let node_ds = pipeline
        .datasets
        .iter()
        .filter(|d| d.properties.contains_key("knut.node.label"))
        .count();
    let edge_ds = pipeline
        .datasets
        .iter()
        .filter(|d| d.properties.contains_key("knut.edge.rel"))
        .count();
    assert_eq!(
        node_ds,
        plan.nodes.len(),
        "thund node datasets == plan nodes"
    );
    assert_eq!(
        edge_ds,
        plan.edges.len(),
        "thund edge datasets == plan edges"
    );
    for n in &plan.nodes {
        assert!(
            pyspark.contains(&n.label),
            "pyspark carries node {}",
            n.label
        );
    }
    for e in &plan.edges {
        assert!(pyspark.contains(&e.rel), "pyspark carries edge {}", e.rel);
    }
    println!(
        "── two projections agree ── {} node(s), {} edge(s); PySpark {} lines, thund pipeline `{}`\n",
        node_ds,
        edge_ds,
        pyspark.lines().count(),
        pipeline.name,
    );
    println!("── thund pipeline.ron (first 12 lines) ──");
    for l in ron.lines().take(12) {
        println!("  {l}");
    }
    println!();

    // 4. RUN the thund pipeline on the native engine over seeded source rows, so
    //    the graph-sink UNWIND…MERGE construction actually executes.
    let mut run = NativeBackend::new()
        .with_input("employees", vec![employees()])
        .with_input("departments", vec![departments()])
        .run(&pipeline)
        .expect("native run of the rel2graph pipeline");

    println!("── native run events ──");
    for ev in run.poll_events().unwrap_or_default() {
        println!("  · {}", ev.message);
    }

    // 5. Report each graph-sink dataset the run resolved (construction-only when
    //    no FalkorDB target is set — never a silent fake).
    println!("\n── graph sinks constructed ──");
    for ds in &pipeline.datasets {
        if let Some(gs) = run.graph_sink(&ds.name) {
            println!(
                "  {:<40} {:?} `{}`  {} row(s)  written={}",
                ds.name, gs.kind, gs.element, gs.rows, gs.written
            );
        }
    }
    println!(
        "\nDone. Set a live target with `.with_falkordb_sink(url, graph)` to MERGE these \
         into FalkorDB (needs a running FalkorDB)."
    );
}