o_sfu_core/engine/room/source_policy/
input.rs1use 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 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 let mut routes = Vec::with_capacity(room.topology.consumer_count());
84 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
119fn 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
166fn rank_room_active_speakers(
175 room: &RoomState,
176 sources: &[ActiveSpeakerSource],
177) -> Vec<ActiveSpeakerSource> {
178 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
212fn 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 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}