#[path = "e2e_support/mod.rs"]
mod support;
use std::time::Duration;
use support::*;
#[test]
#[ignore = "requires Docker"]
fn happy_path_delivers_exactly_and_commits() {
let h = Harness::up();
let params = PipelineParams::defaults("e2e-happy", 19181);
let partitions = 6;
let total = 50_000;
h.create_topic(¶ms.topic, partitions);
h.create_table(¶ms.table);
let produced = events(partitions, total);
h.produce(¶ms.topic, &produced);
let pipeline = h.spawn_pipeline(¶ms);
wait_until(Duration::from_secs(120), "all rows in ClickHouse", || {
h.count(¶ms.table) == total as u64
});
let (status, _) = http_get(pipeline.admin, "/healthz");
assert_eq!(status, 200, "healthy while consuming");
let (status, _) = http_get(pipeline.admin, "/readyz");
assert_eq!(status, 200, "ready: assignment received, sinks probed");
let (status, _) = http_get(pipeline.admin, "/metrics");
assert_eq!(status, 200);
wait_until(Duration::from_secs(30), "operator counted records", || {
let (_, body) = http_get(pipeline.admin, "/metrics");
metric_sum(&body, "spate_operator_records_in_total") >= total as f64
});
wait_until(Duration::from_secs(30), "watermarks committed", || {
let committed = h.committed(¶ms.topic, ¶ms.group, partitions);
committed.iter().sum::<i64>() == total as i64
});
let report = pipeline.stop();
assert!(
matches!(report.state, spate::pipeline::ExitState::Completed),
"clean exit: {report:?}"
);
let mut per_partition = vec![0i64; partitions as usize];
for (p, _) in &produced {
per_partition[*p as usize] += 1;
}
assert_eq!(
h.committed(¶ms.topic, ¶ms.group, partitions),
per_partition,
"every partition committed exactly its produced count"
);
assert_eq!(h.count(¶ms.table), total as u64, "no duplicates");
let name =
h.rt.block_on(
h.ch_client()
.query(&format!(
"SELECT name FROM {} WHERE id = {}",
params.table,
event_id(3, 100)
))
.fetch_one::<String>(),
)
.expect("spot check row");
assert_eq!(name, "evt-3-100");
}