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}