use metrics::{counter, gauge, histogram};
use strum_macros::IntoStaticStr;
pub(super) const FRAMES_ACCEPTED: &str = "price_level_stream_frames_accepted_total";
pub(super) const FRAMES_REJECTED: &str = "price_level_stream_frames_rejected_total";
pub(super) const FRAME_AGE: &str = "price_level_stream_frame_age_seconds";
pub(super) const LAST_SEEN: &str = "price_level_stream_last_seen_timestamp_seconds";
pub(super) const SERVED_COMPONENTS: &str = "price_level_stream_served_components";
pub(super) const STALE_REMOVALS: &str = "price_level_stream_stale_removals_total";
pub(super) const SERVING_STATE: &str = "price_level_stream_serving_state";
pub(super) const RECONNECTS: &str = "price_level_stream_reconnects_total";
pub(super) const UNREGISTERED_PAMM_ENTRIES: &str =
"price_level_stream_unregistered_pamm_entries_total";
#[derive(Clone, Copy, Debug, PartialEq, Eq, IntoStaticStr)]
#[strum(serialize_all = "snake_case")]
pub(super) enum RejectReason {
ParseError,
TooOld,
InFuture,
OutOfOrder,
BlockRegression,
BlockJump,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, IntoStaticStr)]
#[strum(serialize_all = "snake_case")]
pub(super) enum ReconnectReason {
IdleTimeout,
Ended,
Closed,
ReadError,
ConnectFailed,
ConnectTimeout,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(super) enum ServingState {
Unserved = 0,
Serving = 1,
}
pub(super) fn record_frame_accepted() {
counter!(FRAMES_ACCEPTED).increment(1);
}
pub(super) fn record_frame_rejected(reason: RejectReason) {
counter!(FRAMES_REJECTED, "reason" => <&'static str>::from(reason)).increment(1);
}
pub(super) fn record_frame_age(age_seconds: f64) {
histogram!(FRAME_AGE).record(age_seconds);
}
pub(super) fn record_last_seen(venue: &str, unix_seconds: u64) {
gauge!(LAST_SEEN, "venue" => venue.to_string()).set(unix_seconds as f64);
}
pub(super) fn record_served_components(venue: &str, count: usize) {
gauge!(SERVED_COMPONENTS, "venue" => venue.to_string()).set(count as f64);
}
pub(super) fn record_stale_removal(venue: &str) {
counter!(STALE_REMOVALS, "venue" => venue.to_string()).increment(1);
}
pub(super) fn record_serving_state(state: ServingState) {
gauge!(SERVING_STATE).set(state as u8 as f64);
}
pub(super) fn record_reconnect(reason: ReconnectReason) {
counter!(RECONNECTS, "reason" => <&'static str>::from(reason)).increment(1);
}
pub(super) fn record_unregistered_pamm() {
counter!(UNREGISTERED_PAMM_ENTRIES).increment(1);
}
#[cfg(test)]
pub(super) mod recorded {
use std::{
collections::{BTreeMap, HashMap},
future::Future,
};
use metrics_util::{
debugging::{DebugValue, DebuggingRecorder, Snapshot},
MetricKind,
};
pub(in super::super) type SnapshotMap =
HashMap<(MetricKind, String, BTreeMap<String, String>), DebugValue>;
pub(in super::super) fn record_async<T>(future: impl Future<Output = T>) -> (T, SnapshotMap) {
let recorder = DebuggingRecorder::new();
let snapshotter = recorder.snapshotter();
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
let output = metrics::with_local_recorder(&recorder, || runtime.block_on(future));
(output, snapshot_map(snapshotter.snapshot()))
}
pub(in super::super) fn snapshot_map(snapshot: Snapshot) -> SnapshotMap {
snapshot
.into_vec()
.into_iter()
.map(|(composite_key, _unit, _description, value)| {
let name = composite_key.key().name().to_string();
let labels = composite_key
.key()
.labels()
.map(|label| (label.key().to_string(), label.value().to_string()))
.collect();
((composite_key.kind(), name, labels), value)
})
.collect()
}
fn labels(pairs: &[(&str, &str)]) -> BTreeMap<String, String> {
pairs
.iter()
.map(|(key, value)| ((*key).to_string(), (*value).to_string()))
.collect()
}
pub(in super::super) fn counter_value(
snapshot: &SnapshotMap,
name: &str,
label_pairs: &[(&str, &str)],
) -> u64 {
match snapshot.get(&(MetricKind::Counter, name.to_string(), labels(label_pairs))) {
Some(DebugValue::Counter(value)) => *value,
Some(DebugValue::Gauge(_)) | Some(DebugValue::Histogram(_)) | None => 0,
}
}
pub(in super::super) fn gauge_value(
snapshot: &SnapshotMap,
name: &str,
label_pairs: &[(&str, &str)],
) -> f64 {
match snapshot.get(&(MetricKind::Gauge, name.to_string(), labels(label_pairs))) {
Some(DebugValue::Gauge(value)) => value.into_inner(),
Some(DebugValue::Counter(_)) | Some(DebugValue::Histogram(_)) | None => f64::NAN,
}
}
pub(in super::super) fn histogram_values(
snapshot: &SnapshotMap,
name: &str,
label_pairs: &[(&str, &str)],
) -> Vec<f64> {
match snapshot.get(&(MetricKind::Histogram, name.to_string(), labels(label_pairs))) {
Some(DebugValue::Histogram(values)) => values
.iter()
.map(|value| value.into_inner())
.collect(),
Some(DebugValue::Counter(_)) | Some(DebugValue::Gauge(_)) | None => Vec::new(),
}
}
}
#[cfg(test)]
mod tests {
use super::{recorded::*, *};
#[test]
fn every_helper_emits_its_named_metric() {
let ((), snapshot) = record_async(async {
record_frame_accepted();
record_frame_rejected(RejectReason::TooOld);
record_frame_rejected(RejectReason::TooOld);
record_frame_age(0.25);
record_last_seen("fermiswap", 1_700_000_000);
record_served_components("fermiswap", 3);
record_stale_removal("fermiswap");
record_serving_state(ServingState::Serving);
record_reconnect(ReconnectReason::IdleTimeout);
record_unregistered_pamm();
});
assert_eq!(counter_value(&snapshot, "price_level_stream_frames_accepted_total", &[]), 1);
assert_eq!(
counter_value(
&snapshot,
"price_level_stream_frames_rejected_total",
&[("reason", "too_old")]
),
2
);
assert_eq!(
histogram_values(&snapshot, "price_level_stream_frame_age_seconds", &[]),
vec![0.25]
);
assert_eq!(
gauge_value(
&snapshot,
"price_level_stream_last_seen_timestamp_seconds",
&[("venue", "fermiswap")]
),
1_700_000_000.0
);
assert_eq!(
gauge_value(
&snapshot,
"price_level_stream_served_components",
&[("venue", "fermiswap")]
),
3.0
);
assert_eq!(
counter_value(
&snapshot,
"price_level_stream_stale_removals_total",
&[("venue", "fermiswap")]
),
1
);
assert_eq!(gauge_value(&snapshot, "price_level_stream_serving_state", &[]), 1.0);
assert_eq!(
counter_value(
&snapshot,
"price_level_stream_reconnects_total",
&[("reason", "idle_timeout")]
),
1
);
assert_eq!(
counter_value(&snapshot, "price_level_stream_unregistered_pamm_entries_total", &[]),
1
);
}
}