Skip to main content

o_sfu_core/engine/room/media_graph/
subscription.rs

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/// observable receiver route state for room inspection
44#[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        // Keep declined routes committed while planning so the answer turn cannot
158        // immediately recreate the same media lines.
159        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        // Initial admission has no active-speaker snapshot. The following
208        // source-policy turn applies the current speaker ranking.
209        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        // Delivery-active reservations already claim download capacity. Counting
355        // them prevents concurrent setup from oversubscribing the receiver cap.
356        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}