#![allow(clippy::print_stdout, clippy::print_stderr)]
use spate::prelude::*;
use spate::source::LaneId;
use spate_test::{TestDeserializer, TestEncoder, capture_sink, memory_source, wait_until};
use std::time::Duration;
const CONFIG: &str = r#"
pipeline: { name: storefront-orders, threads: 1 }
admin: { listen: none }
checkpoint: { interval: 200ms }
source: { memory: {} }
sink: { capture: {} }
"#;
struct Event<'a> {
kind: &'a str,
order_id: &'a str,
line_items: u64,
amount: f64,
}
fn parse(record: &[u8]) -> Option<Event<'_>> {
let mut fields = std::str::from_utf8(record).ok()?.split(',');
Some(Event {
kind: fields.next()?,
order_id: fields.next()?,
line_items: fields.next()?.parse().ok()?,
amount: fields.next()?.parse().ok()?,
})
}
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 exposition = pipeline.metrics().clone();
let (source, handle) = memory_source();
let (sink, script) = capture_sink(1, 1);
let sink = sink.with_pool_config({
let mut cfg = SinkPoolConfig::default();
cfg.batch.linger = Duration::from_millis(50); cfg
});
let runtime = pipeline
.sink(sink)?
.chains(|ctx| {
let chunk_cfg = ctx.chunk();
let meter = ctx.meter("checkout", "inspect");
let orders = meter.counter("orders_total", &[("channel", "web".into())]);
let refunds = meter.counter("refunds_total", &[("channel", "web".into())]);
let line_items = meter.counter("line_items_total", &[]);
let order_value = meter.histogram("order_value_dollars", &[]);
chain_owned::<Vec<u8>, _>(TestDeserializer::split_on(b'\n'))
.with_metrics(ctx.pipeline, "checkout")
.inspect(move |record: &Vec<u8>| {
let Some(event) = parse(record) else { return };
match event.kind {
"order" => {
orders.increment(1);
line_items.increment(event.line_items);
order_value.record(event.amount);
}
"refund" => refunds.increment(1),
_ => {}
}
})
.filter(|record: &Vec<u8>| parse(record).is_some_and(|event| event.kind == "order"))
.map(|record: Vec<u8>| {
let event = parse(&record).expect("filtered to parseable orders");
format!("{}={:.2}", event.order_id, event.amount).into_bytes()
})
.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 last = handle.push(
p0,
Some(b"storefront"),
concat!(
"order,1042,3,129.99\n",
"order,1043,1,19.50\n",
"refund,1042,1,19.50\n",
"order,1044,2,64.00\n",
"heartbeat"
)
.as_bytes(),
);
wait_until(Duration::from_secs(10), "the offset to commit", || {
handle.last_committed(p0) == Some(last + 1)
});
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();
assert_eq!(
rows.len(),
3,
"three orders; the refund and the junk line went"
);
let scraped = exposition.render();
let labels = r#"pipeline="storefront-orders",component="checkout",component_type="inspect""#;
for series in [
format!(r#"spate_custom_orders_total{{{labels},channel="web"}} 3"#),
format!(r#"spate_custom_refunds_total{{{labels},channel="web"}} 1"#),
format!(r#"spate_custom_line_items_total{{{labels}}} 6"#),
format!(r#"spate_custom_order_value_dollars_count{{{labels}}} 3"#),
] {
assert!(
scraped.contains(&series),
"missing from the exposition: {series}"
);
}
for line in scraped.lines().filter(|l| l.starts_with("spate_custom_")) {
println!("{line}");
}
println!("\npipeline exit: {:?}", report.state);
println!("rows written ({}): {rows:?}", rows.len());
Ok(())
}
#[cfg(test)]
mod tests {
#[test]
fn runs_to_completion() {
super::main().expect("the example must run clean");
}
}