knut-thund 0.1.2

Þ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
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
//! Integration proof for the native backend's **micro-batch streaming** layer.
//!
//! Drives a *synthetic bounded-but-streamed* source (a range/memory source fed
//! in N micro-batches via `NativeBackend::with_stream_input`) through the REAL
//! DataFusion engine — no live Kafka/socket, no heavy run. Asserts:
//!   * incremental batches are produced in order (one increment per trigger),
//!   * the final materialised result equals the batch path over the same data,
//!   * `AvailableNow` collapses to a single bounded pass,
//!   * drop / fail data-quality expectations behave on the streaming result.
//!
//! Gated on `--features native` (the DataFusion tree); a no-op otherwise.
#![cfg(feature = "native")]

use knut_thund::backend::{native::NativeBackend, ExecBackend, RunHandle};
use knut_thund::ir::{
    CdcSpec, Dataset, Expectation, Flow, FlowKind, OnViolation, OutputMode, OutputType, Pipeline,
    ScdType, SourceSpec, Trigger,
};

use datafusion::arrow::array::{Array, Int64Array, RecordBatch, StringArray};
use datafusion::arrow::datatypes::{DataType, Field, Schema};
use std::sync::Arc;

/// An `(id, v)` batch from parallel id/value slices.
fn batch(ids: &[i64], vs: &[i64]) -> RecordBatch {
    let schema = Arc::new(Schema::new(vec![
        Field::new("id", DataType::Int64, false),
        Field::new("v", DataType::Int64, false),
    ]));
    RecordBatch::try_new(
        schema,
        vec![
            Arc::new(Int64Array::from(ids.to_vec())),
            Arc::new(Int64Array::from(vs.to_vec())),
        ],
    )
    .unwrap()
}

/// Three micro-batches of the `events` source.
fn micro_batches() -> Vec<Vec<RecordBatch>> {
    vec![
        vec![batch(&[1, 2], &[10, -5])],
        vec![batch(&[3, 4], &[20, 0])],
        vec![batch(&[5], &[30])],
    ]
}

/// All the same rows, concatenated — for the batch-path comparison.
fn all_rows() -> RecordBatch {
    batch(&[1, 2, 3, 4, 5], &[10, -5, 20, 0, 30])
}

/// A streaming Kafka flow over the seeded `events` topic, running `query` with
/// `output_mode` and `trigger`.
fn stream_pipeline(query: &str, output_mode: OutputMode, trigger: Trigger) -> Pipeline {
    let mut flow = Flow::streaming(
        "f_events",
        "events_out",
        SourceSpec::Kafka {
            bootstrap: "localhost:9092".into(),
            topic: "events".into(),
            format: "json".into(),
        },
    )
    .with_query(query);
    if let FlowKind::Streaming {
        output_mode: om,
        trigger: tr,
        ..
    } = &mut flow.kind
    {
        *om = output_mode;
        *tr = trigger;
    }
    Pipeline::new("stream")
        .with_dataset(Dataset::new("events_out", OutputType::Table).incremental())
        .with_flow(flow)
}

/// The batch-path result of `query` over all rows seen at once.
fn batch_result(query: &str) -> Vec<i64> {
    let p = Pipeline::new("batch")
        .with_dataset(Dataset::new("events_out", OutputType::MaterializedView))
        .with_flow(Flow::batch("f_batch", "events_out", ["events"]).with_query(query));
    let run = NativeBackend::new()
        .with_input("events", vec![all_rows()])
        .run(&p)
        .expect("batch runs");
    sorted_ids(run.output("events_out").expect("batch output"))
}

/// Sorted `id` column across a slice of batches.
fn sorted_ids(batches: &[RecordBatch]) -> Vec<i64> {
    let mut v: Vec<i64> = batches
        .iter()
        .flat_map(|b| {
            b.column(0)
                .as_any()
                .downcast_ref::<Int64Array>()
                .unwrap()
                .values()
                .to_vec()
        })
        .collect();
    v.sort_unstable();
    v
}

/// Append mode over a row-wise filter: each trigger emits exactly the newly
/// arrived rows that pass, in order; the concatenated increments and the final
/// materialised table both equal the batch path.
#[test]
fn append_mode_emits_new_rows_per_trigger_and_matches_batch() {
    let query = "SELECT id, v FROM events WHERE v > 0";
    let p = stream_pipeline(query, OutputMode::Append, Trigger::Continuous);
    let run = NativeBackend::new()
        .with_stream_input("events", micro_batches())
        .run(&p)
        .expect("streaming run completes");

    // One increment per micro-batch, in order.
    assert_eq!(run.trigger_count("events_out"), 3, "one increment per micro-batch");
    let incs = run.increments("events_out");
    // Trigger 0: rows (1,10),(2,-5) -> only id 1 passes v>0.
    assert_eq!(sorted_ids(&incs[0]), vec![1]);
    // Trigger 1: rows (3,20),(4,0) -> only id 3 passes.
    assert_eq!(sorted_ids(&incs[1]), vec![3]);
    // Trigger 2: row (5,30) -> id 5 passes.
    assert_eq!(sorted_ids(&incs[2]), vec![5]);

    // Concatenated increments == final materialised == batch path.
    let concat: Vec<i64> = {
        let mut all: Vec<i64> = incs.iter().flat_map(|i| sorted_ids(i)).collect();
        all.sort_unstable();
        all
    };
    let final_rows = sorted_ids(run.output("events_out").expect("final output"));
    assert_eq!(concat, final_rows, "increments reconstruct the final table");
    assert_eq!(final_rows, batch_result(query), "final == batch path over same data");
}

/// Complete mode over an aggregation: each trigger emits the full running
/// aggregate; the final table equals the batch path.
#[test]
fn complete_mode_emits_full_aggregate_each_trigger_and_matches_batch() {
    let query = "SELECT count(*) AS n, sum(v) AS total FROM events";
    let p = stream_pipeline(query, OutputMode::Complete, Trigger::Continuous);
    let run = NativeBackend::new()
        .with_stream_input("events", micro_batches())
        .run(&p)
        .expect("streaming run completes");

    assert_eq!(run.trigger_count("events_out"), 3, "one increment per micro-batch");
    let incs = run.increments("events_out");
    // Complete mode: every increment is the full result (a single aggregate row).
    for inc in incs {
        let rows: usize = inc.iter().map(|b| b.num_rows()).sum();
        assert_eq!(rows, 1, "complete-mode aggregate emits one row per trigger");
    }
    // Running counts grow monotonically: n = 2, 4, 5.
    let n = |inc: &Vec<RecordBatch>| {
        inc[0].column(0).as_any().downcast_ref::<Int64Array>().unwrap().value(0)
    };
    assert_eq!(n(&incs[0]), 2);
    assert_eq!(n(&incs[1]), 4);
    assert_eq!(n(&incs[2]), 5);

    // Final aggregate: n=5, total=10+(-5)+20+0+30 = 55.
    let out = run.output("events_out").expect("final output");
    let total = out[0].column(1).as_any().downcast_ref::<Int64Array>().unwrap().value(0);
    assert_eq!((n(&out.to_vec()), total), (5, 55));
}

/// `AvailableNow` collapses all currently-available data into ONE bounded pass.
#[test]
fn available_now_is_a_single_bounded_pass() {
    let query = "SELECT id, v FROM events WHERE v > 0";
    let p = stream_pipeline(query, OutputMode::Append, Trigger::AvailableNow);
    let run = NativeBackend::new()
        .with_stream_input("events", micro_batches())
        .run(&p)
        .expect("streaming run completes");

    assert_eq!(run.trigger_count("events_out"), 1, "AvailableNow => single pass");
    let final_rows = sorted_ids(run.output("events_out").expect("final output"));
    assert_eq!(final_rows, batch_result(query), "single pass still equals batch path");
    // The single increment holds every passing row.
    assert_eq!(sorted_ids(&run.increments("events_out")[0]), vec![1, 3, 5]);
}

/// A DROP expectation removes violating rows from the streaming result, just
/// like the batch path.
#[test]
fn drop_expectation_filters_streaming_result() {
    let query = "SELECT id, v FROM events";
    let mut p = stream_pipeline(query, OutputMode::Append, Trigger::AvailableNow);
    p.flows[0]
        .expectations
        .push(Expectation::new("v_positive", "v > 0").on(OnViolation::Drop));
    let mut run = NativeBackend::new()
        .with_stream_input("events", micro_batches())
        .run(&p)
        .expect("streaming run completes");
    // Rows with v>0: ids 1,3,5 (ids 2 (-5) and 4 (0) dropped).
    assert_eq!(sorted_ids(run.output("events_out").expect("final output")), vec![1, 3, 5]);
    assert!(run
        .poll_events()
        .unwrap()
        .iter()
        .any(|e| e.message.contains("DROP")));
}

/// A FAIL expectation aborts the streaming run with a real error.
#[test]
fn fail_expectation_aborts_streaming_run() {
    let query = "SELECT id, v FROM events";
    let mut p = stream_pipeline(query, OutputMode::Append, Trigger::AvailableNow);
    p.flows[0]
        .expectations
        .push(Expectation::new("v_positive", "v > 0").on(OnViolation::Fail));
    let err = NativeBackend::new()
        .with_stream_input("events", micro_batches())
        .run(&p)
        .unwrap_err();
    assert!(matches!(err, knut_thund::ThundError::Backend(m) if m.contains("FAILED")));
}

/// A downstream BATCH flow reads the materialised streaming table by name —
/// proving the streaming result registers into the shared in-process catalog.
#[test]
fn downstream_batch_flow_reads_streaming_result() {
    let mut flow = Flow::streaming(
        "f_events",
        "events_out",
        SourceSpec::Kafka {
            bootstrap: "localhost:9092".into(),
            topic: "events".into(),
            format: "json".into(),
        },
    )
    .with_query("SELECT id, v FROM events WHERE v > 0");
    if let FlowKind::Streaming { trigger, .. } = &mut flow.kind {
        *trigger = Trigger::AvailableNow;
    }
    let p = Pipeline::new("stream_then_batch")
        .with_dataset(Dataset::new("events_out", OutputType::Table).incremental())
        .with_dataset(Dataset::new("summary", OutputType::MaterializedView))
        .with_flow(flow)
        .with_flow(
            Flow::batch("f_sum", "summary", ["events_out"])
                .with_query("SELECT count(*) AS n, sum(v) AS total FROM events_out"),
        );
    let run = NativeBackend::new()
        .with_stream_input("events", micro_batches())
        .run(&p)
        .expect("mixed streaming+batch run completes");
    // Positive rows: v = 10, 20, 30 -> n=3, total=60.
    let out = run.output("summary").expect("summary produced");
    let n = out[0].column(0).as_any().downcast_ref::<Int64Array>().unwrap().value(0);
    let total = out[0].column(1).as_any().downcast_ref::<Int64Array>().unwrap().value(0);
    assert_eq!((n, total), (3, 60));
}

/// With NO seeded stream (and no live infra), a Kafka flow is still deferred
/// honestly — micro-batch streaming is opt-in via `with_stream_input`.
#[test]
fn unseeded_kafka_flow_still_defers() {
    let p = stream_pipeline("SELECT * FROM events", OutputMode::Append, Trigger::Continuous);
    let mut run = NativeBackend::new().run(&p).expect("run completes despite deferral");
    assert_eq!(run.row_count("events_out"), 0, "nothing materialised for an unseeded live source");
    assert!(run.poll_events().unwrap().iter().any(|e| e.message.contains("deferred")));
}

// ── CDC apply-changes / SCD merge ───────────────────────────────────────────
//
// The native backend's CDC merge operator folds an upstream changelog — carrying
// skade & knut-bifrost's shared `_change_type` vocabulary (INSERT / UPDATE_AFTER
// / UPDATE_BEFORE / DELETE) — into the target, key-matched and ordered by the
// flow's `sequence_by`. These prove SCD type 1 (overwrite) and type 2 (history).

/// A changelog batch: parallel `(id, name, seq, _change_type)` slices.
fn change_batch(rows: &[(i64, &str, i64, &str)]) -> RecordBatch {
    let schema = Arc::new(Schema::new(vec![
        Field::new("id", DataType::Int64, false),
        Field::new("name", DataType::Utf8, false),
        Field::new("seq", DataType::Int64, false),
        Field::new("_change_type", DataType::Utf8, false),
    ]));
    let ids: Vec<i64> = rows.iter().map(|r| r.0).collect();
    let names: Vec<&str> = rows.iter().map(|r| r.1).collect();
    let seqs: Vec<i64> = rows.iter().map(|r| r.2).collect();
    let cts: Vec<&str> = rows.iter().map(|r| r.3).collect();
    RecordBatch::try_new(
        schema,
        vec![
            Arc::new(Int64Array::from(ids)),
            Arc::new(StringArray::from(names)),
            Arc::new(Int64Array::from(seqs)),
            Arc::new(StringArray::from(cts)),
        ],
    )
    .unwrap()
}

/// The insert/update/delete changelog seeded as three micro-batches:
///  * t0: INSERT id=1 "a", INSERT id=2 "b"
///  * t1: UPDATE_AFTER id=1 "a2", INSERT id=3 "c"
///  * t2: DELETE id=2
fn cdc_micro_batches() -> Vec<Vec<RecordBatch>> {
    vec![
        vec![change_batch(&[(1, "a", 1, "INSERT"), (2, "b", 2, "INSERT")])],
        vec![change_batch(&[(1, "a2", 3, "UPDATE_AFTER"), (3, "c", 4, "INSERT")])],
        vec![change_batch(&[(2, "b", 5, "DELETE")])],
    ]
}

/// A CDC apply-changes pipeline over the seeded `changes` changelog, keyed by
/// `id` and sequenced by `seq`, with the given SCD type.
fn cdc_pipeline(scd_type: ScdType) -> Pipeline {
    let flow = Flow {
        name: "f_cdc".into(),
        target: "dim".into(),
        reads: vec!["changes".into()],
        kind: FlowKind::Cdc {
            cdc: CdcSpec {
                keys: vec!["id".into()],
                sequence_by: "seq".into(),
                apply_as_deletes: None,
                scd_type,
            },
        },
        query: Some("SELECT * FROM changes".into()),
        expectations: Vec::new(),
    };
    Pipeline::new("cdc")
        .with_dataset(Dataset::new("dim", OutputType::Table).incremental())
        .with_flow(flow)
}

/// Column `name` from a slice of batches as `i64` values (panics if absent/typed
/// otherwise).
fn i64_col(batches: &[RecordBatch], name: &str) -> Vec<i64> {
    batches
        .iter()
        .flat_map(|b| {
            let idx = b.schema().index_of(name).expect("column present");
            b.column(idx)
                .as_any()
                .downcast_ref::<Int64Array>()
                .expect("i64 column")
                .values()
                .to_vec()
        })
        .collect()
}

/// `(id, name)` pairs from the merged target, sorted for a stable assertion.
fn id_name_pairs(batches: &[RecordBatch]) -> Vec<(i64, String)> {
    let mut out = Vec::new();
    for b in batches {
        let ids = b.column(b.schema().index_of("id").unwrap()).as_any().downcast_ref::<Int64Array>().unwrap();
        let names = b.column(b.schema().index_of("name").unwrap()).as_any().downcast_ref::<StringArray>().unwrap();
        for r in 0..b.num_rows() {
            out.push((ids.value(r), names.value(r).to_string()));
        }
    }
    out.sort();
    out
}

/// SCD type 1: the merged target is the latest non-deleted row per key. Over the
/// insert/update/delete changelog: id=1 updated to "a2", id=2 deleted, id=3 new.
#[test]
fn cdc_scd_type1_upserts_and_deletes() {
    let run = NativeBackend::new()
        .with_stream_input("changes", cdc_micro_batches())
        .run(&cdc_pipeline(ScdType::Type1))
        .expect("cdc run completes");

    let out = run.output("dim").expect("merged target produced");
    // id=1 overwritten to "a2", id=3 inserted, id=2 removed by the DELETE.
    assert_eq!(
        id_name_pairs(out),
        vec![(1, "a2".to_string()), (3, "c".to_string())],
        "latest non-deleted row per key"
    );

    // One increment per micro-batch, reflecting the merged snapshot as it grows.
    assert_eq!(run.trigger_count("dim"), 3, "one merge per changelog micro-batch");
    let incs = run.increments("dim");
    // t0: ids {1,2}; t1: id=1 updated + id=3 added -> {1,2,3}; t2: id=2 deleted -> {1,3}.
    let mut t0 = i64_col(&incs[0], "id");
    t0.sort_unstable();
    assert_eq!(t0, vec![1, 2]);
    let mut t2 = i64_col(&incs[2], "id");
    t2.sort_unstable();
    assert_eq!(t2, vec![1, 3], "the DELETE removed id=2 from the target");
}

/// SCD type 2: every version is kept with derived `__start_at` / `__end_at`
/// bounds. id=1 has two versions (a closed, a2 open); id=2's single version is
/// closed by its DELETE (no open version); id=3 has one open version.
#[test]
fn cdc_scd_type2_keeps_version_history() {
    let run = NativeBackend::new()
        .with_stream_input("changes", cdc_micro_batches())
        .run(&cdc_pipeline(ScdType::Type2))
        .expect("cdc run completes");

    let out = run.output("dim").expect("merged target produced");
    let out: Vec<RecordBatch> = out.to_vec();

    // 4 emitted versions: id=1 ×2, id=2 ×1 (closed), id=3 ×1.
    let total: usize = out.iter().map(|b| b.num_rows()).sum();
    assert_eq!(total, 4, "one row per non-delete change (history preserved)");

    // Gather (id, name, start_at, end_at|None) tuples.
    let mut versions: Vec<(i64, String, i64, Option<i64>)> = Vec::new();
    for b in &out {
        let ids = b.column(b.schema().index_of("id").unwrap()).as_any().downcast_ref::<Int64Array>().unwrap();
        let names = b.column(b.schema().index_of("name").unwrap()).as_any().downcast_ref::<StringArray>().unwrap();
        let starts = b.column(b.schema().index_of("__start_at").unwrap()).as_any().downcast_ref::<Int64Array>().unwrap();
        let ends = b.column(b.schema().index_of("__end_at").unwrap()).as_any().downcast_ref::<Int64Array>().unwrap();
        for r in 0..b.num_rows() {
            let end = if ends.is_null(r) { None } else { Some(ends.value(r)) };
            versions.push((ids.value(r), names.value(r).to_string(), starts.value(r), end));
        }
    }
    versions.sort();

    // id=1: "a" valid [1,3), then "a2" open [3,∞).
    assert!(versions.contains(&(1, "a".into(), 1, Some(3))), "id=1 first version closed at seq 3");
    assert!(versions.contains(&(1, "a2".into(), 3, None)), "id=1 current version open");
    // id=2: single version closed by the DELETE at seq 5, no open version.
    assert!(versions.contains(&(2, "b".into(), 2, Some(5))), "id=2 closed by its delete");
    assert!(!versions.iter().any(|v| v.0 == 2 && v.3.is_none()), "id=2 has no open version after delete");
    // id=3: current open version.
    assert!(versions.contains(&(3, "c".into(), 4, None)), "id=3 current version open");
}

/// A CDC flow with no seeded changelog (and no registered source) surfaces the
/// planning failure honestly rather than faking a merge.
#[test]
fn cdc_flow_without_a_changelog_errors() {
    let err = NativeBackend::new()
        .run(&cdc_pipeline(ScdType::Type1))
        .unwrap_err();
    assert!(
        matches!(err, knut_thund::ThundError::Backend(_)),
        "unseeded CDC changelog fails as a backend error, got {err:?}"
    );
}