use crate::events::StorefrontEvent;
use spate_core::metrics::{Counter, Gauge, Meter};
#[derive(Clone, Debug)]
pub(crate) struct LaneCounters {
order_placed: Counter,
payment_captured: Counter,
refund_issued: Counter,
pub(crate) ticks: Counter,
pub(crate) tick_overruns: Counter,
}
pub(crate) fn kind(event: &StorefrontEvent) -> usize {
match event {
StorefrontEvent::OrderPlaced(_) => 0,
StorefrontEvent::PaymentCaptured(_) => 1,
StorefrontEvent::RefundIssued(_) => 2,
}
}
impl LaneCounters {
pub(crate) fn add_generated(&self, tally: [u64; 3]) {
let handles = [
&self.order_placed,
&self.payment_captured,
&self.refund_issued,
];
for (handle, count) in handles.into_iter().zip(tally) {
if count > 0 {
handle.increment(count);
}
}
}
}
#[derive(Debug)]
pub(crate) struct DatagenMetrics {
counters: LaneCounters,
events_remaining: Gauge,
open_orders: Gauge,
committed_offset: Vec<Gauge>,
}
impl DatagenMetrics {
pub(crate) fn new(
meter: &Meter,
partitions: u32,
per_partition_detail: bool,
) -> DatagenMetrics {
let generated = |event: &'static str| {
meter.counter("events_generated_total", &[("event", event.into())])
};
DatagenMetrics {
counters: LaneCounters {
order_placed: generated("order_placed"),
payment_captured: generated("payment_captured"),
refund_issued: generated("refund_issued"),
ticks: meter.counter("ticks_total", &[]),
tick_overruns: meter.counter("tick_overrun_total", &[]),
},
events_remaining: meter.gauge("events_remaining", &[]),
open_orders: meter.gauge("open_orders", &[]),
committed_offset: if per_partition_detail {
(0..partitions)
.map(|p| {
meter.gauge("committed_offset", &[("partition", p.to_string().into())])
})
.collect()
} else {
Vec::new()
},
}
}
pub(crate) fn counters(&self) -> LaneCounters {
self.counters.clone()
}
pub(crate) fn publish(&self, events_remaining: u64, open_orders: u64) {
self.events_remaining.set(events_remaining as f64);
self.open_orders.set(open_orders as f64);
}
pub(crate) fn set_committed(&self, partition: u32, offset: i64) {
if let Some(gauge) = self.committed_offset.get(partition as usize) {
gauge.set(offset as f64);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::events::{OrderPlaced, PaymentCaptured, RefundIssued};
use std::borrow::Cow;
fn render(f: impl FnOnce()) -> String {
let recorder = metrics_exporter_prometheus::PrometheusBuilder::new().build_recorder();
let handle = recorder.handle();
metrics::with_local_recorder(&recorder, f);
handle.run_upkeep();
handle.render()
}
fn meter(component: &'static str) -> Meter {
Meter::with_namespace("datagen", "orders", component, "datagen")
}
fn one_of_each() -> [StorefrontEvent; 3] {
[
StorefrontEvent::OrderPlaced(OrderPlaced {
order_id: 1,
customer_id: 2,
region: Cow::Borrowed("eu-west"),
placed_at: 0,
lines: Vec::new(),
}),
StorefrontEvent::PaymentCaptured(PaymentCaptured {
order_id: 1,
amount_cents: 10,
}),
StorefrontEvent::RefundIssued(RefundIssued {
order_id: 1,
amount_cents: 5,
reason: Cow::Borrowed("damaged"),
}),
]
}
#[test]
fn every_event_kind_counts_against_its_own_pre_registered_handle() {
let rendered = render(|| {
let metrics = DatagenMetrics::new(&meter("events"), 2, false);
let counters = metrics.counters();
let mut tally = [0u64; 3];
for (repeats, event) in one_of_each().iter().enumerate() {
tally[kind(event)] += repeats as u64 + 1;
}
counters.add_generated(tally);
counters.ticks.increment(4);
counters.tick_overruns.increment(1);
});
for (event, count) in [
("order_placed", 1),
("payment_captured", 2),
("refund_issued", 3),
] {
let want = format!(
r#"spate_datagen_events_generated_total{{pipeline="orders",component="events",component_type="datagen",event="{event}"}} {count}"#
);
assert!(rendered.contains(&want), "missing {want} in:\n{rendered}");
}
assert!(rendered.contains("spate_datagen_ticks_total"), "{rendered}");
assert!(
rendered.contains("spate_datagen_tick_overrun_total"),
"{rendered}"
);
}
#[test]
fn the_control_plane_gauges_publish_what_it_was_given() {
let rendered = render(|| {
DatagenMetrics::new(&meter("gauges"), 2, false).publish(97, 12);
});
assert!(
rendered.contains("spate_datagen_events_remaining{") && rendered.contains("} 97"),
"{rendered}"
);
assert!(
rendered.contains("spate_datagen_open_orders{") && rendered.contains("} 12"),
"{rendered}"
);
assert!(!rendered.contains("lag"), "{rendered}");
}
#[test]
fn committed_offset_is_gated_on_per_partition_detail() {
let off =
render(|| DatagenMetrics::new(&meter("detail-off"), 2, false).set_committed(1, 9));
assert!(!off.contains("committed_offset"), "{off}");
let on = render(|| {
let metrics = DatagenMetrics::new(&meter("detail-on"), 2, true);
metrics.set_committed(0, 5);
metrics.set_committed(1, 9);
metrics.set_committed(7, 11);
});
for (partition, offset) in [(0, 5), (1, 9)] {
let want = format!(
r#"spate_datagen_committed_offset{{pipeline="orders",component="detail-on",component_type="datagen",partition="{partition}"}} {offset}"#
);
assert!(on.contains(&want), "missing {want} in:\n{on}");
}
assert!(!on.contains(r#"partition="7""#), "{on}");
}
}