1use std::{
129 collections::{BTreeMap, HashMap, HashSet},
130 sync::{
131 Arc, Weak,
132 atomic::{AtomicU64, Ordering},
133 },
134 time::Duration,
135};
136
137use serde::Serialize;
138use web_async::{Lock, spawn};
139
140use crate::{AsPath, Broadcast, OriginProducer, Path, PathOwned, Track, TrackProducer};
141
142#[derive(Default, Debug)]
152#[non_exhaustive]
153pub struct Counters {
154 pub announced: AtomicU64,
155 pub announced_closed: AtomicU64,
156 pub announced_bytes: AtomicU64,
161 pub subscriptions: AtomicU64,
162 pub subscriptions_closed: AtomicU64,
163 pub broadcasts: AtomicU64,
164 pub broadcasts_closed: AtomicU64,
165 pub bytes: AtomicU64,
166 pub frames: AtomicU64,
167 pub groups: AtomicU64,
168}
169
170impl Counters {
171 fn snapshot(&self) -> RawCounts {
179 let announced_closed = self.announced_closed.load(Ordering::Acquire);
180 let subscriptions_closed = self.subscriptions_closed.load(Ordering::Acquire);
181 let broadcasts_closed = self.broadcasts_closed.load(Ordering::Acquire);
182 let announced = self.announced.load(Ordering::Relaxed);
183 let announced_bytes = self.announced_bytes.load(Ordering::Relaxed);
184 let subscriptions = self.subscriptions.load(Ordering::Relaxed);
185 let broadcasts = self.broadcasts.load(Ordering::Relaxed);
186 let bytes = self.bytes.load(Ordering::Relaxed);
187 let frames = self.frames.load(Ordering::Relaxed);
188 let groups = self.groups.load(Ordering::Relaxed);
189 RawCounts {
190 announced,
191 announced_closed,
192 announced_bytes,
193 broadcasts,
194 broadcasts_closed,
195 subscriptions,
196 subscriptions_closed,
197 bytes,
198 frames,
199 groups,
200 }
201 }
202}
203
204#[derive(Default, Debug)]
208struct SessionCounters {
209 sessions: AtomicU64,
210 sessions_closed: AtomicU64,
211}
212
213impl SessionCounters {
214 fn snapshot(&self) -> (u64, u64) {
218 let closed = self.sessions_closed.load(Ordering::Acquire);
219 let open = self.sessions.load(Ordering::Relaxed);
220 (open, closed)
221 }
222}
223
224#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
226struct RawCounts {
227 announced: u64,
228 announced_closed: u64,
229 announced_bytes: u64,
230 broadcasts: u64,
231 broadcasts_closed: u64,
232 subscriptions: u64,
233 subscriptions_closed: u64,
234 bytes: u64,
235 frames: u64,
236 groups: u64,
237}
238
239#[derive(Copy, Clone, Debug, PartialEq, Eq)]
244pub enum Tier {
245 External,
246 Internal,
247}
248
249impl Tier {
250 fn idx(self) -> usize {
251 match self {
252 Tier::External => 0,
253 Tier::Internal => 1,
254 }
255 }
256
257 pub fn as_str(self) -> &'static str {
264 match self {
265 Tier::External => "",
266 Tier::Internal => "internal",
267 }
268 }
269}
270
271#[derive(Copy, Clone, Debug, PartialEq, Eq)]
275pub enum Role {
276 Publisher,
277 Subscriber,
278}
279
280impl Role {
281 fn idx(self) -> usize {
282 match self {
283 Role::Publisher => 0,
284 Role::Subscriber => 1,
285 }
286 }
287
288 pub fn as_str(self) -> &'static str {
290 match self {
291 Role::Publisher => "publisher",
292 Role::Subscriber => "subscriber",
293 }
294 }
295}
296
297#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
303#[non_exhaustive]
304pub struct CounterTotals {
305 pub announced: u64,
306 pub announced_closed: u64,
307 pub announced_bytes: u64,
308 pub subscriptions: u64,
309 pub subscriptions_closed: u64,
310 pub broadcasts: u64,
311 pub broadcasts_closed: u64,
312 pub bytes: u64,
313 pub frames: u64,
314 pub groups: u64,
315}
316
317impl CounterTotals {
318 fn add(&mut self, raw: RawCounts) {
320 self.announced += raw.announced;
321 self.announced_closed += raw.announced_closed;
322 self.announced_bytes += raw.announced_bytes;
323 self.subscriptions += raw.subscriptions;
324 self.subscriptions_closed += raw.subscriptions_closed;
325 self.broadcasts += raw.broadcasts;
326 self.broadcasts_closed += raw.broadcasts_closed;
327 self.bytes += raw.bytes;
328 self.frames += raw.frames;
329 self.groups += raw.groups;
330 }
331}
332
333#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
337#[non_exhaustive]
338pub struct SessionTotals {
339 pub sessions: u64,
340 pub sessions_closed: u64,
341}
342
343#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
352#[non_exhaustive]
353pub struct StatsSnapshot {
354 traffic: [[CounterTotals; 2]; 2],
356 sessions: [SessionTotals; 2],
358}
359
360impl StatsSnapshot {
361 pub fn traffic(&self) -> [(Tier, Role, CounterTotals); 4] {
364 [
365 (
366 Tier::External,
367 Role::Publisher,
368 self.slot(Tier::External, Role::Publisher),
369 ),
370 (
371 Tier::External,
372 Role::Subscriber,
373 self.slot(Tier::External, Role::Subscriber),
374 ),
375 (
376 Tier::Internal,
377 Role::Publisher,
378 self.slot(Tier::Internal, Role::Publisher),
379 ),
380 (
381 Tier::Internal,
382 Role::Subscriber,
383 self.slot(Tier::Internal, Role::Subscriber),
384 ),
385 ]
386 }
387
388 pub fn sessions(&self) -> [(Tier, SessionTotals); 2] {
390 [
391 (Tier::External, self.sessions[Tier::External.idx()]),
392 (Tier::Internal, self.sessions[Tier::Internal.idx()]),
393 ]
394 }
395
396 fn slot(&self, tier: Tier, role: Role) -> CounterTotals {
397 self.traffic[tier.idx()][role.idx()]
398 }
399}
400
401#[derive(Clone)]
415#[non_exhaustive]
416pub struct StatsConfig {
417 pub origin: Option<OriginProducer>,
420 pub prefix: PathOwned,
424 pub node: Option<PathOwned>,
430 pub interval: Duration,
432 pub depth: usize,
442}
443
444impl StatsConfig {
445 pub fn new() -> Self {
449 Self {
450 origin: None,
451 prefix: PathOwned::from(".stats"),
452 node: None,
453 interval: Duration::from_secs(1),
454 depth: 0,
455 }
456 }
457
458 pub fn with_origin(mut self, origin: impl Into<Option<OriginProducer>>) -> Self {
461 self.origin = origin.into();
462 self
463 }
464
465 pub fn with_prefix(mut self, prefix: impl Into<PathOwned>) -> Self {
467 self.prefix = prefix.into();
468 self
469 }
470
471 pub fn with_interval(mut self, interval: Duration) -> Self {
473 self.interval = interval;
474 self
475 }
476
477 pub fn with_node(mut self, node: impl Into<Option<PathOwned>>) -> Self {
479 self.node = node.into();
480 self
481 }
482
483 pub fn with_depth(mut self, depth: usize) -> Self {
486 self.depth = depth;
487 self
488 }
489}
490
491impl Default for StatsConfig {
492 fn default() -> Self {
493 Self::new()
494 }
495}
496
497#[derive(Clone)]
501pub struct Stats {
502 prefix: PathOwned,
503 shared: Option<Arc<StatsShared>>,
506}
507
508struct StatsShared {
511 origin: OriginProducer,
512 entries: Lock<HashMap<PathOwned, Arc<BroadcastEntry>>>,
513 sessions: [Lock<HashMap<PathOwned, Arc<SessionCounters>>>; 2],
517}
518
519struct BroadcastEntry {
524 publisher: [Counters; 2],
525 subscriber: [Counters; 2],
526}
527
528impl BroadcastEntry {
529 fn new() -> Self {
530 Self {
531 publisher: Default::default(),
532 subscriber: Default::default(),
533 }
534 }
535}
536
537#[derive(Default)]
542struct SlotState {
543 prev_emitted: Option<Snapshot>,
546}
547
548#[derive(Default)]
552struct EntrySnapState {
553 publisher: [SlotState; 2],
554 subscriber: [SlotState; 2],
555}
556
557impl EntrySnapState {
558 fn zip_slots<'a>(&'a mut self, entry: &'a BroadcastEntry) -> [(&'static str, &'a Counters, &'a mut SlotState); 4] {
561 let [pub_ext_state, pub_int_state] = &mut self.publisher;
562 let [sub_ext_state, sub_int_state] = &mut self.subscriber;
563 [
564 ("publisher.json", &entry.publisher[Tier::External.idx()], pub_ext_state),
565 (
566 "subscriber.json",
567 &entry.subscriber[Tier::External.idx()],
568 sub_ext_state,
569 ),
570 (
571 "internal/publisher.json",
572 &entry.publisher[Tier::Internal.idx()],
573 pub_int_state,
574 ),
575 (
576 "internal/subscriber.json",
577 &entry.subscriber[Tier::Internal.idx()],
578 sub_int_state,
579 ),
580 ]
581 }
582}
583
584const NUM_SLOTS: usize = 4;
587
588const TRACK_ORDER: [&str; NUM_SLOTS] = [
591 "publisher.json",
592 "subscriber.json",
593 "internal/publisher.json",
594 "internal/subscriber.json",
595];
596
597const SESSION_TRACK_ORDER: [&str; 2] = ["sessions.json", "internal/sessions.json"];
600
601impl Stats {
602 pub fn new(config: StatsConfig) -> Self {
610 let StatsConfig {
611 origin,
612 prefix,
613 node,
614 interval,
615 depth,
616 } = config;
617 let node = node.filter(|p| !p.is_empty());
622
623 let shared = origin.map(|origin| {
624 let shared = Arc::new(StatsShared {
625 origin,
626 entries: Lock::default(),
627 sessions: Default::default(),
628 });
629 spawn(run_publisher(
630 Arc::downgrade(&shared),
631 prefix.clone(),
632 node.clone(),
633 depth,
634 interval,
635 ));
636 shared
637 });
638
639 Self { prefix, shared }
640 }
641
642 pub fn prefix(&self) -> &Path<'static> {
644 &self.prefix
645 }
646
647 #[cfg(test)]
650 fn shared(&self) -> &Arc<StatsShared> {
651 self.shared.as_ref().expect("enabled stats aggregator")
652 }
653
654 pub fn tier(&self, tier: Tier) -> StatsHandle {
657 StatsHandle {
658 stats: self.clone(),
659 tier,
660 }
661 }
662
663 fn entry(&self, path: impl AsPath) -> Option<Arc<BroadcastEntry>> {
664 let shared = self.shared.as_ref()?;
666 let path = path.as_path();
667 if path.has_prefix(&self.prefix) {
671 return None;
672 }
673 let owned = path.to_owned();
674 let mut entries = shared.entries.lock();
675 Some(
676 entries
677 .entry(owned)
678 .or_insert_with(|| Arc::new(BroadcastEntry::new()))
679 .clone(),
680 )
681 }
682
683 fn session_counters(&self, tier: Tier, root: impl AsPath) -> Option<Arc<SessionCounters>> {
687 let shared = self.shared.as_ref()?;
688 let owned = root.as_path().to_owned();
689 let mut sessions = shared.sessions[tier.idx()].lock();
690 Some(sessions.entry(owned).or_default().clone())
691 }
692
693 pub fn snapshot(&self) -> StatsSnapshot {
701 let mut snap = StatsSnapshot::default();
702 let Some(shared) = self.shared.as_ref() else {
703 return snap;
704 };
705 {
706 let entries = shared.entries.lock();
707 for entry in entries.values() {
708 for tier in [Tier::External, Tier::Internal] {
709 snap.traffic[tier.idx()][Role::Publisher.idx()].add(entry.publisher[tier.idx()].snapshot());
710 snap.traffic[tier.idx()][Role::Subscriber.idx()].add(entry.subscriber[tier.idx()].snapshot());
711 }
712 }
713 }
714 for tier in [Tier::External, Tier::Internal] {
715 let sessions = shared.sessions[tier.idx()].lock();
716 let totals = &mut snap.sessions[tier.idx()];
717 for counters in sessions.values() {
718 let (open, closed) = counters.snapshot();
719 totals.sessions += open;
720 totals.sessions_closed += closed;
721 }
722 }
723 snap
724 }
725}
726
727impl Default for Stats {
728 fn default() -> Self {
729 Self::new(StatsConfig::new())
730 }
731}
732
733#[derive(Clone)]
736pub struct StatsHandle {
737 stats: Stats,
738 tier: Tier,
739}
740
741impl StatsHandle {
742 pub fn parent(&self) -> &Stats {
744 &self.stats
745 }
746
747 pub fn tier(&self) -> Tier {
749 self.tier
750 }
751
752 pub fn broadcast(&self, path: impl AsPath) -> BroadcastStats {
758 BroadcastStats {
759 entry: self.stats.entry(path),
760 tier: self.tier,
761 }
762 }
763
764 pub fn publisher_broadcasts(&self) -> SessionBroadcasts {
769 SessionBroadcasts::new(self.stats.clone(), self.tier, Side::Publisher)
770 }
771
772 pub fn subscriber_broadcasts(&self) -> SessionBroadcasts {
775 SessionBroadcasts::new(self.stats.clone(), self.tier, Side::Subscriber)
776 }
777
778 pub fn session(&self, root: impl AsPath) -> SessionStats {
784 SessionStats::new(self.stats.session_counters(self.tier, root))
785 }
786}
787
788impl Default for StatsHandle {
789 fn default() -> Self {
791 Stats::default().tier(Tier::External)
792 }
793}
794
795#[derive(Clone)]
802pub struct BroadcastStats {
803 entry: Option<Arc<BroadcastEntry>>,
804 tier: Tier,
805}
806
807impl BroadcastStats {
808 pub fn is_empty(&self) -> bool {
812 self.entry.is_none()
813 }
814
815 pub fn publisher(&self) -> PublisherStats {
820 if let Some(entry) = &self.entry {
821 entry.publisher[self.tier.idx()]
822 .announced
823 .fetch_add(1, Ordering::Relaxed);
824 }
825 PublisherStats {
826 entry: self.entry.clone(),
827 tier: self.tier,
828 }
829 }
830
831 pub fn subscriber(&self) -> SubscriberStats {
836 if let Some(entry) = &self.entry {
837 entry.subscriber[self.tier.idx()]
838 .announced
839 .fetch_add(1, Ordering::Relaxed);
840 }
841 SubscriberStats {
842 entry: self.entry.clone(),
843 tier: self.tier,
844 }
845 }
846
847 pub fn publisher_track(&self, _name: &str) -> PublisherTrack {
853 if let Some(entry) = &self.entry {
854 entry.publisher[self.tier.idx()]
855 .subscriptions
856 .fetch_add(1, Ordering::Relaxed);
857 }
858 PublisherTrack {
859 entry: self.entry.clone(),
860 tier: self.tier,
861 }
862 }
863
864 pub fn publisher_announced_bytes(&self, n: u64) {
871 if let Some(entry) = &self.entry {
872 entry.publisher[self.tier.idx()]
873 .announced_bytes
874 .fetch_add(n, Ordering::Relaxed);
875 }
876 }
877
878 pub fn subscriber_announced_bytes(&self, n: u64) {
880 if let Some(entry) = &self.entry {
881 entry.subscriber[self.tier.idx()]
882 .announced_bytes
883 .fetch_add(n, Ordering::Relaxed);
884 }
885 }
886
887 pub fn subscriber_track(&self, _name: &str) -> SubscriberTrack {
889 if let Some(entry) = &self.entry {
890 entry.subscriber[self.tier.idx()]
891 .subscriptions
892 .fetch_add(1, Ordering::Relaxed);
893 }
894 SubscriberTrack {
895 entry: self.entry.clone(),
896 tier: self.tier,
897 }
898 }
899}
900
901#[derive(Copy, Clone)]
903enum Side {
904 Publisher,
905 Subscriber,
906}
907
908impl Side {
909 fn counters(self, entry: &BroadcastEntry, tier: Tier) -> &Counters {
910 match self {
911 Side::Publisher => &entry.publisher[tier.idx()],
912 Side::Subscriber => &entry.subscriber[tier.idx()],
913 }
914 }
915}
916
917#[derive(Clone)]
932pub struct SessionBroadcasts {
933 stats: Stats,
934 tier: Tier,
935 side: Side,
936 counts: Arc<std::sync::Mutex<HashMap<PathOwned, u32>>>,
937}
938
939impl SessionBroadcasts {
940 fn new(stats: Stats, tier: Tier, side: Side) -> Self {
941 Self {
942 stats,
943 tier,
944 side,
945 counts: Arc::new(std::sync::Mutex::new(HashMap::new())),
946 }
947 }
948
949 pub fn subscribe(&self, path: impl AsPath) -> BroadcastSubscription {
954 let path = path.as_path().to_owned();
955 let entry = self.stats.entry(&path);
956 let first = {
957 let mut counts = self.counts.lock().expect("stats refcount poisoned");
958 let n = counts.entry(path.clone()).or_insert(0);
959 let first = *n == 0;
960 *n += 1;
961 first
962 };
963 if first {
964 if let Some(entry) = &entry {
965 self.side
966 .counters(entry, self.tier)
967 .broadcasts
968 .fetch_add(1, Ordering::Relaxed);
969 }
970 }
971 BroadcastSubscription {
972 entry,
973 tier: self.tier,
974 side: self.side,
975 counts: self.counts.clone(),
976 path,
977 }
978 }
979}
980
981#[must_use = "drop the guard to release the subscription"]
984pub struct BroadcastSubscription {
985 entry: Option<Arc<BroadcastEntry>>,
986 tier: Tier,
987 side: Side,
988 counts: Arc<std::sync::Mutex<HashMap<PathOwned, u32>>>,
989 path: PathOwned,
990}
991
992impl Drop for BroadcastSubscription {
993 fn drop(&mut self) {
994 let last = {
995 let mut counts = self.counts.lock().expect("stats refcount poisoned");
996 match counts.get_mut(&self.path) {
997 Some(n) => {
998 *n -= 1;
999 if *n == 0 {
1000 counts.remove(&self.path);
1001 true
1002 } else {
1003 false
1004 }
1005 }
1006 None => false,
1007 }
1008 };
1009 if last {
1010 if let Some(entry) = &self.entry {
1011 self.side
1014 .counters(entry, self.tier)
1015 .broadcasts_closed
1016 .fetch_add(1, Ordering::Release);
1017 }
1018 }
1019 }
1020}
1021
1022#[must_use = "drop the guard to record the session as closed"]
1026pub struct SessionStats {
1027 counters: Option<Arc<SessionCounters>>,
1029}
1030
1031impl SessionStats {
1032 fn new(counters: Option<Arc<SessionCounters>>) -> Self {
1033 if let Some(counters) = &counters {
1034 counters.sessions.fetch_add(1, Ordering::Relaxed);
1035 }
1036 Self { counters }
1037 }
1038}
1039
1040impl Drop for SessionStats {
1041 fn drop(&mut self) {
1042 if let Some(counters) = &self.counters {
1043 counters.sessions_closed.fetch_add(1, Ordering::Release);
1046 }
1047 }
1048}
1049
1050#[must_use = "drop the guard to record the broadcast as closed"]
1052pub struct PublisherStats {
1053 entry: Option<Arc<BroadcastEntry>>,
1054 tier: Tier,
1055}
1056
1057impl PublisherStats {
1058 pub fn track(&self, name: &str) -> PublisherTrack {
1061 BroadcastStats {
1062 entry: self.entry.clone(),
1063 tier: self.tier,
1064 }
1065 .publisher_track(name)
1066 }
1067}
1068
1069impl Drop for PublisherStats {
1070 fn drop(&mut self) {
1071 if let Some(entry) = &self.entry {
1072 entry.publisher[self.tier.idx()]
1076 .announced_closed
1077 .fetch_add(1, Ordering::Release);
1078 }
1079 }
1080}
1081
1082#[must_use = "drop the guard to record the broadcast as closed"]
1084pub struct SubscriberStats {
1085 entry: Option<Arc<BroadcastEntry>>,
1086 tier: Tier,
1087}
1088
1089impl SubscriberStats {
1090 pub fn track(&self, name: &str) -> SubscriberTrack {
1092 BroadcastStats {
1093 entry: self.entry.clone(),
1094 tier: self.tier,
1095 }
1096 .subscriber_track(name)
1097 }
1098}
1099
1100impl Drop for SubscriberStats {
1101 fn drop(&mut self) {
1102 if let Some(entry) = &self.entry {
1103 entry.subscriber[self.tier.idx()]
1105 .announced_closed
1106 .fetch_add(1, Ordering::Release);
1107 }
1108 }
1109}
1110
1111#[must_use = "drop the guard to record the subscription as closed"]
1113pub struct PublisherTrack {
1114 entry: Option<Arc<BroadcastEntry>>,
1115 tier: Tier,
1116}
1117
1118impl PublisherTrack {
1119 pub fn frame(&self) {
1121 if let Some(entry) = &self.entry {
1122 entry.publisher[self.tier.idx()].frames.fetch_add(1, Ordering::Relaxed);
1123 }
1124 }
1125
1126 pub fn bytes(&self, n: u64) {
1128 if let Some(entry) = &self.entry {
1129 entry.publisher[self.tier.idx()].bytes.fetch_add(n, Ordering::Relaxed);
1130 }
1131 }
1132
1133 pub fn group(&self) {
1135 if let Some(entry) = &self.entry {
1136 entry.publisher[self.tier.idx()].groups.fetch_add(1, Ordering::Relaxed);
1137 }
1138 }
1139}
1140
1141impl Drop for PublisherTrack {
1142 fn drop(&mut self) {
1143 if let Some(entry) = &self.entry {
1144 entry.publisher[self.tier.idx()]
1146 .subscriptions_closed
1147 .fetch_add(1, Ordering::Release);
1148 }
1149 }
1150}
1151
1152#[must_use = "drop the guard to record the subscription as closed"]
1154pub struct SubscriberTrack {
1155 entry: Option<Arc<BroadcastEntry>>,
1156 tier: Tier,
1157}
1158
1159impl SubscriberTrack {
1160 pub fn frame(&self) {
1162 if let Some(entry) = &self.entry {
1163 entry.subscriber[self.tier.idx()].frames.fetch_add(1, Ordering::Relaxed);
1164 }
1165 }
1166
1167 pub fn bytes(&self, n: u64) {
1169 if let Some(entry) = &self.entry {
1170 entry.subscriber[self.tier.idx()].bytes.fetch_add(n, Ordering::Relaxed);
1171 }
1172 }
1173
1174 pub fn group(&self) {
1176 if let Some(entry) = &self.entry {
1177 entry.subscriber[self.tier.idx()].groups.fetch_add(1, Ordering::Relaxed);
1178 }
1179 }
1180}
1181
1182impl Drop for SubscriberTrack {
1183 fn drop(&mut self) {
1184 if let Some(entry) = &self.entry {
1185 entry.subscriber[self.tier.idx()]
1187 .subscriptions_closed
1188 .fetch_add(1, Ordering::Release);
1189 }
1190 }
1191}
1192
1193fn process_slot(counters: &Counters, slot_state: &mut SlotState, mut emit: impl FnMut(Snapshot)) {
1197 let raw = counters.snapshot();
1198
1199 let snap = Snapshot {
1200 announced: raw.announced,
1201 announced_closed: raw.announced_closed,
1202 announced_bytes: raw.announced_bytes,
1203 broadcasts: raw.broadcasts,
1204 broadcasts_closed: raw.broadcasts_closed,
1205 subscriptions: raw.subscriptions,
1206 subscriptions_closed: raw.subscriptions_closed,
1207 bytes: raw.bytes,
1208 frames: raw.frames,
1209 groups: raw.groups,
1210 };
1211
1212 let live = snap.announced != snap.announced_closed
1219 || snap.subscriptions != snap.subscriptions_closed
1220 || snap.broadcasts != snap.broadcasts_closed;
1221
1222 let prev_snap = slot_state.prev_emitted.unwrap_or_default();
1233 let changed = snap != prev_snap;
1234 if changed {
1235 slot_state.prev_emitted = Some(snap);
1236 }
1237 if live || changed {
1238 emit(snap);
1239 }
1240}
1241
1242#[derive(Default)]
1245struct SessionSlotState {
1246 prev_emitted: Option<SessionSnapshot>,
1247}
1248
1249fn process_session_slot(
1254 counters: &SessionCounters,
1255 slot_state: &mut SessionSlotState,
1256 mut emit: impl FnMut(SessionSnapshot),
1257) {
1258 let (sessions, sessions_closed) = counters.snapshot();
1259 let snap = SessionSnapshot {
1260 sessions,
1261 sessions_closed,
1262 };
1263
1264 let live = sessions != sessions_closed;
1265 let prev_snap = slot_state.prev_emitted.unwrap_or_default();
1266 let changed = snap != prev_snap;
1267 if changed {
1268 slot_state.prev_emitted = Some(snap);
1269 }
1270 if live || changed {
1271 emit(snap);
1272 }
1273}
1274
1275fn flush_track<T: Serialize>(track: &mut TrackProducer, frame: &T, last: &mut Vec<u8>, name: &str) {
1279 let json = match serde_json::to_vec(frame) {
1280 Ok(b) => b,
1281 Err(err) => {
1282 tracing::debug!(?err, name, "stats: failed to serialize frame");
1283 return;
1284 }
1285 };
1286 if &json == last {
1287 return;
1288 }
1289 if let Err(err) = track.write_frame(json.clone()) {
1290 tracing::debug!(?err, name, "stats: failed to write frame");
1291 return;
1292 }
1293 *last = json;
1294}
1295
1296struct GroupPublisher {
1303 _broadcast: crate::BroadcastProducer,
1306 tracks: Vec<TrackProducer>,
1307 session_tracks: Vec<TrackProducer>,
1308 local: HashMap<PathOwned, EntrySnapState>,
1311 last_payload: [Vec<u8>; NUM_SLOTS],
1312 session_local: [HashMap<PathOwned, SessionSlotState>; 2],
1314 session_last_payload: [Vec<u8>; 2],
1315}
1316
1317type GroupedEntries<'a> = HashMap<PathOwned, Vec<(&'a PathOwned, &'a Arc<BroadcastEntry>)>>;
1318type GroupedSessions<'a> = HashMap<PathOwned, Vec<(&'a PathOwned, &'a Arc<SessionCounters>)>>;
1319
1320impl GroupPublisher {
1321 fn create(origin: &OriginProducer, prefix: &Path, group: &Path, node: Option<&str>) -> Option<Self> {
1325 let mut broadcast = Broadcast::new().produce();
1326
1327 let create = |broadcast: &mut crate::BroadcastProducer, name: &str| match broadcast.create_track(Track {
1329 name: name.into(),
1330 priority: 0,
1331 }) {
1332 Ok(t) => Some(t),
1333 Err(err) => {
1334 tracing::warn!(?err, name, "stats: failed to create track");
1335 None
1336 }
1337 };
1338
1339 let mut tracks: Vec<TrackProducer> = Vec::with_capacity(NUM_SLOTS);
1340 for name in TRACK_ORDER {
1341 tracks.push(create(&mut broadcast, name)?);
1342 }
1343 let mut session_tracks: Vec<TrackProducer> = Vec::with_capacity(SESSION_TRACK_ORDER.len());
1344 for name in SESSION_TRACK_ORDER {
1345 session_tracks.push(create(&mut broadcast, name)?);
1346 }
1347
1348 let advertised = advertised_path(prefix, group, node);
1349 if !origin.publish_broadcast(&advertised, broadcast.consume()) {
1350 tracing::warn!(advertised = %advertised, "stats: origin rejected stats broadcast");
1351 return None;
1352 }
1353 tracing::debug!(advertised = %advertised, "stats: publishing broadcast");
1354
1355 Some(Self {
1356 _broadcast: broadcast,
1357 tracks,
1358 session_tracks,
1359 local: HashMap::new(),
1360 last_payload: Default::default(),
1361 session_local: Default::default(),
1362 session_last_payload: Default::default(),
1363 })
1364 }
1365}
1366
1367fn group_key(path: &str, depth: usize) -> PathOwned {
1371 if depth == 0 {
1372 return Path::empty().to_owned();
1373 }
1374 let mut seen = 0;
1376 let mut end = path.len();
1377 for (i, b) in path.bytes().enumerate() {
1378 if b == b'/' {
1379 seen += 1;
1380 if seen == depth {
1381 end = i;
1382 break;
1383 }
1384 }
1385 }
1386 Path::new(&path[..end]).to_owned()
1387}
1388
1389async fn run_publisher(
1390 weak: Weak<StatsShared>,
1391 prefix: PathOwned,
1392 node: Option<PathOwned>,
1393 depth: usize,
1394 interval: Duration,
1395) {
1396 let node = node.as_ref().map(|p| p.as_str());
1397
1398 let mut groups: HashMap<PathOwned, GroupPublisher> = HashMap::new();
1403 if depth == 0 {
1404 let Some(shared) = weak.upgrade() else {
1405 return;
1406 };
1407 let Some(gp) = GroupPublisher::create(&shared.origin, &prefix, &Path::empty(), node) else {
1408 return;
1409 };
1410 groups.insert(Path::empty().to_owned(), gp);
1411 drop(shared);
1412 }
1413
1414 let mut ticker = web_async::time::interval(interval);
1415 ticker.set_missed_tick_behavior(web_async::time::MissedTickBehavior::Delay);
1416
1417 loop {
1418 ticker.tick().await;
1419
1420 let Some(shared) = weak.upgrade() else {
1421 return;
1422 };
1423
1424 let entries: Vec<(PathOwned, Arc<BroadcastEntry>)> = {
1427 let map = shared.entries.lock();
1428 map.iter().map(|(k, v)| (k.clone(), v.clone())).collect()
1429 };
1430 let session_roots: [Vec<(PathOwned, Arc<SessionCounters>)>; 2] = [
1431 {
1432 let map = shared.sessions[0].lock();
1433 map.iter().map(|(k, v)| (k.clone(), v.clone())).collect()
1434 },
1435 {
1436 let map = shared.sessions[1].lock();
1437 map.iter().map(|(k, v)| (k.clone(), v.clone())).collect()
1438 },
1439 ];
1440
1441 let mut entries_by_group: GroupedEntries<'_> = HashMap::new();
1445 for (path, entry) in &entries {
1446 entries_by_group
1447 .entry(group_key(path.as_str(), depth))
1448 .or_default()
1449 .push((path, entry));
1450 }
1451 let mut roots_by_group: [GroupedSessions<'_>; 2] = Default::default();
1452 for tier_idx in 0..2 {
1453 for (root, counters) in &session_roots[tier_idx] {
1454 roots_by_group[tier_idx]
1455 .entry(group_key(root.as_str(), depth))
1456 .or_default()
1457 .push((root, counters));
1458 }
1459 }
1460
1461 let mut active: HashSet<PathOwned> = HashSet::new();
1464 active.extend(entries_by_group.keys().cloned());
1465 for group_roots in &roots_by_group {
1466 active.extend(group_roots.keys().cloned());
1467 }
1468 if depth == 0 {
1469 active.insert(Path::empty().to_owned());
1470 }
1471
1472 for group in &active {
1473 if !groups.contains_key(group) {
1475 let Some(gp) = GroupPublisher::create(&shared.origin, &prefix, group, node) else {
1476 continue;
1477 };
1478 groups.insert(group.clone(), gp);
1479 }
1480 let gp = groups.get_mut(group).expect("just inserted");
1481
1482 let mut frames: [BTreeMap<String, Snapshot>; NUM_SLOTS] = Default::default();
1484 if let Some(group_entries) = entries_by_group.get(group) {
1485 for &(path, entry) in group_entries {
1486 let snap_state = gp.local.entry(path.clone()).or_default();
1487 for (i, (_track_name, counters, slot_state)) in snap_state.zip_slots(entry).into_iter().enumerate()
1488 {
1489 process_slot(counters, slot_state, |snap| {
1490 frames[i].insert(path.as_str().to_string(), snap);
1491 });
1492 }
1493 }
1494 }
1495 for (i, (frame, last)) in frames.iter().zip(gp.last_payload.iter_mut()).enumerate() {
1496 flush_track(&mut gp.tracks[i], frame, last, TRACK_ORDER[i]);
1497 }
1498
1499 let mut session_frames: [BTreeMap<String, SessionSnapshot>; 2] = Default::default();
1501 for tier_idx in 0..2 {
1502 if let Some(group_roots) = roots_by_group[tier_idx].get(group) {
1503 for &(root, counters) in group_roots {
1504 let state = gp.session_local[tier_idx].entry(root.clone()).or_default();
1505 process_session_slot(counters, state, |snap| {
1506 session_frames[tier_idx].insert(root.as_str().to_string(), snap);
1507 });
1508 }
1509 }
1510 }
1511 for (i, (frame, last)) in session_frames
1512 .iter()
1513 .zip(gp.session_last_payload.iter_mut())
1514 .enumerate()
1515 {
1516 flush_track(&mut gp.session_tracks[i], frame, last, SESSION_TRACK_ORDER[i]);
1517 }
1518 }
1519
1520 drop(entries_by_group);
1523 drop(roots_by_group);
1524 drop(entries);
1525 drop(session_roots);
1526
1527 {
1533 let mut map = shared.entries.lock();
1534 map.retain(|_, entry| Arc::strong_count(entry) > 1);
1535 for gp in groups.values_mut() {
1536 gp.local.retain(|path, _| map.contains_key(path));
1537 }
1538 }
1539 for tier_idx in 0..2 {
1540 let mut map = shared.sessions[tier_idx].lock();
1541 map.retain(|_, counters| Arc::strong_count(counters) > 1);
1542 for gp in groups.values_mut() {
1543 gp.session_local[tier_idx].retain(|root, _| map.contains_key(root));
1544 }
1545 }
1546
1547 groups.retain(|group, _| active.contains(group));
1550
1551 drop(shared);
1552 }
1553}
1554
1555#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize)]
1560#[cfg_attr(test, derive(serde::Deserialize))]
1561struct Snapshot {
1562 announced: u64,
1563 announced_closed: u64,
1564 announced_bytes: u64,
1565 broadcasts: u64,
1566 broadcasts_closed: u64,
1567 subscriptions: u64,
1568 subscriptions_closed: u64,
1569 bytes: u64,
1570 frames: u64,
1571 groups: u64,
1572}
1573
1574#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize)]
1577#[cfg_attr(test, derive(serde::Deserialize))]
1578struct SessionSnapshot {
1579 sessions: u64,
1580 sessions_closed: u64,
1581}
1582
1583fn advertised_path(prefix: &Path, group: &Path, node: Option<&str>) -> PathOwned {
1584 let mut out = prefix.as_str().to_string();
1589 if !group.is_empty() {
1590 out.push('/');
1591 out.push_str(group.as_str());
1592 }
1593 out.push_str("/node");
1594 if let Some(node) = node {
1595 out.push('/');
1596 out.push_str(node);
1597 }
1598 PathOwned::from(out)
1599}
1600
1601#[cfg(test)]
1602mod tests {
1603 use std::{collections::BTreeMap, sync::atomic::Ordering::Relaxed};
1604
1605 use crate::{Origin, Path};
1606
1607 use super::*;
1608
1609 fn test_stats(node: Option<&str>) -> (Stats, OriginProducer) {
1610 let origin = Origin::random().produce();
1611 let stats = Stats::new(
1612 StatsConfig::new()
1613 .with_origin(origin.clone())
1614 .with_node(node.map(|s| PathOwned::from(s.to_string()))),
1615 );
1616 (stats, origin)
1617 }
1618
1619 #[test]
1620 fn advertised_path_with_and_without_node() {
1621 let prefix = Path::new(".stats");
1622 let none = Path::empty();
1623 assert_eq!(advertised_path(&prefix, &none, Some("sjc")).as_str(), ".stats/node/sjc");
1625 assert_eq!(
1626 advertised_path(&prefix, &none, Some("sjc/1")).as_str(),
1627 ".stats/node/sjc/1"
1628 );
1629 assert_eq!(advertised_path(&prefix, &none, None).as_str(), ".stats/node");
1630
1631 let prefix = Path::new("metrics");
1632 assert_eq!(
1633 advertised_path(&prefix, &none, Some("lon")).as_str(),
1634 "metrics/node/lon"
1635 );
1636
1637 let prefix = Path::new(".stats");
1639 let group = Path::new("acme");
1640 assert_eq!(
1641 advertised_path(&prefix, &group, Some("sjc")).as_str(),
1642 ".stats/acme/node/sjc"
1643 );
1644 assert_eq!(advertised_path(&prefix, &group, None).as_str(), ".stats/acme/node");
1645 }
1646
1647 #[test]
1648 fn group_key_takes_leading_segments() {
1649 assert_eq!(group_key("acme/foo/bar", 0).as_str(), "");
1651 assert_eq!(group_key("", 0).as_str(), "");
1652 assert_eq!(group_key("acme/foo/bar", 1).as_str(), "acme");
1654 assert_eq!(group_key("acme", 1).as_str(), "acme");
1655 assert_eq!(group_key("acme/foo/bar", 2).as_str(), "acme/foo");
1657 assert_eq!(group_key("acme", 2).as_str(), "acme");
1659 }
1660
1661 async fn announced_path_for_node(node: &str) -> String {
1665 let origin = Origin::random().produce();
1666 let _stats = Stats::new(
1667 StatsConfig::new()
1668 .with_origin(origin.clone())
1669 .with_node(PathOwned::from(node.to_string())),
1670 );
1671 let mut consumer = origin.consume();
1672 tokio::time::advance(Duration::from_millis(1)).await;
1673 let (path, _broadcast) = consumer.announced().await.expect("expected announce");
1674 path.as_str().to_string()
1675 }
1676
1677 #[tokio::test(start_paused = true)]
1678 async fn new_normalizes_and_drops_empty_node() {
1679 assert_eq!(announced_path_for_node("/sjc//1/").await, ".stats/node/sjc/1");
1680 assert_eq!(announced_path_for_node("///").await, ".stats/node");
1681 }
1682
1683 #[tokio::test(start_paused = true)]
1684 async fn per_broadcast_counters_isolated() {
1685 let (stats, _origin) = test_stats(Some("sjc"));
1687 let bs1 = stats.tier(Tier::External).broadcast("demo/bbb");
1688 let bs2 = stats.tier(Tier::External).broadcast("demo/ccc");
1689 let g1 = bs1.publisher().track("video");
1690 g1.bytes(100);
1691 let g2 = bs2.publisher().track("video");
1692 g2.bytes(7);
1693
1694 let entries = stats.shared().entries.lock();
1695 let e1 = entries.get(&PathOwned::from("demo/bbb")).expect("entry");
1696 let e2 = entries.get(&PathOwned::from("demo/ccc")).expect("entry");
1697 assert_eq!(e1.publisher[Tier::External.idx()].bytes.load(Relaxed), 100);
1698 assert_eq!(e2.publisher[Tier::External.idx()].bytes.load(Relaxed), 7);
1699 }
1700
1701 #[tokio::test(start_paused = true)]
1702 async fn external_and_internal_tiers_are_independent() {
1703 let (stats, _origin) = test_stats(Some("sjc"));
1704 let ext = stats.tier(Tier::External);
1705 let int = stats.tier(Tier::Internal);
1706
1707 let ext_track = ext.broadcast("demo/bbb").publisher().track("video");
1708 ext_track.bytes(100);
1709 let int_track = int.broadcast("demo/bbb").subscriber().track("audio");
1710 int_track.bytes(7);
1711
1712 let entries = stats.shared().entries.lock();
1713 let entry = entries.get(&PathOwned::from("demo/bbb")).expect("entry");
1714 assert_eq!(entry.publisher[Tier::External.idx()].bytes.load(Relaxed), 100);
1715 assert_eq!(entry.subscriber[Tier::External.idx()].bytes.load(Relaxed), 0);
1716 assert_eq!(entry.publisher[Tier::Internal.idx()].bytes.load(Relaxed), 0);
1717 assert_eq!(entry.subscriber[Tier::Internal.idx()].bytes.load(Relaxed), 7);
1718 }
1719
1720 #[tokio::test(start_paused = true)]
1721 async fn snapshot_rolls_up_by_tier_role_and_sessions() {
1722 let (stats, _origin) = test_stats(Some("sjc"));
1723 let ext = stats.tier(Tier::External);
1724 let int = stats.tier(Tier::Internal);
1725
1726 let pub_a = ext.broadcast("demo/aaa").publisher().track("video");
1728 pub_a.bytes(100);
1729 pub_a.frame();
1730 pub_a.group();
1731 let pub_b = ext.broadcast("demo/bbb").publisher().track("video");
1732 pub_b.bytes(50);
1733 let sub_a = int.broadcast("demo/aaa").subscriber().track("audio");
1735 sub_a.bytes(7);
1736
1737 let _s1 = ext.session("acme");
1739 let _s2 = ext.session("acme");
1740 let _s3 = int.session("peer");
1741
1742 let snap = stats.snapshot();
1743
1744 let slot = |tier, role| {
1745 snap.traffic()
1746 .into_iter()
1747 .find(|(t, r, _)| *t == tier && *r == role)
1748 .map(|(_, _, c)| c)
1749 .expect("row present")
1750 };
1751
1752 let ext_pub = slot(Tier::External, Role::Publisher);
1753 assert_eq!(ext_pub.bytes, 150, "external egress bytes sum across broadcasts");
1754 assert_eq!(ext_pub.frames, 1);
1755 assert_eq!(ext_pub.groups, 1);
1756
1757 let int_sub = slot(Tier::Internal, Role::Subscriber);
1758 assert_eq!(int_sub.bytes, 7, "internal ingress isolated by tier/role");
1759 assert_eq!(slot(Tier::External, Role::Subscriber).bytes, 0);
1760 assert_eq!(slot(Tier::Internal, Role::Publisher).bytes, 0);
1761
1762 let sessions = |tier| {
1763 snap.sessions()
1764 .into_iter()
1765 .find(|(t, _)| *t == tier)
1766 .map(|(_, s)| s)
1767 .expect("tier present")
1768 };
1769 let ext_sessions = sessions(Tier::External);
1770 assert_eq!(ext_sessions.sessions, 2, "two external sessions under one root");
1771 assert_eq!(ext_sessions.sessions_closed, 0, "guards still held");
1772 assert_eq!(sessions(Tier::Internal).sessions, 1);
1773 }
1774
1775 #[tokio::test(start_paused = true)]
1776 async fn paths_under_prefix_are_no_op() {
1777 let (stats, _origin) = test_stats(Some("sjc"));
1780 let bs = stats.tier(Tier::External).broadcast(".stats/node/sjc");
1781 assert!(bs.is_empty());
1782 let p = bs.publisher();
1783 let track = p.track("video");
1784 track.bytes(100);
1785 drop(track);
1786 drop(p);
1787 assert!(stats.shared().entries.lock().is_empty());
1788 }
1789
1790 #[tokio::test(start_paused = true)]
1791 async fn disabled_stats_are_noop() {
1792 let stats = Stats::default();
1795 assert!(stats.shared.is_none());
1796 let bs = stats.tier(Tier::External).broadcast("demo/bbb");
1797 assert!(bs.is_empty());
1798 let p = bs.publisher();
1799 let track = p.track("video");
1800 track.bytes(100);
1801 drop(track);
1802 drop(p);
1803 }
1804
1805 #[tokio::test(start_paused = true)]
1806 async fn single_broadcast_path_announced() {
1807 let (stats, origin) = test_stats(Some("sjc/1"));
1810 let mut consumer = origin.consume();
1811
1812 let bs1 = stats.tier(Tier::External).broadcast("foo/bar");
1813 let _t1 = bs1.publisher().track("video");
1814 let bs2 = stats.tier(Tier::External).broadcast("baz/qux");
1815 let _t2 = bs2.publisher().track("video");
1816
1817 tokio::time::advance(Duration::from_millis(1)).await;
1818 let (path, broadcast) = consumer.announced().await.expect("expected announce");
1819 assert!(broadcast.is_some());
1820 assert_eq!(path.as_str(), ".stats/node/sjc/1");
1821 }
1822
1823 #[tokio::test(start_paused = true)]
1824 async fn depth_splits_broadcasts_per_group() {
1825 let origin = Origin::random().produce();
1829 let stats = Stats::new(
1830 StatsConfig::new()
1831 .with_origin(origin.clone())
1832 .with_node(PathOwned::from("sjc".to_string()))
1833 .with_depth(1),
1834 );
1835 let mut consumer = origin.consume();
1836
1837 let _t1 = stats
1838 .tier(Tier::External)
1839 .broadcast("acme/foo")
1840 .publisher()
1841 .track("video");
1842 let _t2 = stats
1843 .tier(Tier::External)
1844 .broadcast("globex/bar")
1845 .publisher()
1846 .track("video");
1847
1848 tokio::time::advance(Duration::from_secs(1)).await;
1851
1852 let mut announced = Vec::new();
1853 for _ in 0..2 {
1854 let (path, broadcast) = consumer.announced().await.expect("expected announce");
1855 assert!(broadcast.is_some());
1856 announced.push(path.as_str().to_string());
1857 }
1858 announced.sort();
1859 assert_eq!(
1860 announced,
1861 vec![".stats/acme/node/sjc".to_string(), ".stats/globex/node/sjc".to_string()]
1862 );
1863 }
1864
1865 #[tokio::test(start_paused = true)]
1866 async fn task_announces_without_node_suffix() {
1867 let origin = Origin::random().produce();
1868 let stats = Stats::new(StatsConfig::new().with_origin(origin.clone()));
1869 let mut consumer = origin.consume();
1870
1871 let bs = stats.tier(Tier::External).broadcast("foo/bar");
1872 let _t = bs.publisher().track("video");
1873
1874 tokio::time::advance(Duration::from_millis(1)).await;
1875 let (path, broadcast) = consumer.announced().await.expect("expected announce");
1876 assert!(broadcast.is_some());
1877 assert_eq!(path.as_str(), ".stats/node");
1878 }
1879
1880 async fn drive_ticks(count: u32) {
1886 for _ in 0..count {
1887 tokio::time::advance(Duration::from_secs(1)).await;
1888 for _ in 0..4 {
1891 tokio::task::yield_now().await;
1892 }
1893 }
1894 }
1895
1896 #[tokio::test(start_paused = true)]
1897 async fn live_entry_kept_while_idle() {
1898 let (stats, _origin) = test_stats(Some("sjc"));
1902 let key = PathOwned::from("foo/bar".to_string());
1903 let bs = stats.tier(Tier::External).broadcast("foo/bar");
1904 let guard = bs.publisher();
1905
1906 drive_ticks(5).await;
1907 assert!(
1908 stats.shared().entries.lock().contains_key(&key),
1909 "announced-but-idle broadcast must stay while the guard is held"
1910 );
1911
1912 drop(guard);
1913 drop(bs);
1914 drive_ticks(1).await;
1917 assert!(
1918 !stats.shared().entries.lock().contains_key(&key),
1919 "entry dropped once the announce guard closes"
1920 );
1921 }
1922
1923 #[tokio::test(start_paused = true)]
1924 async fn entry_dropped_once_fully_closed() {
1925 let (stats, _origin) = test_stats(Some("sjc"));
1928 let key = PathOwned::from("foo/bar".to_string());
1929 let bs = stats.tier(Tier::External).broadcast("foo/bar");
1930 let track = bs.publisher().track("video");
1931
1932 drive_ticks(1).await;
1933 assert!(
1934 stats.shared().entries.lock().contains_key(&key),
1935 "live entry present while the track guard is held"
1936 );
1937
1938 drop(track);
1939 drop(bs);
1940 drive_ticks(1).await;
1941 assert!(
1942 !stats.shared().entries.lock().contains_key(&key),
1943 "fully-closed entry dropped on the next tick"
1944 );
1945 }
1946
1947 #[tokio::test(start_paused = true)]
1948 async fn frame_emits_expected_counters() {
1949 let (stats, origin) = test_stats(Some("sjc"));
1950 let mut consumer = origin.consume();
1951 let bs = stats.tier(Tier::External).broadcast("foo/bar");
1952 let track = bs.publisher().track("video");
1953 track.bytes(42);
1954 track.frame();
1955 let sessions = stats.tier(Tier::External).publisher_broadcasts();
1956 let _sub = sessions.subscribe("foo/bar");
1957
1958 tokio::time::advance(Duration::from_millis(1100)).await;
1959
1960 let (_path, broadcast) = consumer.announced().await.expect("expected announce");
1961 let broadcast = broadcast.expect("active");
1962 let track = broadcast
1963 .subscribe_track(&Track {
1964 name: "publisher.json".into(),
1965 priority: 0,
1966 })
1967 .expect("subscribe");
1968 let frame = read_frame(track).await;
1969 let snap = frame.get("foo/bar").expect("foo/bar entry");
1970 assert_eq!(snap.announced, 1, "publisher() guard bumps announced");
1971 assert_eq!(snap.broadcasts, 1, "one session subscribed");
1972 assert_eq!(snap.subscriptions, 1);
1973 assert_eq!(snap.bytes, 42);
1974 assert_eq!(snap.frames, 1);
1975 }
1976
1977 #[tokio::test(start_paused = true)]
1978 async fn announced_bytes_recorded_per_side() {
1979 let (stats, _origin) = test_stats(Some("sjc"));
1983 let bs = stats.tier(Tier::External).broadcast("foo/bar");
1984 bs.publisher_announced_bytes(40);
1985 bs.publisher_announced_bytes(2);
1986 bs.subscriber_announced_bytes(7);
1987
1988 let entries = stats.shared().entries.lock();
1989 let entry = entries.get(&PathOwned::from("foo/bar")).expect("entry");
1990 let pub_ext = entry.publisher[Tier::External.idx()].snapshot();
1991 let sub_ext = entry.subscriber[Tier::External.idx()].snapshot();
1992 assert_eq!(pub_ext.announced_bytes, 42, "publisher announce bytes accumulate");
1993 assert_eq!(pub_ext.bytes, 0, "announce bytes are not payload bytes");
1994 assert_eq!(sub_ext.announced_bytes, 7, "subscriber side tracked independently");
1995 }
1996
1997 #[tokio::test(start_paused = true)]
1998 async fn announced_bytes_surfaces_in_frame() {
1999 let (stats, origin) = test_stats(Some("sjc"));
2000 let mut consumer = origin.consume();
2001 let bs = stats.tier(Tier::External).broadcast("foo/bar");
2002 let _guard = bs.publisher();
2003 bs.publisher_announced_bytes(123);
2004
2005 tokio::time::advance(Duration::from_millis(1100)).await;
2006
2007 let (_path, broadcast) = consumer.announced().await.expect("announce");
2008 let broadcast = broadcast.expect("active");
2009 let track = broadcast
2010 .subscribe_track(&Track {
2011 name: "publisher.json".into(),
2012 priority: 0,
2013 })
2014 .expect("subscribe");
2015 let frame = read_frame(track).await;
2016 let snap = frame.get("foo/bar").expect("foo/bar entry");
2017 assert_eq!(snap.announced, 1);
2018 assert_eq!(snap.announced_bytes, 123);
2019 }
2020
2021 #[tokio::test(start_paused = true)]
2022 async fn announced_decouples_from_broadcasts() {
2023 let (stats, origin) = test_stats(Some("sjc"));
2026 let mut consumer = origin.consume();
2027 let bs = stats.tier(Tier::External).broadcast("foo/bar");
2028 let _guard = bs.publisher();
2029
2030 tokio::time::advance(Duration::from_millis(1100)).await;
2031
2032 let (_path, broadcast) = consumer.announced().await.expect("announce");
2033 let broadcast = broadcast.expect("active");
2034 let track = broadcast
2035 .subscribe_track(&Track {
2036 name: "publisher.json".into(),
2037 priority: 0,
2038 })
2039 .expect("subscribe");
2040 let frame = read_frame(track).await;
2041 let snap = frame.get("foo/bar").expect("foo/bar entry");
2042 assert_eq!(snap.announced, 1);
2043 assert_eq!(snap.broadcasts, 0, "no subscription, no broadcasts sentinel");
2044 assert_eq!(snap.subscriptions, 0);
2045 }
2046
2047 #[tokio::test(start_paused = true)]
2048 async fn short_lived_sub_is_surfaced() {
2049 let (stats, origin) = test_stats(Some("sjc"));
2055 let mut consumer = origin.consume();
2056 let bs = stats.tier(Tier::External).broadcast("foo/bar");
2057 let sessions = stats.tier(Tier::External).publisher_broadcasts();
2058 {
2059 let track = bs.publisher().track("video");
2060 track.bytes(123);
2061 track.frame();
2062 let _sub = sessions.subscribe("foo/bar");
2063 }
2065
2066 tokio::time::advance(Duration::from_millis(1100)).await;
2067
2068 let (_path, broadcast) = consumer.announced().await.expect("announce");
2069 let broadcast = broadcast.expect("active");
2070 let track = broadcast
2071 .subscribe_track(&Track {
2072 name: "publisher.json".into(),
2073 priority: 0,
2074 })
2075 .expect("subscribe");
2076 let frame = read_frame(track).await;
2077 let snap = frame.get("foo/bar").expect("foo/bar entry");
2078 assert_eq!(snap.subscriptions, 1);
2080 assert_eq!(snap.subscriptions_closed, 1);
2081 assert_eq!(snap.broadcasts, 1, "one session subscribed");
2082 assert_eq!(snap.broadcasts_closed, 1);
2083 assert_eq!(snap.bytes, 123);
2084 assert_eq!(snap.frames, 1);
2085 }
2086
2087 #[tokio::test(start_paused = true)]
2088 async fn multiple_subs_count_as_one_broadcast() {
2089 let (stats, _origin) = test_stats(Some("sjc"));
2094 let bs = stats.tier(Tier::External).broadcast("foo/bar");
2095 let sessions = stats.tier(Tier::External).publisher_broadcasts();
2096 let pub_guard = bs.publisher();
2097 let t1 = pub_guard.track("video");
2098 let t2 = pub_guard.track("audio");
2099 let s1 = sessions.subscribe("foo/bar");
2100 let s2 = sessions.subscribe("foo/bar");
2101
2102 let raw = || {
2103 let entries = stats.shared().entries.lock();
2104 let entry = entries.get(&PathOwned::from("foo/bar")).expect("entry");
2105 entry.publisher[Tier::External.idx()].snapshot()
2106 };
2107
2108 let r = raw();
2109 assert_eq!(r.subscriptions, 2, "two track subs");
2110 assert_eq!(r.subscriptions_closed, 0, "neither dropped yet");
2111 assert_eq!(r.broadcasts, 1, "one session => one broadcast");
2112 assert_eq!(r.broadcasts_closed, 0);
2113
2114 drop(s1);
2115 assert_eq!(raw().broadcasts_closed, 0, "session still has a sub open");
2116
2117 drop(s2);
2118 drop(t1);
2119 drop(t2);
2120 let r = raw();
2121 assert_eq!(r.subscriptions_closed, 2, "both track subs dropped");
2122 assert_eq!(r.broadcasts, 1);
2123 assert_eq!(r.broadcasts_closed, 1, "last sub closed => one broadcasts_closed");
2124
2125 drop(pub_guard);
2126 drop(bs);
2127 }
2128
2129 #[tokio::test(start_paused = true)]
2130 async fn distinct_sessions_count_as_separate_broadcasts() {
2131 let (stats, _origin) = test_stats(Some("sjc"));
2134 let viewer1 = stats.tier(Tier::External).publisher_broadcasts();
2135 let viewer2 = stats.tier(Tier::External).publisher_broadcasts();
2136
2137 let raw = || {
2138 let entries = stats.shared().entries.lock();
2139 let entry = entries.get(&PathOwned::from("foo/bar")).expect("entry");
2140 entry.publisher[Tier::External.idx()].snapshot()
2141 };
2142
2143 let s1 = viewer1.subscribe("foo/bar");
2144 assert_eq!(raw().broadcasts, 1, "one viewer");
2145 let s2 = viewer2.subscribe("foo/bar");
2146 assert_eq!(raw().broadcasts, 2, "two distinct viewers");
2147 assert_eq!(raw().broadcasts_closed, 0);
2148
2149 drop(s1);
2150 let r = raw();
2151 assert_eq!(r.broadcasts, 2, "broadcasts is cumulative");
2152 assert_eq!(r.broadcasts_closed, 1, "one viewer left");
2153 drop(s2);
2156 assert_eq!(raw().broadcasts_closed, 2, "both viewers gone");
2157 }
2158
2159 #[tokio::test(start_paused = true)]
2160 async fn session_counts_by_root() {
2161 let (stats, _origin) = test_stats(Some("sjc"));
2164 let ext = stats.tier(Tier::External);
2165
2166 let snap = |root: &str| {
2167 let map = stats.shared().sessions[Tier::External.idx()].lock();
2168 map.get(&PathOwned::from(root.to_string())).map(|c| c.snapshot())
2169 };
2170
2171 let a1 = ext.session("acme");
2172 let a2 = ext.session("acme");
2173 let b1 = ext.session("globex");
2174 assert_eq!(snap("acme"), Some((2, 0)), "two sessions under one root");
2175 assert_eq!(snap("globex"), Some((1, 0)), "a distinct root is counted separately");
2176
2177 drop(a1);
2178 assert_eq!(snap("acme"), Some((2, 1)));
2179 drop(a2);
2180 drop(b1);
2181 assert_eq!(snap("acme"), Some((2, 2)));
2182 assert_eq!(snap("globex"), Some((1, 1)));
2183 }
2184
2185 #[tokio::test(start_paused = true)]
2186 async fn session_track_surfaces_by_root() {
2187 let (stats, origin) = test_stats(Some("sjc"));
2188 let mut consumer = origin.consume();
2189 let _a = stats.tier(Tier::External).session("acme");
2190 let _b = stats.tier(Tier::External).session("acme");
2191 let _c = stats.tier(Tier::Internal).session("peer");
2192
2193 tokio::time::advance(Duration::from_millis(1100)).await;
2194
2195 let (_path, broadcast) = consumer.announced().await.expect("announce");
2196 let broadcast = broadcast.expect("active");
2197
2198 let track = broadcast
2199 .subscribe_track(&Track {
2200 name: "sessions.json".into(),
2201 priority: 0,
2202 })
2203 .expect("subscribe");
2204 let frame = read_session_frame(track).await;
2205 let snap = frame.get("acme").expect("root entry");
2206 assert_eq!(snap.sessions, 2);
2207 assert_eq!(snap.sessions_closed, 0);
2208 assert!(
2209 !frame.contains_key("peer"),
2210 "internal session must not appear on the external track"
2211 );
2212
2213 let int_track = broadcast
2214 .subscribe_track(&Track {
2215 name: "internal/sessions.json".into(),
2216 priority: 0,
2217 })
2218 .expect("subscribe");
2219 let snap = *read_session_frame(int_track).await.get("peer").expect("internal entry");
2220 assert_eq!(snap.sessions, 1);
2221 }
2222
2223 #[tokio::test(start_paused = true)]
2224 async fn session_root_dropped_when_empty() {
2225 let (stats, _origin) = test_stats(Some("sjc"));
2228 let key = PathOwned::from("acme");
2229 let session = stats.tier(Tier::External).session("acme");
2230
2231 drive_ticks(1).await;
2232 assert!(
2233 stats.shared().sessions[Tier::External.idx()].lock().contains_key(&key),
2234 "root present while a session is connected"
2235 );
2236
2237 drop(session);
2238 drive_ticks(1).await;
2239 assert!(
2240 !stats.shared().sessions[Tier::External.idx()].lock().contains_key(&key),
2241 "root GC'd after the last session leaves"
2242 );
2243 }
2244
2245 #[tokio::test(start_paused = true)]
2246 async fn unused_slots_dont_surface() {
2247 let (stats, origin) = test_stats(Some("sjc"));
2253 let mut consumer = origin.consume();
2254 let bs = stats.tier(Tier::External).broadcast("foo/bar");
2255 let track = bs.publisher().track("video");
2256 track.frame();
2257
2258 drive_ticks(2).await;
2259
2260 let (_path, broadcast) = consumer.announced().await.expect("announce");
2261 let broadcast = broadcast.expect("active");
2262
2263 let pub_track = broadcast
2265 .subscribe_track(&Track {
2266 name: "publisher.json".into(),
2267 priority: 0,
2268 })
2269 .expect("subscribe");
2270 assert!(
2271 read_frame(pub_track).await.contains_key("foo/bar"),
2272 "publisher.json must include the active foo/bar entry"
2273 );
2274
2275 for name in ["subscriber.json", "internal/publisher.json", "internal/subscriber.json"] {
2278 let t = broadcast
2279 .subscribe_track(&Track {
2280 name: name.into(),
2281 priority: 0,
2282 })
2283 .expect("subscribe");
2284 let frame = read_frame(t).await;
2285 assert!(
2286 frame.is_empty(),
2287 "{name} must be empty for an entry with no activity on that slot, got {frame:?}",
2288 );
2289 }
2290 }
2291
2292 #[test]
2293 fn snapshot_reads_closed_before_open() {
2294 let src = include_str!("stats.rs");
2300 let body_start = src
2303 .find("fn snapshot(&self) -> RawCounts")
2304 .expect("snapshot fn present");
2305 let body = &src[body_start..];
2306 let closed_pos = body.find("self.announced_closed.load").expect("announced_closed load");
2307 let open_pos = body.find("self.announced.load(").expect("announced load");
2308 assert!(
2309 closed_pos < open_pos,
2310 "announced_closed must be loaded before announced; reversing breaks the open>=closed invariant",
2311 );
2312 let subs_closed_pos = body
2313 .find("self.subscriptions_closed.load")
2314 .expect("subscriptions_closed load");
2315 let subs_pos = body.find("self.subscriptions.load").expect("subscriptions load");
2316 assert!(
2317 subs_closed_pos < subs_pos,
2318 "subscriptions_closed must be loaded before subscriptions",
2319 );
2320 let bcast_closed_pos = body
2321 .find("self.broadcasts_closed.load")
2322 .expect("broadcasts_closed load");
2323 let bcast_pos = body.find("self.broadcasts.load").expect("broadcasts load");
2324 assert!(
2325 bcast_closed_pos < bcast_pos,
2326 "broadcasts_closed must be loaded before broadcasts",
2327 );
2328 }
2329
2330 #[test]
2331 fn session_snapshot_reads_closed_before_open() {
2332 let src = include_str!("stats.rs");
2336 let body_start = src
2337 .find("fn snapshot(&self) -> (u64, u64)")
2338 .expect("SessionCounters::snapshot fn present");
2339 let body = &src[body_start..];
2340 let closed_pos = body.find("self.sessions_closed.load").expect("sessions_closed load");
2341 let open_pos = body.find("self.sessions.load").expect("sessions load");
2342 assert!(closed_pos < open_pos, "sessions_closed must be loaded before sessions",);
2343 }
2344
2345 async fn read_frame(mut track: crate::TrackConsumer) -> BTreeMap<String, Snapshot> {
2346 let bytes = track.read_frame().await.expect("ok").expect("frame");
2347 serde_json::from_slice(&bytes).expect("json parse")
2348 }
2349
2350 async fn read_session_frame(mut track: crate::TrackConsumer) -> BTreeMap<String, SessionSnapshot> {
2351 let bytes = track.read_frame().await.expect("ok").expect("frame");
2352 serde_json::from_slice(&bytes).expect("json parse")
2353 }
2354}