use std::{
collections::HashMap,
fmt,
sync::{
Arc, Mutex,
atomic::{AtomicU64, Ordering},
},
};
use kio::Lock;
use serde::{Deserialize, Deserializer, Serialize, Serializer};
use crate::{AsPath, PathOwned, Pattern, Patterns};
#[derive(Default, Debug)]
pub(crate) struct Counters {
announces_started: AtomicU64,
announces_ended: AtomicU64,
announced_bytes: AtomicU64,
subscriptions_started: AtomicU64,
subscriptions_ended: AtomicU64,
fetches: AtomicU64,
broadcasts_started: AtomicU64,
broadcasts_ended: AtomicU64,
bytes: AtomicU64,
frames: AtomicU64,
groups: AtomicU64,
datagrams: AtomicU64,
stale: ContentCounters,
}
#[derive(Default, Debug)]
struct ContentCounters {
bytes: AtomicU64,
frames: AtomicU64,
groups: AtomicU64,
datagrams: AtomicU64,
}
impl ContentCounters {
fn snapshot(&self) -> Content {
Content {
bytes: self.bytes.load(Ordering::Relaxed),
frames: self.frames.load(Ordering::Relaxed),
groups: self.groups.load(Ordering::Relaxed),
datagrams: self.datagrams.load(Ordering::Relaxed),
}
}
fn add(&self, content: Content) {
self.bytes.fetch_add(content.bytes, Ordering::Relaxed);
self.frames.fetch_add(content.frames, Ordering::Relaxed);
self.groups.fetch_add(content.groups, Ordering::Relaxed);
self.datagrams.fetch_add(content.datagrams, Ordering::Relaxed);
}
}
impl Counters {
fn snapshot(&self) -> Traffic {
let announces_ended = self.announces_ended.load(Ordering::Acquire);
let subscriptions_ended = self.subscriptions_ended.load(Ordering::Acquire);
let broadcasts_ended = self.broadcasts_ended.load(Ordering::Acquire);
let announces_started = self.announces_started.load(Ordering::Relaxed);
let announced_bytes = self.announced_bytes.load(Ordering::Relaxed);
let subscriptions_started = self.subscriptions_started.load(Ordering::Relaxed);
let fetches = self.fetches.load(Ordering::Relaxed);
let broadcasts_started = self.broadcasts_started.load(Ordering::Relaxed);
let bytes = self.bytes.load(Ordering::Relaxed);
let frames = self.frames.load(Ordering::Relaxed);
let groups = self.groups.load(Ordering::Relaxed);
let datagrams = self.datagrams.load(Ordering::Relaxed);
let stale = self.stale.snapshot();
Traffic {
announces_started,
announces_ended,
announced_bytes,
broadcasts_started,
broadcasts_ended,
subscriptions_started,
subscriptions_ended,
fetches,
bytes,
frames,
groups,
datagrams,
stale,
}
}
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(default)]
#[non_exhaustive]
pub struct Content {
pub bytes: u64,
pub frames: u64,
pub groups: u64,
pub datagrams: u64,
}
impl Content {
pub(crate) fn add(&mut self, other: Self) {
self.bytes += other.bytes;
self.frames += other.frames;
self.groups += other.groups;
self.datagrams += other.datagrams;
}
}
#[derive(Default, Debug)]
struct SessionCounters {
sessions_started: AtomicU64,
sessions_ended: AtomicU64,
}
impl SessionCounters {
fn snapshot(&self) -> Presence {
let sessions_ended = self.sessions_ended.load(Ordering::Acquire);
let sessions_started = self.sessions_started.load(Ordering::Relaxed);
Presence {
sessions_started,
sessions_ended,
}
}
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub struct Traffic {
pub announces_started: u64,
pub announces_ended: u64,
pub announced_bytes: u64,
pub broadcasts_started: u64,
pub broadcasts_ended: u64,
pub subscriptions_started: u64,
pub subscriptions_ended: u64,
pub fetches: u64,
pub bytes: u64,
pub frames: u64,
pub groups: u64,
pub datagrams: u64,
pub stale: Content,
}
#[derive(Default, Clone, Copy)]
struct Edge(Option<u64>);
impl<'de> Deserialize<'de> for Edge {
fn deserialize<D: Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
u64::deserialize(deserializer).map(|v| Self(Some(v)))
}
}
fn counter_edge(canonical: Edge, legacy: Edge) -> u64 {
canonical.0.or(legacy.0).unwrap_or(0)
}
#[derive(Serialize)]
struct TrafficSer {
announces_started: u64,
announced: u64,
announces_ended: u64,
announced_closed: u64,
announced_bytes: u64,
broadcasts_started: u64,
broadcasts: u64,
broadcasts_ended: u64,
broadcasts_closed: u64,
subscriptions_started: u64,
subscriptions: u64,
subscriptions_ended: u64,
subscriptions_closed: u64,
fetches: u64,
bytes: u64,
frames: u64,
groups: u64,
datagrams: u64,
stale: Content,
}
impl From<Traffic> for TrafficSer {
fn from(t: Traffic) -> Self {
Self {
announces_started: t.announces_started,
announced: t.announces_started,
announces_ended: t.announces_ended,
announced_closed: t.announces_ended,
announced_bytes: t.announced_bytes,
broadcasts_started: t.broadcasts_started,
broadcasts: t.broadcasts_started,
broadcasts_ended: t.broadcasts_ended,
broadcasts_closed: t.broadcasts_ended,
subscriptions_started: t.subscriptions_started,
subscriptions: t.subscriptions_started,
subscriptions_ended: t.subscriptions_ended,
subscriptions_closed: t.subscriptions_ended,
fetches: t.fetches,
bytes: t.bytes,
frames: t.frames,
groups: t.groups,
datagrams: t.datagrams,
stale: t.stale,
}
}
}
#[derive(Default, Deserialize)]
#[serde(default)]
struct TrafficDe {
announces_started: Edge,
announced: Edge,
announces_ended: Edge,
announced_closed: Edge,
announced_bytes: u64,
broadcasts_started: Edge,
broadcasts: Edge,
broadcasts_ended: Edge,
broadcasts_closed: Edge,
subscriptions_started: Edge,
subscriptions: Edge,
subscriptions_ended: Edge,
subscriptions_closed: Edge,
fetches: u64,
bytes: u64,
frames: u64,
groups: u64,
datagrams: u64,
stale: Content,
}
impl From<TrafficDe> for Traffic {
fn from(d: TrafficDe) -> Self {
Self {
announces_started: counter_edge(d.announces_started, d.announced),
announces_ended: counter_edge(d.announces_ended, d.announced_closed),
announced_bytes: d.announced_bytes,
broadcasts_started: counter_edge(d.broadcasts_started, d.broadcasts),
broadcasts_ended: counter_edge(d.broadcasts_ended, d.broadcasts_closed),
subscriptions_started: counter_edge(d.subscriptions_started, d.subscriptions),
subscriptions_ended: counter_edge(d.subscriptions_ended, d.subscriptions_closed),
fetches: d.fetches,
bytes: d.bytes,
frames: d.frames,
groups: d.groups,
datagrams: d.datagrams,
stale: d.stale,
}
}
}
impl Serialize for Traffic {
fn serialize<S: Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
TrafficSer::from(*self).serialize(serializer)
}
}
impl<'de> Deserialize<'de> for Traffic {
fn deserialize<D: Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
TrafficDe::deserialize(deserializer).map(Into::into)
}
}
impl Traffic {
pub fn add(&mut self, other: Traffic) {
self.announces_started += other.announces_started;
self.announces_ended += other.announces_ended;
self.announced_bytes += other.announced_bytes;
self.broadcasts_started += other.broadcasts_started;
self.broadcasts_ended += other.broadcasts_ended;
self.subscriptions_started += other.subscriptions_started;
self.subscriptions_ended += other.subscriptions_ended;
self.fetches += other.fetches;
self.bytes += other.bytes;
self.frames += other.frames;
self.groups += other.groups;
self.datagrams += other.datagrams;
self.stale.add(other.stale);
}
pub fn is_announced(&self) -> bool {
self.announces_started > self.announces_ended
}
pub fn active_broadcasts(&self) -> u64 {
self.broadcasts_started.saturating_sub(self.broadcasts_ended)
}
pub fn active_subscriptions(&self) -> u64 {
self.subscriptions_started.saturating_sub(self.subscriptions_ended)
}
pub fn total_bytes(&self) -> u64 {
self.bytes.saturating_add(self.announced_bytes)
}
pub fn is_idle(&self) -> bool {
self.announces_started == self.announces_ended
&& self.subscriptions_started == self.subscriptions_ended
&& self.broadcasts_started == self.broadcasts_ended
}
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub struct Presence {
pub sessions_started: u64,
pub sessions_ended: u64,
}
#[derive(Serialize)]
struct PresenceSer {
sessions_started: u64,
sessions: u64,
sessions_ended: u64,
sessions_closed: u64,
}
impl From<Presence> for PresenceSer {
fn from(p: Presence) -> Self {
Self {
sessions_started: p.sessions_started,
sessions: p.sessions_started,
sessions_ended: p.sessions_ended,
sessions_closed: p.sessions_ended,
}
}
}
#[derive(Default, Deserialize)]
#[serde(default)]
struct PresenceDe {
sessions_started: Edge,
sessions: Edge,
sessions_ended: Edge,
sessions_closed: Edge,
}
impl From<PresenceDe> for Presence {
fn from(d: PresenceDe) -> Self {
Self {
sessions_started: counter_edge(d.sessions_started, d.sessions),
sessions_ended: counter_edge(d.sessions_ended, d.sessions_closed),
}
}
}
impl Serialize for Presence {
fn serialize<S: Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
PresenceSer::from(*self).serialize(serializer)
}
}
impl<'de> Deserialize<'de> for Presence {
fn deserialize<D: Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
PresenceDe::deserialize(deserializer).map(Into::into)
}
}
impl Presence {
pub fn add(&mut self, other: Presence) {
self.sessions_started += other.sessions_started;
self.sessions_ended += other.sessions_ended;
}
pub fn active(&self) -> u64 {
self.sessions_started.saturating_sub(self.sessions_ended)
}
}
#[derive(Clone, Debug, Default, PartialEq, Eq, Hash)]
pub struct Tier(PathOwned);
impl Tier {
pub fn new(label: impl Into<PathOwned>) -> Self {
Self(label.into())
}
pub fn label(&self) -> &PathOwned {
&self.0
}
pub fn is_default(&self) -> bool {
self.0.is_empty()
}
pub fn track_name(&self, name: &str) -> String {
if self.0.is_empty() {
name.to_string()
} else {
format!("{}/{}", self.0.as_str(), name)
}
}
pub fn as_str(&self) -> &str {
self.0.as_str()
}
}
impl fmt::Display for Tier {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
fmt::Display::fmt(&self.0, f)
}
}
#[derive(Copy, Clone, Debug, PartialEq, Eq, Hash)]
pub enum Role {
Publisher,
Subscriber,
}
impl Role {
fn idx(self) -> usize {
match self {
Role::Publisher => 0,
Role::Subscriber => 1,
}
}
pub fn as_str(self) -> &'static str {
match self {
Role::Publisher => "publisher",
Role::Subscriber => "subscriber",
}
}
}
#[derive(Debug, Default, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct Snapshot {
traffic: HashMap<Tier, [Traffic; 2]>,
sessions: HashMap<Tier, Presence>,
}
impl Snapshot {
pub fn traffic(&self) -> Vec<(Tier, Role, Traffic)> {
let mut rows = Vec::with_capacity(self.traffic.len() * 2);
for (tier, roles) in &self.traffic {
rows.push((tier.clone(), Role::Publisher, roles[Role::Publisher.idx()]));
rows.push((tier.clone(), Role::Subscriber, roles[Role::Subscriber.idx()]));
}
rows.sort_by(|a, b| a.0.as_str().cmp(b.0.as_str()).then(a.1.idx().cmp(&b.1.idx())));
rows
}
pub fn sessions(&self) -> Vec<(Tier, Presence)> {
let mut rows: Vec<_> = self.sessions.iter().map(|(tier, s)| (tier.clone(), *s)).collect();
rows.sort_by(|a, b| a.0.as_str().cmp(b.0.as_str()));
rows
}
}
#[derive(Debug, Default, Clone)]
#[non_exhaustive]
pub struct Report {
pub traffic: Vec<TrafficEntry>,
pub sessions: Vec<SessionEntry>,
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct TrafficEntry {
pub path: PathOwned,
pub tier: Tier,
pub publisher: Traffic,
pub subscriber: Traffic,
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct SessionEntry {
pub tier: Tier,
pub root: PathOwned,
pub presence: Presence,
}
#[derive(Clone, Debug, Default)]
#[non_exhaustive]
pub struct Config {
pub exclude: Patterns,
}
impl Config {
pub fn new() -> Self {
Self::default()
}
pub fn with_exclude(mut self, pattern: Pattern) -> Self {
self.exclude.insert(pattern);
self
}
}
#[derive(Clone)]
pub struct Registry {
exclude: Patterns,
shared: Option<Arc<Shared>>,
}
struct Shared {
retired: Lock<Snapshot>,
entries: Lock<HashMap<PathOwned, Arc<BroadcastEntry>>>,
sessions: Lock<HashMap<Tier, HashMap<PathOwned, Arc<SessionCounters>>>>,
}
struct BroadcastEntry {
tiers: Mutex<HashMap<Tier, Arc<TierCounters>>>,
}
impl BroadcastEntry {
fn new() -> Self {
Self {
tiers: Mutex::new(HashMap::new()),
}
}
fn tier(&self, tier: &Tier) -> Arc<TierCounters> {
self.tiers
.lock()
.expect("stats tiers poisoned")
.entry(tier.clone())
.or_default()
.clone()
}
}
#[derive(Default)]
struct TierCounters {
publisher: Counters,
subscriber: Counters,
}
impl Registry {
pub fn new(config: Config) -> Self {
let Config { exclude } = config;
Self {
exclude,
shared: Some(Arc::new(Shared {
retired: Lock::default(),
entries: Lock::default(),
sessions: Default::default(),
})),
}
}
pub fn disabled() -> Self {
Self {
exclude: Patterns::new(),
shared: None,
}
}
pub fn exclude(&self) -> &Patterns {
&self.exclude
}
#[cfg(test)]
fn shared(&self) -> &Arc<Shared> {
self.shared.as_ref().expect("enabled stats registry")
}
pub fn tier(&self, tier: Tier) -> Handle {
Handle {
stats: self.clone(),
tier,
}
}
fn entry(&self, path: impl AsPath) -> Option<Arc<BroadcastEntry>> {
let shared = self.shared.as_ref()?;
let path = path.as_path();
if self.exclude.matches(path.as_str()) {
return None;
}
let owned = path.to_owned();
let mut entries = shared.entries.lock();
Some(
entries
.entry(owned)
.or_insert_with(|| Arc::new(BroadcastEntry::new()))
.clone(),
)
}
fn session_counters(&self, tier: &Tier, root: impl AsPath) -> Option<Arc<SessionCounters>> {
let shared = self.shared.as_ref()?;
let owned = root.as_path().to_owned();
let mut sessions = shared.sessions.lock();
Some(
sessions
.entry(tier.clone())
.or_default()
.entry(owned)
.or_default()
.clone(),
)
}
pub fn snapshot(&self) -> Snapshot {
let Some(shared) = self.shared.as_ref() else {
return Snapshot::default();
};
let retired = shared.retired.lock();
let mut snap = retired.clone();
{
let entries = shared.entries.lock();
for entry in entries.values() {
let tiers = entry.tiers.lock().expect("stats tiers poisoned");
for (tier, counters) in tiers.iter() {
let totals = snap.traffic.entry(tier.clone()).or_default();
totals[Role::Publisher.idx()].add(counters.publisher.snapshot());
totals[Role::Subscriber.idx()].add(counters.subscriber.snapshot());
}
}
}
{
let sessions = shared.sessions.lock();
for (tier, roots) in sessions.iter() {
let totals = snap.sessions.entry(tier.clone()).or_default();
for counters in roots.values() {
totals.add(counters.snapshot());
}
}
}
snap
}
pub fn report(&self, report: &mut Report) {
report.traffic.clear();
report.sessions.clear();
let Some(shared) = self.shared.as_ref() else {
return;
};
let mut retired = shared.retired.lock();
{
let mut entries = shared.entries.lock();
for (path, entry) in entries.iter() {
let tiers = entry.tiers.lock().expect("stats tiers poisoned");
for (tier, counters) in tiers.iter() {
report.traffic.push(TrafficEntry {
path: path.clone(),
tier: tier.clone(),
publisher: counters.publisher.snapshot(),
subscriber: counters.subscriber.snapshot(),
});
}
}
entries.retain(|_, entry| {
if Arc::strong_count(entry) > 1 {
return true;
}
let mut tiers = entry.tiers.lock().expect("stats tiers poisoned");
tiers.retain(|tier, counters| {
if Arc::strong_count(counters) > 1 {
return true;
}
let totals = retired.traffic.entry(tier.clone()).or_default();
totals[Role::Publisher.idx()].add(counters.publisher.snapshot());
totals[Role::Subscriber.idx()].add(counters.subscriber.snapshot());
false
});
!tiers.is_empty()
});
}
{
let mut sessions = shared.sessions.lock();
for (tier, roots) in sessions.iter() {
for (root, counters) in roots.iter() {
report.sessions.push(SessionEntry {
tier: tier.clone(),
root: root.clone(),
presence: counters.snapshot(),
});
}
}
for (tier, roots) in sessions.iter_mut() {
roots.retain(|_, counters| {
if Arc::strong_count(counters) > 1 {
return true;
}
retired
.sessions
.entry(tier.clone())
.or_default()
.add(counters.snapshot());
false
});
}
sessions.retain(|_, roots| !roots.is_empty());
}
}
}
impl Default for Registry {
fn default() -> Self {
Self::disabled()
}
}
#[derive(Clone)]
pub struct Handle {
stats: Registry,
tier: Tier,
}
impl Handle {
pub fn parent(&self) -> &Registry {
&self.stats
}
pub fn tier(&self) -> &Tier {
&self.tier
}
pub fn session(&self, root: impl AsPath) -> Session {
Session::new(self.stats.clone(), self.tier.clone(), root)
}
}
impl Default for Handle {
fn default() -> Self {
Registry::disabled().tier(Tier::default())
}
}
#[derive(Copy, Clone, Default)]
enum Side {
#[default]
Publisher,
Subscriber,
}
impl Side {
fn counters(self, tier: &TierCounters) -> &Counters {
match self {
Side::Publisher => &tier.publisher,
Side::Subscriber => &tier.subscriber,
}
}
}
#[derive(Clone, Default)]
pub struct Session {
inner: Option<Arc<SessionInner>>,
}
struct SessionInner {
registry: Registry,
root: PathOwned,
current: Mutex<Current>,
generation: AtomicU64,
viewers: Mutex<HashMap<PathOwned, Viewer>>,
}
struct Current {
tier: Tier,
presence: Option<Arc<SessionCounters>>,
}
struct Viewer {
subscriptions: u32,
counters: Arc<TierCounters>,
}
impl Session {
fn new(registry: Registry, tier: Tier, root: impl AsPath) -> Self {
let root = root.as_path().to_owned();
let presence = registry.session_counters(&tier, &root);
if let Some(presence) = &presence {
presence.sessions_started.fetch_add(1, Ordering::Relaxed);
}
Self {
inner: Some(Arc::new(SessionInner {
registry,
root,
current: Mutex::new(Current { tier, presence }),
generation: AtomicU64::new(0),
viewers: Mutex::new(HashMap::new()),
})),
}
}
pub fn set_tier(&self, tier: Tier) {
let Some(inner) = &self.inner else { return };
let mut current = inner.current.lock().expect("stats session poisoned");
if current.tier == tier {
return;
}
let presence = inner.registry.session_counters(&tier, &inner.root);
if let Some(presence) = &presence {
presence.sessions_started.fetch_add(1, Ordering::Relaxed);
}
if let Some(old) = std::mem::replace(&mut current.presence, presence) {
old.sessions_ended.fetch_add(1, Ordering::Release);
}
current.tier = tier;
inner.generation.fetch_add(1, Ordering::Release);
}
pub(crate) fn egress(&self, path: impl AsPath) -> Scope {
self.scope(path, Side::Publisher)
}
pub(crate) fn ingress(&self, path: impl AsPath) -> Scope {
self.scope(path, Side::Subscriber)
}
fn scope(&self, path: impl AsPath, side: Side) -> Scope {
let Some(inner) = &self.inner else {
return Scope::default();
};
let path = path.as_path().to_owned();
let resolved = inner.resolve(&path);
Scope {
session: self.clone(),
resolved: Some(Box::new(Mutex::new(resolved))),
side,
path,
}
}
fn viewer_open(&self, path: &PathOwned, counters: &Arc<TierCounters>) {
let Some(inner) = &self.inner else { return };
let mut viewers = inner.viewers.lock().expect("stats viewers poisoned");
let viewer = viewers.entry(path.clone()).or_insert_with(|| {
counters.publisher.broadcasts_started.fetch_add(1, Ordering::Relaxed);
Viewer {
subscriptions: 0,
counters: counters.clone(),
}
});
viewer.subscriptions += 1;
}
fn viewer_close(&self, path: &PathOwned) {
let Some(inner) = &self.inner else { return };
let mut viewers = inner.viewers.lock().expect("stats viewers poisoned");
let Some(viewer) = viewers.get_mut(path) else { return };
viewer.subscriptions -= 1;
if viewer.subscriptions == 0
&& let Some(viewer) = viewers.remove(path)
{
viewer
.counters
.publisher
.broadcasts_ended
.fetch_add(1, Ordering::Release);
}
}
}
impl SessionInner {
fn resolve(&self, path: &PathOwned) -> Resolved {
let (generation, tier) = {
let current = self.current.lock().expect("stats session poisoned");
(self.generation.load(Ordering::Relaxed), current.tier.clone())
};
Resolved {
generation,
counters: self.registry.entry(path).map(|entry| entry.tier(&tier)),
}
}
}
impl Drop for SessionInner {
fn drop(&mut self) {
let current = self.current.get_mut().expect("stats session poisoned");
if let Some(presence) = ¤t.presence {
presence.sessions_ended.fetch_add(1, Ordering::Release);
}
}
}
#[derive(Clone, Default)]
pub(crate) struct Meter {
counters: Option<Arc<TierCounters>>,
side: Side,
}
impl Meter {
fn counters(&self) -> Option<&Counters> {
self.counters.as_ref().map(|c| self.side.counters(c))
}
pub(crate) fn group(&self) {
if let Some(counters) = self.counters() {
counters.groups.fetch_add(1, Ordering::Relaxed);
}
}
pub(crate) fn frames(&self, n: u64) {
if n == 0 {
return;
}
if let Some(counters) = self.counters() {
counters.frames.fetch_add(n, Ordering::Relaxed);
}
}
pub(crate) fn datagram(&self, n: u64) {
if let Some(counters) = self.counters() {
counters.datagrams.fetch_add(1, Ordering::Relaxed);
counters.groups.fetch_add(1, Ordering::Relaxed);
counters.frames.fetch_add(1, Ordering::Relaxed);
counters.bytes.fetch_add(n, Ordering::Relaxed);
}
}
pub(crate) fn is_tracked(&self) -> bool {
self.counters.is_some()
}
pub(crate) fn stale(&self, content: Content) {
if let Some(counters) = self.counters() {
counters.stale.add(content);
}
}
pub(crate) fn bytes(&self, n: u64) {
if n == 0 {
return;
}
if let Some(counters) = self.counters() {
counters.bytes.fetch_add(n, Ordering::Relaxed);
}
}
}
#[derive(Default)]
pub(crate) struct Scope {
session: Session,
resolved: Option<Box<Mutex<Resolved>>>,
side: Side,
path: PathOwned,
}
#[derive(Clone)]
struct Resolved {
generation: u64,
counters: Option<Arc<TierCounters>>,
}
impl Clone for Scope {
fn clone(&self) -> Self {
Self {
session: self.session.clone(),
resolved: self
.resolved
.as_ref()
.map(|r| Box::new(Mutex::new(r.lock().expect("stats scope poisoned").clone()))),
side: self.side,
path: self.path.clone(),
}
}
}
impl Scope {
fn counters(&self) -> Option<Arc<TierCounters>> {
let inner = self.session.inner.as_ref()?;
let generation = inner.generation.load(Ordering::Acquire);
let mut resolved = self.resolved.as_ref()?.lock().expect("stats scope poisoned");
if resolved.generation != generation {
*resolved = inner.resolve(&self.path);
}
resolved.counters.clone()
}
pub(crate) fn meter(&self) -> Meter {
Meter {
counters: self.counters(),
side: self.side,
}
}
pub(crate) fn subscribe(&self) -> Subscription {
let counters = self.counters();
let mut viewer = None;
if let Some(counters) = &counters {
self.side
.counters(counters)
.subscriptions_started
.fetch_add(1, Ordering::Relaxed);
if matches!(self.side, Side::Publisher) {
self.session.viewer_open(&self.path, counters);
viewer = Some((self.session.clone(), self.path.clone()));
}
}
Subscription {
counters,
side: self.side,
viewer,
}
}
pub(crate) fn fetch(&self) {
if let Some(counters) = self.counters() {
self.side.counters(&counters).fetches.fetch_add(1, Ordering::Relaxed);
}
}
pub(crate) fn announce(&self) -> Announce {
let len = self.path.as_str().len() as u64;
let counters = self.counters();
if let Some(counters) = &counters {
let counters = self.side.counters(counters);
counters.announces_started.fetch_add(1, Ordering::Relaxed);
counters.announced_bytes.fetch_add(len, Ordering::Relaxed);
}
Announce {
counters,
side: self.side,
len,
}
}
}
#[derive(Default)]
#[must_use = "drop the guard to record the subscription as closed"]
pub(crate) struct Subscription {
counters: Option<Arc<TierCounters>>,
side: Side,
viewer: Option<(Session, PathOwned)>,
}
impl Drop for Subscription {
fn drop(&mut self) {
if let Some((session, path)) = &self.viewer {
session.viewer_close(path);
}
if let Some(counters) = &self.counters {
self.side
.counters(counters)
.subscriptions_ended
.fetch_add(1, Ordering::Release);
}
}
}
#[must_use = "drop the guard to record the unannounce"]
pub(crate) struct Announce {
counters: Option<Arc<TierCounters>>,
side: Side,
len: u64,
}
impl Drop for Announce {
fn drop(&mut self) {
if let Some(counters) = &self.counters {
let counters = self.side.counters(counters);
counters.announced_bytes.fetch_add(self.len, Ordering::Relaxed);
counters.announces_ended.fetch_add(1, Ordering::Release);
}
}
}
#[cfg(test)]
mod tests {
use std::sync::{Arc, atomic::Ordering::Relaxed};
use super::*;
#[test]
fn default_tier_has_empty_label() {
let tier = Tier::default();
assert_eq!(tier.as_str(), "");
assert_eq!(tier.to_string(), "");
assert_eq!(tier.track_name("publisher.json"), "publisher.json");
}
fn tier_counters(stats: &Registry, path: &str, tier: &Tier) -> Arc<TierCounters> {
stats
.shared()
.entries
.lock()
.get(&PathOwned::from(path.to_string()))
.expect("entry")
.tier(tier)
}
fn session_snapshot(stats: &Registry, tier: &Tier, root: &str) -> Option<Presence> {
stats
.shared()
.sessions
.lock()
.get(tier)
.and_then(|roots| roots.get(&PathOwned::from(root.to_string())).map(|c| c.snapshot()))
}
fn test_stats() -> Registry {
Registry::new(Config::new().with_exclude(Pattern::subtree(".stats").unwrap()))
}
#[test]
fn default_and_named_tiers_are_independent() {
let stats = test_stats();
let default = stats.tier(Tier::default()).session("root");
let regional = stats.tier(Tier::new("region/sjc")).session("root");
default.egress("demo/bbb").meter().bytes(100);
regional.ingress("demo/bbb").meter().bytes(7);
let default_counters = tier_counters(&stats, "demo/bbb", &Tier::default());
let regional_counters = tier_counters(&stats, "demo/bbb", &Tier::new("region/sjc"));
assert_eq!(default_counters.publisher.bytes.load(Relaxed), 100);
assert_eq!(default_counters.subscriber.bytes.load(Relaxed), 0);
assert_eq!(regional_counters.publisher.bytes.load(Relaxed), 0);
assert_eq!(regional_counters.subscriber.bytes.load(Relaxed), 7);
}
#[test]
fn snapshot_rolls_up_by_tier_role_and_sessions() {
let stats = test_stats();
let default = stats.tier(Tier::default());
let regional = stats.tier(Tier::new("region/sjc"));
let s1 = default.session("acme");
let _s2 = default.session("acme");
let s3 = regional.session("peer");
{
let m = s1.egress("demo/aaa").meter();
m.bytes(100);
m.frames(1);
m.group();
}
s1.egress("demo/bbb").meter().bytes(50);
s3.ingress("demo/aaa").meter().bytes(7);
let snap = stats.snapshot();
let slot = |tier, role| {
snap.traffic()
.into_iter()
.find(|(t, r, _)| *t == tier && *r == role)
.map(|(_, _, c)| c)
.expect("row present")
};
let default_publisher = slot(Tier::default(), Role::Publisher);
assert_eq!(
default_publisher.bytes, 150,
"default egress bytes sum across broadcasts"
);
assert_eq!(default_publisher.frames, 1);
assert_eq!(default_publisher.groups, 1);
let regional_subscriber = slot(Tier::new("region/sjc"), Role::Subscriber);
assert_eq!(regional_subscriber.bytes, 7, "regional ingress isolated by tier/role");
assert_eq!(slot(Tier::default(), Role::Subscriber).bytes, 0);
assert_eq!(slot(Tier::new("region/sjc"), Role::Publisher).bytes, 0);
let sessions = |tier| {
snap.sessions()
.into_iter()
.find(|(t, _)| *t == tier)
.map(|(_, s)| s)
.expect("tier present")
};
let default_sessions = sessions(Tier::default());
assert_eq!(
default_sessions.sessions_started, 2,
"two default-tier sessions under one root"
);
assert_eq!(default_sessions.sessions_ended, 0, "guards still held");
assert_eq!(sessions(Tier::new("region/sjc")).sessions_started, 1);
}
fn drain(stats: &Registry) -> Report {
let mut report = Report::default();
stats.report(&mut report);
report
}
#[test]
fn report_reuses_capacity() {
let stats = test_stats();
let ctx = stats.tier(Tier::default()).session("root");
let _scopes: Vec<_> = (0..8).map(|i| ctx.egress(format!("b/{i}").as_str())).collect();
let mut report = Report::default();
stats.report(&mut report);
assert_eq!(report.traffic.len(), 8);
assert_eq!(report.sessions.len(), 1);
let (traffic, sessions) = (report.traffic.as_ptr(), report.sessions.as_ptr());
stats.report(&mut report);
assert_eq!(report.traffic.len(), 8, "refilled, not appended");
assert_eq!(report.sessions.len(), 1);
assert_eq!(report.traffic.as_ptr(), traffic, "traffic buffer reused");
assert_eq!(report.sessions.as_ptr(), sessions, "sessions buffer reused");
}
#[test]
fn report_returns_detail_and_prunes() {
let stats = test_stats();
let key = PathOwned::from("foo/bar");
let ctx = stats.tier(Tier::default()).session("root");
let scope = ctx.egress("foo/bar");
let sub = scope.subscribe();
scope.meter().bytes(42);
let report = drain(&stats);
let row = report
.traffic
.iter()
.find(|row| row.path == key)
.expect("live entry present");
assert_eq!(row.publisher.bytes, 42);
assert_eq!(row.publisher.subscriptions_started, 1);
assert!(!row.publisher.is_idle(), "subscription guard still open");
assert!(
stats.shared().entries.lock().contains_key(&key),
"live entry kept across drains"
);
drop(sub);
drop(scope);
let report = drain(&stats);
let row = report
.traffic
.iter()
.find(|row| row.path == key)
.expect("closing values still reported once");
assert_eq!(row.publisher.subscriptions_ended, 1);
assert!(row.publisher.is_idle());
assert!(
!stats.shared().entries.lock().contains_key(&key),
"fully-closed entry pruned"
);
assert!(drain(&stats).traffic.is_empty(), "nothing left after the prune");
}
#[test]
fn report_keeps_idle_but_announced_entry() {
let stats = test_stats();
let key = PathOwned::from("foo/bar");
let ctx = stats.tier(Tier::default()).session("root");
let scope = ctx.egress("foo/bar");
let guard = scope.announce();
for _ in 0..3 {
let report = drain(&stats);
assert!(
report.traffic.iter().any(|row| row.path == key),
"announced-but-idle broadcast stays while the guard is held"
);
}
drop(guard);
drop(scope);
let report = drain(&stats);
let row = report.traffic.iter().find(|row| row.path == key).expect("final report");
assert!(row.publisher.is_idle());
assert!(!stats.shared().entries.lock().contains_key(&key));
}
#[test]
fn report_prunes_empty_session_roots() {
let stats = test_stats();
let session = stats.tier(Tier::default()).session("acme");
let report = drain(&stats);
let row = report
.sessions
.iter()
.find(|row| row.root.as_str() == "acme")
.expect("root present");
assert_eq!(row.presence.active(), 1);
drop(session);
let report = drain(&stats);
let row = report
.sessions
.iter()
.find(|row| row.root.as_str() == "acme")
.expect("final gauge reported once");
assert_eq!(row.presence.active(), 0);
assert!(drain(&stats).sessions.is_empty(), "root pruned after the last drain");
assert!(session_snapshot(&stats, &Tier::default(), "acme").is_none());
}
#[test]
fn snapshot_preserves_retired_counters() {
let stats = test_stats();
let tier = Tier::default();
let live = stats.tier(tier.clone()).session("live");
let scope = live.egress("live/video");
let _live_sub = scope.subscribe();
for _ in 0..2 {
let session = stats.tier(tier.clone()).session("retired");
let scope = session.egress("retired/video");
let sub = scope.subscribe();
scope.meter().bytes(100);
drop(sub);
drop(scope);
drop(session);
let before = stats.snapshot();
drain(&stats);
assert_eq!(stats.snapshot(), before, "pruning must not reset host counters");
}
let snap = stats.snapshot();
let traffic = snap
.traffic()
.into_iter()
.find(|(_, role, _)| *role == Role::Publisher)
.unwrap()
.2;
assert_eq!(traffic.bytes, 200);
assert_eq!(traffic.subscriptions_started, 3);
assert_eq!(traffic.subscriptions_ended, 2);
let sessions = snap.sessions().into_iter().find(|(label, _)| label == &tier).unwrap().1;
assert_eq!(sessions.sessions_started, 3);
assert_eq!(sessions.sessions_ended, 2);
assert_eq!(stats.shared().entries.lock().len(), 1, "retired paths are still pruned");
assert_eq!(
stats.shared().sessions.lock()[&tier].len(),
1,
"retired roots are still pruned"
);
let retired = stats.shared().retired.lock();
assert_eq!(retired.traffic.len(), 1, "retain only a total per tier");
assert_eq!(retired.sessions.len(), 1);
}
#[cfg(not(target_family = "wasm"))]
#[test]
fn snapshot_and_report_transfer_counters_once() {
let stats = test_stats();
let worker_stats = stats.clone();
let worker = std::thread::spawn(move || {
for i in 0..256 {
let path = format!("root/{i}");
let session = worker_stats.tier(Tier::default()).session(path.as_str());
session.ingress(path.as_str()).meter().bytes(1);
drop(session);
drain(&worker_stats);
}
});
let mut previous = 0;
while !worker.is_finished() {
let bytes: u64 = stats
.snapshot()
.traffic()
.iter()
.map(|(_, _, traffic)| traffic.bytes)
.sum();
assert!(
bytes >= previous,
"retiring an entry must not double-count or lose its bytes"
);
assert!(bytes <= 256);
previous = bytes;
}
worker.join().unwrap();
let snap = stats.snapshot();
assert_eq!(
snap.traffic().iter().map(|(_, _, traffic)| traffic.bytes).sum::<u64>(),
256
);
assert_eq!(snap.sessions()[0].1.sessions_started, 256);
assert_eq!(snap.sessions()[0].1.sessions_ended, 256);
assert!(stats.shared().entries.lock().is_empty());
assert!(stats.shared().sessions.lock().is_empty());
let retired = stats.shared().retired.lock();
assert_eq!(retired.traffic.len(), 1);
assert_eq!(retired.sessions.len(), 1);
}
#[test]
fn paths_under_exclude_are_no_op() {
let stats = test_stats();
let ctx = stats.tier(Tier::default()).session("root");
let scope = ctx.egress(".stats/node/sjc");
scope.meter().bytes(100);
let _guard = scope.announce();
let _sub = scope.subscribe();
assert!(stats.shared().entries.lock().is_empty());
}
#[test]
fn disabled_stats_are_noop() {
let stats = Registry::default();
assert!(stats.shared.is_none());
let ctx = stats.tier(Tier::default()).session("root");
let scope = ctx.egress("demo/bbb");
scope.meter().bytes(100);
let _guard = scope.announce();
let _sub = scope.subscribe();
assert!(drain(&stats).traffic.is_empty());
assert!(stats.snapshot().traffic().is_empty());
}
#[test]
fn session_counts_by_root() {
let stats = test_stats();
let ext = stats.tier(Tier::default());
let snap = |root: &str| {
session_snapshot(&stats, &Tier::default(), root).map(|p| (p.sessions_started, p.sessions_ended))
};
let a1 = ext.session("acme");
let a2 = ext.session("acme");
let b1 = ext.session("globex");
assert_eq!(snap("acme"), Some((2, 0)), "two sessions under one root");
assert_eq!(snap("globex"), Some((1, 0)), "a distinct root is counted separately");
drop(a1);
assert_eq!(snap("acme"), Some((2, 1)));
drop(a2);
drop(b1);
assert_eq!(snap("acme"), Some((2, 2)));
assert_eq!(snap("globex"), Some((1, 1)));
}
#[test]
fn traffic_parses_with_missing_and_unknown_fields() {
let old: Traffic = serde_json::from_str(r#"{"announced":1,"bytes":5}"#).expect("older shape parses");
assert_eq!(old.announces_started, 1);
assert_eq!(old.bytes, 5);
assert_eq!(old.announced_bytes, 0, "missing fields default to zero");
let new: Traffic = serde_json::from_str(r#"{"announces_started":1,"announces_ended":1,"future_counter":9}"#)
.expect("newer shape parses");
assert!(new.is_idle());
}
#[test]
fn snapshot_reads_ended_before_started() {
let src = include_str!("stats.rs");
let body_start = src.find("fn snapshot(&self) -> Traffic").expect("snapshot fn present");
let body = &src[body_start..];
let ended_pos = body.find("self.announces_ended.load").expect("announces_ended load");
let started_pos = body
.find("self.announces_started.load")
.expect("announces_started load");
assert!(
ended_pos < started_pos,
"announces_ended must be loaded before announces_started; reversing breaks the started>=ended invariant",
);
let subs_ended_pos = body
.find("self.subscriptions_ended.load")
.expect("subscriptions_ended load");
let subs_pos = body
.find("self.subscriptions_started.load")
.expect("subscriptions_started load");
assert!(
subs_ended_pos < subs_pos,
"subscriptions_ended must be loaded before subscriptions_started",
);
let bcast_ended_pos = body.find("self.broadcasts_ended.load").expect("broadcasts_ended load");
let bcast_pos = body
.find("self.broadcasts_started.load")
.expect("broadcasts_started load");
assert!(
bcast_ended_pos < bcast_pos,
"broadcasts_ended must be loaded before broadcasts_started",
);
}
#[test]
fn context_presence_closes_on_last_clone() {
let stats = test_stats();
let snap = |root: &str| {
session_snapshot(&stats, &Tier::default(), root).map(|p| (p.sessions_started, p.sessions_ended))
};
let ctx = stats.tier(Tier::default()).session("acme");
assert_eq!(snap("acme"), Some((1, 0)));
let clone = ctx.clone();
assert_eq!(snap("acme"), Some((1, 0)));
drop(ctx);
assert_eq!(snap("acme"), Some((1, 0)));
drop(clone);
assert_eq!(snap("acme"), Some((1, 1)));
}
#[test]
fn set_tier_moves_presence() {
let stats = test_stats();
let gold = Tier::new("gold");
let snap = |tier: &Tier| session_snapshot(&stats, tier, "acme").map(|p| (p.sessions_started, p.sessions_ended));
let ctx = stats.tier(Tier::default()).session("acme");
ctx.set_tier(Tier::default());
assert_eq!(snap(&Tier::default()), Some((1, 0)), "the same tier is a no-op");
ctx.set_tier(gold.clone());
assert_eq!(snap(&Tier::default()), Some((1, 1)));
assert_eq!(snap(&gold), Some((1, 0)));
drop(ctx);
assert_eq!(snap(&gold), Some((1, 1)), "the session closes on its current tier");
}
#[test]
fn set_tier_moves_subsequent_traffic() {
let stats = test_stats();
let gold = Tier::new("gold");
let ctx = stats.tier(Tier::default()).session("acme");
let scope = ctx.egress("demo/bbb");
let clone = scope.clone();
let before = scope.meter();
let sub = scope.subscribe();
let announce = scope.announce();
before.bytes(10);
ctx.set_tier(gold.clone());
before.bytes(1);
scope.meter().bytes(5);
clone.meter().bytes(7);
let sub2 = scope.subscribe();
scope.fetch();
drop(sub);
drop(sub2);
drop(announce);
let old = tier_counters(&stats, "demo/bbb", &Tier::default()).publisher.snapshot();
let new = tier_counters(&stats, "demo/bbb", &gold).publisher.snapshot();
assert_eq!(old.bytes, 11);
assert_eq!(new.bytes, 12);
assert_eq!((old.fetches, new.fetches), (0, 1));
assert_eq!((old.subscriptions_started, old.subscriptions_ended), (1, 1));
assert_eq!((new.subscriptions_started, new.subscriptions_ended), (1, 1));
assert_eq!((old.announces_started, old.announces_ended), (1, 1));
assert_eq!((old.broadcasts_started, old.broadcasts_ended), (1, 1));
assert_eq!((new.broadcasts_started, new.broadcasts_ended), (0, 0));
assert!(old.is_idle() && new.is_idle());
}
#[test]
fn meter_bumps_the_right_side() {
let stats = test_stats();
let ctx = stats.tier(Tier::default()).session("root");
let egress = ctx.egress("demo/bbb").meter();
egress.group();
egress.frames(3);
egress.bytes(100);
let ingress = ctx.ingress("demo/bbb").meter();
ingress.group();
ingress.frames(1);
ingress.bytes(7);
let counters = tier_counters(&stats, "demo/bbb", &Tier::default());
let pub_ = counters.publisher.snapshot();
let sub = counters.subscriber.snapshot();
assert_eq!((pub_.groups, pub_.frames, pub_.bytes), (1, 3, 100));
assert_eq!((sub.groups, sub.frames, sub.bytes), (1, 1, 7));
}
#[test]
fn egress_subscribe_drives_subscriptions_and_viewers() {
let stats = test_stats();
let ctx = stats.tier(Tier::default()).session("root");
let raw = || tier_counters(&stats, "demo/bbb", &Tier::default()).publisher.snapshot();
let scope = ctx.egress("demo/bbb");
let s1 = scope.subscribe();
let s2 = scope.subscribe();
let r = raw();
assert_eq!(r.subscriptions_started, 2, "two track subs");
assert_eq!(r.broadcasts_started, 1, "one context => one viewer");
assert_eq!(r.broadcasts_ended, 0);
drop(s1);
assert_eq!(raw().broadcasts_ended, 0, "context still has a sub open");
drop(s2);
let r = raw();
assert_eq!(r.subscriptions_ended, 2);
assert_eq!(r.broadcasts_ended, 1, "last sub closed => one broadcasts_ended");
}
#[test]
fn distinct_contexts_are_distinct_viewers() {
let stats = test_stats();
let raw = || tier_counters(&stats, "demo/bbb", &Tier::default()).publisher.snapshot();
let v1 = stats.tier(Tier::default()).session("a").egress("demo/bbb").subscribe();
assert_eq!(raw().broadcasts_started, 1);
let v2 = stats.tier(Tier::default()).session("b").egress("demo/bbb").subscribe();
assert_eq!(raw().broadcasts_started, 2, "two distinct contexts => two viewers");
drop(v1);
assert_eq!(raw().active_broadcasts(), 1);
drop(v2);
assert_eq!(raw().broadcasts_ended, 2);
}
#[test]
fn ingress_subscription_has_no_viewer() {
let stats = test_stats();
let ctx = stats.tier(Tier::default()).session("root");
let guard = ctx.ingress("demo/bbb").subscribe();
let sub = tier_counters(&stats, "demo/bbb", &Tier::default())
.subscriber
.snapshot();
assert_eq!(sub.subscriptions_started, 1);
assert_eq!(sub.broadcasts_started, 0, "ingress has no viewer refcount");
drop(guard);
assert_eq!(
tier_counters(&stats, "demo/bbb", &Tier::default())
.subscriber
.snapshot()
.subscriptions_ended,
1
);
}
#[test]
fn fetch_counts_separately_from_subscriptions() {
let stats = test_stats();
let ctx = stats.tier(Tier::default()).session("root");
let scope = ctx.egress("demo/bbb");
scope.fetch();
scope.fetch();
let r = tier_counters(&stats, "demo/bbb", &Tier::default()).publisher.snapshot();
assert_eq!(r.fetches, 2);
assert_eq!(r.subscriptions_started, 0);
assert_eq!(r.broadcasts_started, 0);
}
#[test]
fn announce_guard_records_bytes_on_open_and_close() {
let stats = test_stats();
let ctx = stats.tier(Tier::default()).session("root");
let path_len = "demo/bbb".len() as u64;
let guard = ctx.egress("demo/bbb").announce();
let r = tier_counters(&stats, "demo/bbb", &Tier::default()).publisher.snapshot();
assert_eq!(r.announces_started, 1);
assert_eq!(r.announces_ended, 0);
assert_eq!(r.announced_bytes, path_len);
drop(guard);
let r = tier_counters(&stats, "demo/bbb", &Tier::default()).publisher.snapshot();
assert_eq!(r.announces_ended, 1);
assert_eq!(
r.announced_bytes,
path_len * 2,
"path length recorded on open and close"
);
}
#[test]
fn disabled_context_is_noop() {
let ctx = Session::default();
let scope = ctx.egress("demo/bbb");
scope.meter().bytes(100);
let _guard = scope.announce();
let _sub = scope.subscribe();
scope.fetch();
assert!(ctx.inner.is_none());
}
#[test]
fn fetches_serde_roundtrips() {
let old: Traffic = serde_json::from_str(r#"{"bytes":5}"#).expect("older shape parses");
assert_eq!(old.fetches, 0);
let t = Traffic {
fetches: 9,
..Default::default()
};
let json = serde_json::to_string(&t).unwrap();
let back: Traffic = serde_json::from_str(&json).unwrap();
assert_eq!(back.fetches, 9);
}
#[test]
fn session_snapshot_reads_ended_before_started() {
let src = include_str!("stats.rs");
let body_start = src
.find("fn snapshot(&self) -> Presence")
.expect("SessionCounters::snapshot fn present");
let body = &src[body_start..];
let ended_pos = body.find("self.sessions_ended.load").expect("sessions_ended load");
let started_pos = body.find("self.sessions_started.load").expect("sessions_started load");
assert!(
ended_pos < started_pos,
"sessions_ended must be loaded before sessions_started",
);
}
fn expected_traffic() -> Traffic {
Traffic {
announces_started: 2,
announces_ended: 1,
broadcasts_started: 4,
broadcasts_ended: 3,
subscriptions_started: 6,
subscriptions_ended: 5,
bytes: 9,
..Default::default()
}
}
fn expected_presence() -> Presence {
Presence {
sessions_started: 3,
sessions_ended: 1,
}
}
#[test]
fn traffic_decodes_old_new_and_both_spellings() {
let expected = expected_traffic();
let old = r#"{"announced":2,"announced_closed":1,"broadcasts":4,"broadcasts_closed":3,"subscriptions":6,"subscriptions_closed":5,"bytes":9}"#;
let new = r#"{"announces_started":2,"announces_ended":1,"broadcasts_started":4,"broadcasts_ended":3,"subscriptions_started":6,"subscriptions_ended":5,"bytes":9}"#;
assert_eq!(serde_json::from_str::<Traffic>(old).unwrap(), expected);
assert_eq!(serde_json::from_str::<Traffic>(new).unwrap(), expected);
let both = serde_json::to_string(&expected).unwrap();
assert!(both.contains("\"announces_started\":2"), "{both}");
assert!(both.contains("\"announced\":2"), "{both}");
assert!(both.contains("\"announces_ended\":1"), "{both}");
assert!(both.contains("\"announced_closed\":1"), "{both}");
assert!(both.contains("\"broadcasts_started\":4"), "{both}");
assert!(both.contains("\"broadcasts\":4"), "{both}");
assert!(both.contains("\"subscriptions_started\":6"), "{both}");
assert!(both.contains("\"subscriptions\":6"), "{both}");
assert_eq!(serde_json::from_str::<Traffic>(&both).unwrap(), expected);
}
#[test]
fn presence_decodes_old_new_and_both_spellings() {
let expected = expected_presence();
assert_eq!(
serde_json::from_str::<Presence>(r#"{"sessions":3,"sessions_closed":1}"#).unwrap(),
expected
);
assert_eq!(
serde_json::from_str::<Presence>(r#"{"sessions_started":3,"sessions_ended":1}"#).unwrap(),
expected
);
let both = serde_json::to_string(&expected).unwrap();
assert!(both.contains("\"sessions_started\":3"), "{both}");
assert!(both.contains("\"sessions\":3"), "{both}");
assert!(both.contains("\"sessions_ended\":1"), "{both}");
assert!(both.contains("\"sessions_closed\":1"), "{both}");
assert_eq!(serde_json::from_str::<Presence>(&both).unwrap(), expected);
}
#[test]
fn counter_edge_canonical_wins_when_spellings_disagree() {
let traffic: Traffic =
serde_json::from_str(r#"{"announces_started":9,"announced":1,"announces_ended":8,"announced_closed":0}"#)
.unwrap();
assert_eq!(traffic.announces_started, 9);
assert_eq!(traffic.announces_ended, 8);
let presence: Presence =
serde_json::from_str(r#"{"sessions_started":4,"sessions":0,"sessions_ended":2,"sessions_closed":9}"#)
.unwrap();
assert_eq!(presence.sessions_started, 4);
assert_eq!(presence.sessions_ended, 2);
}
#[test]
fn counter_edge_refuses_null() {
assert!(serde_json::from_str::<Traffic>(r#"{"announces_started":null,"announced":7}"#).is_err());
assert!(serde_json::from_str::<Traffic>(r#"{"subscriptions_closed":null}"#).is_err());
assert!(serde_json::from_str::<Presence>(r#"{"sessions_started":null}"#).is_err());
assert!(serde_json::from_str::<Presence>(r#"{"sessions":null,"sessions_closed":1}"#).is_err());
}
}