use std::collections::{HashMap, HashSet};
use std::sync::{Arc, Weak};
use std::time::Duration;
use std::task::Poll;
use moq_net::stats::{Presence, Registry, Report, Role, Tier, Traffic};
use moq_net::{Path, PathOwned, broadcast, kio, origin, track};
use serde::Serialize;
use web_async::spawn;
use crate::{COMPRESSED_SUFFIX, sessions_track, traffic_track};
#[derive(Clone)]
#[non_exhaustive]
pub struct Config {
pub origin: Option<origin::Producer>,
pub prefix: PathOwned,
pub node: Option<PathOwned>,
pub interval: Duration,
pub depth: usize,
}
impl Config {
pub fn new() -> Self {
Self {
origin: None,
prefix: PathOwned::from(".stats"),
node: None,
interval: Duration::from_secs(1),
depth: 0,
}
}
pub fn with_origin(mut self, origin: impl Into<Option<origin::Producer>>) -> Self {
self.origin = origin.into();
self
}
pub fn with_prefix(mut self, prefix: impl Into<PathOwned>) -> Self {
self.prefix = prefix.into();
self
}
pub fn with_interval(mut self, interval: Duration) -> Self {
self.interval = interval;
self
}
pub fn with_node(mut self, node: impl Into<Option<PathOwned>>) -> Self {
self.node = node.into();
self
}
pub fn with_depth(mut self, depth: usize) -> Self {
self.depth = depth;
self
}
}
impl Default for Config {
fn default() -> Self {
Self::new()
}
}
const MAX_REQUESTED_TRACKS: usize = 64;
const MAX_PARKED_REQUESTS: usize = 256;
struct Keepalive;
#[derive(Clone)]
pub struct Producer {
registry: Registry,
_keepalive: Option<Arc<Keepalive>>,
}
impl Producer {
pub fn new(config: Config) -> Self {
let Config {
origin,
prefix,
node,
interval,
depth,
} = config;
let node = node.filter(|p| !p.is_empty());
let Some(origin) = origin else {
return Self {
registry: Registry::disabled(),
_keepalive: None,
};
};
let exclude = moq_net::Pattern::subtree(prefix.as_str()).expect("the stats prefix is a literal path");
let registry = Registry::new(moq_net::stats::Config::new().with_exclude(exclude));
let keepalive = Arc::new(Keepalive);
let task = Task {
registry: registry.clone(),
origin,
prefix,
node,
depth,
interval,
};
spawn(task.run(Arc::downgrade(&keepalive)));
Self {
registry,
_keepalive: Some(keepalive),
}
}
pub fn registry(&self) -> &Registry {
&self.registry
}
}
struct Task {
registry: Registry,
origin: origin::Producer,
prefix: PathOwned,
node: Option<PathOwned>,
depth: usize,
interval: Duration,
}
impl Task {
async fn run(self, weak: Weak<Keepalive>) {
let interval = self.interval;
let Some(mut drain) = Drain::new(self) else {
return;
};
let mut ticker = web_async::time::interval(interval);
ticker.set_missed_tick_behavior(web_async::time::MissedTickBehavior::Delay);
loop {
ticker.tick().await;
if weak.upgrade().is_none() {
drain.finish();
return;
}
drain.collect();
drain.publish();
}
}
fn node(&self) -> Option<&str> {
self.node.as_ref().map(moq_net::Path::as_str)
}
}
struct Drain {
task: Task,
groups: HashMap<String, GroupPublisher>,
report: Report,
refused: Vec<String>,
tick: u64,
}
impl Drain {
fn new(task: Task) -> Option<Self> {
let mut groups = HashMap::new();
if task.depth == 0 {
let group = GroupPublisher::create(&task.origin, &task.prefix, &Path::empty(), task.node())?;
groups.insert(String::new(), group);
}
Some(Self {
task,
groups,
report: Report::default(),
refused: Vec::new(),
tick: 0,
})
}
fn collect(&mut self) {
self.task.registry.report(&mut self.report);
self.tick += 1;
self.refused.clear();
for group in self.groups.values_mut() {
group.traffic_rows.clear();
group.session_rows.clear();
}
for (i, entry) in self.report.traffic.iter().enumerate() {
let key = group_key(entry.path.as_str(), self.task.depth);
if let Some(group) = Self::group(&self.task, &mut self.groups, &mut self.refused, key) {
group.traffic_rows.push(i);
}
}
for (i, entry) in self.report.sessions.iter().enumerate() {
let key = group_key(entry.root.as_str(), self.task.depth);
if let Some(group) = Self::group(&self.task, &mut self.groups, &mut self.refused, key) {
group.session_rows.push(i);
}
}
for group in self.groups.values_mut() {
group.collect(&self.report, self.tick);
}
}
fn group<'a>(
task: &Task,
groups: &'a mut HashMap<String, GroupPublisher>,
refused: &mut Vec<String>,
key: &str,
) -> Option<&'a mut GroupPublisher> {
if !groups.contains_key(key) {
if refused.iter().any(|name| name == key) {
return None;
}
match GroupPublisher::create(&task.origin, &task.prefix, &Path::new(key), task.node()) {
Some(group) => {
groups.insert(key.to_string(), group);
}
None => {
refused.push(key.to_string());
return None;
}
}
}
groups.get_mut(key)
}
fn publish(&mut self) {
let depth = self.task.depth;
for (_, group) in self
.groups
.extract_if(|_, group| depth > 0 && group.traffic_rows.is_empty() && group.session_rows.is_empty())
{
group.finish();
}
for group in self.groups.values_mut() {
group.flush(self.tick);
group.serve_requests();
}
}
fn finish(self) {
for (_, group) in self.groups {
group.finish();
}
}
}
struct Frame<V> {
entries: Vec<(PathOwned, V)>,
}
impl<V> Default for Frame<V> {
fn default() -> Self {
Self { entries: Vec::new() }
}
}
impl<V: Serialize> Serialize for Frame<V> {
fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
serializer.collect_map(self.entries.iter().map(|(path, value)| (path.as_str(), value)))
}
}
struct TrackPair<V> {
plain: moq_json::snapshot::Producer<Frame<V>>,
compressed: moq_json::snapshot::Producer<Frame<V>>,
frame: Frame<V>,
}
impl<V: Serialize> TrackPair<V> {
fn create(broadcast: &broadcast::Producer, name: &str) -> Result<Self, moq_net::Error> {
let plain_track = broadcast.create_track(name, None)?;
let compressed_track = broadcast.create_track(format!("{name}{COMPRESSED_SUFFIX}").as_str(), None)?;
Ok(Self::from_tracks(plain_track, compressed_track))
}
fn adopt(broadcast: &broadcast::Producer, name: &str, pending: PendingPair) -> Result<Self, moq_net::Error> {
let PendingPair { plain, compressed } = pending;
let plain_track = match plain {
Some(request) => request.accept(None),
None => broadcast.create_track(name, None)?,
};
let compressed_track = match compressed {
Some(request) => request.accept(None),
None => broadcast.create_track(format!("{name}{COMPRESSED_SUFFIX}").as_str(), None)?,
};
Ok(Self::from_tracks(plain_track, compressed_track))
}
fn from_tracks(plain_track: track::Producer, compressed_track: track::Producer) -> Self {
let plain_config = moq_json::snapshot::Config::default().with_delta_ratio(0);
let mut compressed_config = moq_json::snapshot::Config::default();
compressed_config.compression = moq_json::Compression::Deflate;
Self {
plain: moq_json::snapshot::Producer::new(plain_track, plain_config),
compressed: moq_json::snapshot::Producer::new(compressed_track, compressed_config),
frame: Frame::default(),
}
}
fn is_used(&self) -> bool {
self.plain.is_used() || self.compressed.is_used()
}
fn publish(&mut self, name: &str) {
self.frame.entries.sort_unstable_by(|a, b| a.0.cmp(&b.0));
if let Err(err) = self.plain.update(&self.frame) {
tracing::debug!(?err, name, "stats: failed to write frame");
}
if let Err(err) = self.compressed.update(&self.frame) {
tracing::debug!(?err, name, "stats: failed to write compressed frame");
}
self.frame.entries.clear();
}
fn finish(&mut self) {
let _ = self.plain.finish();
let _ = self.compressed.finish();
}
}
#[derive(Default)]
struct PendingPair {
plain: Option<track::Request>,
compressed: Option<track::Request>,
}
impl PendingPair {
fn reject(self, err: moq_net::Error) {
if let Some(request) = self.plain {
request.reject(err.clone());
}
if let Some(request) = self.compressed {
request.reject(err);
}
}
fn is_used(&self, waiter: &kio::Waiter) -> bool {
self.plain
.iter()
.chain(self.compressed.iter())
.any(|request| request.poll_unused(waiter).is_pending())
}
}
struct TrackFamily<V> {
tracks: HashMap<String, TrackPair<V>>,
parked: HashMap<String, PendingPair>,
}
impl<V: Serialize> TrackFamily<V> {
fn new() -> Self {
Self {
tracks: HashMap::new(),
parked: HashMap::new(),
}
}
fn push(
&mut self,
broadcast: &broadcast::Producer,
requested: &mut HashSet<String>,
name: &str,
path: PathOwned,
value: V,
) {
if !self.tracks.contains_key(name) {
let result = match self.parked.remove(name) {
Some(pending) => TrackPair::adopt(broadcast, name, pending),
None => TrackPair::create(broadcast, name),
};
match result {
Ok(pair) => {
self.tracks.insert(name.to_string(), pair);
}
Err(err) => {
tracing::warn!(?err, name, "stats: failed to create track");
return;
}
}
}
if !requested.is_empty() {
requested.remove(name);
}
let pair = self.tracks.get_mut(name).expect("just ensured");
pair.frame.entries.push((path, value));
}
fn flush(&mut self) {
for (name, pair) in self.tracks.iter_mut() {
pair.publish(name);
}
}
fn reclaim(&mut self, requested: &mut HashSet<String>) {
self.tracks.retain(|name, pair| {
if !requested.contains(name) || pair.is_used() {
return true;
}
requested.remove(name);
pair.finish();
false
});
}
fn park(&mut self, plain: String, compressed: bool, request: track::Request, full: bool) {
match self.parked.get_mut(&plain) {
Some(pending) => {
let slot = match compressed {
true => &mut pending.compressed,
false => &mut pending.plain,
};
if slot.is_none() {
*slot = Some(request);
}
}
None if full => request.reject(moq_net::Error::NotFound),
None => {
let mut pending = PendingPair::default();
match compressed {
true => pending.compressed = Some(request),
false => pending.plain = Some(request),
}
self.parked.insert(plain, pending);
}
}
}
fn adopt_parked(&mut self, broadcast: &broadcast::Producer, requested: &mut HashSet<String>) {
let noop = kio::Waiter::noop();
let mut parked = std::mem::take(&mut self.parked);
parked.retain(|plain, pending| {
if !pending.is_used(&noop) {
return false;
}
if requested.len() >= MAX_REQUESTED_TRACKS {
return true;
}
self.adopt_pair(broadcast, requested, plain.clone(), std::mem::take(pending));
false
});
self.parked = parked;
}
fn adopt_pair(
&mut self,
broadcast: &broadcast::Producer,
requested: &mut HashSet<String>,
plain: String,
pending: PendingPair,
) {
if self.tracks.contains_key(&plain) {
pending.reject(moq_net::Error::NotFound);
return;
}
match TrackPair::adopt(broadcast, &plain, pending) {
Ok(mut pair) => {
pair.publish(&plain);
self.tracks.insert(plain.clone(), pair);
requested.insert(plain);
}
Err(err) => tracing::warn!(?err, name = %plain, "stats: failed to adopt requested track"),
}
}
fn finish(&mut self) {
for pair in self.tracks.values_mut() {
pair.finish();
}
}
}
struct GroupPublisher {
broadcast: broadcast::Producer,
dynamic: broadcast::Dynamic,
requested: HashSet<String>,
traffic: TrackFamily<Traffic>,
sessions: TrackFamily<Presence>,
local: HashMap<PathOwned, HashMap<Tier, SideSlots>>,
session_local: HashMap<Tier, HashMap<PathOwned, SessionSlotState>>,
names: HashMap<Tier, TierNames>,
traffic_rows: Vec<usize>,
session_rows: Vec<usize>,
}
struct TierNames {
publisher: String,
subscriber: String,
sessions: String,
}
impl TierNames {
fn new(tier: &Tier) -> Self {
Self {
publisher: traffic_track(tier, Role::Publisher, false),
subscriber: traffic_track(tier, Role::Subscriber, false),
sessions: sessions_track(tier, false),
}
}
}
impl GroupPublisher {
fn create(origin: &origin::Producer, prefix: &Path, group: &Path, node: Option<&str>) -> Option<Self> {
let advertised = advertised_path(prefix, group, node);
let broadcast = match origin.publish(&advertised, origin::Route::default()) {
Ok(broadcast) => broadcast,
Err(err) => {
tracing::warn!(advertised = %advertised, ?err, "stats: origin rejected stats broadcast");
return None;
}
};
tracing::debug!(advertised = %advertised, "stats: publishing broadcast");
let mut traffic = TrackFamily::new();
let mut sessions = TrackFamily::new();
let tier = Tier::default();
for role in [Role::Publisher, Role::Subscriber] {
let name = traffic_track(&tier, role, false);
match TrackPair::create(&broadcast, &name) {
Ok(pair) => {
traffic.tracks.insert(name, pair);
}
Err(err) => {
tracing::warn!(?err, name, "stats: failed to create track");
return None;
}
}
}
let name = sessions_track(&tier, false);
match TrackPair::create(&broadcast, &name) {
Ok(pair) => {
sessions.tracks.insert(name, pair);
}
Err(err) => {
tracing::warn!(?err, name, "stats: failed to create track");
return None;
}
}
let dynamic = broadcast.dynamic();
Some(Self {
broadcast,
dynamic,
requested: HashSet::new(),
traffic,
sessions,
local: HashMap::new(),
session_local: HashMap::new(),
names: HashMap::new(),
traffic_rows: Vec::new(),
session_rows: Vec::new(),
})
}
fn collect(&mut self, report: &Report, tick: u64) {
let Self {
broadcast,
requested,
traffic,
sessions,
local,
session_local,
names,
traffic_rows,
session_rows,
..
} = self;
for &i in traffic_rows.iter() {
let entry = &report.traffic[i];
let names = names
.entry(entry.tier.clone())
.or_insert_with(|| TierNames::new(&entry.tier));
let slots = local
.entry(entry.path.clone())
.or_default()
.entry(entry.tier.clone())
.or_default();
slots.seen = tick;
process_slot(entry.publisher, &mut slots.publisher, |snap| {
traffic.push(broadcast, requested, &names.publisher, entry.path.clone(), snap);
});
process_slot(entry.subscriber, &mut slots.subscriber, |snap| {
traffic.push(broadcast, requested, &names.subscriber, entry.path.clone(), snap);
});
}
for &i in session_rows.iter() {
let entry = &report.sessions[i];
let names = names
.entry(entry.tier.clone())
.or_insert_with(|| TierNames::new(&entry.tier));
let state = session_local
.entry(entry.tier.clone())
.or_default()
.entry(entry.root.clone())
.or_default();
state.seen = tick;
process_session_slot(entry.presence, state, |snap| {
sessions.push(broadcast, requested, &names.sessions, entry.root.clone(), snap);
});
}
}
fn flush(&mut self, tick: u64) {
self.traffic.flush();
self.sessions.flush();
self.local.retain(|_, tiers| {
tiers.retain(|_, slots| slots.seen == tick);
!tiers.is_empty()
});
self.session_local.retain(|_, roots| {
roots.retain(|_, state| state.seen == tick);
!roots.is_empty()
});
}
fn serve_requests(&mut self) {
self.traffic.reclaim(&mut self.requested);
self.sessions.reclaim(&mut self.requested);
let noop = kio::Waiter::noop();
while let Poll::Ready(Ok(request)) = self.dynamic.poll_requested_track(&noop) {
let Some(shape) = requested_track_shape(request.name()) else {
request.reject(moq_net::Error::NotFound);
continue;
};
let full = self.traffic.parked.len() + self.sessions.parked.len() >= MAX_PARKED_REQUESTS;
match shape.sessions {
true => self.sessions.park(shape.plain, shape.compressed, request, full),
false => self.traffic.park(shape.plain, shape.compressed, request, full),
}
}
self.traffic.adopt_parked(&self.broadcast, &mut self.requested);
self.sessions.adopt_parked(&self.broadcast, &mut self.requested);
}
fn finish(mut self) {
self.traffic.finish();
self.sessions.finish();
self.broadcast.finish();
}
}
struct RequestedShape {
plain: String,
compressed: bool,
sessions: bool,
}
fn requested_track_shape(name: &str) -> Option<RequestedShape> {
let (base, compressed) = match name.strip_suffix(COMPRESSED_SUFFIX) {
Some(base) => (base, true),
None => (name, false),
};
let (tier, kind) = match base.rsplit_once('/') {
Some((tier, kind)) => (Some(tier), kind),
None => (None, base),
};
let sessions = match kind {
"publisher.json" | "subscriber.json" => false,
"sessions.json" => true,
_ => return None,
};
if let Some(tier) = tier
&& (tier.is_empty() || tier.starts_with('/') || tier.ends_with('/') || tier.contains("//"))
{
return None;
}
Some(RequestedShape {
plain: base.to_string(),
compressed,
sessions,
})
}
#[derive(Default)]
struct SlotState {
prev_emitted: Option<Traffic>,
}
#[derive(Default)]
struct SideSlots {
publisher: SlotState,
subscriber: SlotState,
seen: u64,
}
#[derive(Default)]
struct SessionSlotState {
prev_emitted: Option<Presence>,
seen: u64,
}
fn process_slot(snap: Traffic, slot_state: &mut SlotState, emit: impl FnOnce(Traffic)) {
let live = !snap.is_idle();
let prev_snap = slot_state.prev_emitted.unwrap_or_default();
let changed = snap != prev_snap;
if changed {
slot_state.prev_emitted = Some(snap);
}
if live || changed {
emit(snap);
}
}
fn process_session_slot(snap: Presence, slot_state: &mut SessionSlotState, emit: impl FnOnce(Presence)) {
let live = snap.active() > 0;
let prev_snap = slot_state.prev_emitted.unwrap_or_default();
let changed = snap != prev_snap;
if changed {
slot_state.prev_emitted = Some(snap);
}
if live || changed {
emit(snap);
}
}
fn group_key(path: &str, depth: usize) -> &str {
if depth == 0 {
return "";
}
match path.match_indices('/').nth(depth - 1) {
Some((end, _)) => &path[..end],
None => path,
}
}
fn advertised_path(prefix: &Path, group: &Path, node: Option<&str>) -> PathOwned {
let mut out = prefix.as_str().to_string();
if !group.is_empty() {
out.push('/');
out.push_str(group.as_str());
}
out.push_str("/node");
if let Some(node) = node {
out.push('/');
out.push_str(node);
}
PathOwned::from(out)
}
#[cfg(test)]
mod tests {
fn produce_origin() -> moq_net::origin::Producer {
let (producer, driver) = moq_net::origin::Producer::new(moq_net::origin::Config::default());
if tokio::runtime::Handle::try_current().is_ok() {
tokio::spawn(moq_net::time::run(driver));
} else {
std::mem::forget(driver);
}
producer
}
use std::collections::BTreeMap;
use moq_net::stats::{Registry, Tier};
use moq_net::{Timestamp, announce, broadcast, track};
use super::*;
fn test_producer(node: Option<&str>) -> (Producer, origin::Producer) {
let origin = produce_origin();
let producer = Producer::new(
Config::new()
.with_origin(origin.clone())
.with_node(node.map(|s| PathOwned::from(s.to_string()))),
);
(producer, origin)
}
#[allow(dead_code)]
struct Feed {
announced: announce::Consumer,
source: broadcast::Producer,
consumer: broadcast::Consumer,
sub: Option<track::Subscriber>,
}
async fn feed(
registry: &Registry,
tier: Tier,
path: &str,
subscribe: bool,
frames: usize,
frame_size: usize,
) -> Feed {
let ctx = registry.tier(tier).session("feed");
let origin = produce_origin();
let egress = origin.consume().with_stats(ctx);
let mut announced = egress.announced();
let source = origin.create_broadcast(path).expect("create_broadcast");
source.announce(origin::Route::default()).expect("announce");
let producer = source.create_track("video", None).expect("create_track");
let update = announced.next().await.expect("announce");
assert!(update.kind.is_active());
let consumer = egress.request_broadcast(path).await.expect("resolve");
let sub = if subscribe {
let mut sub = consumer
.track("video")
.expect("track")
.subscribe(None)
.await
.expect("subscribe");
if frames > 0 {
let mut group = producer.append_group().expect("group");
for _ in 0..frames {
group
.write_frame(Timestamp::ZERO, vec![0u8; frame_size])
.expect("write");
}
group.finish().expect("finish");
let mut group = sub.recv_group().await.expect("recv").expect("group");
while group.read_frame().await.expect("read").is_some() {}
}
Some(sub)
} else {
None
};
Feed {
announced,
source,
consumer,
sub,
}
}
async fn announced(origin: &origin::Producer) -> (String, moq_net::broadcast::Consumer) {
let mut consumer = origin.consume().with_hidden(true).announced();
tokio::time::advance(Duration::from_millis(1)).await;
let update = consumer.next().await.expect("expected announce");
assert!(update.kind.is_active());
let broadcast = origin
.consume()
.request_broadcast(moq_net::Path::new(update.prefix.as_str()))
.await
.expect("resolve");
(update.prefix.as_str().to_string(), broadcast)
}
async fn drive_tick() {
tokio::time::advance(Duration::from_millis(1100)).await;
for _ in 0..4 {
tokio::task::yield_now().await;
}
}
async fn read_frame(broadcast: &moq_net::broadcast::Consumer, name: &str) -> BTreeMap<String, Traffic> {
let mut track = subscribe(broadcast, name).await;
let frame = next_frame(&mut track).await;
serde_json::from_slice(&frame.payload).expect("json parse")
}
async fn read_last_frame(broadcast: &moq_net::broadcast::Consumer, name: &str) -> BTreeMap<String, Traffic> {
let mut track = subscribe(broadcast, name).await;
let mut last = next_frame(&mut track).await;
while let Some(frame) = try_next_frame(&mut track) {
last = frame;
}
serde_json::from_slice(&last.payload).expect("json parse")
}
async fn read_session_frame(broadcast: &moq_net::broadcast::Consumer, name: &str) -> BTreeMap<String, Presence> {
let mut track = subscribe(broadcast, name).await;
let frame = next_frame(&mut track).await;
serde_json::from_slice(&frame.payload).expect("json parse")
}
async fn subscribe(broadcast: &moq_net::broadcast::Consumer, name: &str) -> track::Ordered {
broadcast
.track(name)
.expect("track")
.subscribe(None)
.await
.expect("subscribe")
.ordered()
}
async fn next_frame(track: &mut track::Ordered) -> moq_net::frame::Frame {
let mut group = track.next_group().await.expect("ok").expect("group");
group.read_frame().await.expect("ok").expect("frame")
}
fn try_next_frame(track: &mut track::Ordered) -> Option<moq_net::frame::Frame> {
use futures::FutureExt;
let mut group = track.next_group().now_or_never()?.expect("ok")?;
group.read_frame().now_or_never()?.expect("ok")
}
#[tokio::test(start_paused = true)]
async fn new_normalizes_and_drops_empty_node() {
let (_producer, origin) = test_producer(Some("/sjc//1/"));
assert_eq!(announced(&origin).await.0, ".stats/node/sjc/1");
let (_producer, origin) = test_producer(Some("///"));
assert_eq!(announced(&origin).await.0, ".stats/node");
}
#[tokio::test(start_paused = true)]
async fn single_broadcast_path_announced() {
let (producer, origin) = test_producer(Some("sjc/1"));
let _f1 = feed(producer.registry(), Tier::default(), "foo/bar", true, 1, 8).await;
let _f2 = feed(producer.registry(), Tier::default(), "baz/qux", true, 1, 8).await;
assert_eq!(announced(&origin).await.0, ".stats/node/sjc/1");
}
#[tokio::test(start_paused = true)]
async fn task_announces_without_node_suffix() {
let (producer, origin) = test_producer(None);
let _f = feed(producer.registry(), Tier::default(), "foo/bar", true, 1, 8).await;
assert_eq!(announced(&origin).await.0, ".stats/node");
}
#[tokio::test(start_paused = true)]
async fn frame_emits_expected_counters() {
let (producer, origin) = test_producer(Some("sjc"));
let _f = feed(producer.registry(), Tier::default(), "foo/bar", true, 1, 42).await;
drive_tick().await;
let (_, broadcast) = announced(&origin).await;
let frame = read_last_frame(&broadcast, "publisher.json").await;
let snap = frame.get("foo/bar").expect("foo/bar entry");
assert_eq!(
snap.announces_started, 1,
"egress announce stream bumps announces_started"
);
assert_eq!(snap.broadcasts_started, 1, "one session subscribed");
assert_eq!(snap.subscriptions_started, 1);
assert_eq!(snap.bytes, 42);
assert_eq!(snap.frames, 1);
}
#[tokio::test(start_paused = true)]
async fn announced_bytes_surfaces_in_frame() {
let (producer, origin) = test_producer(Some("sjc"));
let _f = feed(producer.registry(), Tier::default(), "foo/bar", false, 0, 0).await;
drive_tick().await;
let (_, broadcast) = announced(&origin).await;
let frame = read_last_frame(&broadcast, "publisher.json").await;
let snap = frame.get("foo/bar").expect("foo/bar entry");
assert_eq!(snap.announces_started, 1);
assert_eq!(
snap.announced_bytes,
"foo/bar".len() as u64,
"name length recorded on announce"
);
}
#[tokio::test(start_paused = true)]
async fn announced_decouples_from_broadcasts() {
let (producer, origin) = test_producer(Some("sjc"));
let _f = feed(producer.registry(), Tier::default(), "foo/bar", false, 0, 0).await;
drive_tick().await;
let (_, broadcast) = announced(&origin).await;
let frame = read_last_frame(&broadcast, "publisher.json").await;
let snap = frame.get("foo/bar").expect("foo/bar entry");
assert_eq!(snap.announces_started, 1);
assert_eq!(snap.broadcasts_started, 0, "no subscription, no broadcasts sentinel");
assert_eq!(snap.subscriptions_started, 0);
}
#[tokio::test(start_paused = true)]
async fn short_lived_sub_is_surfaced() {
let (producer, origin) = test_producer(Some("sjc"));
{
let _f = feed(producer.registry(), Tier::default(), "foo/bar", true, 1, 123).await;
}
drive_tick().await;
let (_, broadcast) = announced(&origin).await;
let frame = read_last_frame(&broadcast, "publisher.json").await;
let snap = frame.get("foo/bar").expect("foo/bar entry");
assert_eq!(snap.subscriptions_started, 1);
assert_eq!(snap.subscriptions_ended, 1);
assert_eq!(snap.broadcasts_started, 1, "one session subscribed");
assert_eq!(snap.broadcasts_ended, 1);
assert_eq!(snap.bytes, 123);
assert_eq!(snap.frames, 1);
}
#[tokio::test(start_paused = true)]
async fn session_track_surfaces_by_root() {
let (producer, origin) = test_producer(Some("sjc"));
let _a = producer.registry().tier(Tier::default()).session("acme");
let _b = producer.registry().tier(Tier::default()).session("acme");
let _c = producer.registry().tier(Tier::new("region/sjc")).session("peer");
drive_tick().await;
let (_, broadcast) = announced(&origin).await;
let frame = read_session_frame(&broadcast, "sessions.json").await;
let snap = frame.get("acme").expect("root entry");
assert_eq!(snap.sessions_started, 2);
assert_eq!(snap.sessions_ended, 0);
assert!(
!frame.contains_key("peer"),
"regional session must not appear on the default track"
);
let snap = *read_session_frame(&broadcast, "region/sjc/sessions.json")
.await
.get("peer")
.expect("regional entry");
assert_eq!(snap.sessions_started, 1);
}
#[tokio::test(start_paused = true)]
async fn unused_slots_dont_surface() {
let (producer, origin) = test_producer(Some("sjc"));
let _f = feed(producer.registry(), Tier::default(), "foo/bar", true, 1, 8).await;
drive_tick().await;
drive_tick().await;
let (_, broadcast) = announced(&origin).await;
assert!(
read_last_frame(&broadcast, "publisher.json")
.await
.contains_key("foo/bar"),
"publisher.json must include the active foo/bar entry"
);
let frame = read_frame(&broadcast, "subscriber.json").await;
assert!(frame.is_empty(), "subscriber.json must be empty, got {frame:?}");
for name in ["publisher.json.z", "subscriber.json.z", "sessions.json.z"] {
assert!(broadcast.track(name).is_ok(), "{name} must exist");
}
let subscribing = broadcast
.track("region/sjc/publisher.json")
.expect("logical track")
.subscribe(None);
drive_tick().await;
let mut sub = subscribing.await.expect("an idle tier's track is held open").ordered();
let frame = next_frame(&mut sub).await;
let parsed: BTreeMap<String, Traffic> = serde_json::from_slice(&frame.payload).expect("json");
assert!(parsed.is_empty(), "an idle tier serves zeros, got {parsed:?}");
}
#[test]
fn advertised_path_with_and_without_node() {
let prefix = Path::new(".stats");
let empty = Path::empty();
assert_eq!(
advertised_path(&prefix, &empty, Some("sjc")).as_str(),
".stats/node/sjc"
);
assert_eq!(
advertised_path(&prefix, &empty, Some("sjc/1")).as_str(),
".stats/node/sjc/1"
);
assert_eq!(advertised_path(&prefix, &empty, None).as_str(), ".stats/node");
assert_eq!(
advertised_path(&prefix, &Path::new("acme"), Some("sjc")).as_str(),
".stats/acme/node/sjc"
);
let prefix = Path::new("metrics");
assert_eq!(
advertised_path(&prefix, &Path::new("demo/room"), Some("lon")).as_str(),
"metrics/demo/room/node/lon"
);
}
#[test]
fn group_key_uses_leading_segments() {
assert_eq!(group_key("acme/room/cam", 0), "");
assert_eq!(group_key("acme/room/cam", 1), "acme");
assert_eq!(group_key("acme/room/cam", 2), "acme/room");
assert_eq!(group_key("acme/room", 3), "acme/room");
}
#[test]
fn requested_track_shape_classifies() {
let shape = requested_track_shape("rtmp/publisher.json").expect("valid");
assert_eq!(shape.plain, "rtmp/publisher.json");
assert!(!shape.compressed);
assert!(!shape.sessions);
let shape = requested_track_shape("region/sjc/subscriber.json.z").expect("valid");
assert_eq!(shape.plain, "region/sjc/subscriber.json");
assert!(shape.compressed);
assert!(!shape.sessions);
let shape = requested_track_shape("sessions.json").expect("default tier");
assert_eq!(shape.plain, "sessions.json");
assert!(shape.sessions);
assert!(requested_track_shape("bogus.json").is_none());
assert!(requested_track_shape("xpublisher.json").is_none());
assert!(requested_track_shape("/publisher.json").is_none());
assert!(requested_track_shape("rtmp//publisher.json").is_none());
assert!(requested_track_shape("rtmp/publisher.json.z.z").is_none());
}
#[tokio::test(start_paused = true)]
async fn idle_tier_track_resolves_with_zeros() {
let (producer, origin) = test_producer(Some("sjc"));
let _f = feed(producer.registry(), Tier::default(), "foo/bar", true, 1, 42).await;
drive_tick().await;
let (_, broadcast) = announced(&origin).await;
let subscribing = broadcast.track("rtmp/publisher.json").expect("track").subscribe(None);
drive_tick().await;
let mut sub = subscribing.await.expect("held open, not rejected").ordered();
let frame = next_frame(&mut sub).await;
let parsed: BTreeMap<String, Traffic> = serde_json::from_slice(&frame.payload).expect("json");
assert!(parsed.is_empty(), "an idle tier serves zeros");
let _rtmp = feed(producer.registry(), Tier::new("rtmp"), "foo/live", true, 1, 7).await;
drive_tick().await;
let frame = next_frame(&mut sub).await;
let parsed: BTreeMap<String, Traffic> = serde_json::from_slice(&frame.payload).expect("json");
assert_eq!(parsed.get("foo/live").expect("entry").bytes, 7);
}
#[tokio::test(start_paused = true)]
async fn compressed_tier_request_creates_the_pair() {
let (producer, origin) = test_producer(Some("sjc"));
let _f = feed(producer.registry(), Tier::default(), "foo/bar", true, 1, 42).await;
drive_tick().await;
let (_, broadcast) = announced(&origin).await;
let subscribing = broadcast.track("srt/subscriber.json.z").expect("track").subscribe(None);
drive_tick().await;
subscribing.await.expect("compressed flavor held open");
subscribe(&broadcast, "srt/subscriber.json").await;
}
#[tokio::test(start_paused = true)]
async fn idle_tier_sessions_track_resolves_with_zeros() {
let (producer, origin) = test_producer(Some("sjc"));
let _f = feed(producer.registry(), Tier::default(), "foo/bar", true, 1, 42).await;
drive_tick().await;
let (_, broadcast) = announced(&origin).await;
let subscribing = broadcast.track("webrtc/sessions.json").expect("track").subscribe(None);
drive_tick().await;
let mut sub = subscribing.await.expect("held open, not rejected").ordered();
let frame = next_frame(&mut sub).await;
let parsed: BTreeMap<String, Presence> = serde_json::from_slice(&frame.payload).expect("json");
assert!(parsed.is_empty());
}
#[tokio::test(start_paused = true)]
async fn malformed_track_name_rejected() {
let (producer, origin) = test_producer(Some("sjc"));
let _f = feed(producer.registry(), Tier::default(), "foo/bar", true, 1, 42).await;
drive_tick().await;
let (_, broadcast) = announced(&origin).await;
let subscribing = broadcast.track("bogus.json").expect("track").subscribe(None);
drive_tick().await;
assert!(subscribing.await.is_err(), "a non-stats name is rejected");
}
#[tokio::test(start_paused = true)]
async fn request_racing_first_traffic_is_fulfilled() {
let (producer, origin) = test_producer(Some("sjc"));
let _f = feed(producer.registry(), Tier::default(), "foo/bar", true, 1, 42).await;
drive_tick().await;
let (_, broadcast) = announced(&origin).await;
let subscribing = broadcast.track("rtmp/publisher.json").expect("track").subscribe(None);
assert!(subscribing.poll_ok(&moq_net::kio::Waiter::noop()).is_pending());
for _ in 0..8 {
tokio::task::yield_now().await;
}
let _rtmp = feed(producer.registry(), Tier::new("rtmp"), "foo/live", true, 1, 7).await;
drive_tick().await;
let mut sub = subscribing
.await
.expect("fulfilled by the tick's own creation")
.ordered();
let frame = next_frame(&mut sub).await;
let parsed: BTreeMap<String, Traffic> = serde_json::from_slice(&frame.payload).expect("json");
assert_eq!(parsed.get("foo/live").expect("entry").bytes, 7);
}
#[tokio::test(start_paused = true)]
async fn requested_quota_recovers_after_disconnect() {
let (producer, origin) = test_producer(Some("sjc"));
let _f = feed(producer.registry(), Tier::default(), "foo/bar", true, 1, 42).await;
drive_tick().await;
let (_, broadcast) = announced(&origin).await;
let mut held = Vec::new();
for i in 0..MAX_REQUESTED_TRACKS {
let name = format!("junk{i}/publisher.json");
let subscribing = broadcast.track(&name).expect("track").subscribe(None);
drive_tick().await;
held.push(subscribing.await.expect("within the cap"));
}
let subscribing = broadcast.track("real/publisher.json").expect("track").subscribe(None);
assert!(subscribing.poll_ok(&moq_net::kio::Waiter::noop()).is_pending());
drive_tick().await;
assert!(
subscribing.poll_ok(&moq_net::kio::Waiter::noop()).is_pending(),
"an over-quota request parks instead of being rejected"
);
drop(held);
for _ in 0..4 {
tokio::task::yield_now().await;
}
tokio::time::advance(Duration::from_secs(31)).await;
for _ in 0..3 {
drive_tick().await;
}
subscribing
.await
.expect("the parked request is adopted once the quota frees");
}
#[tokio::test(start_paused = true)]
async fn parked_request_is_adopted_by_first_traffic() {
let (producer, origin) = test_producer(Some("sjc"));
let _f = feed(producer.registry(), Tier::default(), "foo/bar", true, 1, 42).await;
drive_tick().await;
let (_, broadcast) = announced(&origin).await;
let mut held = Vec::new();
for i in 0..MAX_REQUESTED_TRACKS {
let name = format!("junk{i}/publisher.json");
let subscribing = broadcast.track(&name).expect("track").subscribe(None);
drive_tick().await;
held.push(subscribing.await.expect("within the cap"));
}
let subscribing = broadcast.track("rt/publisher.json").expect("track").subscribe(None);
assert!(subscribing.poll_ok(&moq_net::kio::Waiter::noop()).is_pending());
drive_tick().await;
assert!(
subscribing.poll_ok(&moq_net::kio::Waiter::noop()).is_pending(),
"parked"
);
let _rt = feed(producer.registry(), Tier::new("rt"), "foo/live", true, 1, 9).await;
drive_tick().await;
let mut sub = subscribing.await.expect("adopted by the flush").ordered();
let frame = next_frame(&mut sub).await;
let parsed: BTreeMap<String, Traffic> = serde_json::from_slice(&frame.payload).expect("json");
assert_eq!(parsed.get("foo/live").expect("entry").bytes, 9);
}
#[test]
fn frame_serializes_like_a_btreemap() {
let traffic = |bytes| {
let mut traffic = Traffic::default();
traffic.bytes = bytes;
traffic
};
let mut frame = Frame::default();
let mut map = BTreeMap::new();
for (path, bytes) in [("room/b", 2), ("room/a", 1), ("other", 3), ("room/a/cam", 4)] {
frame.entries.push((PathOwned::from(path), traffic(bytes)));
map.insert(path.to_string(), traffic(bytes));
}
frame.entries.sort_unstable_by(|a, b| a.0.cmp(&b.0));
assert_eq!(serde_json::to_vec(&frame).unwrap(), serde_json::to_vec(&map).unwrap());
assert_eq!(serde_json::to_vec(&Frame::<Traffic>::default()).unwrap(), b"{}");
}
mod counting {
use std::alloc::{GlobalAlloc, Layout, System};
use std::cell::Cell;
thread_local! {
static ALLOCS: Cell<usize> = const { Cell::new(0) };
}
struct Counting;
unsafe impl GlobalAlloc for Counting {
unsafe fn alloc(&self, layout: Layout) -> *mut u8 {
let _ = ALLOCS.try_with(|n| n.set(n.get() + 1));
unsafe { System.alloc(layout) }
}
unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) {
unsafe { System.dealloc(ptr, layout) }
}
unsafe fn realloc(&self, ptr: *mut u8, layout: Layout, new_size: usize) -> *mut u8 {
let _ = ALLOCS.try_with(|n| n.set(n.get() + 1));
unsafe { System.realloc(ptr, layout, new_size) }
}
}
#[global_allocator]
static GLOBAL: Counting = Counting;
pub fn allocs() -> usize {
ALLOCS.with(Cell::get)
}
}
#[tokio::test(start_paused = true)]
async fn steady_drain_collects_without_allocating() {
for depth in [0, 1] {
for (broadcasts, tiers) in [(1, 1), (16, 1), (1, 4), (16, 4)] {
let registry = Registry::new(moq_net::stats::Config::new());
let mut feeds = Vec::new();
for t in 0..tiers {
let tier = match t {
0 => Tier::default(),
t => Tier::new(format!("tier{t}")),
};
for b in 0..broadcasts {
feeds.push(feed(®istry, tier.clone(), &format!("room{b}/cam"), true, 1, 8).await);
}
}
let mut drain = Drain::new(Task {
registry,
origin: produce_origin(),
prefix: PathOwned::from(".stats"),
node: None,
depth,
interval: Duration::from_secs(1),
})
.expect("drain");
for _ in 0..3 {
drain.collect();
drain.publish();
}
let before = counting::allocs();
drain.collect();
let allocs = counting::allocs() - before;
drain.publish();
let pending: usize = drain
.groups
.values()
.map(|group| group.traffic_rows.len() + group.session_rows.len())
.sum();
assert!(pending > 0, "the drain carried entries");
assert_eq!(allocs, 0, "depth {depth}, {broadcasts} broadcasts x {tiers} tiers");
}
}
}
}