use spate_core::config::PipelineConfig;
use spate_core::deser::Owned;
use spate_core::error::ErrorPolicy;
use spate_core::ops::chain_owned;
use spate_core::pipeline::{ExitState, Pipeline, RuntimeOptions};
use spate_core::record::PartitionId;
use spate_core::sink::KeyHashRouter;
use spate_core::source::LaneId;
use spate_test::{
BytesPassthrough, PipelineRun, TestEncoder, capture_sink, decode_rows, memory_source,
};
use std::time::Duration;
const CONFIG: &str = r#"
pipeline: { name: split-test, threads: 1, io_threads: 1 }
metrics: { exporter: none, listen: "127.0.0.1:0" }
source: { memory: {} }
sinks:
a: { capture: {} }
b: { capture: {} }
"#;
#[test]
fn split_routes_to_named_sinks_and_commits_at_least_once() {
let (source, handle) = memory_source();
let (sink_a, script_a) = capture_sink(1, 1);
let (sink_b, script_b) = capture_sink(1, 1);
let runtime = Pipeline::from_config(PipelineConfig::from_str(CONFIG).expect("config"))
.expect("builder")
.add_sink("a", sink_a)
.expect("sink a")
.add_sink("b", sink_b)
.expect("sink b")
.chains(|ctx| {
let mut split = chain_owned::<Vec<u8>, _>(BytesPassthrough)
.with_metrics(ctx.pipeline.clone(), "main")
.split(ErrorPolicy::Skip);
let a = split.add::<Owned<Vec<u8>>, _, _>(TestEncoder, KeyHashRouter, ctx.sink("a"));
let b = split.add::<Owned<Vec<u8>>, _, _>(TestEncoder, KeyHashRouter, ctx.sink("b"));
split
.route(move |row: Vec<u8>, out| match row.first() {
Some(b'a') => out.emit(a, row),
Some(b'b') => out.emit(b, row),
_ => {}
})
.build()
})
.runtime_options(RuntimeOptions {
handle_signals: false,
..RuntimeOptions::default()
})
.into_runtime(source)
.expect("into_runtime");
let shutdown = runtime.shutdown_handle();
let join = std::thread::spawn(move || runtime.run());
let p = PartitionId(0);
handle.assign_lanes(&[(LaneId(0), p)]);
let mut last = 0;
for payload in [&b"apple"[..], b"banana", b"cherry", b"avocado", b"berry"] {
last = handle.push(p, None, payload);
}
assert!(
handle.wait_committed(p, last + 1, Duration::from_secs(10)),
"commit timeout (last committed: {:?})",
handle.last_committed(p)
);
shutdown.trigger();
let report = join.join().expect("join").expect("run");
assert_eq!(report.exit_code(), 0, "clean drain");
let rows_a: Vec<Vec<u8>> = script_a
.writes()
.iter()
.flat_map(|w| decode_rows(&w.payload))
.collect();
let rows_b: Vec<Vec<u8>> = script_b
.writes()
.iter()
.flat_map(|w| decode_rows(&w.payload))
.collect();
assert_eq!(
rows_a,
vec![b"apple".to_vec(), b"avocado".to_vec()],
"a-prefixed rows routed to sink a, in order"
);
assert_eq!(
rows_b,
vec![b"banana".to_vec(), b"berry".to_vec()],
"b-prefixed rows routed to sink b, in order"
);
}
#[test]
fn split_unmatched_fail_stops_the_pipeline() {
let (source, handle) = memory_source();
let (sink_a, _script_a) = capture_sink(1, 1);
const ONE_SINK: &str = r#"
pipeline: { name: split-fail, threads: 1, io_threads: 1 }
metrics: { exporter: none, listen: "127.0.0.1:0" }
checkpoint: { interval: 200ms, drain_timeout: 2s, stalled_fail_after: 3s }
source: { memory: {} }
sinks:
a: { capture: {} }
"#;
let runtime = Pipeline::from_config(PipelineConfig::from_str(ONE_SINK).expect("config"))
.expect("builder")
.add_sink("a", sink_a)
.expect("sink a")
.chains(|ctx| {
let mut split = chain_owned::<Vec<u8>, _>(BytesPassthrough)
.with_metrics(ctx.pipeline.clone(), "main")
.split(ErrorPolicy::Fail);
let a = split.add::<Owned<Vec<u8>>, _, _>(TestEncoder, KeyHashRouter, ctx.sink("a"));
split
.route(move |row: Vec<u8>, out| {
if row.first() == Some(&b'a') {
out.emit(a, row);
}
})
.build()
})
.runtime_options(RuntimeOptions {
handle_signals: false,
..RuntimeOptions::default()
})
.into_runtime(source)
.expect("into_runtime");
let run = PipelineRun::spawn(move || runtime.run());
let p = PartitionId(0);
handle.assign_lanes(&[(LaneId(0), p)]);
handle.push(p, None, b"apple"); handle.push(p, None, b"zebra");
let report = run
.wait_exit(Duration::from_secs(10))
.expect("pipeline did not stop on its own after a Fail-policy fatal")
.expect("run");
assert!(
matches!(report.state, ExitState::Failed(_)),
"unmatched record under Fail policy must fail the pipeline (got {:?})",
report.state
);
assert_eq!(report.exit_code(), 1);
}