#![deny(rustdoc::broken_intra_doc_links)]
pub mod authoring;
pub mod backend;
pub mod error;
pub mod ir;
pub use backend::{Capabilities, ExecBackend, RunEvent, RunHandle, RunPhase};
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::backend::{native::NativeBackend, spark::SparkBackend, ExecBackend};
use super::ir::*;
use super::ThundError;
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 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");
}
}