Skip to main content

o_sfu_core/engine/media_transport/
mod.rs

1//! Runtime media transport boundary used by core.
2//!
3//! This module contains the service surface that room, server and
4//! [`crate::prelude::SfuCore`] code use when media work must happen outside the
5//! pure router. Callers ask the media transport to create SDP offers, apply
6//! answers, publish or consume media, close sessions, read diagnostics
7//! snapshots and subscribe to source-policy wakeups.
8//!
9//! The boundary exposes the opaque [`MediaTransport`] handle, narrow
10//! construction inputs for the server runtime and transport operations that let
11//! higher layers express intent without knowing about RTC workers, Str0m state,
12//! UDP sockets or worker-local relay routing
13//!
14//! Code above this module should depend on [`MediaTransport`]. Code below this
15//! module, especially the RTC engine, may deal with worker-local state machines
16//! and packet-loop details. Keeping that split explicit prevents room and
17//! signaling code from growing knowledge of the concrete WebRTC
18//! implementation.
19
20mod build;
21mod config;
22mod policy_invalidation;
23mod route_control;
24mod rtc;
25mod teardown;
26#[cfg(any(test, feature = "testing-transport", feature = "internal-benchmarks"))]
27#[path = "TESTS/test_support/mod.rs"]
28pub mod test_support;
29mod types;
30mod workers;
31
32use std::sync::Arc;
33#[cfg(any(test, feature = "testing-transport"))]
34use std::sync::atomic::AtomicUsize;
35
36pub use build::MediaTransportBuildError;
37pub use config::{MediaTransportConfig, MediaTransportDeps};
38use o_sfu_router::{
39    MediaKind,
40    rtp::{MediaCapabilities, MediaStream as RouterRtpParameters},
41};
42pub use policy_invalidation::{SourcePolicySignal, SourcePolicyUpdateSubscription};
43pub(crate) use route_control::{
44    ConsumerRouteControl, ConsumerRouteControlOutcome, MediaControlPlan,
45    TransportSourceActivityEffect,
46};
47pub(in crate::engine::media_transport) use route_control::{
48    ConsumerRouteControlFailure, ProducerRouteControl,
49};
50pub(crate) use teardown::TransportTeardown;
51#[cfg(feature = "internal-benchmarks")]
52pub mod benchmark_support {
53    pub use super::rtc::benchmark_support::*;
54}
55#[cfg(any(test, fuzzing))]
56pub mod fuzz_support {
57    pub use super::rtc::{
58        client_rtp_capabilities_from_answer, fuzz_support::route_packet_loop_ingress_demux,
59    };
60}
61use rtc::{
62    ParsedSessionAnswer, RtcSessionOffer, RtcWorker, RtcWorkerCommand, RtcWorkerResponse,
63    RtpProfile,
64};
65use tracing::warn;
66pub use types::{
67    ActiveSpeakerActivityReason, ActiveSpeakerActivityState, ActiveSpeakerSource,
68    ActiveSpeakerSourceDiagnostic, AppliedProducer, AppliedSessionAnswer, ConsumerActivity,
69    ProducerActivity, ReceiverBandwidthSnapshot, ReceiverBweTargetUpdate, RelayRouteActivity,
70    SessionOffer, SessionUploadEncoding, SessionUploadSlot, SourcePacketGate,
71    TransportAdapterError, TransportBitrateSnapshot, TransportConsumerRoute,
72    TransportHealthSnapshot, TransportMediaId, TransportQualitySample, TransportQualitySnapshot,
73    TransportRelayRouteAction, TransportRelayRouteEffect, TransportResult, TransportRidActivity,
74    TransportSessionHealth, TransportSessionKey, TransportSourceActivity,
75    TransportSourceDiagnosticsSnapshot, TransportSourceKey, TransportWorkerPressureSnapshot,
76};
77pub(crate) use types::{SourceActivityRevision, SourceActivityUpdate};
78
79use self::workers::signaling_to_str0m_media_kind;
80use crate::engine::metrics::RuntimeMetrics;
81
82/// Opaque runtime media transport handle.
83///
84/// [`MediaTransport`] is the handle server code and [`crate::prelude::SfuCore`]
85/// should hold.
86/// It hides the production RTC worker topology. Callers express intent through
87/// inherent methods and must not branch on concrete worker internals.
88///
89/// The handle also centralizes warning logs for failed transport effects. Inner
90/// backends return typed errors, while this boundary adds stable diagnostic
91/// context such as session keys, media ids and SDP lengths.
92#[derive(Debug, Clone)]
93pub struct MediaTransport {
94    workers: Arc<[RtcWorker]>,
95    profile: Arc<RtpProfile>,
96    metrics: Arc<RuntimeMetrics>,
97    #[cfg(test)]
98    media_control_batches: test_support::MediaControlBatchLog,
99    #[cfg(any(test, feature = "testing-transport"))]
100    source_diagnostics_requests: Arc<AtomicUsize>,
101    /// Coalesces transport observations for room policy.
102    source_policy_signal: SourcePolicySignal,
103}
104
105#[derive(Debug, Clone, Copy)]
106enum TransportCommandOp {
107    CreateInitialSessionOffer,
108    CreateSessionRenegotiationOffer,
109    ApplySessionAnswer,
110    PublishMedia,
111    ConsumeMedia,
112}
113
114fn warn_session_command_failed(
115    session_key: &TransportSessionKey,
116    op: TransportCommandOp,
117    error: TransportAdapterError,
118) {
119    warn!(
120        ?session_key,
121        ?op,
122        ?error,
123        "media transport worker command failed"
124    );
125}
126
127impl MediaTransport {
128    /// Cancels every RTC worker without waiting for termination.
129    pub fn cancel(&self) {
130        for worker in self.workers.iter() {
131            worker.cancel();
132        }
133    }
134
135    /// Cancels every RTC worker and waits for its thread to terminate.
136    pub async fn shutdown(&self) {
137        self.cancel();
138        for worker in self.workers.iter() {
139            worker.wait_for_shutdown().await;
140        }
141    }
142
143    /// returns the router capability snapshot compiled from the RTC wire profile
144    #[must_use]
145    pub fn router_rtp_capabilities(&self) -> MediaCapabilities {
146        self.profile.router_capabilities()
147    }
148
149    async fn request_session_command<T>(
150        &self,
151        session_key: &TransportSessionKey,
152        build: impl FnOnce(RtcWorkerResponse<T>) -> RtcWorkerCommand,
153        log_error: impl FnOnce(TransportAdapterError),
154    ) -> Result<T, TransportAdapterError> {
155        let result = match self.require_worker_for_user(session_key) {
156            Ok(worker) => worker.request_worker(build).await,
157            Err(error) => Err(error),
158        };
159        if let Err(error) = &result {
160            log_error(*error);
161        }
162        result
163    }
164
165    /// Creates the first SDP offer for a transport session.
166    ///
167    /// The session must already be assigned to a transport worker by its
168    /// [`TransportSessionKey`]. The returned offer is transport state that
169    /// callers should send to the browser unchanged.
170    ///
171    /// # Errors
172    ///
173    /// Returns [`TransportAdapterError`] when the session cannot be addressed or
174    /// the active backend cannot create the offer.
175    pub async fn create_initial_session_offer(
176        &self,
177        room_id: &str,
178        session_key: &TransportSessionKey,
179    ) -> Result<SessionOffer, TransportAdapterError> {
180        let room_id: Arc<str> = Arc::from(room_id);
181        self.request_session_command(
182            session_key,
183            move |response| RtcWorkerCommand::CreateInitialSessionOffer {
184                room_id,
185                session_key: session_key.clone(),
186                response,
187            },
188            |error| {
189                warn_session_command_failed(
190                    session_key,
191                    TransportCommandOp::CreateInitialSessionOffer,
192                    error,
193                );
194            },
195        )
196        .await
197        .map(RtcSessionOffer::into_session_offer)
198    }
199
200    /// Creates a new SDP offer after transport media state changed.
201    ///
202    /// Backends reject sessions that cannot renegotiate in their current state.
203    /// Callers should treat the returned offer as replacing any older pending
204    /// transport offer for the same session.
205    ///
206    /// # Errors
207    ///
208    /// Returns [`TransportAdapterError`] when the session cannot be addressed or
209    /// the active backend cannot create a renegotiation offer.
210    pub async fn create_session_renegotiation_offer(
211        &self,
212        session_key: &TransportSessionKey,
213    ) -> Result<SessionOffer, TransportAdapterError> {
214        self.request_session_command(
215            session_key,
216            |response| RtcWorkerCommand::CreateSessionRenegotiationOffer {
217                session_key: session_key.clone(),
218                response,
219            },
220            |error| {
221                warn_session_command_failed(
222                    session_key,
223                    TransportCommandOp::CreateSessionRenegotiationOffer,
224                    error,
225                );
226            },
227        )
228        .await
229        .map(RtcSessionOffer::into_session_offer)
230    }
231
232    /// Applies a browser SDP answer to a pending transport offer.
233    ///
234    /// The answer can reveal transport-derived producer facts such as mapped
235    /// RTP parameters. Those facts are returned so room code can commit staged
236    /// media using values observed by the transport.
237    ///
238    /// # Errors
239    ///
240    /// Returns [`TransportAdapterError`] when the session cannot be addressed,
241    /// the answer is invalid or the active backend cannot apply it.
242    pub async fn apply_session_answer(
243        &self,
244        session_key: &TransportSessionKey,
245        answer_sdp: &str,
246    ) -> Result<AppliedSessionAnswer, TransportAdapterError> {
247        async {
248            let worker = self.require_worker_for_user(session_key)?;
249            let answer = ParsedSessionAnswer::parse(answer_sdp)?;
250            worker
251                .request_worker(|response| RtcWorkerCommand::ApplySessionAnswer {
252                    session_key: session_key.clone(),
253                    answer,
254                    response,
255                })
256                .await
257        }
258        .await
259        .inspect_err(|error| {
260            warn!(
261                ?session_key,
262                op = ?TransportCommandOp::ApplySessionAnswer,
263                answer_len = answer_sdp.len(),
264                ?error,
265                "media transport failed to apply session answer"
266            );
267        })
268    }
269
270    /// Declares a new producer on a transport session.
271    ///
272    /// `rtp_parameters` must come from router media state accepted by the core.
273    /// The returned [`TransportMediaId`] addresses the backend-local producer
274    /// for later route, activity and cleanup operations.
275    ///
276    /// # Errors
277    ///
278    /// Returns [`TransportAdapterError`] when the session cannot be addressed,
279    /// the RTP parameters are invalid or the active backend cannot create the
280    /// producer.
281    pub async fn publish_media(
282        &self,
283        session_key: &TransportSessionKey,
284        media_kind: MediaKind,
285        rtp_parameters: &RouterRtpParameters,
286    ) -> Result<TransportMediaId, TransportAdapterError> {
287        self.request_session_command(
288            session_key,
289            |response| RtcWorkerCommand::AddRecvMedia {
290                session_key: session_key.clone(),
291                media_kind: signaling_to_str0m_media_kind(media_kind),
292                rtp_parameters: rtp_parameters.clone(),
293                response,
294            },
295            |error| {
296                warn!(
297                    ?session_key,
298                    op = ?TransportCommandOp::PublishMedia,
299                    ?media_kind,
300                    mid = rtp_parameters.mid(),
301                    ?error,
302                    "media transport worker command failed"
303                );
304            },
305        )
306        .await
307    }
308
309    /// Declares a new consumer route from a source session to a consumer session.
310    ///
311    /// The source and consumer sessions must belong to the same room instance.
312    /// Cross-worker routing is an implementation detail hidden behind this
313    /// method. The initial activity controls whether the packet loop can
314    /// forward packets before later room policy updates arrive.
315    ///
316    /// # Errors
317    ///
318    /// Returns [`TransportAdapterError`] when either session cannot be
319    /// addressed, the source media id is unknown or the active backend cannot
320    /// create the consumer.
321    pub async fn consume_media(
322        &self,
323        consumer_session_key: &TransportSessionKey,
324        media_kind: MediaKind,
325        source_session_key: &TransportSessionKey,
326        source_media_id: TransportMediaId,
327        consumer_rtp_parameters: &RouterRtpParameters,
328        initial_activity: ConsumerActivity,
329    ) -> Result<TransportMediaId, TransportAdapterError> {
330        let result = async {
331            Self::ensure_same_room(consumer_session_key, source_session_key)?;
332            let relay_route =
333                self.relay_registration_workers(consumer_session_key, source_session_key)?;
334            let remote_source_control =
335                relay_route
336                    .as_ref()
337                    .map(|(source_worker, consumer_worker)| {
338                        source_worker.remote_source_control(consumer_worker)
339                    });
340            self.require_worker_for_user(consumer_session_key)?
341                .request_worker(|response| RtcWorkerCommand::AddSendMedia {
342                    consumer_key: consumer_session_key.clone(),
343                    media_kind: signaling_to_str0m_media_kind(media_kind),
344                    source: TransportSourceKey::new(source_session_key.clone(), source_media_id),
345                    remote_source_control,
346                    consumer_rtp_parameters: consumer_rtp_parameters.clone(),
347                    active: initial_activity.is_active(),
348                    response,
349                })
350                .await
351        }
352        .await;
353        if let Err(error) = &result {
354            warn!(
355                ?consumer_session_key,
356                ?source_session_key,
357                ?source_media_id,
358                op = ?TransportCommandOp::ConsumeMedia,
359                ?media_kind,
360                mid = consumer_rtp_parameters.mid(),
361                initial_active = initial_activity.is_active(),
362                ?error,
363                "media transport worker command failed"
364            );
365        }
366        result
367    }
368
369    /// Applies a room relay route mutation.
370    ///
371    /// The transport may install, release or gate the packet-loop relay target.
372    /// Room state remains the lifecycle owner.
373    ///
374    /// # Errors
375    ///
376    /// Returns [`TransportAdapterError`] when the referenced route cannot be
377    /// addressed or the active backend cannot apply the relay mutation.
378    pub(crate) async fn apply_relay_route_effect(
379        &self,
380        effect: &TransportRelayRouteEffect,
381    ) -> Result<(), TransportAdapterError> {
382        self.execute_relay_route_effect(effect)
383            .await
384            .inspect_err(|error| {
385                warn!(
386                source = ?effect.source,
387                target_media_worker_id = effect.target_media_worker_id.as_usize(),
388                action = ?effect.action,
389                ?error,
390                "media transport failed to apply relay route effect"
391                );
392            })
393    }
394
395    /// Applies one activity revision to a remote source.
396    ///
397    /// # Errors
398    ///
399    /// Returns [`TransportAdapterError`] when the target worker is unavailable
400    /// or rejects the update.
401    pub(crate) async fn apply_remote_source_activity_effect(
402        &self,
403        effect: &TransportSourceActivityEffect,
404    ) -> Result<(), TransportAdapterError> {
405        self.execute_remote_source_activity_effect(effect)
406            .await
407            .inspect_err(|error| {
408                warn!(
409                    source = ?effect.source,
410                    target_media_worker_id = effect.target_media_worker_id.as_usize(),
411                    update = ?effect.update,
412                    ?error,
413                    "media transport failed to apply remote source activity"
414                );
415            })
416    }
417
418    /// Subscribes to source-policy invalidation signals emitted by transport
419    /// workers.
420    #[must_use]
421    pub fn source_policy_subscription(&self) -> SourcePolicyUpdateSubscription {
422        self.source_policy_signal.subscribe()
423    }
424}
425
426#[cfg(test)]
427#[allow(non_snake_case, reason = "test modules map to local TESTS directories")]
428mod TESTS;