#![cfg(feature = "sql")]
use knut_thund::authoring::sql;
use knut_thund::ir::{FlowKind, OnViolation, OutputType};
const SCRIPT: &str = r#"
CREATE MATERIALIZED VIEW bronze_events (
event_id BIGINT NOT NULL,
event_ts TIMESTAMP,
amount DOUBLE,
CONSTRAINT amount_non_negative EXPECT (amount >= 0) ON VIOLATION DROP ROW
)
PARTITIONED BY (months(event_ts))
COMMENT 'raw events, month-partitioned'
TBLPROPERTIES ('owner' = 'data-eng')
AS SELECT event_id, event_ts, amount FROM landing_events;
CREATE MATERIALIZED VIEW silver_events AS
SELECT event_id, event_ts, amount FROM bronze_events WHERE amount > 0;
CREATE STREAMING TABLE gold_live AS
SELECT event_id, amount FROM STREAM(silver_events);
CREATE FLOW gold_backfill AS
INSERT INTO gold_live BY NAME
SELECT event_id, amount FROM STREAM(bronze_events);
"#;
#[test]
fn parses_every_clause_into_the_ir() {
let p = sql::from_sql("medallion", SCRIPT).unwrap();
assert_eq!(p.datasets.len(), 3);
assert_eq!(p.flows.len(), 4);
let bronze = p.dataset("bronze_events").unwrap();
assert_eq!(bronze.output_type, OutputType::MaterializedView);
assert_eq!(
bronze.comment.as_deref(),
Some("raw events, month-partitioned")
);
assert_eq!(bronze.partition_cols, vec!["months(event_ts)".to_string()]);
assert_eq!(bronze.schema.fields.len(), 3);
assert_eq!(bronze.schema.fields[0].name, "event_id");
assert_eq!(bronze.schema.fields[0].arrow_type, "Int64");
assert!(!bronze.schema.fields[0].nullable, "BIGINT NOT NULL");
assert!(bronze.schema.fields[1].nullable, "TIMESTAMP is nullable");
assert_eq!(
bronze.schema.fields[1].arrow_type,
"Timestamp(Microsecond, None)"
);
let bf = p
.flows
.iter()
.find(|f| f.target == "bronze_events")
.unwrap();
assert_eq!(bf.expectations.len(), 1);
assert_eq!(bf.expectations[0].name, "amount_non_negative");
assert_eq!(bf.expectations[0].on_violation, OnViolation::Drop);
assert_eq!(
p.dataset("gold_live").unwrap().output_type,
OutputType::Table
);
}
#[test]
fn stream_marks_the_flow_not_the_query() {
let p = sql::from_sql("medallion", SCRIPT).unwrap();
let live = p.flows.iter().find(|f| f.target == "gold_live").unwrap();
assert!(matches!(live.kind, FlowKind::Streaming { .. }));
assert_eq!(live.reads, vec!["silver_events".to_string()]);
let q = live.query.as_ref().unwrap();
assert!(
!q.to_uppercase().contains("STREAM"),
"query still says STREAM: {q}"
);
assert!(q.contains("silver_events"));
}
#[test]
fn dag_is_inferred_from_reads() {
let p = sql::from_sql("medallion", SCRIPT).unwrap();
let order = p.topo_order().expect("acyclic");
let pos = |n: &str| order.iter().position(|x| x == n).unwrap();
assert!(pos("bronze_events") < pos("silver_events"));
assert!(pos("silver_events") < pos("gold_live"));
}
#[test]
fn standalone_create_flow_appends_to_an_existing_target() {
let p = sql::from_sql("medallion", SCRIPT).unwrap();
let backfill = p.flows.iter().find(|f| f.name == "gold_backfill").unwrap();
assert_eq!(backfill.target, "gold_live");
assert_eq!(backfill.reads, vec!["bronze_events".to_string()]);
assert_eq!(
p.flows.iter().filter(|f| f.target == "gold_live").count(),
2
);
}
#[test]
fn sql_round_trips_through_itself() {
let p1 = sql::from_sql("medallion", SCRIPT).unwrap();
let rendered = sql::to_sql(&p1).unwrap();
let p2 = sql::from_sql("medallion", &rendered)
.unwrap_or_else(|e| panic!("re-parse failed: {e}\n--- rendered ---\n{rendered}"));
assert_eq!(p1.datasets.len(), p2.datasets.len());
for d1 in &p1.datasets {
let d2 = p2
.dataset(&d1.name)
.expect("dataset survives the round trip");
assert_eq!(d1.output_type, d2.output_type, "{}", d1.name);
assert_eq!(d1.partition_cols, d2.partition_cols, "{}", d1.name);
assert_eq!(d1.comment, d2.comment, "{}", d1.name);
assert_eq!(d1.schema, d2.schema, "{}", d1.name);
assert_eq!(d1.properties, d2.properties, "{}", d1.name);
}
}
#[cfg(feature = "dsl")]
#[test]
fn sql_to_ron_to_sql_is_the_same_pipeline() {
use knut_thund::authoring::dsl;
let p1 = sql::from_sql("medallion", SCRIPT).unwrap();
let ron = dsl::to_ron(&p1).unwrap();
let p2 = dsl::from_ron(&ron).unwrap();
assert_eq!(p1, p2, "RON round-trip must be lossless");
assert_eq!(sql::to_sql(&p1).unwrap(), sql::to_sql(&p2).unwrap());
}
#[test]
fn partition_transforms_are_preserved_for_the_backend_to_judge() {
let p = sql::from_sql(
"p",
"CREATE MATERIALIZED VIEW t
PARTITIONED BY (months(event_ts), bucket(16, user_id), truncate(10, name))
AS SELECT * FROM src;",
)
.unwrap();
assert_eq!(
p.dataset("t").unwrap().partition_cols,
vec![
"months(event_ts)".to_string(),
"bucket(16, user_id)".to_string(),
"truncate(10, name)".to_string(),
]
);
}
#[test]
fn tblproperties_are_captured_and_reach_sdp() {
let p = sql::from_sql("medallion", SCRIPT).unwrap();
let bronze = p.dataset("bronze_events").unwrap();
assert_eq!(
bronze.properties.get("owner").map(String::as_str),
Some("data-eng")
);
let sdp = p.to_sdp();
let lowered = sdp
.datasets
.iter()
.find(|d| d.name == "bronze_events")
.unwrap();
assert_eq!(
lowered.properties.get("owner").map(String::as_str),
Some("data-eng")
);
}
#[test]
fn tblproperties_accept_dotted_bare_and_empty_forms() {
let p = sql::from_sql(
"p",
"CREATE MATERIALIZED VIEW a
TBLPROPERTIES ('delta.appendOnly' = 'true', bare = 42, 'q' = 'it''s')
AS SELECT 1 AS x;
CREATE MATERIALIZED VIEW b TBLPROPERTIES () AS SELECT 1 AS x;",
)
.unwrap();
let a = p.dataset("a").unwrap();
assert_eq!(
a.properties.get("delta.appendOnly").map(String::as_str),
Some("true")
);
assert_eq!(a.properties.get("bare").map(String::as_str), Some("42"));
assert_eq!(a.properties.get("q").map(String::as_str), Some("it's"));
assert!(p.dataset("b").unwrap().properties.is_empty());
}
#[test]
fn tblproperties_round_trip_through_sql() {
let p1 = sql::from_sql("medallion", SCRIPT).unwrap();
let p2 = sql::from_sql("medallion", &sql::to_sql(&p1).unwrap()).unwrap();
assert_eq!(
p1.dataset("bronze_events").unwrap().properties,
p2.dataset("bronze_events").unwrap().properties
);
}
#[test]
fn a_cyclic_script_is_rejected() {
let err = sql::from_sql(
"p",
"CREATE MATERIALIZED VIEW a AS SELECT * FROM b;
CREATE MATERIALIZED VIEW b AS SELECT * FROM a;",
)
.unwrap_err();
assert!(
matches!(err, knut_thund::ThundError::Cyclic),
"expected a cycle error, got {err}"
);
}
#[test]
fn non_definition_sql_is_rejected_with_a_useful_message() {
let err = sql::from_sql("p", "SELECT 1;").unwrap_err();
let msg = err.to_string();
assert!(msg.contains("CREATE"), "unhelpful message: {msg}");
}
#[test]
fn expect_syntax_spark_rejects_is_accepted_here() {
let p = sql::from_sql(
"p",
"CREATE MATERIALIZED VIEW t (
x INT,
CONSTRAINT x_pos EXPECT (x > 0) ON VIOLATION FAIL UPDATE
) AS SELECT x FROM src;",
)
.unwrap();
let f = p.flows.iter().find(|f| f.target == "t").unwrap();
assert_eq!(f.expectations[0].on_violation, OnViolation::Fail);
}