#![deny(rustdoc::broken_intra_doc_links)]
pub mod authoring;
pub mod backend;
pub mod error;
pub mod ir;
#[cfg(feature = "rel2graph")]
pub mod rel2graph;
pub use backend::{
Capabilities, ExecBackend, RunEvent, RunHandle, RunPhase, SdpFlowDrop, SdpLoweringReport,
};
pub use error::{Result, ThundError};
pub use ir::Pipeline;
#[inline]
pub fn functional_status(component: &str, check: &str, ok: bool, detail: &str) {
#[cfg(feature = "testmatrix")]
nornir_testmatrix::functional_status(component, check, ok, detail);
#[cfg(not(feature = "testmatrix"))]
{
let _ = (component, check, ok, detail);
}
}
#[cfg(test)]
mod tests {
use super::ThundError;
use super::backend::{ExecBackend, native::NativeBackend, spark::SparkBackend};
use super::ir::*;
fn sample_streaming() -> Pipeline {
Pipeline::new("clicks")
.with_dataset(
Dataset::new("click_counts", OutputType::Table)
.incremental()
.with_schema(
DatasetSchema::new()
.field("window_end", "Timestamp(Microsecond, None)", false)
.field("n", "Int64", false),
),
)
.with_flow(
Flow::streaming(
"f_clicks",
"click_counts",
SourceSpec::Kafka {
bootstrap: "localhost:9092".into(),
topic: "clicks".into(),
format: "json".into(),
},
)
.with_watermark(Watermark {
event_time_column: "ts".into(),
allowed_lateness_ms: 5_000,
idle_timeout_ms: 30_000,
})
.with_query("SELECT window_end, count(*) n FROM clicks GROUP BY TUMBLE(ts, '10s')")
.expect(Expectation::new("n_positive", "n > 0").on(OnViolation::Drop)),
)
}
#[test]
fn validate_and_topo_order() {
let p = Pipeline::new("p")
.with_dataset(Dataset::new("a", OutputType::Table))
.with_dataset(Dataset::new("b", OutputType::MaterializedView))
.with_flow(Flow::batch("fa", "a", Vec::<String>::new()))
.with_flow(Flow::batch("fb", "b", ["a"]));
p.validate().expect("valid");
assert_eq!(
p.topo_order().unwrap(),
vec!["a".to_string(), "b".to_string()]
);
}
#[test]
fn dangling_flow_is_rejected() {
let p = Pipeline::new("p").with_flow(Flow::batch("f", "ghost", Vec::<String>::new()));
let err = p.validate().unwrap_err();
assert!(
matches!(err, ThundError::DanglingFlow { .. }),
"got {err:?}"
);
}
#[test]
fn cycle_is_rejected() {
let p = Pipeline::new("p")
.with_dataset(Dataset::new("a", OutputType::Table))
.with_dataset(Dataset::new("b", OutputType::Table))
.with_flow(Flow::batch("fa", "a", ["b"]))
.with_flow(Flow::batch("fb", "b", ["a"]));
assert!(matches!(p.validate().unwrap_err(), ThundError::Cyclic));
assert!(p.topo_order().is_none());
}
#[test]
fn streaming_pipeline_detected_and_lowers_to_sdp() {
let p = sample_streaming();
p.validate().expect("valid");
assert!(p.is_streaming(), "kafka flow makes it streaming");
let sdp = p.to_sdp();
sdp.validate().expect("lowered SDP is valid");
assert_eq!(sdp.datasets.len(), 1);
assert_eq!(
sdp.datasets[0].output_type,
knut_pipelines::OutputType::Table
);
assert_eq!(
sdp.datasets[0].schema.as_deref(),
Some("window_end Timestamp(Microsecond, None), n Int64")
);
assert_eq!(sdp.flows.len(), 1);
assert_eq!(
sdp.flows[0].query.as_deref(),
Some("SELECT window_end, count(*) n FROM clicks GROUP BY TUMBLE(ts, '10s')")
);
}
#[test]
fn spark_backend_rejects_expectations_via_capabilities() {
let p = sample_streaming(); let spark = SparkBackend::new("sc://localhost:15002");
let err = spark.check(&p).unwrap_err();
assert!(
matches!(err, ThundError::Unsupported { ref what, .. } if what.contains("expectations")),
"got {err:?}"
);
}
#[test]
fn partition_transforms_are_refused_up_front_by_both_backends() {
use crate::ir::{Dataset, OutputType, Pipeline};
let p = Pipeline::new("p").with_dataset(
Dataset::new("events", OutputType::MaterializedView)
.with_partition_cols(["months(event_ts)"]),
);
for err in [
SparkBackend::new("sc://localhost:15002")
.check(&p)
.unwrap_err(),
NativeBackend::new().check(&p).unwrap_err(),
] {
assert!(
matches!(err, ThundError::Unsupported { ref what, .. }
if what.contains("months(event_ts)")),
"got {err:?}"
);
}
let ok = Pipeline::new("p").with_dataset(
Dataset::new("events", OutputType::MaterializedView)
.with_partition_cols(["event_month"]),
);
SparkBackend::new("sc://localhost:15002")
.check(&ok)
.unwrap();
NativeBackend::new().check(&ok).unwrap();
}
#[test]
fn graph_sink_pipelines_are_refused_up_front_by_the_spark_backend() {
use crate::ir::{Dataset, OutputType, Pipeline};
let p = Pipeline::new("rel2graph").with_dataset(
Dataset::new("node__User", OutputType::Sink)
.with_format("falkordb-node")
.with_properties([("knut.node.label", "User")]),
);
let spark = SparkBackend::new("sc://localhost:15002");
assert!(!spark.capabilities().graph_sinks);
let err = spark.check(&p).unwrap_err();
assert!(
matches!(err, ThundError::Unsupported { backend, ref what }
if backend == "spark-connect-sdp"
&& what.contains("graph-sink")
&& what.contains("SparkBackend")
&& what.contains("NativeBackend")
&& what.contains("PySpark")),
"got {err:?}"
);
let native = NativeBackend::new();
if native.capabilities().graph_sinks {
native
.check(&p)
.expect("native (falkordb) accepts a graph-sink pipeline");
} else {
assert!(
matches!(
native.check(&p).unwrap_err(),
ThundError::Unsupported { .. }
),
"native without the falkordb feature refuses graph sinks"
);
}
let plain = Pipeline::new("p")
.with_dataset(Dataset::new("audit", OutputType::Sink).with_format("parquet"));
SparkBackend::new("sc://x:15002")
.check(&plain)
.expect("a non-graph Sink is not caught by the graph-sink gate");
}
#[test]
fn native_backend_accepts_everything_in_the_ir() {
let p = sample_streaming();
let native = NativeBackend::new();
native
.check(&p)
.expect("native accepts streaming + expectations + watermark");
let caps = native.capabilities();
assert!(caps.streaming && caps.event_time && caps.cdc && caps.expectations);
}
#[test]
fn capabilities_differ_between_backends() {
let n = NativeBackend::new().capabilities();
let s = SparkBackend::new("sc://x:15002").capabilities();
assert!(
n.event_time && !s.event_time,
"native has event-time, spark does not"
);
assert_eq!(s.output_modes, vec!["append".to_string()]);
assert!(n.output_modes.contains(&"complete".to_string()));
}
#[cfg(feature = "dsl")]
#[test]
fn dsl_roundtrips_through_ir() {
let p = sample_streaming();
let ron = crate::authoring::dsl::to_ron(&p).expect("serialize");
let back = crate::authoring::dsl::from_ron(&ron).expect("parse");
assert_eq!(p, back, "RON DSL is a lossless view of the IR");
}
#[cfg(feature = "dsl")]
#[test]
fn dsl_roundtrips_declarative_file_sink() {
let p = Pipeline::new("rollup_to_disk")
.with_dataset(
Dataset::new("by_region", OutputType::Sink)
.with_path("out/by_region")
.with_format("parquet")
.with_partition_cols(["region"]),
)
.with_flow(
Flow::batch("agg", "by_region", ["orders"])
.with_query("SELECT region, SUM(amount) AS total FROM orders GROUP BY region"),
);
let ron = crate::authoring::dsl::to_ron(&p).expect("serialize");
let back = crate::authoring::dsl::from_ron(&ron).expect("parse");
assert_eq!(
p, back,
"the declarative sink path/format/partitioning round-trips"
);
let d = back.dataset("by_region").expect("sink dataset present");
assert_eq!(
d.path.as_deref(),
Some("out/by_region"),
"the sink path survived"
);
assert_eq!(d.partition_cols, vec!["region".to_string()]);
}
#[cfg(feature = "dsl")]
#[test]
fn dsl_roundtrips_typed_projection_filter() {
let p = Pipeline::new("prune")
.with_dataset(Dataset::new("slim", OutputType::MaterializedView))
.with_flow(
Flow::batch("proj", "slim", ["orders"])
.with_projection(["customer", "amount"])
.with_filter("amount > 100"),
);
let ron = crate::authoring::dsl::to_ron(&p).expect("serialize");
let back = crate::authoring::dsl::from_ron(&ron).expect("parse");
assert_eq!(
p, back,
"the typed projection + filter round-trip losslessly"
);
let f = &back.flows[0];
assert_eq!(
f.projection,
vec!["customer".to_string(), "amount".to_string()]
);
assert_eq!(f.filter.as_deref(), Some("amount > 100"));
assert_eq!(
f.to_sdp().query.as_deref(),
Some("SELECT customer, amount FROM orders WHERE amount > 100"),
);
}
#[cfg(feature = "dsl")]
#[test]
fn dsl_roundtrips_write_mode() {
use crate::ir::WriteMode;
let p = Pipeline::new("refresh")
.with_dataset(
Dataset::new("by_region", OutputType::Sink)
.with_path("out/by_region")
.with_partition_cols(["region"])
.overwrite(),
)
.with_dataset(Dataset::new("plain_sink", OutputType::Sink).with_path("out/plain_sink"))
.with_flow(
Flow::batch("agg", "by_region", ["orders"])
.with_query("SELECT region, SUM(amount) AS total FROM orders GROUP BY region"),
);
let ron = crate::authoring::dsl::to_ron(&p).expect("serialize");
let back = crate::authoring::dsl::from_ron(&ron).expect("parse");
assert_eq!(p, back, "the sink write_mode round-trips losslessly");
assert_eq!(
back.dataset("by_region").unwrap().write_mode,
WriteMode::Overwrite,
"the declared Overwrite mode survived",
);
assert_eq!(
back.dataset("plain_sink").unwrap().write_mode,
WriteMode::Append,
"a sink with no declared write_mode defaults to Append (L2-additive)",
);
assert!(
!ron.contains("append"),
"the default Append mode is not serialized (skip_serializing_if) — wire stays additive",
);
}
#[cfg(feature = "dsl")]
#[test]
fn from_ron_full_pipeline_lowers_to_the_dataflow_graph_korp_renders() {
let authored = Pipeline::new("medallion")
.with_dataset(Dataset::new("bronze", OutputType::Table))
.with_dataset(Dataset::new("silver", OutputType::MaterializedView))
.with_dataset(Dataset::new("gold", OutputType::MaterializedView))
.with_flow(
Flow::batch("refine", "silver", ["bronze"])
.with_query("SELECT * FROM bronze WHERE ok"),
)
.with_flow(
Flow::batch("rollup", "gold", ["silver"])
.with_query("SELECT region, SUM(amount) total FROM silver GROUP BY region"),
);
let text = crate::authoring::dsl::to_ron(&authored).expect("render RON");
let parsed = crate::authoring::dsl::from_ron(&text).expect("parse RON");
let g = parsed.to_sdp();
g.validate().expect("the lowered SDP graph is valid");
assert_eq!(
g.datasets
.iter()
.map(|d| d.name.as_str())
.collect::<Vec<_>>(),
["bronze", "silver", "gold"],
);
assert_eq!(g.datasets[0].output_type, knut_pipelines::OutputType::Table);
assert_eq!(
g.datasets[1].output_type,
knut_pipelines::OutputType::MaterializedView
);
assert_eq!(
g.datasets[2].output_type,
knut_pipelines::OutputType::MaterializedView
);
assert_eq!(g.flows.len(), 2);
let refine = g
.flows
.iter()
.find(|f| f.name == "refine")
.expect("refine flow");
assert_eq!(refine.target, "silver");
assert_eq!(
refine.query.as_deref(),
Some("SELECT * FROM bronze WHERE ok")
);
let rollup = g
.flows
.iter()
.find(|f| f.name == "rollup")
.expect("rollup flow");
assert_eq!(rollup.target, "gold");
assert_eq!(
rollup.query.as_deref(),
Some("SELECT region, SUM(amount) total FROM silver GROUP BY region"),
);
}
}