Skip to main content

o_sfu_core/engine/room/media_graph/
consumer_setup.rs

1use std::mem;
2
3use o_sfu_router::{
4    MediaKind as RouterMediaKind, rtp::MediaStream as RouterRtpParameters,
5    topology::RoutedProducerId,
6};
7use tracing::warn;
8
9use super::{
10    super::{
11        outbound::{OutboundSender, VersionedRemoteTrackSnapshot},
12        state::RoomState,
13    },
14    ConsumerId, ConsumerRouteTarget, PublishedSource, SubscriptionKey,
15    route_graph::{ConsumerRouteReservation, RelayRouteKey},
16};
17use crate::engine::{
18    MediaWorkerId,
19    media_transport::{
20        ConsumerActivity, MediaTransport, ProducerActivity, SourceActivityUpdate,
21        TransportConsumerRoute, TransportMediaId, TransportRelayRouteEffect, TransportSessionKey,
22        TransportSourceActivityEffect, TransportSourceKey,
23    },
24    source_model::{PublishedSourceId, UserStreamId},
25};
26
27#[derive(Debug)]
28pub struct ConsumerSetupTarget {
29    pub session: TransportSessionKey,
30    pub source: TransportSourceKey,
31    pub source_id: PublishedSourceId,
32    pub stream: UserStreamId,
33    pub kind: RouterMediaKind,
34    pub routed: RoutedProducerId,
35}
36
37/// Consumer-route reservation awaiting transport declaration.
38///
39/// It carries graph and optional relay ownership across work performed without
40/// the room lock. Callers must commit or release it after declaration.
41#[derive(Debug)]
42#[must_use = "pending consumer setups reserve route graph state and must be committed or released"]
43pub struct PendingConsumerSetup {
44    pub(super) target: ConsumerSetupTarget,
45    pub(super) consumer: ConsumerId,
46    pub(super) reservation: ConsumerRouteReservation,
47    pub(super) sender: OutboundSender,
48    pub(super) rtp: RouterRtpParameters,
49    pub(super) relays: Vec<TransportRelayRouteEffect>,
50}
51
52pub struct DeclaredConsumerSetup {
53    pub(super) pending: PendingConsumerSetup,
54    pub(super) route: TransportConsumerRoute,
55    pub(super) mid: Option<String>,
56}
57
58pub struct CommittedConsumerSetup {
59    pub(super) target: ConsumerSetupTarget,
60    pub(super) route: TransportConsumerRoute,
61    pub(super) sender: OutboundSender,
62    pub(super) transport_activity_update: Option<bool>,
63}
64
65#[allow(
66    clippy::large_enum_variant,
67    reason = "consumer setup outcomes are returned and matched immediately so boxing the committed setup would allocate on every successful consumer setup"
68)]
69#[derive(Debug)]
70pub enum ConsumerSetupOutcome {
71    Committed {
72        target: ConsumerSetupTarget,
73        route: TransportConsumerRoute,
74        sender: OutboundSender,
75        track_snapshot: VersionedRemoteTrackSnapshot,
76        remote_source_activity: Option<TransportSourceActivityEffect>,
77        transport_activity_update: Option<bool>,
78        readiness_keyframe: Option<ConsumerRouteTarget>,
79    },
80    Released(TransportConsumerRoute, Vec<TransportRelayRouteEffect>),
81}
82
83#[derive(Debug, Clone, Copy)]
84pub enum ConsumerSetupOrigin {
85    Readiness,
86    Publish,
87    Subscribe,
88}
89
90impl ConsumerSetupOrigin {
91    pub const fn as_diagnostic_str(self) -> &'static str {
92        match self {
93            Self::Readiness => "readiness",
94            Self::Publish => "publish",
95            Self::Subscribe => "subscribe",
96        }
97    }
98}
99
100impl RoomState {
101    /// Commits a transport-declared consumer route.
102    ///
103    /// Returns [`ConsumerSetupOutcome::Released`] when receiver or source
104    /// identity changed during declaration or the router rejects the dependency.
105    pub fn commit_declared_consumer_setup(
106        &mut self,
107        setup: DeclaredConsumerSetup,
108        origin: ConsumerSetupOrigin,
109    ) -> ConsumerSetupOutcome {
110        let target = &setup.pending.target;
111        let session = &target.session;
112        if self
113            .user_for_connection(session.user_id(), session.connection_id())
114            .is_some_and(|user| user.parsed_client_rtp_capabilities.is_some())
115            && let Some((source_active, source_activity_revision)) = self
116                .topology
117                .published_source(target.source_id)
118                .filter(|source| target.matches_identity(source))
119                .map(|source| (source.active, source.activity_revision))
120        {
121            // Transport declaration ran without the room lock. Recompute from
122            // current intent and policy before committing the pending route.
123            let selection = self.setup_selection(target, source_active);
124            let delivery_active = selection.delivery_active();
125            match self.topology.commit_consumer_setup(setup, selection) {
126                Ok(commit) => {
127                    // A consumer on another media worker must be told the source
128                    // activity because that worker does not host the producer.
129                    let remote_source_activity =
130                        (commit.route.source().session_key().media_worker_id()
131                            != commit.route.consumer_session_key().media_worker_id())
132                        .then(|| TransportSourceActivityEffect {
133                            source: commit.route.source().clone(),
134                            target_media_worker_id: commit
135                                .route
136                                .consumer_session_key()
137                                .media_worker_id(),
138                            update: SourceActivityUpdate::new(
139                                ProducerActivity::from_active(source_active),
140                                source_activity_revision,
141                            ),
142                        });
143                    let track_snapshot =
144                        self.remote_track_snapshot_for_user(commit.target.session.user_id(), true);
145                    // A pending activity correction to active already requests a
146                    // keyframe so avoid a duplicate readiness refresh.
147                    let readiness_keyframe = match origin {
148                        ConsumerSetupOrigin::Readiness
149                            if delivery_active
150                                && commit.transport_activity_update != Some(true)
151                                && commit.target.kind == RouterMediaKind::Video =>
152                        {
153                            Some(commit.target.route_target(commit.route.clone()))
154                        }
155                        ConsumerSetupOrigin::Readiness
156                        | ConsumerSetupOrigin::Publish
157                        | ConsumerSetupOrigin::Subscribe => None,
158                    };
159                    ConsumerSetupOutcome::Committed {
160                        target: commit.target,
161                        route: commit.route,
162                        sender: commit.sender,
163                        track_snapshot,
164                        remote_source_activity,
165                        transport_activity_update: commit.transport_activity_update,
166                        readiness_keyframe,
167                    }
168                }
169                Err((route, relays)) => ConsumerSetupOutcome::Released(route, relays),
170            }
171        } else {
172            let DeclaredConsumerSetup { pending, route, .. } = setup;
173            ConsumerSetupOutcome::Released(route, self.topology.release_consumer_setup(pending))
174        }
175    }
176
177    pub fn release_pending_consumer_setup(
178        &mut self,
179        setup: PendingConsumerSetup,
180    ) -> Vec<TransportRelayRouteEffect> {
181        self.topology.release_consumer_setup(setup)
182    }
183}
184
185impl PendingConsumerSetup {
186    pub(in crate::engine::room) fn take_relays(&mut self) -> Vec<TransportRelayRouteEffect> {
187        mem::take(&mut self.relays)
188    }
189
190    pub(in crate::engine::room) async fn declare(
191        self,
192        media_transport: &MediaTransport,
193        origin: ConsumerSetupOrigin,
194    ) -> Result<DeclaredConsumerSetup, Self> {
195        let activity =
196            ConsumerActivity::from_active(self.reservation.selection().delivery_active());
197        match media_transport
198            .consume_media(
199                &self.target.session,
200                self.target.kind,
201                self.target.source.session_key(),
202                self.target.source.transport_media_id(),
203                &self.rtp,
204                activity,
205            )
206            .await
207        {
208            Ok(media) => {
209                let mid = media_transport
210                    .transport_media_mid(&self.target.session, media)
211                    .await;
212                Ok(DeclaredConsumerSetup {
213                    route: self.target.transport_consumer_route(media),
214                    pending: self,
215                    mid,
216                })
217            }
218            Err(error) => {
219                warn!(
220                    consumer_user_id = ?self.target.session.user_id(),
221                    consumer_connection_id = ?self.target.session.connection_id(),
222                    producer_user_id = ?self.target.source.session_key().user_id(),
223                    producer_connection_id = ?self.target.source.session_key().connection_id(),
224                    source_transport_media_id = ?self.target.source.transport_media_id(),
225                    error = ?error,
226                    consumer_mid = self.rtp.mid(),
227                    ?origin,
228                    "media transport rejected consume media declaration"
229                );
230                Err(self)
231            }
232        }
233    }
234}
235
236impl ConsumerSetupTarget {
237    pub fn new(session: TransportSessionKey, source: &PublishedSource) -> Self {
238        Self {
239            session,
240            source: source.transport.clone(),
241            source_id: source.descriptor.source_id(),
242            stream: source.descriptor.stream_id().clone(),
243            kind: source.descriptor.media_kind(),
244            routed: source.routed,
245        }
246    }
247
248    pub(super) fn subscription_key(&self) -> SubscriptionKey {
249        SubscriptionKey::new(
250            self.session.user_id(),
251            self.source.session_key().user_id(),
252            &self.stream,
253        )
254    }
255
256    pub(super) fn transport_consumer_route(
257        &self,
258        consumer_media: TransportMediaId,
259    ) -> TransportConsumerRoute {
260        TransportConsumerRoute::new(self.session.clone(), consumer_media, self.source.clone())
261    }
262
263    fn route_target(&self, route: TransportConsumerRoute) -> ConsumerRouteTarget {
264        ConsumerRouteTarget::new(route, self.stream.clone(), self.kind)
265    }
266
267    pub(super) fn relay_route_key(&self, target_worker: MediaWorkerId) -> RelayRouteKey {
268        RelayRouteKey {
269            source_user: self.source.session_key().user_id().clone(),
270            source_connection: self.source.session_key().connection_id(),
271            source_media: self.source.transport_media_id(),
272            target_worker,
273        }
274    }
275
276    pub(super) fn matches_identity(&self, source: &PublishedSource) -> bool {
277        source.descriptor.source_id() == self.source_id
278            && source.transport == self.source
279            && source.descriptor.stream_id() == &self.stream
280            && source.descriptor.media_kind() == self.kind
281            && source.routed == self.routed
282    }
283}