Skip to main content

o_sfu_core/engine/room/media_graph/
producer.rs

1use itertools::Itertools;
2use o_sfu_rfc::rtp::{Mid, Rid, Ssrc};
3use o_sfu_router::{
4    RouterError,
5    rtp::{MediaFormat, MediaStream as RouterRtpParameters},
6};
7use tracing::{error, warn};
8
9use super::{
10    super::{
11        outbound::{OutboundSender, VersionedRemoteTrackSnapshot},
12        state::{PresenceCommit, RoomState},
13    },
14    ReceiverRouteWork,
15    source_index::PublishedSources,
16    subscription::ReceiverRouteScope,
17};
18use crate::{
19    Bitrate,
20    engine::{
21        ConnectionId, UserId, UserInfo,
22        media_transport::{
23            ProducerActivity, SessionUploadEncoding, SourceActivityUpdate, TransportMediaId,
24            TransportSessionKey, TransportSourceActivityEffect, TransportSourceKey,
25        },
26        source_model::{
27            PublishedSourceDescriptor, PublishedSourceDescriptorParts, PublishedSourceId,
28            PublishedSourceOwner, SourceEncodingDescriptor, SourceEncodingDescriptorParts,
29            SourceModelError, SourcePublishIntent, UploadLayerPolicyRole, UserStreamId,
30        },
31    },
32};
33
34#[derive(Debug, Clone)]
35pub struct ValidatedPublish {
36    pub session_key: TransportSessionKey,
37    pub intent: SourcePublishIntent,
38}
39
40#[derive(Debug)]
41pub struct PublishCommit {
42    pub receiver_route_work: ReceiverRouteWork,
43    pub presence: Option<PresenceCommit>,
44}
45
46#[derive(Debug)]
47pub enum PublishIntentPlan {
48    Activate(ProducerActivityCommit),
49    Noop,
50    Queue,
51    Stage(ValidatedPublish),
52}
53
54#[derive(Debug)]
55pub struct ProducerActivityCommit {
56    pub source: TransportSourceKey,
57    pub stream_id: UserStreamId,
58    pub update: SourceActivityUpdate,
59    pub remote_activity_effects: Vec<TransportSourceActivityEffect>,
60    pub track_snapshots: Vec<(OutboundSender, VersionedRemoteTrackSnapshot)>,
61    pub presence: Option<PresenceCommit>,
62}
63
64#[derive(Debug)]
65pub enum ProducerActivityRejection {
66    MissingPublication,
67    StalePublication,
68}
69
70#[derive(Debug, thiserror::Error)]
71pub(in crate::engine::room) enum PublicationCommitError {
72    #[error(transparent)]
73    Source(#[from] SourceModelError),
74    #[error(transparent)]
75    Router(#[from] RouterError),
76}
77
78impl RoomState {
79    pub fn apply_publish_intent(
80        &mut self,
81        user_id: &UserId,
82        publisher_connection_id: ConnectionId,
83        intent: &SourcePublishIntent,
84        can_stage: bool,
85    ) -> PublishIntentPlan {
86        let Some(user) = self.users.get(user_id) else {
87            warn!(
88                ?user_id,
89                publisher_connection_id = ?publisher_connection_id,
90                stream_id = %intent.stream_id(),
91                "cannot start publish because the user is missing from room state"
92            );
93            return PublishIntentPlan::Noop;
94        };
95        if user.connection_id != publisher_connection_id {
96            warn!(
97                ?user_id,
98                publisher_connection_id = ?publisher_connection_id,
99                current_connection_id = ?user.connection_id,
100                stream_id = %intent.stream_id(),
101                "cannot start publish because the connection is stale"
102            );
103            return PublishIntentPlan::Noop;
104        }
105        if self
106            .published_source_id(user_id, publisher_connection_id, intent.stream_id())
107            .is_some()
108        {
109            return self
110                .apply_publication_activity(
111                    user_id,
112                    publisher_connection_id,
113                    intent.stream_id(),
114                    true,
115                    intent.presence(),
116                )
117                .map_or_else(|_| PublishIntentPlan::Noop, PublishIntentPlan::Activate);
118        }
119        if !can_stage {
120            return PublishIntentPlan::Queue;
121        }
122        self.validate_publish(user_id, publisher_connection_id, intent)
123            .map_or(PublishIntentPlan::Noop, PublishIntentPlan::Stage)
124    }
125
126    pub fn validate_publish(
127        &self,
128        user_id: &UserId,
129        publisher_connection_id: ConnectionId,
130        intent: &SourcePublishIntent,
131    ) -> Option<ValidatedPublish> {
132        let Some(user) = self.users.get(user_id) else {
133            warn!(
134                ?user_id,
135                publisher_connection_id = ?publisher_connection_id,
136                stream_id = %intent.stream_id(),
137                "cannot prepare negotiated publish because the user is missing from room state"
138            );
139            return None;
140        };
141        if user.connection_id != publisher_connection_id {
142            warn!(
143                ?user_id,
144                publisher_connection_id = ?publisher_connection_id,
145                current_connection_id = ?user.connection_id,
146                stream_id = %intent.stream_id(),
147                "cannot prepare negotiated publish because the connection is stale"
148            );
149            return None;
150        }
151        if user.parsed_client_rtp_capabilities.is_none() {
152            warn!(
153                ?user_id,
154                publisher_connection_id = ?publisher_connection_id,
155                stream_id = %intent.stream_id(),
156                "cannot prepare negotiated publish because the user is not publish-ready"
157            );
158            return None;
159        }
160        Some(ValidatedPublish {
161            session_key: self.transport_user_key(user_id, publisher_connection_id),
162            intent: intent.clone(),
163        })
164    }
165
166    pub fn commit_publish_reservation(
167        &mut self,
168        publish: ValidatedPublish,
169        consumable_rtp_parameters: RouterRtpParameters,
170        upload_encodings: &[SessionUploadEncoding],
171        transport_media_id: TransportMediaId,
172    ) -> Option<PublishCommit> {
173        self.validate_publish_commit(&publish, transport_media_id)?;
174        let owner_user_id = publish.session_key.user_id().clone();
175        let owner_connection_id = publish.session_key.connection_id();
176        let stream_id = publish.intent.stream_id().clone();
177        let presence = publish.intent.presence().cloned();
178        let source_id = match self.topology.commit_publication(
179            publish,
180            consumable_rtp_parameters,
181            upload_encodings,
182            transport_media_id,
183        ) {
184            Ok(source_id) => source_id,
185            Err(error) => {
186                error!(
187                    user_id = ?owner_user_id,
188                    ?owner_connection_id,
189                    stream_id = %stream_id,
190                    ?transport_media_id,
191                    ?error,
192                    "failed to commit negotiated publication"
193                );
194                return None;
195            }
196        };
197        let receiver_route_work =
198            self.plan_missing_receiver_routes(ReceiverRouteScope::Source(source_id));
199        let presence = presence.and_then(|info| {
200            self.apply_presence_update(&owner_user_id, owner_connection_id, &info)
201        });
202        Some(PublishCommit {
203            receiver_route_work,
204            presence,
205        })
206    }
207
208    pub(in crate::engine::room) fn validate_publish_commit(
209        &self,
210        publish: &ValidatedPublish,
211        transport_media_id: TransportMediaId,
212    ) -> Option<()> {
213        let user_id = publish.session_key.user_id();
214        let connection_id = publish.session_key.connection_id();
215        let Some(user) = self.users.get(user_id) else {
216            warn!(
217                ?user_id,
218                ?connection_id,
219                stream_id = %publish.intent.stream_id(),
220                ?transport_media_id,
221                "cannot commit negotiated publish because the user is missing from room state"
222            );
223            return None;
224        };
225        let publish_ready = user.parsed_client_rtp_capabilities.is_some();
226        if user.connection_id != connection_id || !publish_ready {
227            warn!(
228                ?user_id,
229                ?connection_id,
230                current_connection_id = ?user.connection_id,
231                publish_ready,
232                stream_id = %publish.intent.stream_id(),
233                ?transport_media_id,
234                "cannot commit negotiated publish because the user state changed before commit"
235            );
236            return None;
237        }
238        if self
239            .topology
240            .source_id_for_owner_stream(user_id, publish.intent.stream_id())
241            .is_some()
242        {
243            warn!(
244                ?user_id,
245                ?connection_id,
246                stream_id = %publish.intent.stream_id(),
247                ?transport_media_id,
248                "cannot commit negotiated publish because a source already exists for this stream"
249            );
250            return None;
251        }
252        Some(())
253    }
254
255    #[must_use]
256    pub fn producer_stream_id_for_transport_media_id(
257        &self,
258        transport_media_id: TransportMediaId,
259    ) -> Option<UserStreamId> {
260        self.topology
261            .source_for_transport_media(transport_media_id)
262            .map(|source| source.descriptor.stream_id().clone())
263    }
264
265    #[must_use]
266    pub fn published_source_id(
267        &self,
268        owner: &UserId,
269        connection: ConnectionId,
270        stream: &UserStreamId,
271    ) -> Option<PublishedSourceId> {
272        self.topology.published_source_id(owner, connection, stream)
273    }
274
275    #[cfg(any(test, feature = "testing-transport"))]
276    pub fn published_source_id_for_user(
277        &self,
278        user: &UserId,
279        stream: &UserStreamId,
280    ) -> Option<PublishedSourceId> {
281        self.published_source_id(user, self.user_connection_id(user)?, stream)
282    }
283
284    pub fn apply_publication_activity(
285        &mut self,
286        user_id: &UserId,
287        connection_id: ConnectionId,
288        stream_id: &UserStreamId,
289        active: bool,
290        presence: Option<&UserInfo>,
291    ) -> Result<ProducerActivityCommit, ProducerActivityRejection> {
292        let source_id = self
293            .published_source_id(user_id, connection_id, stream_id)
294            .ok_or(ProducerActivityRejection::MissingPublication)?;
295        let revision = self
296            .topology
297            .set_published_source_activity(source_id, connection_id, active)
298            .ok_or(ProducerActivityRejection::StalePublication)?;
299        let source_recipients = self
300            .topology
301            .committed_consumer_user_ids_for_source(source_id);
302        let source = self
303            .topology
304            .published_source(source_id)
305            .ok_or(ProducerActivityRejection::StalePublication)?
306            .transport
307            .clone();
308        let update = SourceActivityUpdate::new(ProducerActivity::from_active(active), revision);
309        let remote_activity_effects = self.topology.source_activity_effects(&source, update);
310        let presence =
311            presence.and_then(|info| self.apply_presence_update(user_id, connection_id, info));
312        Ok(ProducerActivityCommit {
313            source,
314            stream_id: stream_id.clone(),
315            update,
316            remote_activity_effects,
317            track_snapshots: self.remote_track_snapshots_for_users(source_recipients, false),
318            presence,
319        })
320    }
321}
322
323pub(super) fn allocate_source_descriptor(
324    sources: &mut PublishedSources,
325    publish: &ValidatedPublish,
326    consumable_rtp_parameters: &RouterRtpParameters,
327    upload_encodings: &[SessionUploadEncoding],
328) -> Result<PublishedSourceDescriptor, SourceModelError> {
329    let source_id = sources.allocate_id();
330    let encodings = consumable_rtp_parameters
331        .bindings()
332        .map(|binding| {
333            let upload_profile = upload_profile_for_rid(upload_encodings, binding.rid());
334            SourceEncodingDescriptor::new(SourceEncodingDescriptorParts {
335                encoding_id: sources.allocate_encoding_id(),
336                source_id,
337                rid: binding.rid().map(Rid::new),
338                primary_ssrc: binding.ssrc().map(Ssrc::new),
339                repair_ssrc: None,
340                max_bitrate: binding
341                    .max_bitrate()
342                    .map(Bitrate::from_bps)
343                    .or_else(|| upload_profile.and_then(|(_, encoding)| encoding.max_bitrate)),
344                resolution_scale: upload_profile
345                    .and_then(|(_, encoding)| encoding.resolution_scale),
346                max_framerate: upload_profile.and_then(|(_, encoding)| encoding.max_framerate),
347                policy_role: upload_profile.map(|(rank, _)| {
348                    upload_layer_policy_role_for_rank(rank, upload_encodings.len())
349                }),
350                negotiated_format: negotiated_format_for_binding(
351                    consumable_rtp_parameters,
352                    binding.payload_type(),
353                ),
354            })
355        })
356        .collect::<Vec<_>>();
357    PublishedSourceDescriptor::new(PublishedSourceDescriptorParts {
358        source_id,
359        owner: PublishedSourceOwner::new(publish.session_key.user_id().clone()),
360        stream_id: publish.intent.stream_id().clone(),
361        media_kind: publish.intent.media_kind(),
362        policy: publish.intent.policy(),
363        mid: consumable_rtp_parameters.mid().map(Mid::new),
364        encodings,
365    })
366}
367
368fn upload_profile_for_rid<'a>(
369    upload_encodings: &'a [SessionUploadEncoding],
370    rid: Option<&str>,
371) -> Option<(usize, &'a SessionUploadEncoding)> {
372    let rid = rid?;
373    upload_encodings
374        .iter()
375        .enumerate()
376        .find(|(_rank, encoding)| encoding.rid == rid)
377}
378
379const fn upload_layer_policy_role_for_rank(
380    rank: usize,
381    layer_count: usize,
382) -> UploadLayerPolicyRole {
383    if layer_count > 1 && rank + 1 == layer_count {
384        UploadLayerPolicyRole::Featured
385    } else if rank == 0 && layer_count > 2 {
386        UploadLayerPolicyRole::DegradedThumbnail
387    } else {
388        UploadLayerPolicyRole::Thumbnail
389    }
390}
391
392fn negotiated_format_for_binding(
393    parameters: &RouterRtpParameters,
394    payload_type: Option<u8>,
395) -> Option<MediaFormat> {
396    if let Some(payload_type) = payload_type
397        && let Some(format) = parameters
398            .formats()
399            .find(|format| format.payload_type() == payload_type)
400    {
401        return Some(format.clone());
402    }
403    parameters
404        .formats()
405        .find_or_first(|format| !format.codec().is_rtx())
406        .cloned()
407}