Skip to main content

moq_stats/
produce.rs

1//! The publishing half: drain a [`Registry`] on an interval into stats tracks.
2
3use 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/// Settings for a [`Producer`]. Construct with [`ProducerConfig::new`] and chain
15/// the `with_*` setters (e.g.
16/// `ProducerConfig::new().with_origin(origin).with_prefix(".foo")`), then hand it
17/// to [`Producer::new`].
18///
19/// With no origin set the resulting producer is a no-op: its registry is
20/// disabled (bumps are dropped) and no task spawns. Call
21/// [`ProducerConfig::with_origin`] to publish.
22#[derive(Clone)]
23#[non_exhaustive]
24pub struct ProducerConfig {
25	/// Origin the stats broadcasts are created on.
26	/// When `None`, [`Producer::new`] spawns no task and publishes nothing.
27	pub origin: Option<origin::Producer>,
28	/// Top-level path stats are published under (default `.stats`). The full
29	/// advertised path is `<prefix>/node/<node>` (or `<prefix>/node` when
30	/// `node` is unset). Also the registry's exclude prefix, so serving a
31	/// stats broadcast doesn't generate more stats.
32	pub prefix: PathOwned,
33	/// Node suffix that disambiguates broadcasts from different relays sharing a
34	/// cluster origin. Set this on every node in multi-relay deployments. May be
35	/// multi-segment (e.g. `sjc/1`, `sjc/2`) so a region with multiple hosts can
36	/// nest under a shared region key. An empty path is treated as unset.
37	/// Default none.
38	pub node: Option<PathOwned>,
39	/// How long the publish task waits between drains. Default 1s.
40	pub interval: Duration,
41	/// How many leading broadcast-path segments to use as a grouping key.
42	///
43	/// Default `0` publishes one `<prefix>/node/<node>` broadcast carrying every
44	/// path. `1` publishes one broadcast per first segment at
45	/// `<prefix>/<group>/node/<node>`, and larger values include more leading
46	/// segments. Group broadcasts are announced while their group has live traffic;
47	/// at depth `0`, the single broadcast stays announced for the producer's life.
48	pub depth: usize,
49}
50
51impl ProducerConfig {
52	/// A config with default settings: no origin (no-op), `.stats` prefix, 1s
53	/// interval, and no node suffix. Call [`Self::with_origin`] to actually
54	/// publish.
55	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	/// Set the origin to publish the stats broadcast on. Without this the
66	/// producer is a no-op.
67	pub fn with_origin(mut self, origin: impl Into<Option<origin::Producer>>) -> Self {
68		self.origin = origin.into();
69		self
70	}
71
72	/// Override the top-level prefix (default `.stats`).
73	pub fn with_prefix(mut self, prefix: impl Into<PathOwned>) -> Self {
74		self.prefix = prefix.into();
75		self
76	}
77
78	/// Override the publish interval (default 1s).
79	pub fn with_interval(mut self, interval: Duration) -> Self {
80		self.interval = interval;
81		self
82	}
83
84	/// Set the node suffix (default none). An empty path is treated as unset.
85	pub fn with_node(mut self, node: impl Into<Option<PathOwned>>) -> Self {
86		self.node = node.into();
87		self
88	}
89
90	/// Set the grouping depth (default 0, a single broadcast). See [`Self::depth`].
91	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
103/// Keeps the publish task alive: the task holds only a `Weak` to this, so it
104/// exits once the last [`Producer`] clone drops.
105struct Keepalive;
106
107/// Publishes a [`Registry`]'s counters as stats broadcasts. Cheap to clone.
108///
109/// [`Producer::new`] builds the registry itself (wiring the config's prefix as
110/// its exclude prefix) and spawns the publish task; hand sessions tier-scoped
111/// handles via [`Registry::tier`] on [`Producer::registry`]. The task drains
112/// the registry every interval and writes a frame per changed track, running
113/// until the last [`Producer`] clone is dropped.
114#[derive(Clone)]
115pub struct Producer {
116	registry: Registry,
117	/// `None` for a no-op producer (config had no origin): no task was spawned
118	/// and the registry is disabled.
119	_keepalive: Option<Arc<Keepalive>>,
120}
121
122impl Producer {
123	/// Build a producer from `config`.
124	///
125	/// When `config` has an origin, this spawns the publish task immediately
126	/// and announces the stats broadcast; the task runs until the last
127	/// [`Producer`] clone is dropped. With no origin the producer is a no-op
128	/// (its registry is disabled, nothing is published) and no task spawns, so
129	/// it's safe to build outside an async runtime.
130	pub fn new(config: ProducerConfig) -> Self {
131		let ProducerConfig {
132			origin,
133			prefix,
134			node,
135			interval,
136			depth,
137		} = config;
138		// An empty path after normalization is indistinguishable from "no node
139		// set"; collapse it so downstream code only sees a single representation.
140		// We do this here (not in `with_node`) so a directly-assigned
141		// `config.node` is normalized too.
142		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	/// The registry this producer drains. Hand sessions tier-scoped handles via
170	/// [`Registry::tier`]; read node totals back with [`Registry::snapshot`].
171	/// Disabled (all bumps no-op) for a no-op producer.
172	pub fn registry(&self) -> &Registry {
173		&self.registry
174	}
175}
176
177/// Everything the publish task owns.
178struct 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	/// Publishes stats broadcasts and writes a frame per drain. Runs until
189	/// every [`Producer`] clone is dropped (`weak.upgrade()` returns `None`).
190	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			// Drain the registry: current per-broadcast values, with dead
215			// entries pruned (their final values are still in this report).
216			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			// Drop change-detection state for entries the report no longer
297			// carries (they were pruned on a previous drain).
298			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			// Deliberate unpublish: finish evicted broadcasts rather than dropping
314			// them, so there is no dropped-without-finish warning.
315			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
329/// A plain track and its `.z` sibling, kept in lockstep. The plain side runs
330/// moq-json with deltas and compression off, which is wire-identical to
331/// writing each frame as its own single-frame group; the compressed side uses
332/// merge-patch deltas inside a shared DEFLATE window.
333struct 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	/// Publish `frame` on both flavors; moq-json skips unchanged values.
353	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
363/// Ensure a track pair exists for every frame this drain produced, then push
364/// each pair its frame (an empty one when the drain had nothing for it, so a
365/// track whose last entry closed transitions to `{}` exactly once).
366fn 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
388/// One group stats broadcast and its change-detection state.
389struct 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		// The default tier's tracks always exist, even while idle.
413		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/// Change-detection state for one `(path, tier, side)` slot, owned by the
448/// publish task. The task is single-threaded so this needs no atomics.
449#[derive(Default)]
450struct SlotState {
451	/// Last [`Traffic`] we emitted for this slot, used to detect changes that
452	/// warrant re-emission.
453	prev_emitted: Option<Traffic>,
454}
455
456/// Change-detection state for one `(path, tier)`: a [`SlotState`] per side.
457#[derive(Default)]
458struct SideSlots {
459	publisher: SlotState,
460	subscriber: SlotState,
461}
462
463/// Change-detection state for one session-track root, mirroring [`SlotState`].
464#[derive(Default)]
465struct SessionSlotState {
466	prev_emitted: Option<Presence>,
467}
468
469/// Per-drain work for a single `(side, tier)` slot: update the slot's
470/// `prev_emitted` and hand `snap` to `emit` iff the slot is live or changed
471/// this drain.
472fn process_slot(snap: Traffic, slot_state: &mut SlotState, emit: impl FnOnce(Traffic)) {
473	// A slot is live while any open counter still exceeds its `*_closed`
474	// counterpart: a guard is held, so a subscription could begin at any
475	// moment. Live slots are emitted every drain so a downstream "currently
476	// active" view always sees the full set. Once every pair is equal no
477	// traffic can flow and the entry is on its way out (the registry pruned
478	// it as soon as the last guard released its handle).
479	let live = !snap.is_idle();
480
481	// Include the entry whenever it's live OR its snapshot changed this
482	// drain. Change-driven inclusion catches bumps since the previous drain
483	// (incl. sub-interval flickers) and emits the final close snapshot on the
484	// drain a slot transitions to fully closed.
485	//
486	// `None` (slot never emitted) is treated as the default Traffic so a
487	// first-drain all-zeros snap on an unused tier-side slot doesn't count
488	// as a "change". Without this, every entry would surface in all four
489	// tracks with zeros on the drain after creation even if only one slot
490	// is actually in use.
491	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
501/// Per-drain work for one session-track root: same live-or-changed rule as
502/// [`process_slot`].
503fn 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	// `<prefix>/<group>/node/<node>`. The group segment is empty at depth 0.
536	// The fixed `node` category leaves room for sibling categories (e.g.
537	// `<top-prefix>/<group>/cluster` for relay-mesh stats) under the same prefix.
538	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	/// Kept-alive handles from [`feed`]: dropping them closes the subscription and the
571	/// announce (bumping the `_closed` counters).
572	#[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	/// Drive a tagged egress broadcast so `registry` records publisher-side traffic on
581	/// `path` under `tier`. The local publisher (ingress) is left untagged, so only the
582	/// egress (publisher) counters move, matching how a relay bills read-out traffic.
583	///
584	/// Announces the broadcast; if `subscribe`, opens one subscription and reads a group
585	/// of `frames` frames of `frame_size` bytes each (so `bytes`/`frames`/`groups` move).
586	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		// Egress (publisher side) is tagged; the local publisher stays untagged.
597		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		// Let the origin's source watcher attach and announce (paused time advances
606		// instantly and yields to the spawned tasks).
607		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	/// Awaits the stats announce and returns its broadcast.
647	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	/// Advance past one publish interval so the task drains and writes frames.
655	async fn drive_tick() {
656		tokio::time::advance(Duration::from_millis(1100)).await;
657		// Yield several times to let the task wake, drain the registry, write
658		// the frames, and re-await the next tick.
659		for _ in 0..4 {
660			tokio::task::yield_now().await;
661		}
662	}
663
664	/// Reads the first frame off a plain track as raw JSON, pinning the plain
665	/// wire format (a full JSON object per frame, no compression).
666	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	/// Read the latest buffered traffic frame off a track. The producer emits an
673	/// immediate first (often empty) frame at time zero, so a test that records
674	/// traffic asynchronously reads the accumulated state rather than that stale one.
675	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	/// The advertised path normalizes a messy node suffix and drops an
701	/// all-empty one. Observed through the announced path, since the task
702	/// announces at construction.
703	#[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		// No matter how many broadcasts get bumped, exactly one stats
715		// broadcast is announced (the per-node aggregate).
716		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		// One announced broadcast, one subscription, one 42-byte frame read out.
735		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		// Announce only: the guard records the broadcast-name length once on open.
753		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		// An announce with no subscription should bump announced but NOT broadcasts
771		// (which only counts sessions with an active sub).
772		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		// A subscription that opens AND closes within a single drain window
788		// must still surface as a complete broadcasts open/close cycle. The
789		// cumulative counters retain broadcasts=1/broadcasts_closed=1, and the
790		// change-driven inclusion surfaces the entry even though it's net-idle
791		// by drain time.
792		let (producer, origin) = test_producer(Some("sjc"));
793		{
794			// Subscribe, read one 123-byte frame, then drop everything within the
795			// first interval so the open and close both land before the drain.
796			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		// One session opened then closed a subscription within the drain.
805		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		// A broadcast that only sees default-tier publisher traffic must NOT
842		// surface on its sibling default-tier subscriber track, and a tier
843		// with no traffic gets no tracks at all.
844		let (producer, origin) = test_producer(Some("sjc"));
845		// Only the egress (publisher) side is tagged, so `foo/bar` gets publisher
846		// traffic and no subscriber traffic.
847		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		// Default-tier publisher slot SHOULD include foo/bar.
855		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		// The default-tier subscriber slot had zero activity; its first frame
863		// must be `{}`, not `{"foo/bar": {all zeros}}`.
864		let frame = read_frame(&broadcast, "subscriber.json").await;
865		assert!(frame.is_empty(), "subscriber.json must be empty, got {frame:?}");
866
867		// The compressed siblings of the default tracks always exist.
868		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		// The regional tier never saw traffic, so its tracks were never created.
873		// The announced broadcast is an origin-owned splice, so `track()` always
874		// hands back a logical track; only the subscription reveals absence.
875		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}