use std::sync::Arc;
use std::sync::Mutex;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, Instant};
use crate::{Error, Result};
use arc_swap::ArcSwap;
use tokio::sync::broadcast;
use zenoh::Session;
use zenoh::sample::SampleKind;
use crate::model::retain::{Retention, RetentionBudget, RetentionStats};
use crate::model::stats::StatsTable;
use crate::model::tree::KeyTreeSnapshot;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct SampleSource {
pub zid: zenoh::config::ZenohId,
pub eid: u32,
pub sn: u32,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum StampProvenance {
SelfStamped,
Foreign { stamper: zenoh::time::TimestampId },
Unattributable { stamper: zenoh::time::TimestampId },
}
impl StampProvenance {
pub fn stamper(self) -> Option<zenoh::time::TimestampId> {
match self {
StampProvenance::SelfStamped => None,
StampProvenance::Foreign { stamper } | StampProvenance::Unattributable { stamper } => {
Some(stamper)
}
}
}
fn of(timestamp: &zenoh::time::Timestamp, source: Option<SampleSource>) -> StampProvenance {
let stamper = *timestamp.get_id();
match source {
Some(src) if stamper == zenoh::time::TimestampId::from(src.zid) => {
StampProvenance::SelfStamped
}
Some(_) => StampProvenance::Foreign { stamper },
None => StampProvenance::Unattributable { stamper },
}
}
}
#[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>,
pub stamped_by: Option<StampProvenance>,
pub attachment: Option<zenoh::bytes::ZBytes>,
pub priority: zenoh::qos::Priority,
pub congestion_control: zenoh::qos::CongestionControl,
pub reliability: zenoh::qos::Reliability,
pub express: bool,
pub source: Option<SampleSource>,
pub received: Instant,
}
impl SampleView {
pub fn of(sample: &zenoh::sample::Sample) -> SampleView {
let source = sample.source_info().map(|si| SampleSource {
zid: si.source_id().zid(),
eid: si.source_id().eid(),
sn: si.source_sn(),
});
let timestamp = sample.timestamp().copied();
SampleView {
key: sample.key_expr().as_str().to_string(),
payload: sample.payload().clone(),
encoding: sample.encoding().to_string(),
kind: sample.kind(),
stamped_by: timestamp.as_ref().map(|t| StampProvenance::of(t, source)),
timestamp,
attachment: sample.attachment().cloned(),
priority: sample.priority(),
congestion_control: sample.congestion_control(),
reliability: sample.reliability(),
express: sample.express(),
source,
received: Instant::now(),
}
}
pub fn qos_matches(&self, profile: zenkey::qos::QosProfile) -> bool {
self.priority == profile.priority()
&& self.congestion_control == profile.congestion_control()
&& self.reliability == profile.reliability()
&& self.express == profile.express()
}
}
#[derive(Debug, Clone)]
pub enum FleetEvent {
Sample(Arc<SampleView>),
NodeUp(String),
NodeDown(String),
StatsTick,
WatchChanged,
WatchSeeded {
id: WatchId,
coverage: crate::report::SeedCoverage,
},
}
#[derive(Debug, Clone)]
pub struct MonitorSpec {
pub selectors: Vec<String>,
pub liveliness: Vec<String>,
pub stats_tick: Duration,
pub capacity: usize,
pub max_keys: usize,
}
impl Default for MonitorSpec {
fn default() -> Self {
MonitorSpec {
selectors: Vec::new(),
liveliness: Vec::new(),
stats_tick: Duration::from_millis(250),
capacity: 1024,
max_keys: crate::model::bounded::DEFAULT_MAX_KEYS,
}
}
}
pub struct MonitorCore {
tx: broadcast::Sender<FleetEvent>,
stats: Mutex<StatsTable>,
tree: ArcSwap<KeyTreeSnapshot>,
dropped: AtomicU64,
retain: Mutex<Retention>,
tick_seq: AtomicU64,
published_seq: Mutex<u64>,
}
impl MonitorCore {
pub fn new(capacity: usize) -> Arc<MonitorCore> {
MonitorCore::bounded(capacity, crate::model::bounded::DEFAULT_MAX_KEYS)
}
pub fn bounded(capacity: usize, max_keys: usize) -> Arc<MonitorCore> {
let (tx, _) = broadcast::channel(capacity.max(2));
Arc::new(MonitorCore {
tx,
stats: Mutex::new(StatsTable::with_capacity(max_keys)),
tree: ArcSwap::from_pointee(KeyTreeSnapshot::default()),
dropped: AtomicU64::new(0),
retain: Mutex::new(Retention::new(RetentionBudget::default())),
tick_seq: AtomicU64::new(0),
published_seq: Mutex::new(0),
})
}
pub fn ingest(&self, view: SampleView, sn: Option<u32>) {
self.ingest_at(
Arc::new(view),
sn,
Instant::now(),
std::time::SystemTime::now(),
);
}
pub fn ingest_at(
&self,
view: Arc<SampleView>,
sn: Option<u32>,
now: Instant,
wall: std::time::SystemTime,
) {
{
let latency = view.timestamp.as_ref().map(|t| {
let stamped = t.get_time().to_system_time();
let us = match wall.duration_since(stamped) {
Ok(d) => i64::try_from(d.as_micros()).unwrap_or(i64::MAX),
Err(e) => -i64::try_from(e.duration().as_micros()).unwrap_or(i64::MAX),
};
let class = match view.stamped_by {
Some(StampProvenance::SelfStamped) => {
crate::model::stats::StampClass::SelfStamped
}
Some(StampProvenance::Foreign { .. }) => {
crate::model::stats::StampClass::Foreign
}
_ => crate::model::stats::StampClass::Unattributable,
};
(us, class)
});
let stamper = view.stamped_by.and_then(StampProvenance::stamper);
let mut stats = self.stats.lock().expect("stats lock");
stats.record(&view.key, view.payload.len(), sn, now, latency, stamper);
}
self.retain
.lock()
.expect("retain lock")
.push(Arc::clone(&view), now);
let _ = self.tx.send(FleetEvent::Sample(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 (rows, seq) = self.stats_rows();
self.publish(KeyTreeSnapshot::fold(rows), seq);
}
pub async fn tick_off_runtime(&self) {
let (rows, seq) = self.stats_rows();
match tokio::task::spawn_blocking(move || KeyTreeSnapshot::fold(rows)).await {
Ok(snapshot) => self.publish(snapshot, seq),
Err(e) => tracing::warn!("key-tree fold: {e}"),
}
}
fn stats_rows(&self) -> (crate::model::tree::TreeRows, u64) {
let stats = self.stats.lock().expect("stats lock");
let seq = self.tick_seq.fetch_add(1, Ordering::Relaxed) + 1;
(stats.rows(), seq)
}
fn publish(&self, snapshot: KeyTreeSnapshot, seq: u64) {
let mut published = self.published_seq.lock().expect("tick seq lock");
if seq <= *published {
return;
}
*published = seq;
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 with_stats_mut<R>(&self, f: impl FnOnce(&mut StatsTable) -> R) -> R {
f(&mut self.stats.lock().expect("stats lock"))
}
pub fn keys_unwatched(&self) -> u64 {
self.with_stats(|s| s.unwatched())
}
pub fn dropped(&self) -> u64 {
self.dropped.load(Ordering::Relaxed)
}
pub fn keys_evicted(&self) -> u64 {
self.with_stats(|s| s.evicted())
}
pub fn retained(&self) -> Arc<[Arc<SampleView>]> {
let parts = self
.retain
.lock()
.expect("retain lock")
.parts(Instant::now());
parts.flatten()
}
pub fn retention(&self) -> RetentionStats {
self.retain
.lock()
.expect("retain lock")
.stats(Instant::now())
}
pub fn set_retention_budget(&self, budget: RetentionBudget) {
self.retain.lock().expect("retain lock").set_budget(budget);
}
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 fn into_stream(self) -> impl futures_core::Stream<Item = StreamItem> + Send {
futures_util::stream::unfold(self, |mut events| async move {
events.recv().await.map(|item| (item, events))
})
}
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,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub struct WatchId(u64);
struct WatchEntry {
selector: String,
subscriber: zenoh::pubsub::Subscriber<()>,
seed_task: Option<tokio::task::JoinHandle<()>>,
}
pub struct Monitor {
core: Arc<MonitorCore>,
session: Session,
watches: tokio::sync::Mutex<std::collections::HashMap<WatchId, WatchEntry>>,
next_watch: AtomicU64,
tasks: Vec<tokio::task::JoinHandle<()>>,
}
impl std::fmt::Debug for Monitor {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Monitor").finish_non_exhaustive()
}
}
impl Monitor {
pub async fn start(session: &Session, spec: MonitorSpec) -> Result<Monitor> {
let core = MonitorCore::bounded(spec.capacity, spec.max_keys);
let mut tasks = Vec::new();
for liveliness_sel in &spec.liveliness {
let subscriber = crate::bus::teardown::declared(
"liveliness subscribe",
liveliness_sel,
session
.liveliness()
.declare_subscriber(liveliness_sel)
.history(true),
)
.await?;
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_off_runtime().await;
}
}));
}
let monitor = Monitor {
core,
session: session.clone(),
watches: tokio::sync::Mutex::new(std::collections::HashMap::new()),
next_watch: AtomicU64::new(0),
tasks,
};
for selector in &spec.selectors {
monitor.watch(selector).await?;
}
Ok(monitor)
}
pub async fn watch(&self, selector: &str) -> Result<WatchId> {
let core = Arc::clone(&self.core);
let subscriber = crate::bus::teardown::declared(
"subscribe",
selector,
self.session
.declare_subscriber(selector)
.callback(move |sample| {
let view = SampleView::of(&sample);
let sn = view.source.map(|s| s.sn);
core.ingest(view, sn);
}),
)
.await?;
let id = WatchId(self.next_watch.fetch_add(1, Ordering::Relaxed));
self.watches.lock().await.insert(
id,
WatchEntry {
selector: selector.to_string(),
subscriber,
seed_task: None,
},
);
let _ = self.core.tx.send(FleetEvent::WatchChanged);
Ok(id)
}
pub async fn watching<S: AsRef<str>>(
self,
selectors: impl IntoIterator<Item = S>,
) -> Result<Monitor> {
for selector in selectors {
if let Err(declare) = self.watch(selector.as_ref()).await {
if let Err(teardown) = self.shutdown().await {
tracing::warn!("after a failed watch: {teardown}");
}
return Err(declare);
}
}
Ok(self)
}
pub async fn watch_seeded(
&self,
selector: &str,
policy: crate::bus::seed::SeedPolicy,
) -> Result<WatchId> {
use crate::bus::seed::{Merge, cache_selector, seed_get, view_of};
use crate::report::SeedCoverage;
let gate: Arc<arc_swap::ArcSwapOption<Merge>> =
Arc::new(arc_swap::ArcSwapOption::from_pointee(Merge::new()));
let core = Arc::clone(&self.core);
let cb_gate = Arc::clone(&gate);
let subscriber = crate::bus::teardown::declared(
"seeded subscribe",
selector,
self.session
.declare_subscriber(selector)
.callback(move |sample| {
let sn = sample.source_info().map(|si| si.source_sn());
let view = view_of(&sample);
if let Some(merge) = cb_gate.load_full()
&& !merge.admit(&view)
{
return;
}
core.ingest(view, sn);
}),
)
.await?;
let id = WatchId(self.next_watch.fetch_add(1, Ordering::Relaxed));
self.watches.lock().await.insert(
id,
WatchEntry {
selector: selector.to_string(),
subscriber,
seed_task: None,
},
);
let seed_task = {
let session = self.session.clone();
let core = Arc::clone(&self.core);
let selector = selector.to_string();
tokio::spawn(async move {
let merge = gate
.load_full()
.expect("gate holds the merge while seeding");
let history = async {
if policy.history {
let sel = cache_selector(&selector);
Some(
seed_get(&session, &sel, policy.timeout, &merge, |view| {
core.ingest(view, None);
})
.await,
)
} else {
None
}
};
let storage = async {
if policy.storage {
Some(
seed_get(&session, &selector, policy.timeout, &merge, |view| {
core.ingest(view, None);
})
.await,
)
} else {
None
}
};
let (history_replies, storage_replies) = tokio::join!(history, storage);
let coverage = SeedCoverage {
history_replies,
storage_replies,
superseded: merge.superseded(),
};
gate.store(None);
core.tick();
let _ = core.tx.send(FleetEvent::WatchSeeded { id, coverage });
})
};
match self.watches.lock().await.get_mut(&id) {
Some(entry) => entry.seed_task = Some(seed_task),
None => seed_task.abort(),
}
let _ = self.core.tx.send(FleetEvent::WatchChanged);
Ok(id)
}
pub async fn unwatch(&self, id: WatchId) -> Result<()> {
let mut entry = {
let mut watches = self.watches.lock().await;
watches.remove(&id).ok_or_else(|| {
Error::unaskable(
format!("watch id {id:?}"),
"is not a watch this monitor holds",
)
})?
};
if let Some(task) = entry.seed_task.take() {
task.abort();
}
entry
.subscriber
.undeclare()
.await
.map_err(|e| Error::bus("undeclare", &entry.selector, e))?;
let kept: Vec<String> = {
let watches = self.watches.lock().await;
watches.values().map(|w| w.selector.clone()).collect()
};
self.core.with_stats_mut(|stats| {
stats.retire_unwatched(&entry.selector, &kept);
});
self.core.tick();
let _ = self.core.tx.send(FleetEvent::WatchChanged);
Ok(())
}
pub async fn watched(&self) -> Vec<(WatchId, String)> {
let watches = self.watches.lock().await;
let mut v: Vec<(WatchId, String)> = watches
.iter()
.map(|(id, w)| (*id, w.selector.clone()))
.collect();
v.sort();
v
}
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) {
drop(self);
}
pub async fn shutdown(self) -> Result<()> {
let drained: Vec<WatchEntry> = {
let mut watches = self.watches.lock().await;
watches.drain().map(|(_, entry)| entry).collect()
};
let mut failed = Vec::new();
for mut entry in drained {
if let Some(task) = entry.seed_task.take() {
task.abort();
}
if let Err(e) = entry.subscriber.undeclare().await {
failed.push(format!("{}: {e}", entry.selector));
}
}
drop(self);
if failed.is_empty() {
Ok(())
} else {
Err(Error::bus(
"undeclare",
failed.join("; "),
"one or more handles refused",
))
}
}
}
impl Drop for Monitor {
fn drop(&mut self) {
for t in &self.tasks {
t.abort();
}
for entry in self.watches.get_mut().values_mut() {
if let Some(task) = entry.seed_task.take() {
task.abort();
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn an_older_fold_never_walks_the_tree_backwards() {
let core = MonitorCore::new(8);
let (empty_rows, first) = core.stats_rows();
core.ingest(view("v1/h-3fa9c2d41b7e/telemetry/p/x", 4), None);
let (seeded_rows, second) = core.stats_rows();
assert!(second > first, "the sequence orders the two takes");
core.publish(KeyTreeSnapshot::fold(seeded_rows), second);
assert_eq!(core.tree().keys, 1, "the boundary tick published");
core.publish(KeyTreeSnapshot::fold(empty_rows), first);
assert_eq!(
core.tree().keys,
1,
"an older fold overwrote a newer snapshot — the tree walked backwards"
);
}
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,
stamped_by: None,
attachment: None,
priority: zenoh::qos::Priority::DEFAULT,
congestion_control: zenoh::qos::CongestionControl::DEFAULT,
reliability: zenoh::qos::Reliability::DEFAULT,
express: false,
source: None,
received: Instant::now(),
}
}
#[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);
}
#[test]
fn injected_wall_clock_drives_the_latency_fold_deterministically() {
let stamp_epoch = Duration::from_secs(1_000_000);
let ts = zenoh::time::Timestamp::new(
zenoh::time::NTP64::from(stamp_epoch),
zenoh::time::TimestampId::rand(),
);
let key = "zs/v1/h-a/telemetry/x/m";
let stamped_view = || {
let mut v = view(key, 4);
v.timestamp = Some(ts);
v.stamped_by = Some(StampProvenance::Unattributable {
stamper: *ts.get_id(),
});
v
};
let now = Instant::now();
let wall = ts.get_time().to_system_time() + Duration::from_millis(5);
let fold = || {
let core = MonitorCore::new(8);
core.ingest_at(Arc::new(stamped_view()), None, now, wall);
core.with_stats(|s| s.get(key).expect("recorded").latency())
.expect("a stamped sample has a latency window")
};
let a = fold();
let summary = a.unattributable.expect("unattributable population");
assert_eq!(summary.samples, 1);
assert_eq!(
summary.median_us, 5_000,
"the latency is wall − HLC, on the injected wall clock"
);
assert_eq!(a, fold());
}
#[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");
};
}
#[tokio::test]
async fn the_ring_sees_what_a_lagging_receiver_missed() {
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(_)) = slow.recv().await else {
panic!("the broadcast lagged");
};
let window = core.retained();
assert_eq!(window.len(), 10, "the ring is upstream of the lag");
assert_eq!(window[0].key, "zs/v1/h-a/telemetry/x/m0");
assert_eq!(window[9].key, "zs/v1/h-a/telemetry/x/m9");
}
#[test]
fn the_tick_holds_the_ingest_lock_only_for_the_row_copy() {
const KEYS: usize = 5_000;
let core = MonitorCore::bounded(2, KEYS * 2);
let now = Instant::now();
core.with_stats_mut(|stats| {
for i in 0..KEYS {
stats.record(
&format!(
"zs/v1/h-{:04}/telemetry/proc-{i}/group/sub/leaf/m{i}",
i % 97
),
64,
None,
now,
None,
None,
);
}
});
let t0 = Instant::now();
let (rows, _seq) = core.stats_rows();
let copy = t0.elapsed();
assert_eq!(rows.rows.len(), KEYS);
let copied_rows = rows.rows.len();
let t1 = Instant::now();
let snapshot = KeyTreeSnapshot::fold(rows);
let fold = t1.elapsed();
assert_eq!(snapshot.keys, KEYS);
assert_eq!(copied_rows, snapshot.keys, "one row in, one key out");
assert!(
copy < fold * 100,
"a sane machine folds slower than it copies; copy {copy:?}, fold {fold:?}"
);
}
#[test]
fn the_split_fold_is_the_same_snapshot() {
let core = MonitorCore::new(8);
for i in 0..50 {
core.ingest(view(&format!("zs/v1/h-a/telemetry/x/m{i}"), 4), None);
}
core.tick();
let ticked = core.tree();
let direct = core.with_stats(KeyTreeSnapshot::build);
assert_eq!(ticked.keys, direct.keys);
assert_eq!(ticked.root.subtree_count, direct.root.subtree_count);
assert_eq!(ticked.root.subtree_bytes, direct.root.subtree_bytes);
assert_eq!(
ticked
.node(&["zs", "v1", "h-a", "telemetry", "x"])
.unwrap()
.subtree_keys,
50
);
}
#[test]
fn a_retained_read_holds_the_ingest_lock_for_chunk_pointers_only() {
const SAMPLES: usize = 40_000;
let core = MonitorCore::new(2);
let now = Instant::now();
for i in 0..SAMPLES {
core.ingest_at(
Arc::new(view(&format!("zs/v1/h-a/telemetry/x/m{i}"), 8)),
None,
now,
std::time::SystemTime::now(),
);
}
let t0 = Instant::now();
let parts = core
.retain
.lock()
.expect("retain lock")
.parts(Instant::now());
let under_lock = t0.elapsed();
let held_chunks = parts.chunks();
let t1 = Instant::now();
let window = parts.flatten();
let flatten = t1.elapsed();
assert_eq!(window.len(), SAMPLES, "the whole window, unchanged");
assert_eq!(window[0].key, "zs/v1/h-a/telemetry/x/m0");
assert!(
held_chunks <= SAMPLES / crate::model::retain::CHUNK + 2,
"the critical section holds chunk pointers, not samples: {held_chunks} chunks for {SAMPLES} samples"
);
let _ = (under_lock, flatten);
}
#[tokio::test]
async fn the_eviction_populations_are_never_folded() {
const SAMPLES: usize = 100;
let core = MonitorCore::bounded(2, 8);
core.set_retention_budget(crate::model::retain::RetentionBudget {
max_bytes: 1100,
max_age: Duration::from_secs(3600),
});
let mut slow = core.events();
for i in 0..SAMPLES {
core.ingest(view(&format!("zs/v1/h-a/telemetry/x/m{i}"), 64), None);
}
let Some(StreamItem::Dropped(lagged)) = slow.recv().await else {
panic!("the broadcast lagged");
};
assert_eq!(core.dropped(), lagged, "lag counts only broadcast lag");
let table_kept = core.with_stats(StatsTable::len);
assert_eq!(
table_kept as u64 + core.keys_evicted(),
SAMPLES as u64,
"every key is in the table or in its eviction count"
);
let r = core.retention();
assert_eq!(
r.retained as u64 + r.evicted,
SAMPLES as u64,
"every sample is in the ring or in its eviction count"
);
assert_eq!(r.expired, 0, "nothing aged out in this window");
assert_eq!(core.keys_unwatched(), 0, "nothing was unwatched");
assert_ne!(r.evicted, core.keys_evicted());
assert_ne!(r.evicted, core.dropped());
}
}