Skip to main content

o_sfu_core/engine/room/source_policy/
input.rs

1use std::{
2    cmp::Reverse,
3    collections::{BTreeMap, BTreeSet},
4};
5
6use o_sfu_router::MediaKind;
7
8use super::action::FeaturedUserUpdate;
9use crate::{
10    Bitrate, RoomMediaLimits, VideoAdaptationTuning,
11    engine::{
12        ConnectionId, UserId,
13        media_transport::{
14            ActiveSpeakerSource, ReceiverBandwidthSnapshot, TransportBitrateSnapshot,
15            TransportMediaId,
16        },
17        room::{
18            media_graph::ConsumerRouteView,
19            state::{ActiveUser, RoomState},
20        },
21    },
22};
23
24const ACTIVE_SPEAKER_FEATURED_CLEAR_LIMIT: usize = 5;
25
26#[derive(Debug)]
27pub(super) struct SourcePolicySnapshot<'a> {
28    pub(super) routes: Vec<ConsumerRouteView<'a>>,
29    pub(super) receiver_bandwidth_by_connection: BTreeMap<ConnectionId, Bitrate>,
30    pub(super) source_bitrate_by_media: BTreeMap<TransportMediaId, Bitrate>,
31    pub(super) limited_audio_media_ids: BTreeSet<TransportMediaId>,
32    pub(super) deaf_receiver_connection_ids: BTreeSet<ConnectionId>,
33    pub(super) featured_source_user_ids: BTreeSet<UserId>,
34    pub(super) active_speaker_rank_by_user: BTreeMap<UserId, usize>,
35    pub(super) featured_user_updates: Vec<FeaturedUserUpdate>,
36    pub(super) user_count: usize,
37    pub(super) media_limits: RoomMediaLimits,
38    pub(super) video_adaptation_tuning: VideoAdaptationTuning,
39    pub(super) audio_reserve_by_connection: BTreeMap<ConnectionId, Bitrate>,
40}
41
42impl<'a> SourcePolicySnapshot<'a> {
43    pub(super) fn from_state(
44        room: &'a RoomState,
45        active_speaker_sources: &[ActiveSpeakerSource],
46        receiver_bandwidth_snapshot: &ReceiverBandwidthSnapshot,
47        source_bitrate_snapshot: &TransportBitrateSnapshot,
48    ) -> Self {
49        let ranked_sources = rank_room_active_speakers(room, active_speaker_sources);
50        let media_limits = room.media_limits;
51        let tuning = room.video_adaptation_tuning;
52        let admitted_audio_speakers = admitted_audio_media_ids(
53            room,
54            &ranked_sources,
55            media_limits.max_active_audio_speakers(),
56        );
57        let deaf_receiver_connection_ids = deaf_receiver_connection_ids(room);
58        let mut featured_source_user_ids = BTreeSet::new();
59        let mut active_speaker_rank_by_user = BTreeMap::new();
60        let mut desired_featured_user_id = None;
61        for (index, user_id) in ranked_sources
62            .iter()
63            .filter_map(|source| {
64                room.topology
65                    .active_speaker_detector_owner(source.transport_media_id())
66            })
67            .enumerate()
68        {
69            if index == 0 {
70                desired_featured_user_id = Some(user_id.clone());
71            }
72            // The clear limit counts eligible sources before owner deduplication.
73            if index < ACTIVE_SPEAKER_FEATURED_CLEAR_LIMIT {
74                featured_source_user_ids.insert(user_id.clone());
75            }
76            let next_rank = active_speaker_rank_by_user.len();
77            active_speaker_rank_by_user
78                .entry(user_id)
79                .or_insert(next_rank);
80        }
81        let featured_user_updates = featured_user_updates(room, desired_featured_user_id.as_ref());
82        // Reserve the committed-route bound to avoid copying growing snapshots.
83        let mut routes = Vec::with_capacity(room.topology.consumer_count());
84        // Include policy-paused routes so later turns can resume them. Filtering
85        // on `delivery_active()` would make a policy pause self-perpetuating.
86        routes.extend(
87            room.committed_consumer_routes()
88                .filter(|route| route.source.active && route.selection.active()),
89        );
90        let audio_reserve_by_connection = audio_reserve_by_connection(
91            &routes,
92            &admitted_audio_speakers,
93            &deaf_receiver_connection_ids,
94            tuning.audio_reserve_per_speaker,
95        );
96        Self {
97            routes,
98            receiver_bandwidth_by_connection: receiver_bandwidth_by_connection(
99                receiver_bandwidth_snapshot,
100            ),
101            source_bitrate_by_media: source_bitrate_snapshot.per_media.iter().copied().collect(),
102            limited_audio_media_ids: ranked_sources
103                .iter()
104                .map(|source| source.transport_media_id())
105                .filter(|media_id| !admitted_audio_speakers.contains(media_id))
106                .collect(),
107            deaf_receiver_connection_ids,
108            featured_source_user_ids,
109            active_speaker_rank_by_user,
110            featured_user_updates,
111            user_count: room.user_count(),
112            media_limits,
113            video_adaptation_tuning: tuning,
114            audio_reserve_by_connection,
115        }
116    }
117}
118
119/// Bandwidth reserved for admitted audio before video budgeting, per receiver
120/// connection.
121///
122/// Each receiver reserves `per_speaker` for every admitted audio route it
123/// actually consumes, so a receiver that disabled audio, deafened itself (or a
124/// publisher with no consumer routes) reserves nothing and keeps its full video
125/// budget. The reserve is fixed per route, so it is deterministic and
126/// independent of policy-turn cadence. A zero per-speaker rate disables the
127/// reservation and returns an empty map.
128fn audio_reserve_by_connection(
129    routes: &[ConsumerRouteView<'_>],
130    admitted_audio_media_ids: &BTreeSet<TransportMediaId>,
131    deaf_receiver_connection_ids: &BTreeSet<ConnectionId>,
132    per_speaker: Bitrate,
133) -> BTreeMap<ConnectionId, Bitrate> {
134    if per_speaker.as_bps() == 0 {
135        return BTreeMap::new();
136    }
137    let mut reserve_by_connection = BTreeMap::new();
138    for route in routes {
139        if route.source.descriptor.media_kind() != MediaKind::Audio
140            || !admitted_audio_media_ids.contains(&route.route.source_transport_media_id())
141        {
142            continue;
143        }
144        let connection_id = route.route.consumer_session_key().connection_id();
145        if deaf_receiver_connection_ids.contains(&connection_id) {
146            continue;
147        }
148        let reserve = reserve_by_connection
149            .entry(connection_id)
150            .or_insert_with(Bitrate::zero);
151        *reserve = reserve.saturating_add(per_speaker);
152    }
153    reserve_by_connection
154}
155
156fn receiver_bandwidth_by_connection(
157    snapshot: &ReceiverBandwidthSnapshot,
158) -> BTreeMap<ConnectionId, Bitrate> {
159    snapshot
160        .per_session
161        .iter()
162        .map(|(session, estimate)| (session.connection_id(), *estimate))
163        .collect()
164}
165
166/// Filters and ranks a list of active speaker sources for a room.
167///
168/// **Ranking Criteria:**
169/// 1. **Recency:** Most recently active first (highest `observed_at`).
170/// 2. **Loudness:** Highest audio level first (`last_audio_level_dbov`).
171/// 3. **Tie-breaker:** Transport media ID.
172///
173/// Only retains sources that are active and present in the current room topology.
174fn rank_room_active_speakers(
175    room: &RoomState,
176    sources: &[ActiveSpeakerSource],
177) -> Vec<ActiveSpeakerSource> {
178    // Reserve the snapshot bound to keep filtering to at most one allocation.
179    let mut ranked = Vec::with_capacity(sources.len());
180    ranked.extend(
181        sources
182            .iter()
183            .filter(|source| {
184                room.topology
185                    .source_for_transport_media(source.transport_media_id())
186                    .is_some_and(|source| source.active)
187            })
188            .copied(),
189    );
190    ranked.sort_unstable_by_key(|source| {
191        (
192            Reverse(source.observed_at()),
193            Reverse(source.last_audio_level_dbov().unwrap_or(i8::MIN)),
194            source.transport_media_id().as_u64(),
195        )
196    });
197    ranked
198}
199
200fn user_for_source<'a>(
201    room: &'a RoomState,
202    source: &ActiveSpeakerSource,
203) -> Option<&'a ActiveUser> {
204    room.topology
205        .source_for_transport_media(source.transport_media_id())
206        .and_then(|published_source| {
207            room.users
208                .get(published_source.descriptor.owner().user_id())
209        })
210}
211
212/// Takes audio media IDs from `sources` up to the provided `limit`,
213/// prioritizing participants who are currently screen sharing.
214fn admitted_audio_media_ids(
215    room: &RoomState,
216    sources: &[ActiveSpeakerSource],
217    limit: usize,
218) -> BTreeSet<TransportMediaId> {
219    let mut admitted = BTreeSet::new();
220    let mut deferred = Vec::with_capacity(limit);
221    for source in sources {
222        if admitted.len() == limit {
223            break;
224        }
225        let media_id = source.transport_media_id();
226        // prioritize participants who are currently screen sharing
227        if user_for_source(room, source).is_some_and(ActiveUser::is_screensharing) {
228            admitted.insert(media_id);
229        } else if deferred.len() < limit - admitted.len() {
230            deferred.push(media_id);
231        }
232    }
233    admitted.extend(deferred.into_iter().take(limit - admitted.len()));
234    admitted
235}
236
237fn deaf_receiver_connection_ids(room: &RoomState) -> BTreeSet<ConnectionId> {
238    room.users
239        .values()
240        .filter(|user| user.is_deaf())
241        .map(|user| user.connection_id)
242        .collect()
243}
244
245fn featured_user_updates(
246    room: &RoomState,
247    desired_featured_user_id: Option<&UserId>,
248) -> Vec<FeaturedUserUpdate> {
249    if desired_featured_user_id.is_none()
250        && !room.users.values().any(|user| user.featured().is_some())
251    {
252        return Vec::new();
253    }
254    room.users
255        .iter()
256        .filter_map(|(user_id, user)| {
257            let current_featured = user.featured();
258            let desired_featured = match desired_featured_user_id {
259                Some(featured_user_id) => Some(featured_user_id == user_id),
260                None if current_featured.is_some() => Some(false),
261                None => None,
262            };
263            (desired_featured != current_featured).then(|| FeaturedUserUpdate {
264                user_id: user_id.clone(),
265                connection_id: user.connection_id,
266                featured: desired_featured,
267            })
268        })
269        .collect()
270}