1use std::collections::{HashMap, HashSet};
4use std::sync::{Arc, Weak};
5use std::time::Duration;
6
7use moq_net::stats::{Presence, Registry, Role, Tier, Traffic};
8use moq_net::{Path, PathOwned, broadcast, origin};
9use serde::Serialize;
10use web_async::spawn;
11
12use crate::{COMPRESSED_SUFFIX, SessionsFrame, TrafficFrame, sessions_track, traffic_track};
13
14#[derive(Clone)]
23#[non_exhaustive]
24pub struct ProducerConfig {
25 pub origin: Option<origin::Producer>,
28 pub prefix: PathOwned,
33 pub node: Option<PathOwned>,
39 pub interval: Duration,
41 pub depth: usize,
49}
50
51impl ProducerConfig {
52 pub fn new() -> Self {
56 Self {
57 origin: None,
58 prefix: PathOwned::from(".stats"),
59 node: None,
60 interval: Duration::from_secs(1),
61 depth: 0,
62 }
63 }
64
65 pub fn with_origin(mut self, origin: impl Into<Option<origin::Producer>>) -> Self {
68 self.origin = origin.into();
69 self
70 }
71
72 pub fn with_prefix(mut self, prefix: impl Into<PathOwned>) -> Self {
74 self.prefix = prefix.into();
75 self
76 }
77
78 pub fn with_interval(mut self, interval: Duration) -> Self {
80 self.interval = interval;
81 self
82 }
83
84 pub fn with_node(mut self, node: impl Into<Option<PathOwned>>) -> Self {
86 self.node = node.into();
87 self
88 }
89
90 pub fn with_depth(mut self, depth: usize) -> Self {
92 self.depth = depth;
93 self
94 }
95}
96
97impl Default for ProducerConfig {
98 fn default() -> Self {
99 Self::new()
100 }
101}
102
103struct Keepalive;
106
107#[derive(Clone)]
115pub struct Producer {
116 registry: Registry,
117 _keepalive: Option<Arc<Keepalive>>,
120}
121
122impl Producer {
123 pub fn new(config: ProducerConfig) -> Self {
131 let ProducerConfig {
132 origin,
133 prefix,
134 node,
135 interval,
136 depth,
137 } = config;
138 let node = node.filter(|p| !p.is_empty());
143
144 let Some(origin) = origin else {
145 return Self {
146 registry: Registry::disabled(),
147 _keepalive: None,
148 };
149 };
150
151 let registry = Registry::new(moq_net::stats::Config::new().with_exclude(prefix.clone()));
152 let keepalive = Arc::new(Keepalive);
153 let task = Task {
154 registry: registry.clone(),
155 origin,
156 prefix,
157 node,
158 depth,
159 interval,
160 };
161 spawn(task.run(Arc::downgrade(&keepalive)));
162
163 Self {
164 registry,
165 _keepalive: Some(keepalive),
166 }
167 }
168
169 pub fn registry(&self) -> &Registry {
173 &self.registry
174 }
175}
176
177struct Task {
179 registry: Registry,
180 origin: origin::Producer,
181 prefix: PathOwned,
182 node: Option<PathOwned>,
183 depth: usize,
184 interval: Duration,
185}
186
187impl Task {
188 async fn run(self, weak: Weak<Keepalive>) {
191 let node = self.node.as_ref().map(|p| p.as_str());
192 let mut groups: HashMap<PathOwned, GroupPublisher> = HashMap::new();
193
194 if self.depth == 0 {
195 let Some(group) = GroupPublisher::create(&self.origin, &self.prefix, &Path::empty(), node) else {
196 return;
197 };
198 groups.insert(Path::empty().to_owned(), group);
199 }
200
201 let mut ticker = web_async::time::interval(self.interval);
202 ticker.set_missed_tick_behavior(web_async::time::MissedTickBehavior::Delay);
203
204 loop {
205 ticker.tick().await;
206
207 if weak.upgrade().is_none() {
208 for (_, mut publisher) in groups.drain() {
209 publisher.broadcast.finish();
210 }
211 return;
212 }
213
214 let report = self.registry.report();
217
218 let mut entries_by_group: HashMap<PathOwned, Vec<&moq_net::stats::TrafficEntry>> = HashMap::new();
219 for entry in &report.traffic {
220 entries_by_group
221 .entry(group_key(entry.path.as_str(), self.depth))
222 .or_default()
223 .push(entry);
224 }
225
226 let mut sessions_by_group: HashMap<PathOwned, Vec<&moq_net::stats::SessionEntry>> = HashMap::new();
227 for entry in &report.sessions {
228 sessions_by_group
229 .entry(group_key(entry.root.as_str(), self.depth))
230 .or_default()
231 .push(entry);
232 }
233
234 let mut active: HashSet<PathOwned> = HashSet::new();
235 active.extend(entries_by_group.keys().cloned());
236 active.extend(sessions_by_group.keys().cloned());
237 if self.depth == 0 {
238 active.insert(Path::empty().to_owned());
239 }
240
241 for group in &active {
242 if !groups.contains_key(group) {
243 let Some(publisher) = GroupPublisher::create(&self.origin, &self.prefix, group, node) else {
244 continue;
245 };
246 groups.insert(group.clone(), publisher);
247 }
248 let publisher = groups.get_mut(group).expect("just inserted");
249
250 let mut frames: HashMap<String, TrafficFrame> = HashMap::new();
251 if let Some(group_entries) = entries_by_group.get(group) {
252 for entry in group_entries {
253 let slots = publisher
254 .local
255 .entry(entry.path.clone())
256 .or_default()
257 .entry(entry.tier.clone())
258 .or_default();
259 process_slot(entry.publisher, &mut slots.publisher, |snap| {
260 frames
261 .entry(traffic_track(&entry.tier, Role::Publisher, false))
262 .or_default()
263 .insert(entry.path.as_str().to_string(), snap);
264 });
265 process_slot(entry.subscriber, &mut slots.subscriber, |snap| {
266 frames
267 .entry(traffic_track(&entry.tier, Role::Subscriber, false))
268 .or_default()
269 .insert(entry.path.as_str().to_string(), snap);
270 });
271 }
272 }
273
274 let mut session_frames: HashMap<String, SessionsFrame> = HashMap::new();
275 if let Some(group_sessions) = sessions_by_group.get(group) {
276 for entry in group_sessions {
277 let state = publisher
278 .session_local
279 .entry(entry.tier.clone())
280 .or_default()
281 .entry(entry.root.clone())
282 .or_default();
283 process_session_slot(entry.presence, state, |snap| {
284 session_frames
285 .entry(sessions_track(&entry.tier, false))
286 .or_default()
287 .insert(entry.root.as_str().to_string(), snap);
288 });
289 }
290 }
291
292 flush_dynamic(&mut publisher.broadcast, &mut publisher.traffic_tracks, &frames);
293 flush_dynamic(&mut publisher.broadcast, &mut publisher.session_tracks, &session_frames);
294 }
295
296 let reported: HashSet<(&PathOwned, &Tier)> =
299 report.traffic.iter().map(|entry| (&entry.path, &entry.tier)).collect();
300 let reported_sessions: HashSet<(&Tier, &PathOwned)> =
301 report.sessions.iter().map(|entry| (&entry.tier, &entry.root)).collect();
302 for publisher in groups.values_mut() {
303 publisher.local.retain(|path, tiers| {
304 tiers.retain(|tier, _| reported.contains(&(path, tier)));
305 !tiers.is_empty()
306 });
307 publisher.session_local.retain(|tier, roots| {
308 roots.retain(|root, _| reported_sessions.contains(&(tier, root)));
309 !roots.is_empty()
310 });
311 }
312
313 let evicted: Vec<PathOwned> = groups
316 .keys()
317 .filter(|group| !active.contains(*group))
318 .cloned()
319 .collect();
320 for group in evicted {
321 if let Some(mut publisher) = groups.remove(&group) {
322 publisher.broadcast.finish();
323 }
324 }
325 }
326 }
327}
328
329struct TrackPair<T> {
334 plain: moq_json::snapshot::Producer<T>,
335 compressed: moq_json::snapshot::Producer<T>,
336}
337
338impl<T: Serialize> TrackPair<T> {
339 fn create(broadcast: &mut broadcast::Producer, name: &str) -> Result<Self, moq_net::Error> {
340 let plain_track = broadcast.create_track(name, None)?;
341 let compressed_track = broadcast.create_track(format!("{name}{COMPRESSED_SUFFIX}").as_str(), None)?;
342
343 let plain_config = moq_json::snapshot::ProducerConfig::default().with_delta_ratio(0);
344 let compressed_config = moq_json::snapshot::ProducerConfig::default().with_compression(true);
345
346 Ok(Self {
347 plain: moq_json::snapshot::Producer::new(plain_track, plain_config),
348 compressed: moq_json::snapshot::Producer::new(compressed_track, compressed_config),
349 })
350 }
351
352 fn update(&mut self, name: &str, frame: &T) {
354 if let Err(err) = self.plain.update(frame) {
355 tracing::debug!(?err, name, "stats: failed to write frame");
356 }
357 if let Err(err) = self.compressed.update(frame) {
358 tracing::debug!(?err, name, "stats: failed to write compressed frame");
359 }
360 }
361}
362
363fn flush_dynamic<T: Serialize + Default>(
367 broadcast: &mut broadcast::Producer,
368 tracks: &mut HashMap<String, TrackPair<T>>,
369 frames: &HashMap<String, T>,
370) {
371 for name in frames.keys() {
372 if !tracks.contains_key(name) {
373 match TrackPair::create(broadcast, name) {
374 Ok(pair) => {
375 tracks.insert(name.clone(), pair);
376 }
377 Err(err) => tracing::warn!(?err, name, "stats: failed to create track"),
378 }
379 }
380 }
381
382 let empty = T::default();
383 for (name, pair) in tracks.iter_mut() {
384 pair.update(name, frames.get(name).unwrap_or(&empty));
385 }
386}
387
388struct GroupPublisher {
390 broadcast: broadcast::Producer,
391 traffic_tracks: HashMap<String, TrackPair<TrafficFrame>>,
392 session_tracks: HashMap<String, TrackPair<SessionsFrame>>,
393 local: HashMap<PathOwned, HashMap<Tier, SideSlots>>,
394 session_local: HashMap<Tier, HashMap<PathOwned, SessionSlotState>>,
395}
396
397impl GroupPublisher {
398 fn create(origin: &origin::Producer, prefix: &Path, group: &Path, node: Option<&str>) -> Option<Self> {
399 let advertised = advertised_path(prefix, group, node);
400 let mut broadcast = match origin.create_broadcast(&advertised, broadcast::Route::new().with_announce(true)) {
401 Ok(broadcast) => broadcast,
402 Err(err) => {
403 tracing::warn!(advertised = %advertised, ?err, "stats: origin rejected stats broadcast");
404 return None;
405 }
406 };
407 tracing::debug!(advertised = %advertised, "stats: publishing broadcast");
408
409 let mut traffic_tracks = HashMap::new();
410 let mut session_tracks = HashMap::new();
411
412 let tier = Tier::default();
414 for role in [Role::Publisher, Role::Subscriber] {
415 let name = traffic_track(&tier, role, false);
416 match TrackPair::create(&mut broadcast, &name) {
417 Ok(pair) => {
418 traffic_tracks.insert(name, pair);
419 }
420 Err(err) => {
421 tracing::warn!(?err, name, "stats: failed to create track");
422 return None;
423 }
424 }
425 }
426 let name = sessions_track(&tier, false);
427 match TrackPair::create(&mut broadcast, &name) {
428 Ok(pair) => {
429 session_tracks.insert(name, pair);
430 }
431 Err(err) => {
432 tracing::warn!(?err, name, "stats: failed to create track");
433 return None;
434 }
435 }
436
437 Some(Self {
438 broadcast,
439 traffic_tracks,
440 session_tracks,
441 local: HashMap::new(),
442 session_local: HashMap::new(),
443 })
444 }
445}
446
447#[derive(Default)]
450struct SlotState {
451 prev_emitted: Option<Traffic>,
454}
455
456#[derive(Default)]
458struct SideSlots {
459 publisher: SlotState,
460 subscriber: SlotState,
461}
462
463#[derive(Default)]
465struct SessionSlotState {
466 prev_emitted: Option<Presence>,
467}
468
469fn process_slot(snap: Traffic, slot_state: &mut SlotState, emit: impl FnOnce(Traffic)) {
473 let live = !snap.is_idle();
480
481 let prev_snap = slot_state.prev_emitted.unwrap_or_default();
492 let changed = snap != prev_snap;
493 if changed {
494 slot_state.prev_emitted = Some(snap);
495 }
496 if live || changed {
497 emit(snap);
498 }
499}
500
501fn process_session_slot(snap: Presence, slot_state: &mut SessionSlotState, emit: impl FnOnce(Presence)) {
504 let live = snap.active() > 0;
505 let prev_snap = slot_state.prev_emitted.unwrap_or_default();
506 let changed = snap != prev_snap;
507 if changed {
508 slot_state.prev_emitted = Some(snap);
509 }
510 if live || changed {
511 emit(snap);
512 }
513}
514
515fn group_key(path: &str, depth: usize) -> PathOwned {
516 if depth == 0 {
517 return Path::empty().to_owned();
518 }
519
520 let mut seen = 0;
521 let mut end = path.len();
522 for (i, b) in path.bytes().enumerate() {
523 if b == b'/' {
524 seen += 1;
525 if seen == depth {
526 end = i;
527 break;
528 }
529 }
530 }
531 Path::new(&path[..end]).to_owned()
532}
533
534fn advertised_path(prefix: &Path, group: &Path, node: Option<&str>) -> PathOwned {
535 let mut out = prefix.as_str().to_string();
539 if !group.is_empty() {
540 out.push('/');
541 out.push_str(group.as_str());
542 }
543 out.push_str("/node");
544 if let Some(node) = node {
545 out.push('/');
546 out.push_str(node);
547 }
548 PathOwned::from(out)
549}
550
551#[cfg(test)]
552mod tests {
553 use std::collections::BTreeMap;
554
555 use moq_net::stats::{Registry, Tier};
556 use moq_net::{Origin, Timestamp, announce, broadcast, track};
557
558 use super::*;
559
560 fn test_producer(node: Option<&str>) -> (Producer, origin::Producer) {
561 let origin = Origin::random().produce();
562 let producer = Producer::new(
563 ProducerConfig::new()
564 .with_origin(origin.clone())
565 .with_node(node.map(|s| PathOwned::from(s.to_string()))),
566 );
567 (producer, origin)
568 }
569
570 #[allow(dead_code)]
573 struct Feed {
574 announced: announce::Consumer,
575 source: broadcast::Producer,
576 consumer: broadcast::Consumer,
577 sub: Option<track::Subscriber>,
578 }
579
580 async fn feed(
587 registry: &Registry,
588 tier: Tier,
589 path: &str,
590 subscribe: bool,
591 frames: usize,
592 frame_size: usize,
593 ) -> Feed {
594 let ctx = registry.tier(tier).session("feed");
595 let origin = Origin::random().produce();
596 let egress = origin.consume().with_stats(ctx);
598
599 let mut announced = egress.announced();
600 let mut source = origin
601 .create_broadcast(path, broadcast::Route::announced())
602 .expect("create_broadcast");
603 let mut producer = source.create_track("video", None).expect("create_track");
604
605 tokio::time::sleep(Duration::from_millis(1)).await;
608 tokio::time::sleep(Duration::from_millis(1)).await;
609
610 let announce::Update { broadcast, .. } = announced.next().await.expect("announce");
611 let consumer = broadcast.expect("active");
612
613 let sub = if subscribe {
614 let mut sub = consumer
615 .track("video")
616 .expect("track")
617 .subscribe(None)
618 .await
619 .expect("subscribe");
620
621 if frames > 0 {
622 let mut group = producer.append_group().expect("group");
623 for _ in 0..frames {
624 group
625 .write_frame(Timestamp::ZERO, vec![0u8; frame_size])
626 .expect("write");
627 }
628 group.finish().expect("finish");
629
630 let mut group = sub.recv_group().await.expect("recv").expect("group");
631 while group.read_frame().await.expect("read").is_some() {}
632 }
633 Some(sub)
634 } else {
635 None
636 };
637
638 Feed {
639 announced,
640 source,
641 consumer,
642 sub,
643 }
644 }
645
646 async fn announced(origin: &origin::Producer) -> (String, moq_net::broadcast::Consumer) {
648 let mut consumer = origin.consume().announced();
649 tokio::time::advance(Duration::from_millis(1)).await;
650 let announce::Update { path, broadcast } = consumer.next().await.expect("expected announce");
651 (path.as_str().to_string(), broadcast.expect("active"))
652 }
653
654 async fn drive_tick() {
656 tokio::time::advance(Duration::from_millis(1100)).await;
657 for _ in 0..4 {
660 tokio::task::yield_now().await;
661 }
662 }
663
664 async fn read_frame(broadcast: &moq_net::broadcast::Consumer, name: &str) -> BTreeMap<String, Traffic> {
667 let mut track = subscribe(broadcast, name).await;
668 let frame = track.read_frame().await.expect("ok").expect("frame");
669 serde_json::from_slice(&frame.payload).expect("json parse")
670 }
671
672 async fn read_last_frame(broadcast: &moq_net::broadcast::Consumer, name: &str) -> BTreeMap<String, Traffic> {
676 use futures::FutureExt;
677 let mut track = subscribe(broadcast, name).await;
678 let mut last = track.read_frame().await.expect("ok").expect("frame");
679 while let Some(Ok(Some(frame))) = track.read_frame().now_or_never() {
680 last = frame;
681 }
682 serde_json::from_slice(&last.payload).expect("json parse")
683 }
684
685 async fn read_session_frame(broadcast: &moq_net::broadcast::Consumer, name: &str) -> BTreeMap<String, Presence> {
686 let mut track = subscribe(broadcast, name).await;
687 let frame = track.read_frame().await.expect("ok").expect("frame");
688 serde_json::from_slice(&frame.payload).expect("json parse")
689 }
690
691 async fn subscribe(broadcast: &moq_net::broadcast::Consumer, name: &str) -> track::Subscriber {
692 broadcast
693 .track(name)
694 .expect("track")
695 .subscribe(None)
696 .await
697 .expect("subscribe")
698 }
699
700 #[tokio::test(start_paused = true)]
704 async fn new_normalizes_and_drops_empty_node() {
705 let (_producer, origin) = test_producer(Some("/sjc//1/"));
706 assert_eq!(announced(&origin).await.0, ".stats/node/sjc/1");
707
708 let (_producer, origin) = test_producer(Some("///"));
709 assert_eq!(announced(&origin).await.0, ".stats/node");
710 }
711
712 #[tokio::test(start_paused = true)]
713 async fn single_broadcast_path_announced() {
714 let (producer, origin) = test_producer(Some("sjc/1"));
717
718 let _f1 = feed(producer.registry(), Tier::default(), "foo/bar", true, 1, 8).await;
719 let _f2 = feed(producer.registry(), Tier::default(), "baz/qux", true, 1, 8).await;
720
721 assert_eq!(announced(&origin).await.0, ".stats/node/sjc/1");
722 }
723
724 #[tokio::test(start_paused = true)]
725 async fn task_announces_without_node_suffix() {
726 let (producer, origin) = test_producer(None);
727 let _f = feed(producer.registry(), Tier::default(), "foo/bar", true, 1, 8).await;
728 assert_eq!(announced(&origin).await.0, ".stats/node");
729 }
730
731 #[tokio::test(start_paused = true)]
732 async fn frame_emits_expected_counters() {
733 let (producer, origin) = test_producer(Some("sjc"));
734 let _f = feed(producer.registry(), Tier::default(), "foo/bar", true, 1, 42).await;
736
737 drive_tick().await;
738
739 let (_, broadcast) = announced(&origin).await;
740 let frame = read_last_frame(&broadcast, "publisher.json").await;
741 let snap = frame.get("foo/bar").expect("foo/bar entry");
742 assert_eq!(snap.announced, 1, "egress announce stream bumps announced");
743 assert_eq!(snap.broadcasts, 1, "one session subscribed");
744 assert_eq!(snap.subscriptions, 1);
745 assert_eq!(snap.bytes, 42);
746 assert_eq!(snap.frames, 1);
747 }
748
749 #[tokio::test(start_paused = true)]
750 async fn announced_bytes_surfaces_in_frame() {
751 let (producer, origin) = test_producer(Some("sjc"));
752 let _f = feed(producer.registry(), Tier::default(), "foo/bar", false, 0, 0).await;
754
755 drive_tick().await;
756
757 let (_, broadcast) = announced(&origin).await;
758 let frame = read_last_frame(&broadcast, "publisher.json").await;
759 let snap = frame.get("foo/bar").expect("foo/bar entry");
760 assert_eq!(snap.announced, 1);
761 assert_eq!(
762 snap.announced_bytes,
763 "foo/bar".len() as u64,
764 "name length recorded on announce"
765 );
766 }
767
768 #[tokio::test(start_paused = true)]
769 async fn announced_decouples_from_broadcasts() {
770 let (producer, origin) = test_producer(Some("sjc"));
773 let _f = feed(producer.registry(), Tier::default(), "foo/bar", false, 0, 0).await;
774
775 drive_tick().await;
776
777 let (_, broadcast) = announced(&origin).await;
778 let frame = read_last_frame(&broadcast, "publisher.json").await;
779 let snap = frame.get("foo/bar").expect("foo/bar entry");
780 assert_eq!(snap.announced, 1);
781 assert_eq!(snap.broadcasts, 0, "no subscription, no broadcasts sentinel");
782 assert_eq!(snap.subscriptions, 0);
783 }
784
785 #[tokio::test(start_paused = true)]
786 async fn short_lived_sub_is_surfaced() {
787 let (producer, origin) = test_producer(Some("sjc"));
793 {
794 let _f = feed(producer.registry(), Tier::default(), "foo/bar", true, 1, 123).await;
797 }
798
799 drive_tick().await;
800
801 let (_, broadcast) = announced(&origin).await;
802 let frame = read_last_frame(&broadcast, "publisher.json").await;
803 let snap = frame.get("foo/bar").expect("foo/bar entry");
804 assert_eq!(snap.subscriptions, 1);
806 assert_eq!(snap.subscriptions_closed, 1);
807 assert_eq!(snap.broadcasts, 1, "one session subscribed");
808 assert_eq!(snap.broadcasts_closed, 1);
809 assert_eq!(snap.bytes, 123);
810 assert_eq!(snap.frames, 1);
811 }
812
813 #[tokio::test(start_paused = true)]
814 async fn session_track_surfaces_by_root() {
815 let (producer, origin) = test_producer(Some("sjc"));
816 let _a = producer.registry().tier(Tier::default()).session("acme");
817 let _b = producer.registry().tier(Tier::default()).session("acme");
818 let _c = producer.registry().tier(Tier::new("region/sjc")).session("peer");
819
820 drive_tick().await;
821
822 let (_, broadcast) = announced(&origin).await;
823 let frame = read_session_frame(&broadcast, "sessions.json").await;
824 let snap = frame.get("acme").expect("root entry");
825 assert_eq!(snap.sessions, 2);
826 assert_eq!(snap.sessions_closed, 0);
827 assert!(
828 !frame.contains_key("peer"),
829 "regional session must not appear on the default track"
830 );
831
832 let snap = *read_session_frame(&broadcast, "region/sjc/sessions.json")
833 .await
834 .get("peer")
835 .expect("regional entry");
836 assert_eq!(snap.sessions, 1);
837 }
838
839 #[tokio::test(start_paused = true)]
840 async fn unused_slots_dont_surface() {
841 let (producer, origin) = test_producer(Some("sjc"));
845 let _f = feed(producer.registry(), Tier::default(), "foo/bar", true, 1, 8).await;
848
849 drive_tick().await;
850 drive_tick().await;
851
852 let (_, broadcast) = announced(&origin).await;
853
854 assert!(
856 read_last_frame(&broadcast, "publisher.json")
857 .await
858 .contains_key("foo/bar"),
859 "publisher.json must include the active foo/bar entry"
860 );
861
862 let frame = read_frame(&broadcast, "subscriber.json").await;
865 assert!(frame.is_empty(), "subscriber.json must be empty, got {frame:?}");
866
867 for name in ["publisher.json.z", "subscriber.json.z", "sessions.json.z"] {
869 assert!(broadcast.track(name).is_ok(), "{name} must exist");
870 }
871
872 for name in ["region/sjc/publisher.json", "region/sjc/publisher.json.z"] {
876 let track = broadcast.track(name).expect("logical track");
877 assert!(
878 track.subscribe(None).await.is_err(),
879 "{name} must not exist for a tier with no traffic",
880 );
881 }
882 }
883
884 #[test]
885 fn advertised_path_with_and_without_node() {
886 let prefix = Path::new(".stats");
887 let empty = Path::empty();
888 assert_eq!(
889 advertised_path(&prefix, &empty, Some("sjc")).as_str(),
890 ".stats/node/sjc"
891 );
892 assert_eq!(
893 advertised_path(&prefix, &empty, Some("sjc/1")).as_str(),
894 ".stats/node/sjc/1"
895 );
896 assert_eq!(advertised_path(&prefix, &empty, None).as_str(), ".stats/node");
897 assert_eq!(
898 advertised_path(&prefix, &Path::new("acme"), Some("sjc")).as_str(),
899 ".stats/acme/node/sjc"
900 );
901
902 let prefix = Path::new("metrics");
903 assert_eq!(
904 advertised_path(&prefix, &Path::new("demo/room"), Some("lon")).as_str(),
905 "metrics/demo/room/node/lon"
906 );
907 }
908
909 #[test]
910 fn group_key_uses_leading_segments() {
911 assert_eq!(group_key("acme/room/cam", 0), Path::empty().to_owned());
912 assert_eq!(group_key("acme/room/cam", 1), Path::new("acme").to_owned());
913 assert_eq!(group_key("acme/room/cam", 2), Path::new("acme/room").to_owned());
914 assert_eq!(group_key("acme/room", 3), Path::new("acme/room").to_owned());
915 }
916}