#![allow(clippy::print_stdout, clippy::print_stderr)]
use spate::backpressure::InflightBudget;
use spate::metrics::{ComponentLabels, MetricsSettings, SinkShardMetrics, install};
use spate::pipeline::{PipelineRuntime, SinkRuntime, metrics_settings};
use spate::prelude::*;
use spate::sink::{SinkDrainFn, SinkPool, shard_queues};
use spate::source::LaneId;
use spate::telemetry;
use spate_test::{TestDeserializer, TestEncoder, capture_sink, memory_source, wait_until};
use std::sync::Arc;
use std::time::Duration;
const CONFIG: &str = r#"
pipeline: { name: manual-assembly-demo, threads: 1, io_threads: 1 }
admin: { listen: none }
checkpoint: { interval: 200ms }
source: { memory: {} }
sink: { capture: {} }
"#;
fn main() -> Result<(), Box<dyn std::error::Error>> {
let config = PipelineConfig::from_str(CONFIG)?;
let pipeline_name = config.pipeline.name.clone();
telemetry::init(telemetry::LogFormat::Pretty, "info");
let settings: MetricsSettings = metrics_settings(&config);
let metrics = install(&settings)?;
let io = tokio::runtime::Builder::new_multi_thread()
.worker_threads(config.pipeline.io_threads)
.thread_name("spate-io")
.enable_all()
.build()?;
let budget = Arc::new(InflightBudget::new());
let sink_options = SinkOptions::default().with_queue_capacity(16);
let (sink, script) = capture_sink(1, 1);
let sink = sink.with_pool_config({
let mut pool = SinkPoolConfig::default();
pool.batch.linger = Duration::from_millis(50); pool
});
let parts = sink.into_parts();
let shards = parts.shard_endpoints.len();
let component_type = parts.component_type.clone();
let replica_labels = parts.effective_replica_labels();
let (queues, receivers) = shard_queues(shards, sink_options.queue_capacity);
let labels = ComponentLabels::new(pipeline_name.clone(), "sink", component_type);
let shard_metrics: Vec<SinkShardMetrics> = replica_labels
.iter()
.enumerate()
.map(|(shard, replicas)| {
SinkShardMetrics::try_new(
&labels,
u32::try_from(shard).unwrap_or(u32::MAX),
replicas,
settings.e2e_basis,
)
})
.collect::<Result<_, _>>()?;
let pool = SinkPool::spawn(
Arc::new(parts.writer),
parts.shard_endpoints,
receivers,
parts.pool,
Arc::clone(&budget),
shard_metrics,
&pipeline_name,
io.handle(),
);
let drain: SinkDrainFn =
Box::new(move |deadline| Box::pin(async move { pool.drain(deadline).await }));
let sink_runtime = SinkRuntime {
queues: vec![queues.clone()],
drain,
probe: parts.probe,
};
let (source, handle) = memory_source();
let chain_name = pipeline_name.clone();
let chain_budget = Arc::clone(&budget);
let chains = move |_thread: usize| {
chain_owned::<Vec<u8>, _>(TestDeserializer::split_on(b','))
.with_metrics(chain_name.clone(), "main")
.filter(|order_id: &Vec<u8>| !order_id.is_empty())
.map(|order_id: Vec<u8>| order_id.to_ascii_uppercase())
.sink(
TestEncoder,
KeyHashRouter,
ChunkConfig::default(),
queues.clone(),
Arc::clone(&chain_budget),
)
.build()
};
let runtime = PipelineRuntime::new(config, source, chains, sink_runtime, budget)
.with_options(RuntimeOptions {
handle_signals: false, ..RuntimeOptions::default()
})
.with_io_runtime(io);
let shutdown = runtime.shutdown_handle();
let join = std::thread::spawn(move || runtime.run());
let orders = PartitionId(0);
handle.assign_lanes(&[(LaneId(0), orders)]);
let mut last = 0;
for payload in [&b"order-1,order-2"[..], b"order-3,,order-4", b"order-5"] {
last = handle.push(orders, Some(b"eu-west"), payload);
}
wait_until(Duration::from_secs(10), "the last offset to commit", || {
handle.last_committed(orders) == 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(|row| String::from_utf8_lossy(&row).into_owned())
.collect();
assert_eq!(rows.len(), 5, "six fields minus one filtered empty");
assert!(rows.contains(&"ORDER-1".to_string()));
let exposition = metrics.render();
assert!(
exposition.contains("spate_sink_"),
"the sink shard handles must render through the exporter installed by hand"
);
println!("\npipeline exit: {:?}", report.state);
println!("final watermarks: {:?}", report.final_watermarks);
println!("rows written ({}): {rows:?}", rows.len());
Ok(())
}
#[cfg(test)]
mod tests {
#[test]
fn runs_to_completion() {
super::main().expect("the example must run clean");
}
}