1use std::collections::{BTreeMap, BTreeSet};
2
3use o_sfu_router::{MediaKind as RouterMediaKind, negotiation::negotiate_consumer_rtp_parameters};
4
5use super::{
6 super::state::{ActiveUser, RoomState},
7 ConsumerId, ConsumerRouteTarget, SubscriptionKey,
8 consumer_setup::{ConsumerSetupTarget, PendingConsumerSetup},
9};
10use crate::engine::{
11 ConnectionId, UserId,
12 media_transport::{TransportMediaId, TransportRelayRouteEffect, TransportTeardown},
13 room::{
14 outbound::{OutboundSender, VersionedRemoteTrackSnapshot},
15 source_policy::VideoAdmissionRank,
16 },
17 source_model::{
18 ConsumerSourceSelection, PolicyPauseReason, PublishedSourceId, SourceRoutePriority,
19 SourceSubscriptionIntent, UserStreamId,
20 },
21};
22
23#[derive(Debug, Clone, PartialEq, Eq)]
24pub struct ReceiverRouteActivity {
25 target: ConsumerRouteTarget,
26 active: bool,
27}
28
29impl ReceiverRouteActivity {
30 pub const fn new(target: ConsumerRouteTarget, active: bool) -> Self {
31 Self { target, active }
32 }
33
34 pub const fn target(&self) -> &ConsumerRouteTarget {
35 &self.target
36 }
37
38 pub const fn active(&self) -> bool {
39 self.active
40 }
41}
42
43#[cfg(any(test, feature = "testing-transport"))]
45#[derive(Debug, Clone, Copy, PartialEq, Eq)]
46pub enum ConsumerRouteState {
47 Absent,
48 Inactive,
49 Active,
50}
51
52#[derive(Debug, Default)]
53pub struct ReceiverRouteWork {
54 pub(in crate::engine::room) activities: Vec<ReceiverRouteActivity>,
55 pub(in crate::engine::room) setups: Vec<PendingConsumerSetup>,
56 pub(in crate::engine::room) relays: Vec<TransportRelayRouteEffect>,
57 pub(in crate::engine::room) teardown: Vec<TransportTeardown>,
58 pub(in crate::engine::room) evictions: usize,
59}
60
61#[derive(Debug)]
62pub struct ReceiverRouteCommit {
63 pub(in crate::engine::room) work: ReceiverRouteWork,
64 pub(in crate::engine::room) track_snapshots:
65 Vec<(OutboundSender, VersionedRemoteTrackSnapshot)>,
66}
67
68#[derive(Clone, Copy)]
69pub(super) enum ReceiverRouteScope<'a> {
70 Source(PublishedSourceId),
71 Receiver(&'a UserId, ConnectionId),
72 SourceUser(&'a UserId, ConnectionId, &'a UserId),
73}
74
75impl RoomState {
76 pub fn apply_receiver_intent(
77 &mut self,
78 user_id: &UserId,
79 connection_id: ConnectionId,
80 target_user_id: &UserId,
81 intents: &BTreeMap<UserStreamId, SourceSubscriptionIntent>,
82 ) -> Option<ReceiverRouteCommit> {
83 self.user_for_connection(user_id, connection_id)?;
84 let work =
85 self.plan_receiver_intent_change(user_id, connection_id, target_user_id, intents);
86 Some(ReceiverRouteCommit {
87 work,
88 track_snapshots: Vec::new(),
89 })
90 }
91
92 #[cfg(test)]
93 pub fn plan_receiver_route_work(
94 &mut self,
95 user_id: &UserId,
96 connection_id: ConnectionId,
97 target_user_id: &UserId,
98 intents: &BTreeMap<UserStreamId, SourceSubscriptionIntent>,
99 ) -> ReceiverRouteWork {
100 if self.user_for_connection(user_id, connection_id).is_none() {
101 return ReceiverRouteWork::default();
102 }
103 self.plan_receiver_intent_change(user_id, connection_id, target_user_id, intents)
104 }
105
106 fn plan_receiver_intent_change(
107 &mut self,
108 user_id: &UserId,
109 connection_id: ConnectionId,
110 target_user_id: &UserId,
111 intents: &BTreeMap<UserStreamId, SourceSubscriptionIntent>,
112 ) -> ReceiverRouteWork {
113 let mut activities = Vec::new();
114 let mut relays = Vec::new();
115 let mut evictions = 0;
116 let receiver_deafened = self
117 .user_for_connection(user_id, connection_id)
118 .is_some_and(ActiveUser::is_deaf);
119 for (stream_id, intent) in intents {
120 let commit = self.topology.apply_subscription_intent(
121 &SubscriptionKey::new(user_id, target_user_id, stream_id),
122 connection_id,
123 *intent,
124 receiver_deafened,
125 self.users.contains_key(target_user_id),
126 );
127 evictions += commit.evictions;
128 relays.extend(commit.relay_effects);
129 activities.extend(commit.update);
130 }
131 let ReceiverRouteWork { setups, .. } = self.plan_missing_receiver_routes(
132 ReceiverRouteScope::SourceUser(user_id, connection_id, target_user_id),
133 );
134 ReceiverRouteWork {
135 activities,
136 setups,
137 relays,
138 evictions,
139 ..Default::default()
140 }
141 }
142
143 pub fn refresh_consumer_readiness(
144 &mut self,
145 user_id: &UserId,
146 connection_id: ConnectionId,
147 declined_consumers: &[TransportMediaId],
148 ) -> Option<ReceiverRouteCommit> {
149 let sender = {
150 let user = self.user_for_connection(user_id, connection_id)?;
151 user.parsed_client_rtp_capabilities.as_ref()?;
152 (!declined_consumers.is_empty()).then(|| user.sender.clone())
153 };
154 let session = self
155 .topology
156 .transport_user_key(user_id.clone(), connection_id);
157 let mut work =
160 self.plan_missing_receiver_routes(ReceiverRouteScope::Receiver(user_id, connection_id));
161 let (relays, teardown, detached) = self
162 .topology
163 .detach_declined_consumers(&session, declined_consumers);
164 work.relays.extend(relays);
165 work.teardown.extend(teardown);
166 let track_snapshots = sender
167 .filter(|_| detached)
168 .map(|sender| (sender, self.remote_track_snapshot_for_user(user_id, false)))
169 .into_iter()
170 .collect();
171 Some(ReceiverRouteCommit {
172 work,
173 track_snapshots,
174 })
175 }
176
177 #[cfg(test)]
178 pub fn plan_missing_consumers(
179 &mut self,
180 user_id: &UserId,
181 connection_id: ConnectionId,
182 ) -> Option<Vec<PendingConsumerSetup>> {
183 let user = self.user_for_connection(user_id, connection_id)?;
184 if user.parsed_client_rtp_capabilities.is_none() {
185 return Some(Vec::new());
186 }
187 let ReceiverRouteWork { setups, .. } =
188 self.plan_missing_receiver_routes(ReceiverRouteScope::Receiver(user_id, connection_id));
189 Some(setups)
190 }
191
192 pub(super) fn plan_missing_receiver_routes(
193 &mut self,
194 scope: ReceiverRouteScope<'_>,
195 ) -> ReceiverRouteWork {
196 let targets = self.missing_receiver_route_targets(scope);
197 ReceiverRouteWork {
198 setups: self.plan_consumers(targets),
199 ..Default::default()
200 }
201 }
202
203 fn plan_consumers(
204 &mut self,
205 mut targets: Vec<ConsumerSetupTarget>,
206 ) -> Vec<PendingConsumerSetup> {
207 let active_speakers = BTreeSet::new();
210 targets.sort_by_key(|target| self.setup_rank(target, &active_speakers));
211 targets
212 .into_iter()
213 .filter_map(|target| self.plan_consumer(target))
214 .collect()
215 }
216
217 fn setup_rank(
218 &self,
219 target: &ConsumerSetupTarget,
220 active_speakers: &BTreeSet<UserId>,
221 ) -> VideoAdmissionRank {
222 if target.kind != RouterMediaKind::Video {
223 return VideoAdmissionRank::new(
224 SourceRoutePriority::PinnedOrFeatured,
225 None,
226 target.source_id,
227 );
228 }
229 let Some(source) = self.topology.source_descriptor(target.source_id) else {
230 return VideoAdmissionRank::new(
231 SourceRoutePriority::HiddenOrOverflow,
232 None,
233 target.source_id,
234 );
235 };
236 VideoAdmissionRank::new(
237 self.receiver_video_layout_role(target.session.user_id(), source, active_speakers)
238 .priority(),
239 None,
240 target.source_id,
241 )
242 }
243
244 fn missing_receiver_route_targets(
245 &mut self,
246 scope: ReceiverRouteScope<'_>,
247 ) -> Vec<ConsumerSetupTarget> {
248 match scope {
249 ReceiverRouteScope::Source(source_id) => {
250 let users = &self.users;
251 self.topology.missing_consumer_targets_for_source(
252 source_id,
253 users
254 .iter()
255 .map(|(user, state)| (user, state.connection_id)),
256 )
257 }
258 ReceiverRouteScope::Receiver(user, connection) => self
259 .topology
260 .missing_consumer_targets(user, connection, |_| true),
261 ReceiverRouteScope::SourceUser(user, connection, source_user_id) => self
262 .topology
263 .missing_consumer_targets(user, connection, |source| {
264 source.descriptor.owner().user_id() == source_user_id
265 }),
266 }
267 }
268
269 fn plan_consumer(&mut self, target: ConsumerSetupTarget) -> Option<PendingConsumerSetup> {
270 let (sender, client_caps) = {
271 let user = self.users.get(target.session.user_id())?;
272 if user.connection_id != target.session.connection_id() {
273 return None;
274 }
275 (
276 user.sender.clone(),
277 user.parsed_client_rtp_capabilities.as_ref()?,
278 )
279 };
280 let source = self.topology.published_source(target.source_id)?;
281 if !target.matches_identity(source) {
282 return None;
283 }
284 let selection = self.setup_selection(&target, source.active);
285 let rtp = negotiate_consumer_rtp_parameters(&source.rtp, client_caps).ok()?;
286 let consumer = ConsumerId::allocate(&mut self.next_consumer_id);
287 self.topology
288 .reserve_consumer_setup(target, consumer, selection, sender, rtp)
289 }
290
291 pub(super) fn setup_selection(
292 &self,
293 target: &ConsumerSetupTarget,
294 source_active: bool,
295 ) -> ConsumerSourceSelection {
296 let key = target.subscription_key();
297 let selection = self
298 .topology
299 .consumer_source_selection(&key, target.source_id)
300 .unwrap_or_else(|| {
301 ConsumerSourceSelection::open(
302 self.topology
303 .subscription_intent(&key)
304 .active()
305 .unwrap_or(true),
306 )
307 });
308 let selection = self.apply_initial_video_download_cap(target, source_active, selection);
309 self.apply_initial_receiver_deafened(target, selection)
310 }
311
312 fn apply_initial_receiver_deafened(
313 &self,
314 target: &ConsumerSetupTarget,
315 mut selection: ConsumerSourceSelection,
316 ) -> ConsumerSourceSelection {
317 if target.kind == RouterMediaKind::Audio
318 && selection.delivery_active()
319 && self
320 .user_for_connection(target.session.user_id(), target.session.connection_id())
321 .is_some_and(ActiveUser::is_deaf)
322 {
323 selection.set_policy_pause_reason(Some(PolicyPauseReason::ReceiverDeafened));
324 }
325 selection
326 }
327
328 fn apply_initial_video_download_cap(
329 &self,
330 target: &ConsumerSetupTarget,
331 source_active: bool,
332 mut selection: ConsumerSourceSelection,
333 ) -> ConsumerSourceSelection {
334 if target.kind != RouterMediaKind::Video
335 || !source_active
336 || !selection.delivery_active()
337 || self.active_video_count(target.session.user_id())
338 < self.media_limits.max_video_downloads_per_receiver()
339 {
340 return selection;
341 }
342 selection.set_policy_pause_reason(Some(PolicyPauseReason::VideoDownloadLimit));
343 selection
344 }
345
346 fn active_video_count(&self, consumer_user_id: &UserId) -> usize {
347 let committed = self
348 .topology
349 .committed_consumer_routes_for_user(consumer_user_id)
350 .filter(|route| route.source.descriptor.media_kind() == RouterMediaKind::Video)
351 .filter(|route| route.source.active)
352 .filter(|route| route.selection.delivery_active())
353 .count();
354 let pending = self
357 .topology
358 .pending_consumer_routes_for_user(consumer_user_id)
359 .filter(|route| route.source.descriptor.media_kind() == RouterMediaKind::Video)
360 .filter(|route| route.source.active)
361 .filter(|route| route.selection.delivery_active())
362 .count();
363 committed + pending
364 }
365
366 #[cfg(any(test, feature = "testing-transport"))]
367 pub fn consumer_route_state(
368 &self,
369 consumer_user_id: &UserId,
370 producer_user_id: &UserId,
371 stream_id: &UserStreamId,
372 ) -> Option<ConsumerRouteState> {
373 self.users.get(consumer_user_id)?;
374 let key = SubscriptionKey::new(consumer_user_id, producer_user_id, stream_id);
375 let Some(route) = self.topology.committed_consumer_route_for_key(&key) else {
376 return Some(ConsumerRouteState::Absent);
377 };
378 let route_active = route.source.active && route.selection.delivery_active();
379 Some(if route_active {
380 ConsumerRouteState::Active
381 } else {
382 ConsumerRouteState::Inactive
383 })
384 }
385}