1use std::{
107 collections::HashMap,
108 fmt,
109 sync::{
110 Arc, Mutex,
111 atomic::{AtomicU64, Ordering},
112 },
113};
114
115use serde::{Deserialize, Serialize};
116use web_async::Lock;
117
118use crate::{AsPath, PathOwned};
119
120#[derive(Default, Debug)]
133pub(crate) struct Counters {
134 announced: AtomicU64,
135 announced_closed: AtomicU64,
136 announced_bytes: AtomicU64,
141 subscriptions: AtomicU64,
142 subscriptions_closed: AtomicU64,
143 fetches: AtomicU64,
147 broadcasts: AtomicU64,
148 broadcasts_closed: AtomicU64,
149 bytes: AtomicU64,
150 frames: AtomicU64,
151 groups: AtomicU64,
152 datagrams: AtomicU64,
154}
155
156impl Counters {
157 fn snapshot(&self) -> Traffic {
165 let announced_closed = self.announced_closed.load(Ordering::Acquire);
166 let subscriptions_closed = self.subscriptions_closed.load(Ordering::Acquire);
167 let broadcasts_closed = self.broadcasts_closed.load(Ordering::Acquire);
168 let announced = self.announced.load(Ordering::Relaxed);
169 let announced_bytes = self.announced_bytes.load(Ordering::Relaxed);
170 let subscriptions = self.subscriptions.load(Ordering::Relaxed);
171 let fetches = self.fetches.load(Ordering::Relaxed);
172 let broadcasts = self.broadcasts.load(Ordering::Relaxed);
173 let bytes = self.bytes.load(Ordering::Relaxed);
174 let frames = self.frames.load(Ordering::Relaxed);
175 let groups = self.groups.load(Ordering::Relaxed);
176 let datagrams = self.datagrams.load(Ordering::Relaxed);
177 Traffic {
178 announced,
179 announced_closed,
180 announced_bytes,
181 broadcasts,
182 broadcasts_closed,
183 subscriptions,
184 subscriptions_closed,
185 fetches,
186 bytes,
187 frames,
188 groups,
189 datagrams,
190 }
191 }
192}
193
194#[derive(Default, Debug)]
198struct SessionCounters {
199 sessions: AtomicU64,
200 sessions_closed: AtomicU64,
201}
202
203impl SessionCounters {
204 fn snapshot(&self) -> Presence {
208 let sessions_closed = self.sessions_closed.load(Ordering::Acquire);
209 let sessions = self.sessions.load(Ordering::Relaxed);
210 Presence {
211 sessions,
212 sessions_closed,
213 }
214 }
215}
216
217#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
226#[serde(default)]
227#[non_exhaustive]
228pub struct Traffic {
229 pub announced: u64,
231 pub announced_closed: u64,
233 pub announced_bytes: u64,
236 pub broadcasts: u64,
239 pub broadcasts_closed: u64,
241 pub subscriptions: u64,
243 pub subscriptions_closed: u64,
245 pub fetches: u64,
250 pub bytes: u64,
252 pub frames: u64,
254 pub groups: u64,
256 pub datagrams: u64,
260}
261
262impl Traffic {
263 pub fn add(&mut self, other: Traffic) {
265 self.announced += other.announced;
266 self.announced_closed += other.announced_closed;
267 self.announced_bytes += other.announced_bytes;
268 self.broadcasts += other.broadcasts;
269 self.broadcasts_closed += other.broadcasts_closed;
270 self.subscriptions += other.subscriptions;
271 self.subscriptions_closed += other.subscriptions_closed;
272 self.fetches += other.fetches;
273 self.bytes += other.bytes;
274 self.frames += other.frames;
275 self.groups += other.groups;
276 self.datagrams += other.datagrams;
277 }
278
279 pub fn is_announced(&self) -> bool {
281 self.announced > self.announced_closed
282 }
283
284 pub fn active_broadcasts(&self) -> u64 {
286 self.broadcasts.saturating_sub(self.broadcasts_closed)
287 }
288
289 pub fn active_subscriptions(&self) -> u64 {
291 self.subscriptions.saturating_sub(self.subscriptions_closed)
292 }
293
294 pub fn total_bytes(&self) -> u64 {
298 self.bytes.saturating_add(self.announced_bytes)
299 }
300
301 pub fn is_idle(&self) -> bool {
304 self.announced == self.announced_closed
305 && self.subscriptions == self.subscriptions_closed
306 && self.broadcasts == self.broadcasts_closed
307 }
308}
309
310#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
317#[serde(default)]
318#[non_exhaustive]
319pub struct Presence {
320 pub sessions: u64,
322 pub sessions_closed: u64,
324}
325
326impl Presence {
327 pub fn add(&mut self, other: Presence) {
329 self.sessions += other.sessions;
330 self.sessions_closed += other.sessions_closed;
331 }
332
333 pub fn active(&self) -> u64 {
335 self.sessions.saturating_sub(self.sessions_closed)
336 }
337}
338
339#[derive(Clone, Debug, Default, PartialEq, Eq, Hash)]
350pub struct Tier(PathOwned);
351
352impl Tier {
353 pub fn new(label: impl Into<PathOwned>) -> Self {
355 Self(label.into())
356 }
357
358 pub fn label(&self) -> &PathOwned {
360 &self.0
361 }
362
363 pub fn is_default(&self) -> bool {
365 self.0.is_empty()
366 }
367
368 pub fn track_name(&self, name: &str) -> String {
371 if self.0.is_empty() {
372 name.to_string()
373 } else {
374 format!("{}/{}", self.0.as_str(), name)
375 }
376 }
377
378 pub fn as_str(&self) -> &str {
383 self.0.as_str()
384 }
385}
386
387impl fmt::Display for Tier {
388 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
390 fmt::Display::fmt(&self.0, f)
391 }
392}
393
394#[derive(Copy, Clone, Debug, PartialEq, Eq, Hash)]
401pub enum Role {
402 Publisher,
404 Subscriber,
406}
407
408impl Role {
409 fn idx(self) -> usize {
410 match self {
411 Role::Publisher => 0,
412 Role::Subscriber => 1,
413 }
414 }
415
416 pub fn as_str(self) -> &'static str {
418 match self {
419 Role::Publisher => "publisher",
420 Role::Subscriber => "subscriber",
421 }
422 }
423}
424
425#[derive(Debug, Default, Clone, PartialEq, Eq)]
436#[non_exhaustive]
437pub struct Snapshot {
438 traffic: HashMap<Tier, [Traffic; 2]>,
441 sessions: HashMap<Tier, Presence>,
443}
444
445impl Snapshot {
446 pub fn traffic(&self) -> Vec<(Tier, Role, Traffic)> {
449 let mut rows = Vec::with_capacity(self.traffic.len() * 2);
450 for (tier, roles) in &self.traffic {
451 rows.push((tier.clone(), Role::Publisher, roles[Role::Publisher.idx()]));
452 rows.push((tier.clone(), Role::Subscriber, roles[Role::Subscriber.idx()]));
453 }
454 rows.sort_by(|a, b| a.0.as_str().cmp(b.0.as_str()).then(a.1.idx().cmp(&b.1.idx())));
455 rows
456 }
457
458 pub fn sessions(&self) -> Vec<(Tier, Presence)> {
461 let mut rows: Vec<_> = self.sessions.iter().map(|(tier, s)| (tier.clone(), *s)).collect();
462 rows.sort_by(|a, b| a.0.as_str().cmp(b.0.as_str()));
463 rows
464 }
465}
466
467#[derive(Debug, Default, Clone)]
471#[non_exhaustive]
472pub struct Report {
473 pub traffic: Vec<TrafficEntry>,
475 pub sessions: Vec<SessionEntry>,
477}
478
479#[derive(Debug, Clone)]
481#[non_exhaustive]
482pub struct TrafficEntry {
483 pub path: PathOwned,
485 pub tier: Tier,
487 pub publisher: Traffic,
489 pub subscriber: Traffic,
491}
492
493#[derive(Debug, Clone)]
495#[non_exhaustive]
496pub struct SessionEntry {
497 pub tier: Tier,
499 pub root: PathOwned,
501 pub presence: Presence,
503}
504
505#[derive(Clone, Debug, Default)]
511#[non_exhaustive]
512pub struct Config {
513 pub exclude: Vec<PathOwned>,
518}
519
520impl Config {
521 pub fn new() -> Self {
523 Self::default()
524 }
525
526 pub fn with_exclude(mut self, prefix: impl Into<PathOwned>) -> Self {
529 self.exclude.push(prefix.into());
530 self
531 }
532}
533
534#[derive(Clone)]
539pub struct Registry {
540 exclude: Vec<PathOwned>,
543 shared: Option<Arc<Shared>>,
545}
546
547struct Shared {
549 entries: Lock<HashMap<PathOwned, Arc<BroadcastEntry>>>,
550 sessions: Lock<HashMap<Tier, HashMap<PathOwned, Arc<SessionCounters>>>>,
554}
555
556struct BroadcastEntry {
561 tiers: Mutex<HashMap<Tier, Arc<TierCounters>>>,
562}
563
564impl BroadcastEntry {
565 fn new() -> Self {
566 Self {
567 tiers: Mutex::new(HashMap::new()),
568 }
569 }
570
571 fn tier(&self, tier: &Tier) -> Arc<TierCounters> {
573 self.tiers
574 .lock()
575 .expect("stats tiers poisoned")
576 .entry(tier.clone())
577 .or_default()
578 .clone()
579 }
580}
581
582#[derive(Default)]
586struct TierCounters {
587 publisher: Counters,
588 subscriber: Counters,
589}
590
591impl Registry {
592 pub fn new(config: Config) -> Self {
594 let Config { exclude } = config;
595 Self {
596 exclude,
597 shared: Some(Arc::new(Shared {
598 entries: Lock::default(),
599 sessions: Default::default(),
600 })),
601 }
602 }
603
604 pub fn disabled() -> Self {
606 Self {
607 exclude: Vec::new(),
608 shared: None,
609 }
610 }
611
612 pub fn exclude(&self) -> &[PathOwned] {
614 &self.exclude
615 }
616
617 #[cfg(test)]
620 fn shared(&self) -> &Arc<Shared> {
621 self.shared.as_ref().expect("enabled stats registry")
622 }
623
624 pub fn tier(&self, tier: Tier) -> Handle {
627 Handle {
628 stats: self.clone(),
629 tier,
630 }
631 }
632
633 fn entry(&self, path: impl AsPath) -> Option<Arc<BroadcastEntry>> {
634 let shared = self.shared.as_ref()?;
636 let path = path.as_path();
637 if self.exclude.iter().any(|prefix| path.has_prefix(prefix)) {
641 return None;
642 }
643 let owned = path.to_owned();
644 let mut entries = shared.entries.lock();
645 Some(
646 entries
647 .entry(owned)
648 .or_insert_with(|| Arc::new(BroadcastEntry::new()))
649 .clone(),
650 )
651 }
652
653 fn session_counters(&self, tier: &Tier, root: impl AsPath) -> Option<Arc<SessionCounters>> {
657 let shared = self.shared.as_ref()?;
658 let owned = root.as_path().to_owned();
659 let mut sessions = shared.sessions.lock();
660 Some(
661 sessions
662 .entry(tier.clone())
663 .or_default()
664 .entry(owned)
665 .or_default()
666 .clone(),
667 )
668 }
669
670 pub fn snapshot(&self) -> Snapshot {
678 let mut snap = Snapshot::default();
679 let Some(shared) = self.shared.as_ref() else {
680 return snap;
681 };
682 {
683 let entries = shared.entries.lock();
684 for entry in entries.values() {
685 let tiers = entry.tiers.lock().expect("stats tiers poisoned");
686 for (tier, counters) in tiers.iter() {
687 let totals = snap.traffic.entry(tier.clone()).or_default();
688 totals[Role::Publisher.idx()].add(counters.publisher.snapshot());
689 totals[Role::Subscriber.idx()].add(counters.subscriber.snapshot());
690 }
691 }
692 }
693 {
694 let sessions = shared.sessions.lock();
695 for (tier, roots) in sessions.iter() {
696 let totals = snap.sessions.entry(tier.clone()).or_default();
697 for counters in roots.values() {
698 totals.add(counters.snapshot());
699 }
700 }
701 }
702 snap
703 }
704
705 pub fn report(&self) -> Report {
715 let mut report = Report::default();
716 let Some(shared) = self.shared.as_ref() else {
717 return report;
718 };
719 {
720 let mut entries = shared.entries.lock();
721 for (path, entry) in entries.iter() {
722 let tiers = entry.tiers.lock().expect("stats tiers poisoned");
723 for (tier, counters) in tiers.iter() {
724 report.traffic.push(TrafficEntry {
725 path: path.clone(),
726 tier: tier.clone(),
727 publisher: counters.publisher.snapshot(),
728 subscriber: counters.subscriber.snapshot(),
729 });
730 }
731 }
732 entries.retain(|_, entry| {
737 if Arc::strong_count(entry) > 1 {
738 return true;
739 }
740 let mut tiers = entry.tiers.lock().expect("stats tiers poisoned");
741 tiers.retain(|_, counters| Arc::strong_count(counters) > 1);
742 !tiers.is_empty()
743 });
744 }
745 {
746 let mut sessions = shared.sessions.lock();
747 for (tier, roots) in sessions.iter() {
748 for (root, counters) in roots.iter() {
749 report.sessions.push(SessionEntry {
750 tier: tier.clone(),
751 root: root.clone(),
752 presence: counters.snapshot(),
753 });
754 }
755 }
756 for roots in sessions.values_mut() {
757 roots.retain(|_, counters| Arc::strong_count(counters) > 1);
758 }
759 sessions.retain(|_, roots| !roots.is_empty());
760 }
761 report
762 }
763}
764
765impl Default for Registry {
766 fn default() -> Self {
768 Self::disabled()
769 }
770}
771
772#[derive(Clone)]
775pub struct Handle {
776 stats: Registry,
777 tier: Tier,
778}
779
780impl Handle {
781 pub fn parent(&self) -> &Registry {
783 &self.stats
784 }
785
786 pub fn tier(&self) -> &Tier {
788 &self.tier
789 }
790
791 pub fn session(&self, root: impl AsPath) -> Session {
797 Session::new(self.clone(), self.stats.session_counters(&self.tier, root))
798 }
799}
800
801impl Default for Handle {
802 fn default() -> Self {
804 Registry::disabled().tier(Tier::default())
805 }
806}
807
808#[derive(Copy, Clone, Default)]
812enum Side {
813 #[default]
814 Publisher,
815 Subscriber,
816}
817
818impl Side {
819 fn counters(self, tier: &TierCounters) -> &Counters {
820 match self {
821 Side::Publisher => &tier.publisher,
822 Side::Subscriber => &tier.subscriber,
823 }
824 }
825}
826
827#[derive(Clone, Default)]
845pub struct Session {
846 inner: Option<Arc<SessionInner>>,
848}
849
850struct SessionInner {
853 handle: Handle,
855 presence: Option<Arc<SessionCounters>>,
857 viewers: Mutex<HashMap<PathOwned, u32>>,
861}
862
863impl Session {
864 fn new(handle: Handle, presence: Option<Arc<SessionCounters>>) -> Self {
865 if let Some(presence) = &presence {
866 presence.sessions.fetch_add(1, Ordering::Relaxed);
867 }
868 Self {
869 inner: Some(Arc::new(SessionInner {
870 handle,
871 presence,
872 viewers: Mutex::new(HashMap::new()),
873 })),
874 }
875 }
876
877 pub(crate) fn egress(&self, path: impl AsPath) -> Scope {
880 self.scope(path, Side::Publisher)
881 }
882
883 pub(crate) fn ingress(&self, path: impl AsPath) -> Scope {
885 self.scope(path, Side::Subscriber)
886 }
887
888 fn scope(&self, path: impl AsPath, side: Side) -> Scope {
889 let Some(inner) = &self.inner else {
890 return Scope::default();
891 };
892 let path = path.as_path().to_owned();
893 let counters = inner
894 .handle
895 .stats
896 .entry(&path)
897 .map(|entry| entry.tier(&inner.handle.tier));
898 Scope {
899 session: self.clone(),
900 counters,
901 side,
902 path,
903 }
904 }
905
906 fn viewer_open(&self, path: &PathOwned) -> bool {
909 let Some(inner) = &self.inner else { return false };
910 let mut viewers = inner.viewers.lock().expect("stats viewers poisoned");
911 let n = viewers.entry(path.clone()).or_insert(0);
912 let first = *n == 0;
913 *n += 1;
914 first
915 }
916
917 fn viewer_close(&self, path: &PathOwned) -> bool {
920 let Some(inner) = &self.inner else { return false };
921 let mut viewers = inner.viewers.lock().expect("stats viewers poisoned");
922 match viewers.get_mut(path) {
923 Some(n) => {
924 *n -= 1;
925 if *n == 0 {
926 viewers.remove(path);
927 true
928 } else {
929 false
930 }
931 }
932 None => false,
933 }
934 }
935}
936
937impl Drop for SessionInner {
938 fn drop(&mut self) {
939 if let Some(presence) = &self.presence {
940 presence.sessions_closed.fetch_add(1, Ordering::Release);
943 }
944 }
945}
946
947#[derive(Clone, Default)]
961pub(crate) struct Meter {
962 counters: Option<Arc<TierCounters>>,
963 side: Side,
964}
965
966impl Meter {
967 fn counters(&self) -> Option<&Counters> {
968 self.counters.as_ref().map(|c| self.side.counters(c))
969 }
970
971 pub(crate) fn group(&self) {
973 if let Some(counters) = self.counters() {
974 counters.groups.fetch_add(1, Ordering::Relaxed);
975 }
976 }
977
978 pub(crate) fn frames(&self, n: u64) {
980 if n == 0 {
981 return;
982 }
983 if let Some(counters) = self.counters() {
984 counters.frames.fetch_add(n, Ordering::Relaxed);
985 }
986 }
987
988 pub(crate) fn datagram(&self, n: u64) {
992 if let Some(counters) = self.counters() {
993 counters.datagrams.fetch_add(1, Ordering::Relaxed);
994 counters.groups.fetch_add(1, Ordering::Relaxed);
995 counters.frames.fetch_add(1, Ordering::Relaxed);
996 counters.bytes.fetch_add(n, Ordering::Relaxed);
997 }
998 }
999
1000 pub(crate) fn bytes(&self, n: u64) {
1002 if n == 0 {
1003 return;
1004 }
1005 if let Some(counters) = self.counters() {
1006 counters.bytes.fetch_add(n, Ordering::Relaxed);
1007 }
1008 }
1009}
1010
1011#[derive(Clone, Default)]
1016pub(crate) struct Scope {
1017 session: Session,
1019 counters: Option<Arc<TierCounters>>,
1021 side: Side,
1022 path: PathOwned,
1025}
1026
1027impl Scope {
1028 fn counters(&self) -> Option<&Counters> {
1029 self.counters.as_ref().map(|c| self.side.counters(c))
1030 }
1031
1032 pub(crate) fn meter(&self) -> Meter {
1034 Meter {
1035 counters: self.counters.clone(),
1036 side: self.side,
1037 }
1038 }
1039
1040 pub(crate) fn subscribe(&self) -> Subscription {
1044 if let Some(counters) = self.counters() {
1045 counters.subscriptions.fetch_add(1, Ordering::Relaxed);
1046 }
1047 let viewer = if matches!(self.side, Side::Publisher) && self.counters.is_some() {
1050 if self.session.viewer_open(&self.path)
1051 && let Some(counters) = self.counters()
1052 {
1053 counters.broadcasts.fetch_add(1, Ordering::Relaxed);
1054 }
1055 Some((self.session.clone(), self.path.clone()))
1056 } else {
1057 None
1058 };
1059 Subscription {
1060 counters: self.counters.clone(),
1061 side: self.side,
1062 viewer,
1063 }
1064 }
1065
1066 pub(crate) fn fetch(&self) {
1068 if let Some(counters) = self.counters() {
1069 counters.fetches.fetch_add(1, Ordering::Relaxed);
1070 }
1071 }
1072
1073 pub(crate) fn open_subscription(&self) {
1078 if let Some(counters) = self.counters() {
1079 counters.subscriptions.fetch_add(1, Ordering::Relaxed);
1080 }
1081 }
1082
1083 pub(crate) fn close_subscription(&self) {
1085 if let Some(counters) = self.counters() {
1086 counters.subscriptions_closed.fetch_add(1, Ordering::Release);
1088 }
1089 }
1090
1091 pub(crate) fn announce(&self) -> Announce {
1096 let len = self.path.as_str().len() as u64;
1097 if let Some(counters) = self.counters() {
1098 counters.announced.fetch_add(1, Ordering::Relaxed);
1099 counters.announced_bytes.fetch_add(len, Ordering::Relaxed);
1100 }
1101 Announce {
1102 counters: self.counters.clone(),
1103 side: self.side,
1104 len,
1105 }
1106 }
1107}
1108
1109#[derive(Default)]
1112#[must_use = "drop the guard to record the subscription as closed"]
1113pub(crate) struct Subscription {
1114 counters: Option<Arc<TierCounters>>,
1115 side: Side,
1116 viewer: Option<(Session, PathOwned)>,
1118}
1119
1120impl Drop for Subscription {
1121 fn drop(&mut self) {
1122 if let Some((session, path)) = &self.viewer
1123 && session.viewer_close(path)
1124 && let Some(counters) = &self.counters
1125 {
1126 self.side
1128 .counters(counters)
1129 .broadcasts_closed
1130 .fetch_add(1, Ordering::Release);
1131 }
1132 if let Some(counters) = &self.counters {
1133 self.side
1135 .counters(counters)
1136 .subscriptions_closed
1137 .fetch_add(1, Ordering::Release);
1138 }
1139 }
1140}
1141
1142#[must_use = "drop the guard to record the unannounce"]
1144pub(crate) struct Announce {
1145 counters: Option<Arc<TierCounters>>,
1146 side: Side,
1147 len: u64,
1148}
1149
1150impl Drop for Announce {
1151 fn drop(&mut self) {
1152 if let Some(counters) = &self.counters {
1153 let counters = self.side.counters(counters);
1154 counters.announced_bytes.fetch_add(self.len, Ordering::Relaxed);
1155 counters.announced_closed.fetch_add(1, Ordering::Release);
1157 }
1158 }
1159}
1160
1161#[cfg(test)]
1162mod tests {
1163 use std::sync::{Arc, atomic::Ordering::Relaxed};
1164
1165 use super::*;
1166
1167 #[test]
1168 fn default_tier_has_empty_label() {
1169 let tier = Tier::default();
1170 assert_eq!(tier.as_str(), "");
1171 assert_eq!(tier.to_string(), "");
1172 assert_eq!(tier.track_name("publisher.json"), "publisher.json");
1173 }
1174
1175 fn tier_counters(stats: &Registry, path: &str, tier: &Tier) -> Arc<TierCounters> {
1177 stats
1178 .shared()
1179 .entries
1180 .lock()
1181 .get(&PathOwned::from(path.to_string()))
1182 .expect("entry")
1183 .tier(tier)
1184 }
1185
1186 fn session_snapshot(stats: &Registry, tier: &Tier, root: &str) -> Option<Presence> {
1188 stats
1189 .shared()
1190 .sessions
1191 .lock()
1192 .get(tier)
1193 .and_then(|roots| roots.get(&PathOwned::from(root.to_string())).map(|c| c.snapshot()))
1194 }
1195
1196 fn test_stats() -> Registry {
1197 Registry::new(Config::new().with_exclude(".stats"))
1198 }
1199
1200 #[test]
1201 fn default_and_named_tiers_are_independent() {
1202 let stats = test_stats();
1203 let default = stats.tier(Tier::default()).session("root");
1204 let regional = stats.tier(Tier::new("region/sjc")).session("root");
1205
1206 default.egress("demo/bbb").meter().bytes(100);
1207 regional.ingress("demo/bbb").meter().bytes(7);
1208
1209 let default_counters = tier_counters(&stats, "demo/bbb", &Tier::default());
1210 let regional_counters = tier_counters(&stats, "demo/bbb", &Tier::new("region/sjc"));
1211 assert_eq!(default_counters.publisher.bytes.load(Relaxed), 100);
1212 assert_eq!(default_counters.subscriber.bytes.load(Relaxed), 0);
1213 assert_eq!(regional_counters.publisher.bytes.load(Relaxed), 0);
1214 assert_eq!(regional_counters.subscriber.bytes.load(Relaxed), 7);
1215 }
1216
1217 #[test]
1218 fn snapshot_rolls_up_by_tier_role_and_sessions() {
1219 let stats = test_stats();
1220 let default = stats.tier(Tier::default());
1221 let regional = stats.tier(Tier::new("region/sjc"));
1222
1223 let s1 = default.session("acme");
1225 let _s2 = default.session("acme");
1226 let s3 = regional.session("peer");
1227
1228 {
1230 let m = s1.egress("demo/aaa").meter();
1231 m.bytes(100);
1232 m.frames(1);
1233 m.group();
1234 }
1235 s1.egress("demo/bbb").meter().bytes(50);
1236 s3.ingress("demo/aaa").meter().bytes(7);
1238
1239 let snap = stats.snapshot();
1240
1241 let slot = |tier, role| {
1242 snap.traffic()
1243 .into_iter()
1244 .find(|(t, r, _)| *t == tier && *r == role)
1245 .map(|(_, _, c)| c)
1246 .expect("row present")
1247 };
1248
1249 let default_publisher = slot(Tier::default(), Role::Publisher);
1250 assert_eq!(
1251 default_publisher.bytes, 150,
1252 "default egress bytes sum across broadcasts"
1253 );
1254 assert_eq!(default_publisher.frames, 1);
1255 assert_eq!(default_publisher.groups, 1);
1256
1257 let regional_subscriber = slot(Tier::new("region/sjc"), Role::Subscriber);
1258 assert_eq!(regional_subscriber.bytes, 7, "regional ingress isolated by tier/role");
1259 assert_eq!(slot(Tier::default(), Role::Subscriber).bytes, 0);
1260 assert_eq!(slot(Tier::new("region/sjc"), Role::Publisher).bytes, 0);
1261
1262 let sessions = |tier| {
1263 snap.sessions()
1264 .into_iter()
1265 .find(|(t, _)| *t == tier)
1266 .map(|(_, s)| s)
1267 .expect("tier present")
1268 };
1269 let default_sessions = sessions(Tier::default());
1270 assert_eq!(default_sessions.sessions, 2, "two default-tier sessions under one root");
1271 assert_eq!(default_sessions.sessions_closed, 0, "guards still held");
1272 assert_eq!(sessions(Tier::new("region/sjc")).sessions, 1);
1273 }
1274
1275 #[test]
1276 fn report_returns_detail_and_prunes() {
1277 let stats = test_stats();
1281 let key = PathOwned::from("foo/bar");
1282 let ctx = stats.tier(Tier::default()).session("root");
1283 let scope = ctx.egress("foo/bar");
1284 let sub = scope.subscribe();
1285 scope.meter().bytes(42);
1286
1287 let report = stats.report();
1288 let row = report
1289 .traffic
1290 .iter()
1291 .find(|row| row.path == key)
1292 .expect("live entry present");
1293 assert_eq!(row.publisher.bytes, 42);
1294 assert_eq!(row.publisher.subscriptions, 1);
1295 assert!(!row.publisher.is_idle(), "subscription guard still open");
1296 assert!(
1297 stats.shared().entries.lock().contains_key(&key),
1298 "live entry kept across drains"
1299 );
1300
1301 drop(sub);
1302 drop(scope);
1303
1304 let report = stats.report();
1307 let row = report
1308 .traffic
1309 .iter()
1310 .find(|row| row.path == key)
1311 .expect("closing values still reported once");
1312 assert_eq!(row.publisher.subscriptions_closed, 1);
1313 assert!(row.publisher.is_idle());
1314 assert!(
1315 !stats.shared().entries.lock().contains_key(&key),
1316 "fully-closed entry pruned"
1317 );
1318 assert!(stats.report().traffic.is_empty(), "nothing left after the prune");
1319 }
1320
1321 #[test]
1322 fn report_keeps_idle_but_announced_entry() {
1323 let stats = test_stats();
1327 let key = PathOwned::from("foo/bar");
1328 let ctx = stats.tier(Tier::default()).session("root");
1329 let scope = ctx.egress("foo/bar");
1330 let guard = scope.announce();
1331
1332 for _ in 0..3 {
1333 let report = stats.report();
1334 assert!(
1335 report.traffic.iter().any(|row| row.path == key),
1336 "announced-but-idle broadcast stays while the guard is held"
1337 );
1338 }
1339
1340 drop(guard);
1341 drop(scope);
1342 let report = stats.report();
1343 let row = report.traffic.iter().find(|row| row.path == key).expect("final report");
1344 assert!(row.publisher.is_idle());
1345 assert!(!stats.shared().entries.lock().contains_key(&key));
1346 }
1347
1348 #[test]
1349 fn report_prunes_empty_session_roots() {
1350 let stats = test_stats();
1353 let session = stats.tier(Tier::default()).session("acme");
1354
1355 let report = stats.report();
1356 let row = report
1357 .sessions
1358 .iter()
1359 .find(|row| row.root.as_str() == "acme")
1360 .expect("root present");
1361 assert_eq!(row.presence.active(), 1);
1362
1363 drop(session);
1364 let report = stats.report();
1365 let row = report
1366 .sessions
1367 .iter()
1368 .find(|row| row.root.as_str() == "acme")
1369 .expect("final gauge reported once");
1370 assert_eq!(row.presence.active(), 0);
1371 assert!(stats.report().sessions.is_empty(), "root pruned after the last drain");
1372 assert!(session_snapshot(&stats, &Tier::default(), "acme").is_none());
1373 }
1374
1375 #[test]
1376 fn paths_under_exclude_are_no_op() {
1377 let stats = test_stats();
1380 let ctx = stats.tier(Tier::default()).session("root");
1381 let scope = ctx.egress(".stats/node/sjc");
1382 scope.meter().bytes(100);
1383 let _guard = scope.announce();
1384 let _sub = scope.subscribe();
1385 assert!(stats.shared().entries.lock().is_empty());
1386 }
1387
1388 #[test]
1389 fn disabled_stats_are_noop() {
1390 let stats = Registry::default();
1393 assert!(stats.shared.is_none());
1394 let ctx = stats.tier(Tier::default()).session("root");
1395 let scope = ctx.egress("demo/bbb");
1396 scope.meter().bytes(100);
1397 let _guard = scope.announce();
1398 let _sub = scope.subscribe();
1399 assert!(stats.report().traffic.is_empty());
1400 assert!(stats.snapshot().traffic().is_empty());
1401 }
1402
1403 #[test]
1404 fn session_counts_by_root() {
1405 let stats = test_stats();
1408 let ext = stats.tier(Tier::default());
1409
1410 let snap =
1411 |root: &str| session_snapshot(&stats, &Tier::default(), root).map(|p| (p.sessions, p.sessions_closed));
1412
1413 let a1 = ext.session("acme");
1414 let a2 = ext.session("acme");
1415 let b1 = ext.session("globex");
1416 assert_eq!(snap("acme"), Some((2, 0)), "two sessions under one root");
1417 assert_eq!(snap("globex"), Some((1, 0)), "a distinct root is counted separately");
1418
1419 drop(a1);
1420 assert_eq!(snap("acme"), Some((2, 1)));
1421 drop(a2);
1422 drop(b1);
1423 assert_eq!(snap("acme"), Some((2, 2)));
1424 assert_eq!(snap("globex"), Some((1, 1)));
1425 }
1426
1427 #[test]
1428 fn traffic_parses_with_missing_and_unknown_fields() {
1429 let old: Traffic = serde_json::from_str(r#"{"announced":1,"bytes":5}"#).expect("older shape parses");
1432 assert_eq!(old.announced, 1);
1433 assert_eq!(old.bytes, 5);
1434 assert_eq!(old.announced_bytes, 0, "missing fields default to zero");
1435
1436 let new: Traffic = serde_json::from_str(r#"{"announced":1,"announced_closed":1,"future_counter":9}"#)
1437 .expect("newer shape parses");
1438 assert!(new.is_idle());
1439 }
1440
1441 #[test]
1442 fn snapshot_reads_closed_before_open() {
1443 let src = include_str!("stats.rs");
1448 let body_start = src.find("fn snapshot(&self) -> Traffic").expect("snapshot fn present");
1451 let body = &src[body_start..];
1452 let closed_pos = body.find("self.announced_closed.load").expect("announced_closed load");
1453 let open_pos = body.find("self.announced.load(").expect("announced load");
1454 assert!(
1455 closed_pos < open_pos,
1456 "announced_closed must be loaded before announced; reversing breaks the open>=closed invariant",
1457 );
1458 let subs_closed_pos = body
1459 .find("self.subscriptions_closed.load")
1460 .expect("subscriptions_closed load");
1461 let subs_pos = body.find("self.subscriptions.load").expect("subscriptions load");
1462 assert!(
1463 subs_closed_pos < subs_pos,
1464 "subscriptions_closed must be loaded before subscriptions",
1465 );
1466 let bcast_closed_pos = body
1467 .find("self.broadcasts_closed.load")
1468 .expect("broadcasts_closed load");
1469 let bcast_pos = body.find("self.broadcasts.load").expect("broadcasts load");
1470 assert!(
1471 bcast_closed_pos < bcast_pos,
1472 "broadcasts_closed must be loaded before broadcasts",
1473 );
1474 }
1475
1476 #[test]
1477 fn context_presence_closes_on_last_clone() {
1478 let stats = test_stats();
1481 let snap =
1482 |root: &str| session_snapshot(&stats, &Tier::default(), root).map(|p| (p.sessions, p.sessions_closed));
1483
1484 let ctx = stats.tier(Tier::default()).session("acme");
1485 assert_eq!(snap("acme"), Some((1, 0)));
1486
1487 let clone = ctx.clone();
1488 assert_eq!(snap("acme"), Some((1, 0)));
1490 drop(ctx);
1491 assert_eq!(snap("acme"), Some((1, 0)));
1492 drop(clone);
1493 assert_eq!(snap("acme"), Some((1, 1)));
1494 }
1495
1496 #[test]
1497 fn meter_bumps_the_right_side() {
1498 let stats = test_stats();
1500 let ctx = stats.tier(Tier::default()).session("root");
1501
1502 let egress = ctx.egress("demo/bbb").meter();
1503 egress.group();
1504 egress.frames(3);
1505 egress.bytes(100);
1506
1507 let ingress = ctx.ingress("demo/bbb").meter();
1508 ingress.group();
1509 ingress.frames(1);
1510 ingress.bytes(7);
1511
1512 let counters = tier_counters(&stats, "demo/bbb", &Tier::default());
1513 let pub_ = counters.publisher.snapshot();
1514 let sub = counters.subscriber.snapshot();
1515 assert_eq!((pub_.groups, pub_.frames, pub_.bytes), (1, 3, 100));
1516 assert_eq!((sub.groups, sub.frames, sub.bytes), (1, 1, 7));
1517 }
1518
1519 #[test]
1520 fn egress_subscribe_drives_subscriptions_and_viewers() {
1521 let stats = test_stats();
1524 let ctx = stats.tier(Tier::default()).session("root");
1525 let raw = || tier_counters(&stats, "demo/bbb", &Tier::default()).publisher.snapshot();
1526
1527 let scope = ctx.egress("demo/bbb");
1528 let s1 = scope.subscribe();
1529 let s2 = scope.subscribe();
1530 let r = raw();
1531 assert_eq!(r.subscriptions, 2, "two track subs");
1532 assert_eq!(r.broadcasts, 1, "one context => one viewer");
1533 assert_eq!(r.broadcasts_closed, 0);
1534
1535 drop(s1);
1536 assert_eq!(raw().broadcasts_closed, 0, "context still has a sub open");
1537 drop(s2);
1538 let r = raw();
1539 assert_eq!(r.subscriptions_closed, 2);
1540 assert_eq!(r.broadcasts_closed, 1, "last sub closed => one broadcasts_closed");
1541 }
1542
1543 #[test]
1544 fn distinct_contexts_are_distinct_viewers() {
1545 let stats = test_stats();
1547 let raw = || tier_counters(&stats, "demo/bbb", &Tier::default()).publisher.snapshot();
1548
1549 let v1 = stats.tier(Tier::default()).session("a").egress("demo/bbb").subscribe();
1550 assert_eq!(raw().broadcasts, 1);
1551 let v2 = stats.tier(Tier::default()).session("b").egress("demo/bbb").subscribe();
1552 assert_eq!(raw().broadcasts, 2, "two distinct contexts => two viewers");
1553
1554 drop(v1);
1555 assert_eq!(raw().active_broadcasts(), 1);
1556 drop(v2);
1557 assert_eq!(raw().broadcasts_closed, 2);
1558 }
1559
1560 #[test]
1561 fn ingress_subscription_has_no_viewer() {
1562 let stats = test_stats();
1565 let ctx = stats.tier(Tier::default()).session("root");
1566 let scope = ctx.ingress("demo/bbb");
1567 scope.open_subscription();
1568 let sub = tier_counters(&stats, "demo/bbb", &Tier::default())
1569 .subscriber
1570 .snapshot();
1571 assert_eq!(sub.subscriptions, 1);
1572 assert_eq!(sub.broadcasts, 0, "ingress has no viewer refcount");
1573 scope.close_subscription();
1574 assert_eq!(
1575 tier_counters(&stats, "demo/bbb", &Tier::default())
1576 .subscriber
1577 .snapshot()
1578 .subscriptions_closed,
1579 1
1580 );
1581 }
1582
1583 #[test]
1584 fn fetch_counts_separately_from_subscriptions() {
1585 let stats = test_stats();
1587 let ctx = stats.tier(Tier::default()).session("root");
1588 let scope = ctx.egress("demo/bbb");
1589 scope.fetch();
1590 scope.fetch();
1591 let r = tier_counters(&stats, "demo/bbb", &Tier::default()).publisher.snapshot();
1592 assert_eq!(r.fetches, 2);
1593 assert_eq!(r.subscriptions, 0);
1594 assert_eq!(r.broadcasts, 0);
1595 }
1596
1597 #[test]
1598 fn announce_guard_records_bytes_on_open_and_close() {
1599 let stats = test_stats();
1602 let ctx = stats.tier(Tier::default()).session("root");
1603 let path_len = "demo/bbb".len() as u64;
1604
1605 let guard = ctx.egress("demo/bbb").announce();
1606 let r = tier_counters(&stats, "demo/bbb", &Tier::default()).publisher.snapshot();
1607 assert_eq!(r.announced, 1);
1608 assert_eq!(r.announced_closed, 0);
1609 assert_eq!(r.announced_bytes, path_len);
1610
1611 drop(guard);
1612 let r = tier_counters(&stats, "demo/bbb", &Tier::default()).publisher.snapshot();
1613 assert_eq!(r.announced_closed, 1);
1614 assert_eq!(
1615 r.announced_bytes,
1616 path_len * 2,
1617 "path length recorded on open and close"
1618 );
1619 }
1620
1621 #[test]
1622 fn disabled_context_is_noop() {
1623 let ctx = Session::default();
1625 let scope = ctx.egress("demo/bbb");
1626 scope.meter().bytes(100);
1627 let _guard = scope.announce();
1628 let _sub = scope.subscribe();
1629 scope.fetch();
1630 assert!(ctx.inner.is_none());
1632 }
1633
1634 #[test]
1635 fn fetches_serde_roundtrips() {
1636 let old: Traffic = serde_json::from_str(r#"{"bytes":5}"#).expect("older shape parses");
1639 assert_eq!(old.fetches, 0);
1640
1641 let t = Traffic {
1642 fetches: 9,
1643 ..Default::default()
1644 };
1645 let json = serde_json::to_string(&t).unwrap();
1646 let back: Traffic = serde_json::from_str(&json).unwrap();
1647 assert_eq!(back.fetches, 9);
1648 }
1649
1650 #[test]
1651 fn session_snapshot_reads_closed_before_open() {
1652 let src = include_str!("stats.rs");
1656 let body_start = src
1657 .find("fn snapshot(&self) -> Presence")
1658 .expect("SessionCounters::snapshot fn present");
1659 let body = &src[body_start..];
1660 let closed_pos = body.find("self.sessions_closed.load").expect("sessions_closed load");
1661 let open_pos = body.find("self.sessions.load").expect("sessions load");
1662 assert!(closed_pos < open_pos, "sessions_closed must be loaded before sessions",);
1663 }
1664}