use std::time::Instant;
use crate::fanout::fanout;
use crate::ids::SfuRid;
use crate::propagate::{ClientId, Propagated};
use super::Registry;
#[cfg(feature = "pacer")]
fn apply_pacer_action(
action: crate::bwe::PacerAction,
peer_id: ClientId,
) -> (Option<Propagated>, Option<bool>) {
use crate::bwe::PacerAction;
match action {
PacerAction::GoAudioOnly => (
Some(Propagated::AudioOnlyMode {
peer_id,
audio_only: true,
}),
None,
),
PacerAction::RestoreVideo => (
Some(Propagated::AudioOnlyMode {
peer_id,
audio_only: false,
}),
None,
),
PacerAction::SuspendVideo => (
Some(Propagated::SuspendVideo {
peer_id,
suspended: true,
}),
Some(true),
),
PacerAction::RestoreAudio => (
Some(Propagated::SuspendVideo {
peer_id,
suspended: false,
}),
Some(false),
),
PacerAction::ChangeLayer(_) | PacerAction::NoChange => (None, None),
}
}
#[cfg(feature = "pacer")]
const _PACER_SINGLE_DRIVER_GUARD: () = {
let poll_all_drives = cfg!(all(feature = "pacer", not(feature = "kalman-bwe")));
let update_pacer_layers_drives = cfg!(all(feature = "pacer", feature = "kalman-bwe"));
assert!(
poll_all_drives ^ update_pacer_layers_drives,
"exactly one pacer drive site must be live per feature combo (ADR-3+4)"
);
};
impl Registry {
pub fn poll_all(&mut self, now: Instant) -> Instant {
let mut deadline = now + std::time::Duration::from_millis(100);
#[cfg(all(feature = "kalman-bwe", feature = "googcc-bwe"))]
let mut last_video_timing: Option<(f64, f64)> = None;
for client in self.clients.iter_mut() {
loop {
if !client.is_alive() {
break;
}
match client.poll_output() {
Propagated::Timeout(t) => {
deadline = deadline.min(t);
break;
}
Propagated::Noop => continue,
Propagated::BandwidthEstimate {
peer_id,
ref estimate,
} => {
self.metrics.update_peer_bwe(*peer_id, estimate.bps);
self.to_propagate.push_back(Propagated::BandwidthEstimate {
peer_id,
estimate: *estimate,
});
#[cfg(all(feature = "kalman-bwe", feature = "pacer"))]
{
self.bandwidth
.record_native_estimate(peer_id, estimate.bps as f64);
}
#[cfg(all(feature = "pacer", not(feature = "kalman-bwe")))]
{
let (event, suspend) =
apply_pacer_action(client.drive_pacer(estimate.bps), peer_id);
if let Some(suspended) = suspend {
client.set_suspended(suspended);
self.metrics.inc_suspend_video(if suspended {
"enter"
} else {
"exit"
});
}
if let Some(event) = event {
self.to_propagate.push_back(event);
}
}
}
Propagated::RtcpStats { peer_id, ref stats } => {
self.metrics.update_peer_rtcp(
*peer_id,
stats.fraction_lost,
stats.rtt.as_secs_f64() * 1000.0,
stats.jitter.as_secs_f64() * 1000.0,
);
self.to_propagate.push_back(Propagated::RtcpStats {
peer_id,
stats: *stats,
});
}
other => {
#[cfg(feature = "active-speaker")]
if let Propagated::MediaData(ref origin, ref data) = other {
if let Some(raw) = data.audio_level_raw() {
if !client.is_relay() {
let level = (-(raw as i16)).clamp(0, 127) as u8;
let now_ms = self.detector_epoch.elapsed().as_millis() as u64;
self.detector.record_level(**origin, level, now_ms);
}
}
}
#[cfg(all(feature = "kalman-bwe", feature = "googcc-bwe"))]
if let Propagated::MediaData(_, ref data) = other {
if data.is_video() {
let arrival_ms =
now.saturating_duration_since(self.bwe_epoch).as_secs_f64()
* 1000.0;
last_video_timing = Some((arrival_ms, data.rtp_send_ms()));
}
}
self.to_propagate.push_back(other);
}
}
}
}
#[cfg(all(feature = "kalman-bwe", feature = "googcc-bwe"))]
if let Some((arrival_ms, send_ms)) = last_video_timing {
for client in &self.clients {
if let Some(gcc) = self.bandwidth.googcc_for_subscriber_mut(client.id) {
gcc.on_receive(arrival_ms, send_ms, 0.0);
}
}
}
deadline
}
#[cfg(feature = "active-speaker")]
#[cfg_attr(docsrs, doc(cfg(feature = "active-speaker")))]
pub fn tick_active_speaker(&mut self, now: Instant) {
let now_ms = now
.saturating_duration_since(self.detector_epoch)
.as_millis() as u64;
if let Some(change) = self.detector.tick(now_ms) {
self.metrics.inc_dominant_speaker_changes();
self.to_propagate
.push_back(Propagated::ActiveSpeakerChanged {
peer_id: change.peer_id,
confidence: change.c2_margin,
});
}
}
#[cfg(all(feature = "active-speaker", feature = "metrics-prometheus"))]
#[cfg_attr(
docsrs,
doc(cfg(all(feature = "active-speaker", feature = "metrics-prometheus")))
)]
pub fn tick_speaker_scores(&mut self) {
for (peer_id, imm, med, lng) in self.detector.peer_scores() {
self.metrics
.update_peer_speaker_scores(peer_id, imm, med, lng);
}
}
pub fn tick(&mut self, now: Instant) {
for client in self.clients.iter_mut() {
client.handle_timeout(now);
}
}
pub fn fanout_pending(&mut self) {
#[cfg(feature = "kalman-bwe")]
let now = Instant::now();
while let Some(p) = self.to_propagate.pop_front() {
#[cfg(feature = "kalman-bwe")]
if let Propagated::ClientBudgetHint(subscriber_id, bps) = &p {
self.bandwidth.record_client_hint(*subscriber_id, *bps, now);
continue;
}
#[cfg(all(feature = "kalman-bwe", feature = "pacer"))]
if let Propagated::MediaData(origin, _) = &p {
self.update_pacer_layers(*origin, now);
}
fanout(&p, &mut self.clients);
}
}
pub fn emit_publisher_layer_hints(&mut self) {
use crate::client::layer;
use std::collections::HashMap;
let mut max_per_publisher: HashMap<ClientId, SfuRid> = HashMap::new();
for subscriber in &self.clients {
let sub_desired = subscriber.desired_layer();
for track_out in &subscriber.tracks_out {
if let Some(track_in) = track_out.track_in.upgrade() {
let publisher_id = track_in.origin;
let entry = max_per_publisher.entry(publisher_id).or_insert(layer::LOW);
let rank = |r: SfuRid| -> u8 {
if r == SfuRid::LOW {
0
} else if r == SfuRid::MEDIUM {
1
} else {
2
}
};
if rank(sub_desired) > rank(*entry) {
*entry = sub_desired;
}
}
}
}
for (publisher_id, max_rid) in max_per_publisher {
let is_relay = self
.clients
.iter()
.any(|c| c.id == publisher_id && c.is_relay());
if is_relay {
self.to_propagate
.push_back(Propagated::PublisherLayerHintForUpstream {
publisher_relay_id: publisher_id,
max_rid,
});
} else {
self.to_propagate.push_back(Propagated::PublisherLayerHint {
publisher_id,
max_rid,
});
}
}
}
}
#[cfg(all(test, feature = "pacer"))]
mod tests {
use super::*;
use crate::bwe::PacerAction;
const PEER: ClientId = ClientId(7);
#[test]
fn go_audio_only_emits_audio_only_mode_true_no_suspend_change() {
let (event, suspend) = apply_pacer_action(PacerAction::GoAudioOnly, PEER);
match event {
Some(Propagated::AudioOnlyMode {
peer_id,
audio_only,
}) => {
assert_eq!(peer_id, PEER);
assert!(audio_only);
}
other => panic!("expected AudioOnlyMode, got {other:?}"),
}
assert_eq!(suspend, None);
}
#[test]
fn restore_video_emits_audio_only_mode_false_no_suspend_change() {
let (event, suspend) = apply_pacer_action(PacerAction::RestoreVideo, PEER);
match event {
Some(Propagated::AudioOnlyMode {
peer_id,
audio_only,
}) => {
assert_eq!(peer_id, PEER);
assert!(!audio_only);
}
other => panic!("expected AudioOnlyMode, got {other:?}"),
}
assert_eq!(suspend, None);
}
#[test]
fn suspend_video_emits_suspend_video_true_and_suspend_change_true() {
let (event, suspend) = apply_pacer_action(PacerAction::SuspendVideo, PEER);
match event {
Some(Propagated::SuspendVideo { peer_id, suspended }) => {
assert_eq!(peer_id, PEER);
assert!(suspended);
}
other => panic!("expected SuspendVideo, got {other:?}"),
}
assert_eq!(suspend, Some(true));
}
#[test]
fn restore_audio_emits_suspend_video_false_and_suspend_change_false() {
let (event, suspend) = apply_pacer_action(PacerAction::RestoreAudio, PEER);
match event {
Some(Propagated::SuspendVideo { peer_id, suspended }) => {
assert_eq!(peer_id, PEER);
assert!(!suspended);
}
other => panic!("expected SuspendVideo, got {other:?}"),
}
assert_eq!(suspend, Some(false));
}
#[test]
fn change_layer_and_no_change_emit_nothing() {
let (event, suspend) =
apply_pacer_action(PacerAction::ChangeLayer(crate::ids::SfuRid::MEDIUM), PEER);
assert!(event.is_none());
assert_eq!(suspend, None);
let (event, suspend) = apply_pacer_action(PacerAction::NoChange, PEER);
assert!(event.is_none());
assert_eq!(suspend, None);
}
#[test]
fn exactly_one_pacer_driver() {
let poll_all_drives = cfg!(all(feature = "pacer", not(feature = "kalman-bwe")));
let update_pacer_layers_drives = cfg!(all(feature = "pacer", feature = "kalman-bwe"));
assert!(
poll_all_drives ^ update_pacer_layers_drives,
"exactly one pacer drive site must be live (poll_all XOR update_pacer_layers)"
);
}
}
#[cfg(all(test, feature = "pacer", feature = "kalman-bwe"))]
mod min_tick_floor_tests {
use super::*;
use crate::bwe::PACER_MIN_TICK_INTERVAL;
use crate::client::test_seed::new_client;
use std::time::{Duration, Instant};
fn suspend_video_enters(evs: &[Propagated]) -> usize {
evs.iter()
.filter(|e| {
matches!(
e,
Propagated::SuspendVideo {
suspended: true,
..
}
)
})
.count()
}
#[test]
fn min_tick_floor_prevents_burst_suspend() {
let mut reg = Registry::new_for_tests();
let pub_id = ClientId(900);
let sub_id = ClientId(901);
reg.insert(new_client(pub_id));
reg.insert(new_client(sub_id));
reg.bandwidth_mut_for_tests()
.record_native_estimate(sub_id, 5_000.0);
let t0 = Instant::now();
reg.update_pacer_layers(pub_id, t0);
assert_eq!(
suspend_video_enters(®.drain_propagated_for_tests()),
0,
"first below-threshold tick must not suspend (debounce)"
);
reg.update_pacer_layers(pub_id, t0 + Duration::from_millis(10));
assert_eq!(
suspend_video_enters(®.drain_propagated_for_tests()),
0,
"burst tick inside the min-tick floor must not advance the FSM"
);
reg.update_pacer_layers(pub_id, t0 + PACER_MIN_TICK_INTERVAL);
assert_eq!(
suspend_video_enters(®.drain_propagated_for_tests()),
1,
"a tick at/after the floor must advance the FSM and suspend"
);
}
#[test]
fn first_tick_records_drive_instant() {
let mut reg = Registry::new_for_tests();
let pub_id = ClientId(910);
let sub_id = ClientId(911);
reg.insert(new_client(pub_id));
reg.insert(new_client(sub_id));
reg.bandwidth_mut_for_tests()
.record_native_estimate(sub_id, 1_000_000.0);
reg.update_pacer_layers(pub_id, Instant::now());
let recorded = reg
.clients_mut_for_tests()
.iter()
.find(|c| c.id == sub_id)
.map(|c| c.last_pacer_drive.is_some())
.expect("subscriber present");
assert!(
recorded,
"first update_pacer_layers must record a drive instant for the subscriber"
);
}
}
#[cfg(all(feature = "kalman-bwe", feature = "pacer"))]
const DEFAULT_SIMULCAST_LADDER: &[crate::ids::SfuRid] = &[
crate::ids::SfuRid::LOW,
crate::ids::SfuRid::MEDIUM,
crate::ids::SfuRid::HIGH,
];
#[cfg(all(feature = "kalman-bwe", feature = "pacer"))]
impl Registry {
pub fn update_pacer_layers(&mut self, origin: crate::propagate::ClientId, now: Instant) {
let publisher_rids: Vec<crate::ids::SfuRid> = self
.clients
.iter()
.find(|c| c.id == origin)
.map(|c| c.active_rids())
.unwrap_or_default();
let _available: &[crate::ids::SfuRid] = if publisher_rids.is_empty() {
DEFAULT_SIMULCAST_LADDER
} else {
&publisher_rids
};
for client in self.clients.iter_mut() {
if client.id == origin {
continue;
}
let sub_id = client.id;
let Some((budget, term)) = self.bandwidth.estimate_with_term(sub_id, now) else {
continue;
};
if !client.pacer_tick_ready(now, crate::bwe::PACER_MIN_TICK_INTERVAL) {
self.metrics.inc_pacer_tick_throttled();
continue;
}
#[cfg(feature = "metrics-prometheus")]
{
self.metrics.update_peer_bwe(*sub_id, budget);
self.metrics
.update_peer_binding_term(*sub_id, term.as_str());
}
let (event, suspend) = apply_pacer_action(client.drive_pacer(budget), sub_id);
if let Some(suspended) = suspend {
client.set_suspended(suspended);
self.metrics
.inc_suspend_video(if suspended { "enter" } else { "exit" });
}
if let Some(event) = event {
self.to_propagate.push_back(event);
}
}
}
}