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
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
//! **rel2graph → thund**: project a bifröst [`GraphPlan`] into a real thund
//! [`Pipeline`] that any [`ExecBackend`] runs.
//!
//! This is the *second* projection of the two-projections invariant (the first,
//! the PySpark job, lives in `knut_bifrost::codegen`). bifröst maps a relational
//! schema to a graph **once** — the canonical [`GraphPlan`] (labels, merge keys,
//! edge directions, props, batch size). Here that ONE plan becomes a thund
//! [`Pipeline`]: one [`Dataset`] per node/edge (an [`OutputType::Sink`] carrying
//! the graph target in its `knut.*` [`Dataset::properties`]) fed by a projection
//! [`Flow`] over the source table. The resulting pipeline is a normal thund IR —
//! it `validate()`s, round-trips through the RON DSL, and is handed to
//! [`SparkBackend`](crate::backend::spark) (→ SDP) or
//! [`NativeBackend`](crate::backend::native) (→ DataFusion). "bifrost mapping →
//! thund → Spark|Native" is one path; thund picks the tech.
//!
//! **What is a projection vs a worklist item.** The graph *shape* (which nodes,
//! which merge keys, which directed edges, which columns, in canonical order)
//! projects fully and losslessly, and that shape is exactly what must agree with
//! the PySpark job — pinned by [`tests::two_projections_agree`]. The graph
//! *sink write itself* (a backend that turns an `OutputType::Sink` +
//! `knut.node.*`/`knut.edge.*` properties into `UNWIND … MERGE`) is a backend
//! capability not yet built — see the design doc's worklist. Until then the
//! Pipeline is the runnable, reviewable, diffable **description**; the PySpark
//! job is the executable Spark path.
//!
//! [`GraphPlan`]: knut_bifrost::GraphPlan
//! [`ExecBackend`]: crate::backend::ExecBackend

use std::collections::BTreeMap;

use knut_bifrost::{GraphPlan, PlanMode};

use crate::error::Result;
use crate::ir::{Dataset, Flow, Materialization, OutputType, Pipeline};

/// The property key namespace the graph target is carried under on each
/// [`Dataset`] — read by a (future) graph-sink `ExecBackend`, and by the
/// two-projections-agree test.
pub mod prop {
    /// Node label.
    pub const NODE_LABEL: &str = "knut.node.label";
    /// Node merge-key column.
    pub const NODE_KEY: &str = "knut.node.key";
    /// Node property columns (comma-joined, ordinal).
    pub const NODE_PROPS: &str = "knut.node.props";
    /// Source table the node rows come from.
    pub const NODE_TABLE: &str = "knut.node.table";
    /// Edge relationship type.
    pub const EDGE_REL: &str = "knut.edge.rel";
    /// Edge tail (child) label.
    pub const EDGE_FROM: &str = "knut.edge.from";
    /// Edge head (parent) label.
    pub const EDGE_TO: &str = "knut.edge.to";
    /// The tail node's **merge-key property** — the property the `from` node is
    /// MERGEd by (its `id_column`). NOT necessarily the column the value is read
    /// from; see [`EDGE_FROM_COL`].
    pub const EDGE_FROM_KEY: &str = "knut.edge.from_key";
    /// The head node's **merge-key property** — the property the `to` node is
    /// MERGEd by (its `id_column`). Differs from [`EDGE_TO_COL`] for a self-ref /
    /// renamed FK (MERGE on `employee_id`, value read from `manager_id`).
    pub const EDGE_TO_KEY: &str = "knut.edge.to_key";
    /// The column in the edge table the **tail** endpoint VALUE is read from
    /// (`src.column`). The projection flow projects this column (renamed to
    /// `src` by the graph-sink backend). Usually equals [`EDGE_FROM_KEY`].
    pub const EDGE_FROM_COL: &str = "knut.edge.from_col";
    /// The column in the edge table the **head** endpoint VALUE is read from
    /// (the FK column, `dst.column`). The projection flow projects THIS column,
    /// not [`EDGE_TO_KEY`] — the two differ for a self-ref / renamed FK.
    pub const EDGE_TO_COL: &str = "knut.edge.to_col";
    /// Source table the edge rows come from.
    pub const EDGE_TABLE: &str = "knut.edge.table";
    /// Rows per `UNWIND … MERGE` batch (mirrors the PySpark `BATCH_ROWS`).
    pub const BATCH_ROWS: &str = "knut.graph.batch_rows";
    /// The resolution tier (`T1`…`T4`) the element was resolved at.
    pub const TIER: &str = "knut.graph.tier";
    /// The graph sink token (`falkordb` / `graphar`).
    pub const SINK: &str = "knut.graph.sink";
    /// The source plane token (`postgres` / `iceberg`).
    pub const SOURCE: &str = "knut.graph.source";
}

/// **Project a [`GraphPlan`] into a thund [`Pipeline`].**
///
/// One [`Dataset`] per node (`node__<Label>`) and per edge
/// (`edge__<From>__<Rel>__<To>`), each an [`OutputType::Sink`] whose
/// `knut.*` properties carry the graph target, fed by a projection [`Flow`]
/// reading the source table's needed columns. Nodes precede edges (declaration
/// order follows the plan's canonical order). Streaming plans mark datasets
/// [`Materialization::Incremental`]; batch plans [`Materialization::Full`].
///
/// The returned pipeline is `validate()`d. T4 proposals are **not** projected —
/// they are never applied (they remain on [`GraphPlan::proposals`] for a
/// reviewer, and are emitted commented-out only by the PySpark projection).
pub fn to_thund_pipeline(plan: &GraphPlan) -> Result<Pipeline> {
    let mut p = Pipeline::new("rel2graph");
    p.storage = Some(plan.checkpoint.clone());

    let mat = match plan.mode {
        PlanMode::Stream => Materialization::Incremental,
        PlanMode::Batch => Materialization::Full,
    };
    let sink_fmt = plan.sink.token();
    let batch = plan.batch_rows.to_string();

    for n in &plan.nodes {
        let ds_name = node_dataset(&n.label);
        let mut props: BTreeMap<String, String> = BTreeMap::new();
        props.insert(prop::NODE_LABEL.into(), n.label.clone());
        props.insert(prop::NODE_KEY.into(), n.key.clone());
        props.insert(prop::NODE_PROPS.into(), n.props.join(","));
        props.insert(prop::NODE_TABLE.into(), n.table.clone());
        props.insert(prop::BATCH_ROWS.into(), batch.clone());
        props.insert(prop::TIER.into(), n.tier.token().into());
        props.insert(prop::SINK.into(), sink_fmt.into());
        props.insert(prop::SOURCE.into(), plan.source.token().into());

        let ds = Dataset::new(ds_name.clone(), OutputType::Sink)
            .with_format(format!("{sink_fmt}-node"))
            .with_properties(props);
        let ds = with_materialization(ds, mat);
        p = p.with_dataset(ds);

        // Projection flow: read the source table, prune to key + props.
        let mut cols = vec![n.key.clone()];
        cols.extend(n.props.iter().cloned());
        p = p.with_flow(
            Flow::batch(format!("flow_{ds_name}"), ds_name, [n.table.clone()])
                .with_projection(cols),
        );
    }

    for e in &plan.edges {
        let ds_name = edge_dataset(&e.from_label, &e.rel, &e.to_label);
        let mut props: BTreeMap<String, String> = BTreeMap::new();
        props.insert(prop::EDGE_REL.into(), e.rel.clone());
        props.insert(prop::EDGE_FROM.into(), e.from_label.clone());
        props.insert(prop::EDGE_TO.into(), e.to_label.clone());
        props.insert(prop::EDGE_FROM_KEY.into(), e.from_key.clone());
        props.insert(prop::EDGE_TO_KEY.into(), e.to_key.clone());
        props.insert(prop::EDGE_FROM_COL.into(), e.from_col.clone());
        props.insert(prop::EDGE_TO_COL.into(), e.to_col.clone());
        props.insert(prop::EDGE_TABLE.into(), e.table.clone());
        props.insert(prop::BATCH_ROWS.into(), batch.clone());
        props.insert(prop::TIER.into(), e.tier.token().into());
        props.insert(prop::SINK.into(), sink_fmt.into());
        props.insert(prop::SOURCE.into(), plan.source.token().into());

        let ds = Dataset::new(ds_name.clone(), OutputType::Sink)
            .with_format(format!("{sink_fmt}-edge"))
            .with_properties(props);
        let ds = with_materialization(ds, mat);
        p = p.with_dataset(ds);

        // Projection flow: read the edge table, prune to the two endpoint
        // VALUE columns (the FK columns `from_col`/`to_col`, NOT the merge-key
        // property names) + edge props. For a self-ref / renamed FK the value
        // lives in a different column than the node key (`manager_id` read, node
        // MERGEd on `employee_id`); projecting the merge key here would read the
        // wrong column and collapse the edge to a self-loop.
        let mut cols = vec![e.from_col.clone(), e.to_col.clone()];
        cols.extend(e.props.iter().cloned());
        p = p.with_flow(
            Flow::batch(format!("flow_{ds_name}"), ds_name, [e.table.clone()])
                .with_projection(cols),
        );
    }

    p.validate()?;
    Ok(p)
}

/// Project a [`GraphPlan`] into a thund [`Pipeline`] rendered as **RON** — the
/// `pipeline.ron` artifact, a first-class thund DSL spec (requires the `dsl`
/// feature, which `rel2graph` enables). Round-trips through
/// [`crate::authoring::dsl::from_ron`].
pub fn to_pipeline_ron(plan: &GraphPlan) -> Result<String> {
    let p = to_thund_pipeline(plan)?;
    crate::authoring::dsl::to_ron(&p)
}

// ── helpers ─────────────────────────────────────────────────────────────────

fn with_materialization(ds: Dataset, mat: Materialization) -> Dataset {
    match mat {
        Materialization::Incremental => ds.incremental(),
        Materialization::Full => ds, // Full is the default
    }
}

/// A stable, identifier-safe dataset name for a node label.
fn node_dataset(label: &str) -> String {
    format!("node__{}", sanitize(label))
}

/// A stable, identifier-safe dataset name for a directed edge.
fn edge_dataset(from: &str, rel: &str, to: &str) -> String {
    format!(
        "edge__{}__{}__{}",
        sanitize(from),
        sanitize(rel),
        sanitize(to)
    )
}

fn sanitize(s: &str) -> String {
    s.chars()
        .map(|c| if c.is_ascii_alphanumeric() { c } else { '_' })
        .collect()
}

// ═══════════════════════════════════════════════════════════════════════════
// Tests — THE two-projections-must-agree invariant lives here, where BOTH
// projections are visible: the PySpark job (bifröst) and the thund Pipeline.
// ═══════════════════════════════════════════════════════════════════════════

#[cfg(test)]
mod tests {
    use super::*;
    use knut_bifrost::MappingSpec;
    use knut_bifrost::codegen::CodegenOptions;

    fn sample_plan(opts: &CodegenOptions) -> GraphPlan {
        // Declaration order is deliberately unsorted so the canonical ordering
        // is exercised on both projections identically.
        let yaml = r#"
version: 1
tables:
  - table: lake.social.users
    pipeline:
      - vertex: { label: User, id_column: user_id, properties: [name, age] }
      - fk_edges:
          type: WORKS_AT
          src: { label: User,    column: user_id }
          dst: { label: Company, column: company_id }
          properties: [since]
  - table: lake.social.companies
    pipeline:
      - vertex: { label: Company, id_column: company_id, properties: [name] }
"#;
        let spec = MappingSpec::from_yaml(yaml).unwrap();
        GraphPlan::from_mapping(&spec, opts).unwrap()
    }

    /// (label, key) node facts parsed back out of the PySpark job's startup
    /// index statements — the Spark projection's ground truth.
    fn nodes_from_pyspark(py: &str) -> Vec<(String, String)> {
        let mut out = Vec::new();
        for line in py.lines() {
            // ... CREATE INDEX FOR (n:`Label`) ON (n.`key`) ...
            if let Some(rest) = line.split_once("CREATE INDEX FOR (n:`") {
                let (label, tail) = rest.1.split_once("`) ON (n.`").unwrap();
                let key = tail.split_once("`)").unwrap().0;
                out.push((label.to_string(), key.to_string()));
            }
        }
        out.sort();
        out
    }

    /// (from, rel, to) edge facts parsed back out of the PySpark job's
    /// `_write_edges("From", "REL", "To", …)` calls.
    fn edges_from_pyspark(py: &str) -> Vec<(String, String, String)> {
        let mut out = Vec::new();
        for line in py.lines() {
            // Skip the writer's *definition* (`def _write_edges(fl, rel, tl…)`),
            // whose args are unquoted params, not the call sites.
            if line.trim_start().starts_with("def ") {
                continue;
            }
            if let Some((_, rest)) = line.split_once("_write_edges(") {
                // Args are Rust {:?} string literals: "From", "REL", "To", ...
                let args: Vec<&str> = rest.split(", ").collect();
                let unq = |s: &str| s.trim().trim_matches('"').to_string();
                out.push((unq(args[0]), unq(args[1]), unq(args[2])));
            }
        }
        out.sort();
        out
    }

    fn batch_from_pyspark(py: &str) -> String {
        py.lines()
            .find_map(|l| l.strip_prefix("BATCH_ROWS = "))
            .unwrap()
            .to_string()
    }

    fn nodes_from_pipeline(p: &Pipeline) -> Vec<(String, String)> {
        let mut out: Vec<(String, String)> = p
            .datasets
            .iter()
            .filter_map(|d| {
                Some((
                    d.properties.get(prop::NODE_LABEL)?.clone(),
                    d.properties.get(prop::NODE_KEY)?.clone(),
                ))
            })
            .collect();
        out.sort();
        out
    }

    fn edges_from_pipeline(p: &Pipeline) -> Vec<(String, String, String)> {
        let mut out: Vec<(String, String, String)> = p
            .datasets
            .iter()
            .filter_map(|d| {
                Some((
                    d.properties.get(prop::EDGE_FROM)?.clone(),
                    d.properties.get(prop::EDGE_REL)?.clone(),
                    d.properties.get(prop::EDGE_TO)?.clone(),
                ))
            })
            .collect();
        out.sort();
        out
    }

    fn batch_from_pipeline(p: &Pipeline) -> String {
        p.datasets
            .iter()
            .find_map(|d| d.properties.get(prop::BATCH_ROWS).cloned())
            .unwrap()
    }

    /// **THE invariant.** The PySpark job and the thund Pipeline are two
    /// projections of the SAME [`GraphPlan`], and they must agree on labels,
    /// merge keys, edge directions, and batch size. We derive each set
    /// INDEPENDENTLY — parsed from the generated Python on one side, read from
    /// the Pipeline's `knut.*` properties on the other — and assert equality.
    /// A divergence in either emitter (a renamed label, a flipped edge, a
    /// different batch) fails here.
    #[test]
    fn two_projections_agree() {
        let plan = sample_plan(&CodegenOptions::default());
        let py = plan.to_pyspark();
        let pipe = to_thund_pipeline(&plan).unwrap();

        assert_eq!(
            nodes_from_pyspark(&py),
            nodes_from_pipeline(&pipe),
            "node labels+keys agree across the two projections"
        );
        assert_eq!(
            edges_from_pyspark(&py),
            edges_from_pipeline(&pipe),
            "edge directions (from,rel,to) agree across the two projections"
        );
        assert_eq!(
            batch_from_pyspark(&py),
            batch_from_pipeline(&pipe),
            "batch size agrees across the two projections"
        );

        // And both agree with the plan's own canonical contract.
        let contract = plan.contract();
        assert!(contract.contains("node User=user_id"));
        assert!(contract.contains("node Company=company_id"));
        assert!(contract.contains("edge User-WORKS_AT->Company"));
        assert!(contract.contains("batch_rows=1000"));

        // Non-empty, so the assertions above are not vacuously true.
        assert_eq!(nodes_from_pyspark(&py).len(), 2);
        assert_eq!(edges_from_pyspark(&py).len(), 1);
    }

    /// The thund Pipeline projection is a REAL, valid thund IR: it validates,
    /// its graph target rides in `knut.*` properties, nodes precede edges, and
    /// it round-trips through the RON DSL losslessly (so `pipeline.ron` is a
    /// first-class thund spec an `ExecBackend` consumes).
    #[test]
    fn pipeline_is_real_thund_ir_and_round_trips() {
        let plan = sample_plan(&CodegenOptions::default());
        let pipe = to_thund_pipeline(&plan).unwrap();
        pipe.validate()
            .expect("the projected pipeline is valid thund IR");

        // Nodes-before-edges in declaration order.
        let names: Vec<&str> = pipe.datasets.iter().map(|d| d.name.as_str()).collect();
        assert_eq!(
            names,
            [
                "node__Company",
                "node__User",
                "edge__User__WORKS_AT__Company"
            ]
        );
        // Every dataset is a graph Sink carrying its target.
        assert!(
            pipe.datasets
                .iter()
                .all(|d| d.output_type == OutputType::Sink)
        );

        // The projection flows prune to exactly the needed columns.
        let user_flow = pipe
            .flows
            .iter()
            .find(|f| f.target == "node__User")
            .unwrap();
        assert_eq!(user_flow.projection, vec!["user_id", "name", "age"]);
        assert_eq!(user_flow.reads, vec!["lake.social.users"]);

        // Round-trip through the RON DSL (the pipeline.ron artifact).
        let ron = to_pipeline_ron(&plan).unwrap();
        let back = crate::authoring::dsl::from_ron(&ron).unwrap();
        assert_eq!(back, pipe, "pipeline.ron round-trips losslessly");
        // The graph target genuinely reached the RON.
        assert!(ron.contains("knut.node.label"));
        assert!(ron.contains("WORKS_AT"));
    }

    /// A streaming plan projects to incremental datasets (and still validates /
    /// round-trips), matching the PySpark `foreachBatch` streaming path.
    #[test]
    fn streaming_plan_projects_incremental() {
        use knut_bifrost::{PlanMode, SourceKind};
        let opts = CodegenOptions {
            source: SourceKind::Iceberg,
            mode: PlanMode::Stream,
            ..Default::default()
        };
        let plan = sample_plan(&opts);
        let pipe = to_thund_pipeline(&plan).unwrap();
        assert!(
            pipe.datasets
                .iter()
                .all(|d| d.materialization == Materialization::Incremental),
            "streaming plan → incremental datasets"
        );
        // Still agrees with the Spark projection's streaming form.
        assert!(plan.to_pyspark().contains("foreachBatch("));
    }
}