use std::collections::BTreeMap;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use crate::bus::lock::lock;
#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
pub enum RuntimeDirection {
Publish,
Subscribe,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
pub enum RuntimeBufferKind {
Outbound,
Latest,
Subscriber,
}
#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord)]
pub struct RuntimeMetricKey {
pub topic: String,
pub direction: RuntimeDirection,
pub buffer_kind: RuntimeBufferKind,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct RuntimeMetricSnapshot {
pub key: RuntimeMetricKey,
pub count: u64,
pub drops: u64,
pub latest_overwrites: u64,
pub bounded_evictions: u64,
pub capacity: u64,
pub current_depth: u64,
pub high_water_depth: u64,
pub decode_errors: u64,
pub timeline_filtered: u64,
}
#[derive(Debug, Default)]
struct Counters {
count: AtomicU64,
drops: AtomicU64,
latest_overwrites: AtomicU64,
bounded_evictions: AtomicU64,
capacity: AtomicU64,
current_depth: AtomicU64,
high_water_depth: AtomicU64,
decode_errors: AtomicU64,
timeline_filtered: AtomicU64,
}
#[derive(Clone, Debug)]
pub(crate) struct RuntimeMetricHandle {
counters: Arc<Counters>,
local_depth: Option<Arc<AtomicU64>>,
}
impl RuntimeMetricHandle {
pub(crate) fn record_message(&self) {
self.counters.count.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn record_drop(&self) {
self.counters.drops.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn record_decode_error(&self) {
self.counters.decode_errors.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn record_timeline_filtered(&self, count: u64) {
self.counters
.timeline_filtered
.fetch_add(count, Ordering::Relaxed);
}
pub(crate) fn record_latest(&self, overwrote: bool) {
self.record_message();
if overwrote {
self.record_latest_overwrite();
}
self.set_inbound_depth(1);
}
pub(crate) fn record_pending(&self) {
self.record_message();
}
pub(crate) fn record_latest_depth(&self, occupied: bool) {
self.set_inbound_depth(u64::from(occupied));
}
pub(crate) fn record_subscriber(&self, evicted: bool, current_depth: usize) {
self.record_message();
if evicted {
self.counters
.bounded_evictions
.fetch_add(1, Ordering::Relaxed);
self.record_drop();
}
self.set_inbound_depth(u64::try_from(current_depth).unwrap_or(u64::MAX));
}
pub(crate) fn record_subscriber_pop(&self, current_depth: usize) {
self.set_inbound_depth(u64::try_from(current_depth).unwrap_or(u64::MAX));
}
pub(crate) fn record_latest_overwrite(&self) {
self.counters
.latest_overwrites
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn record_bounded_eviction(&self) {
self.counters
.bounded_evictions
.fetch_add(1, Ordering::Relaxed);
self.record_drop();
}
pub(crate) fn enqueue_started(&self) {
let current = self
.counters
.current_depth
.fetch_add(1, Ordering::Relaxed)
.saturating_add(1);
update_max(&self.counters.high_water_depth, current);
}
pub(crate) fn enqueue_finished(&self) {
let _ = self.counters.current_depth.fetch_update(
Ordering::Relaxed,
Ordering::Relaxed,
|value| Some(value.saturating_sub(1)),
);
}
fn set_inbound_depth(&self, depth: u64) {
let Some(local) = self.local_depth.as_ref() else {
debug_assert!(
false,
"outbound runtime metric cannot update an inbound depth gauge"
);
return;
};
let previous = local.swap(depth, Ordering::Relaxed);
let current = if depth >= previous {
self.counters
.current_depth
.fetch_add(depth - previous, Ordering::Relaxed)
.saturating_add(depth - previous)
} else {
self.counters
.current_depth
.fetch_sub(previous - depth, Ordering::Relaxed)
.saturating_sub(previous - depth)
};
update_max(&self.counters.high_water_depth, current);
}
}
#[derive(Debug, Default)]
pub(crate) struct RuntimeMetrics {
rows: Mutex<BTreeMap<RuntimeMetricKey, Arc<Counters>>>,
}
impl RuntimeMetrics {
pub(crate) fn register_outbound(&self, topic: &str, capacity: usize) -> RuntimeMetricHandle {
self.register(
RuntimeMetricKey {
topic: topic.to_string(),
direction: RuntimeDirection::Publish,
buffer_kind: RuntimeBufferKind::Outbound,
},
capacity,
false,
)
}
pub(crate) fn register_latest(&self, topic: &str) -> RuntimeMetricHandle {
self.register(
RuntimeMetricKey {
topic: topic.to_string(),
direction: RuntimeDirection::Subscribe,
buffer_kind: RuntimeBufferKind::Latest,
},
1,
true,
)
}
pub(crate) fn register_subscriber(&self, topic: &str, capacity: usize) -> RuntimeMetricHandle {
self.register(
RuntimeMetricKey {
topic: topic.to_string(),
direction: RuntimeDirection::Subscribe,
buffer_kind: RuntimeBufferKind::Subscriber,
},
capacity,
true,
)
}
fn register(
&self,
key: RuntimeMetricKey,
capacity: usize,
additive_capacity: bool,
) -> RuntimeMetricHandle {
let mut rows = lock(&self.rows);
let counters = rows.entry(key).or_default();
let capacity = u64::try_from(capacity).unwrap_or(u64::MAX);
if additive_capacity {
counters.capacity.fetch_add(capacity, Ordering::Relaxed);
} else {
update_max(&counters.capacity, capacity);
}
RuntimeMetricHandle {
counters: Arc::clone(counters),
local_depth: additive_capacity.then(|| Arc::new(AtomicU64::new(0))),
}
}
pub(crate) fn take(&self) -> Vec<RuntimeMetricSnapshot> {
let rows = lock(&self.rows);
rows.iter()
.map(|(key, counters)| {
let current_depth = counters.current_depth.load(Ordering::Relaxed);
RuntimeMetricSnapshot {
key: key.clone(),
count: counters.count.swap(0, Ordering::Relaxed),
drops: counters.drops.swap(0, Ordering::Relaxed),
latest_overwrites: counters.latest_overwrites.swap(0, Ordering::Relaxed),
bounded_evictions: counters.bounded_evictions.swap(0, Ordering::Relaxed),
capacity: counters.capacity.load(Ordering::Relaxed),
current_depth,
high_water_depth: counters
.high_water_depth
.swap(current_depth, Ordering::Relaxed),
decode_errors: counters.decode_errors.swap(0, Ordering::Relaxed),
timeline_filtered: counters.timeline_filtered.swap(0, Ordering::Relaxed),
}
})
.collect()
}
}
fn update_max(target: &AtomicU64, value: u64) {
let _ = target.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
(value > current).then_some(value)
});
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
use serial_test::serial;
use zenoh::bytes::Encoding;
use zenoh::key_expr::OwnedKeyExpr;
use crate::bus::abi::CodecId;
use crate::bus::handle::publisher::StatePublisher;
use crate::bus::handle::subscriber::{Latest, Subscriber};
use crate::bus::session::BusOwner;
use crate::bus::test_support::{
TARGET_TOPIC, Target, bound, metadata, participant_config, step,
};
#[test]
fn every_direction_and_buffer_kind_is_accounted_for_on_the_wire() {
const DIRECTIONS: [RuntimeDirection; 2] =
[RuntimeDirection::Publish, RuntimeDirection::Subscribe];
const BUFFER_KINDS: [RuntimeBufferKind; 3] = [
RuntimeBufferKind::Outbound,
RuntimeBufferKind::Latest,
RuntimeBufferKind::Subscriber,
];
for direction in DIRECTIONS {
let name = match direction {
RuntimeDirection::Publish => "publish",
RuntimeDirection::Subscribe => "subscribe",
};
assert!(!name.is_empty());
}
for buffer_kind in BUFFER_KINDS {
let name = match buffer_kind {
RuntimeBufferKind::Outbound => "outbound",
RuntimeBufferKind::Latest => "latest",
RuntimeBufferKind::Subscriber => "subscriber",
};
assert!(!name.is_empty());
}
assert_eq!(
DIRECTIONS.len(),
2,
"a new RuntimeDirection needs a wire variant in the runtime contract \
family and an arm in the runtime performance rollup"
);
assert_eq!(
BUFFER_KINDS.len(),
3,
"a new RuntimeBufferKind needs a wire variant in the runtime contract \
family and an arm in the runtime performance rollup"
);
}
#[test]
fn identical_keys_aggregate_and_quiet_rows_persist() {
let metrics = RuntimeMetrics::default();
let first = metrics.register_subscriber("robot/drive/state", 4);
let second = metrics.register_subscriber("robot/drive/state", 8);
first.record_subscriber(false, 1);
second.record_subscriber(true, 8);
let rows = metrics.take();
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].count, 2);
assert_eq!(rows[0].bounded_evictions, 1);
assert_eq!(rows[0].drops, 1);
assert_eq!(rows[0].capacity, 12);
assert_eq!(rows[0].current_depth, 9);
assert_eq!(rows[0].high_water_depth, 9);
let quiet = metrics.take();
assert_eq!(quiet.len(), 1);
assert_eq!(quiet[0].count, 0);
assert_eq!(quiet[0].capacity, 12);
assert_eq!(quiet[0].current_depth, 9);
}
#[test]
fn outbound_capacity_is_a_non_additive_view_of_one_shared_queue() {
let metrics = RuntimeMetrics::default();
let first = metrics.register_outbound("robot/drive/target", 1_024);
let second = metrics.register_outbound("robot/drive/target", 1_024);
let _other = metrics.register_outbound("robot/motion/target", 1_024);
first.enqueue_started();
second.enqueue_started();
let rows = metrics.take();
assert_eq!(rows.len(), 2);
assert!(rows.iter().all(|row| row.capacity == 1_024));
assert_eq!(rows[0].current_depth, 2);
}
#[test]
fn fixed_setup_rows_persist_after_the_declaring_handle_is_dropped() {
let metrics = RuntimeMetrics::default();
{
let _declared = metrics.register_latest("robot/drive/state");
}
let rows = metrics.take();
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].capacity, 1);
assert_eq!(rows[0].current_depth, 0);
assert_eq!(metrics.take().len(), 1);
}
#[serial]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn runtime_metrics_cover_quiet_latest_overwrite_eviction_and_decode_error_rows() {
let (owner, bus) = BusOwner::open(participant_config("metrics")).await.unwrap();
let pub_topic = bound::<Target>(TARGET_TOPIC).owner();
let sub_topic = bound::<Target>(TARGET_TOPIC).client();
let publisher = StatePublisher::<Target>::new(bus.clone(), &pub_topic).unwrap();
let latest = Latest::<Target>::new(&bus, &sub_topic).await.unwrap();
let subscriber = Subscriber::<Target>::new(&bus, &sub_topic).await.unwrap();
let quiet = bus.take_runtime_metrics().unwrap();
assert_eq!(quiet.len(), 3);
assert!(quiet.iter().all(|row| row.count == 0));
for value in [1.0_f32, 2.0, 3.0] {
publisher
.publish(
&step(1, value as u64),
Target {
linear_x_mps: value,
angular_z_radps: 0.0,
},
)
.unwrap();
for _ in 0..50 {
if latest
.latest()
.is_some_and(|sample| sample.linear_x_mps == value)
{
break;
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
}
for _ in 0..50 {
if latest
.latest()
.is_some_and(|sample| sample.linear_x_mps == 3.0)
&& bus
.health()
.inbound_drops
.load(std::sync::atomic::Ordering::Relaxed)
>= 2
{
break;
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
bus.session()
.unwrap()
.put(
OwnedKeyExpr::new(bus.full_key(TARGET_TOPIC)).unwrap(),
vec![0xc1_u8],
)
.encoding(Encoding::from(CodecId::MessagePack.encoding_string()))
.attachment(metadata().encode().expect("test metadata encodes"))
.await
.unwrap();
for _ in 0..50 {
if bus
.health()
.decode_errors
.load(std::sync::atomic::Ordering::Relaxed)
>= 2
{
break;
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
let rows = bus.take_runtime_metrics().unwrap();
let outbound = rows
.iter()
.find(|row| row.key.direction == RuntimeDirection::Publish)
.unwrap();
assert_eq!(outbound.key.buffer_kind, RuntimeBufferKind::Outbound);
assert_eq!(outbound.key.topic, TARGET_TOPIC);
assert_eq!(outbound.count, 3);
let latest_row = rows
.iter()
.find(|row| row.key.buffer_kind == RuntimeBufferKind::Latest)
.unwrap();
assert_eq!(latest_row.count, 3);
assert_eq!(latest_row.latest_overwrites, 2);
assert_eq!(latest_row.capacity, 1);
assert_eq!(latest_row.current_depth, 1);
assert_eq!(latest_row.decode_errors, 1);
let subscriber_row = rows
.iter()
.find(|row| row.key.buffer_kind == RuntimeBufferKind::Subscriber)
.unwrap();
assert_eq!(subscriber_row.count, 3);
assert_eq!(subscriber_row.bounded_evictions, 2);
assert_eq!(subscriber_row.drops, 2);
assert_eq!(subscriber_row.current_depth, 1);
assert_eq!(subscriber_row.high_water_depth, 1);
assert_eq!(subscriber_row.decode_errors, 1);
drop(subscriber);
owner.close().await;
}
}