#![allow(clippy::print_stdout, clippy::print_stderr)]
use spate::prelude::*;
use spate::source::LaneId;
use spate_test::{TestDeserializer, TestEncoder, capture_sink, memory_source};
use std::time::{Duration, Instant};
const CONFIG: &str = r#"
pipeline: { name: memory-demo, threads: 1 }
checkpoint: { interval: 200ms }
source: { memory: {} }
sink: { capture: {} }
"#;
fn main() -> Result<(), Box<dyn std::error::Error>> {
spate::telemetry::init(spate::telemetry::LogFormat::Pretty, "info");
let pipeline = Pipeline::from_config(PipelineConfig::from_str(CONFIG)?)?;
let (source, handle) = memory_source();
let (sink, script) = capture_sink(1, 1);
let pool_cfg = {
let mut cfg = SinkPoolConfig::default();
cfg.batch.linger = Duration::from_millis(50); cfg
};
let sink = sink.with_pool_config(pool_cfg);
let runtime = pipeline
.sink(sink)?
.chains(|ctx| {
let chunk_cfg = ctx.chunk();
chain_owned::<Vec<u8>, _>(TestDeserializer::split_on(b','))
.with_metrics(ctx.pipeline, "main")
.filter(|word: &Vec<u8>| !word.is_empty())
.map(|word: Vec<u8>| word.to_ascii_uppercase())
.sink(
TestEncoder,
KeyHashRouter,
chunk_cfg,
ctx.queues,
ctx.budget,
)
.build()
})
.runtime_options(RuntimeOptions {
handle_signals: false, ..RuntimeOptions::default()
})
.into_runtime(source)?;
let shutdown = runtime.shutdown_handle();
let join = std::thread::spawn(move || runtime.run());
let p0 = PartitionId(0);
handle.assign_lanes(&[(LaneId(0), p0)]);
let mut last = 0;
for payload in [
&b"alpha,beta,gamma"[..],
b"delta,,epsilon",
b"zeta,eta,theta",
] {
last = handle.push(p0, Some(b"demo"), payload);
}
let deadline = Instant::now() + Duration::from_secs(10);
while handle.last_committed(p0) != Some(last + 1) {
assert!(Instant::now() < deadline, "commit not observed in time");
std::thread::sleep(Duration::from_millis(20));
}
shutdown.trigger();
let report = join.join().expect("pipeline thread")?;
let rows: Vec<String> = script
.writes()
.iter()
.flat_map(|w| spate_test::decode_rows(&w.payload))
.map(|r| String::from_utf8_lossy(&r).into_owned())
.collect();
println!("\npipeline exit: {:?}", report.state);
println!("final watermarks: {:?}", report.final_watermarks);
println!("rows written ({}): {rows:?}", rows.len());
assert_eq!(rows.len(), 8, "nine words minus one filtered empty");
assert!(rows.contains(&"ALPHA".to_string()));
Ok(())
}