use std::sync::Arc;
use std::sync::Mutex;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, Instant};
use anyhow::{Result, anyhow};
use arc_swap::ArcSwap;
use tokio::sync::broadcast;
use zenoh::Session;
use zenoh::sample::SampleKind;
use crate::stats::StatsTable;
use crate::tree::KeyTreeSnapshot;
#[derive(Debug, Clone)]
pub struct SampleView {
pub key: String,
pub payload: zenoh::bytes::ZBytes,
pub encoding: String,
pub kind: SampleKind,
pub timestamp: Option<zenoh::time::Timestamp>,
}
#[derive(Debug, Clone)]
pub enum FleetEvent {
Sample(Arc<SampleView>),
NodeUp(String),
NodeDown(String),
StatsTick,
}
#[derive(Debug, Clone)]
pub struct MonitorSpec {
pub selectors: Vec<String>,
pub liveliness: Option<String>,
pub stats_tick: Duration,
pub capacity: usize,
}
impl Default for MonitorSpec {
fn default() -> Self {
MonitorSpec {
selectors: Vec::new(),
liveliness: None,
stats_tick: Duration::from_millis(250),
capacity: 1024,
}
}
}
pub struct MonitorCore {
tx: broadcast::Sender<FleetEvent>,
stats: Mutex<StatsTable>,
tree: ArcSwap<KeyTreeSnapshot>,
dropped: AtomicU64,
}
impl MonitorCore {
pub fn new(capacity: usize) -> Arc<MonitorCore> {
let (tx, _) = broadcast::channel(capacity.max(2));
Arc::new(MonitorCore {
tx,
stats: Mutex::new(StatsTable::new()),
tree: ArcSwap::from_pointee(KeyTreeSnapshot::default()),
dropped: AtomicU64::new(0),
})
}
pub fn ingest(&self, view: SampleView, sn: Option<u32>) {
{
let mut stats = self.stats.lock().expect("stats lock");
stats.record(&view.key, view.payload.len(), sn, Instant::now());
}
let _ = self.tx.send(FleetEvent::Sample(Arc::new(view)));
}
pub fn node_event(&self, key: String, up: bool) {
let _ = self.tx.send(if up {
FleetEvent::NodeUp(key)
} else {
FleetEvent::NodeDown(key)
});
}
pub fn tick(&self) {
let snapshot = {
let stats = self.stats.lock().expect("stats lock");
KeyTreeSnapshot::build(&stats)
};
self.tree.store(Arc::new(snapshot));
let _ = self.tx.send(FleetEvent::StatsTick);
}
pub fn tree(&self) -> Arc<KeyTreeSnapshot> {
self.tree.load_full()
}
pub fn with_stats<R>(&self, f: impl FnOnce(&StatsTable) -> R) -> R {
f(&self.stats.lock().expect("stats lock"))
}
pub fn dropped(&self) -> u64 {
self.dropped.load(Ordering::Relaxed)
}
pub fn events(self: &Arc<Self>) -> EventStream {
EventStream {
rx: self.tx.subscribe(),
core: Arc::clone(self),
}
}
}
pub struct EventStream {
rx: broadcast::Receiver<FleetEvent>,
core: Arc<MonitorCore>,
}
#[derive(Debug, Clone)]
pub enum StreamItem {
Event(FleetEvent),
Dropped(u64),
}
impl EventStream {
pub async fn recv(&mut self) -> Option<StreamItem> {
match self.rx.recv().await {
Ok(ev) => Some(StreamItem::Event(ev)),
Err(broadcast::error::RecvError::Lagged(n)) => {
self.core.dropped.fetch_add(n, Ordering::Relaxed);
Some(StreamItem::Dropped(n))
}
Err(broadcast::error::RecvError::Closed) => None,
}
}
}
pub struct Monitor {
core: Arc<MonitorCore>,
tasks: Vec<tokio::task::JoinHandle<()>>,
}
impl Monitor {
pub async fn start(session: &Session, spec: MonitorSpec) -> Result<Monitor> {
let core = MonitorCore::new(spec.capacity);
let mut tasks = Vec::new();
for selector in &spec.selectors {
let subscriber = session
.declare_subscriber(selector)
.await
.map_err(|e| anyhow!("subscribe {selector}: {e}"))?;
let core = Arc::clone(&core);
tasks.push(tokio::spawn(async move {
while let Ok(sample) = subscriber.recv_async().await {
let sn = sample.source_info().map(|si| si.source_sn());
core.ingest(
SampleView {
key: sample.key_expr().as_str().to_string(),
payload: sample.payload().clone(),
encoding: sample.encoding().to_string(),
kind: sample.kind(),
timestamp: sample.timestamp().copied(),
},
sn,
);
}
}));
}
if let Some(liveliness_sel) = &spec.liveliness {
let subscriber = session
.liveliness()
.declare_subscriber(liveliness_sel)
.history(true)
.await
.map_err(|e| anyhow!("liveliness subscribe {liveliness_sel}: {e}"))?;
let core = Arc::clone(&core);
tasks.push(tokio::spawn(async move {
while let Ok(sample) = subscriber.recv_async().await {
let key = sample.key_expr().as_str().to_string();
core.node_event(key, sample.kind() == SampleKind::Put);
}
}));
}
{
let core = Arc::clone(&core);
let period = spec.stats_tick;
tasks.push(tokio::spawn(async move {
let mut interval = tokio::time::interval(period);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
interval.tick().await;
core.tick();
}
}));
}
Ok(Monitor { core, tasks })
}
pub fn core(&self) -> &Arc<MonitorCore> {
&self.core
}
pub fn events(&self) -> EventStream {
self.core.events()
}
pub fn tree(&self) -> Arc<KeyTreeSnapshot> {
self.core.tree()
}
pub fn stop(self) {
for t in &self.tasks {
t.abort();
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn view(key: &str, len: usize) -> SampleView {
SampleView {
key: key.to_string(),
payload: zenoh::bytes::ZBytes::from(vec![0u8; len]),
encoding: "zenoh/bytes".to_string(),
kind: SampleKind::Put,
timestamp: None,
}
}
#[tokio::test]
async fn events_flow_and_snapshots_rebuild_on_tick() {
let core = MonitorCore::new(8);
let mut events = core.events();
core.ingest(view("zs/v1/h-a/telemetry/x/m", 4), None);
core.tick();
let Some(StreamItem::Event(FleetEvent::Sample(s))) = events.recv().await else {
panic!("expected sample");
};
assert_eq!(s.key, "zs/v1/h-a/telemetry/x/m");
assert_eq!(s.payload.len(), 4);
let Some(StreamItem::Event(FleetEvent::StatsTick)) = events.recv().await else {
panic!("expected tick");
};
let snap = core.tree();
assert_eq!(snap.keys, 1);
assert_eq!(snap.root.subtree_count, 1);
}
#[tokio::test]
async fn overflow_surfaces_as_dropped_counts() {
let core = MonitorCore::new(2);
let mut slow = core.events();
for i in 0..10 {
core.ingest(view(&format!("zs/v1/h-a/telemetry/x/m{i}"), 1), None);
}
let Some(StreamItem::Dropped(n)) = slow.recv().await else {
panic!("expected a dropped count first");
};
assert!(n >= 8, "missed at least 8, reported {n}");
assert_eq!(core.dropped(), n);
let Some(StreamItem::Event(FleetEvent::Sample(_))) = slow.recv().await else {
panic!("expected a sample after the gap report");
};
}
}