use spate_core::config::PipelineConfig;
use spate_core::metrics::{Exporter, MetricsSettings, install};
use spate_core::ops::chain_owned;
use spate_core::pipeline::{Pipeline, RuntimeOptions};
use spate_core::record::PartitionId;
use spate_core::sink::KeyHashRouter;
use spate_core::source::LaneId;
use spate_test::{BytesPassthrough, TestEncoder, capture_sink, decode_rows, memory_source};
use std::net::{Ipv4Addr, SocketAddr};
use std::time::Duration;
const CONFIG: &str = r#"
pipeline: { name: seam-test, threads: 1, io_threads: 1 }
metrics: { exporter: prometheus, listen: "127.0.0.1:0" }
source: { memory: {} }
sink: { capture: {} }
"#;
#[test]
fn source_and_sink_receive_role_scoped_meters() {
let (source, handle) = memory_source();
let (sink, script) = capture_sink(1, 1);
let runtime = Pipeline::from_config(PipelineConfig::from_str(CONFIG).expect("config"))
.expect("builder")
.sink(sink)
.expect("sink")
.chains(|ctx| {
let chunk_cfg = ctx.chunk();
chain_owned::<Vec<u8>, _>(BytesPassthrough)
.with_metrics(ctx.pipeline, "main")
.sink(
TestEncoder,
KeyHashRouter,
chunk_cfg,
ctx.queues,
ctx.budget,
)
.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"alpha"[..], b"beta"] {
last = handle.push(p, Some(b"key"), payload);
}
assert!(
handle.wait_committed(p, last + 1, Duration::from_secs(10)),
"timed out waiting for commit (last: {:?})",
handle.last_committed(p)
);
shutdown.trigger();
let report = join.join().expect("join").expect("run");
assert_eq!(report.exit_code(), 0, "clean drain");
let rows: Vec<Vec<u8>> = script
.writes()
.iter()
.flat_map(|w| decode_rows(&w.payload))
.collect();
assert_eq!(rows, vec![b"alpha".to_vec(), b"beta".to_vec()]);
let rendered = install(&MetricsSettings {
exporter: Exporter::Prometheus,
listen: SocketAddr::from((Ipv4Addr::LOCALHOST, 0)),
..MetricsSettings::default()
})
.expect("reuse installed exporter")
.render();
assert!(
rendered.contains(
r#"spate_memory_source_opens_total{pipeline="seam-test",component="source",component_type="memory"}"#
),
"source family missing:\n{rendered}"
);
assert!(
rendered.contains(
r#"spate_capture_sink_attaches_total{pipeline="seam-test",component="sink",component_type="capture"}"#
),
"sink family missing:\n{rendered}"
);
}