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
32#[cfg(any(test, feature = "testing-transport"))]
33use std::sync::atomic::AtomicUsize;
34use std::{ptr, sync::Arc};
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
79pub(crate) use self::workers::WorkerPlacementState;
80use self::workers::signaling_to_str0m_media_kind;
81use crate::engine::metrics::RuntimeMetrics;
82
83/// Opaque runtime media transport handle.
84///
85/// [`MediaTransport`] is the handle server code and [`crate::prelude::SfuCore`]
86/// should hold.
87/// It hides the production RTC worker topology. Callers express intent through
88/// inherent methods and must not branch on concrete worker internals.
89///
90/// The handle also centralizes warning logs for failed transport effects. Inner
91/// backends return typed errors, while this boundary adds stable diagnostic
92/// context such as session keys, media ids and SDP lengths.
93#[derive(Debug, Clone)]
94pub struct MediaTransport {
95    workers: Arc<[RtcWorker]>,
96    profile: Arc<RtpProfile>,
97    metrics: Arc<RuntimeMetrics>,
98    #[cfg(test)]
99    media_control_batches: test_support::MediaControlBatchLog,
100    #[cfg(any(test, feature = "testing-transport"))]
101    source_diagnostics_requests: Arc<AtomicUsize>,
102    /// Coalesces transport observations for room policy.
103    source_policy_signal: SourcePolicySignal,
104}
105
106#[derive(Debug, Clone, Copy)]
107enum TransportCommandOp {
108    CreateInitialSessionOffer,
109    CreateSessionRenegotiationOffer,
110    ApplySessionAnswer,
111    PublishMedia,
112    ConsumeMedia,
113}
114
115fn warn_session_command_failed(
116    session_key: &TransportSessionKey,
117    op: TransportCommandOp,
118    error: TransportAdapterError,
119) {
120    warn!(
121        ?session_key,
122        ?op,
123        ?error,
124        "media transport worker command failed"
125    );
126}
127
128impl MediaTransport {
129    /// Cancels every RTC worker without waiting for termination.
130    pub fn cancel(&self) {
131        for worker in self.workers.iter() {
132            worker.cancel();
133        }
134    }
135
136    /// Cancels every RTC worker and waits for its thread to terminate.
137    pub async fn shutdown(&self) {
138        self.cancel();
139        for worker in self.workers.iter() {
140            worker.wait_for_shutdown().await;
141        }
142    }
143
144    /// returns the router capability snapshot compiled from the RTC wire profile
145    #[must_use]
146    pub fn router_rtp_capabilities(&self) -> MediaCapabilities {
147        self.profile.router_capabilities()
148    }
149
150    async fn request_session_command<T>(
151        &self,
152        session_key: &TransportSessionKey,
153        build: impl FnOnce(RtcWorkerResponse<T>) -> RtcWorkerCommand,
154        log_error: impl FnOnce(TransportAdapterError),
155    ) -> Result<T, TransportAdapterError> {
156        let result = match self.require_worker_for_user(session_key) {
157            Ok(worker) => worker.request_worker(build).await,
158            Err(error) => Err(error),
159        };
160        if let Err(error) = &result {
161            log_error(*error);
162        }
163        result
164    }
165
166    /// Creates the first SDP offer for a transport session.
167    ///
168    /// The session must already be assigned to a transport worker by its
169    /// [`TransportSessionKey`]. The returned offer is transport state that
170    /// callers should send to the browser unchanged.
171    ///
172    /// # Errors
173    ///
174    /// Returns [`TransportAdapterError`] when the session cannot be addressed or
175    /// the active backend cannot create the offer.
176    pub async fn create_initial_session_offer(
177        &self,
178        room_id: &str,
179        session_key: &TransportSessionKey,
180    ) -> Result<SessionOffer, TransportAdapterError> {
181        let room_id: Arc<str> = Arc::from(room_id);
182        self.request_session_command(
183            session_key,
184            move |response| RtcWorkerCommand::CreateInitialSessionOffer {
185                room_id,
186                session_key: session_key.clone(),
187                response,
188            },
189            |error| {
190                warn_session_command_failed(
191                    session_key,
192                    TransportCommandOp::CreateInitialSessionOffer,
193                    error,
194                );
195            },
196        )
197        .await
198        .map(RtcSessionOffer::into_session_offer)
199    }
200
201    /// Creates a new SDP offer after transport media state changed.
202    ///
203    /// Backends reject sessions that cannot renegotiate in their current state.
204    /// Callers should treat the returned offer as replacing any older pending
205    /// transport offer for the same session.
206    ///
207    /// # Errors
208    ///
209    /// Returns [`TransportAdapterError`] when the session cannot be addressed or
210    /// the active backend cannot create a renegotiation offer.
211    pub async fn create_session_renegotiation_offer(
212        &self,
213        session_key: &TransportSessionKey,
214    ) -> Result<SessionOffer, TransportAdapterError> {
215        self.request_session_command(
216            session_key,
217            |response| RtcWorkerCommand::CreateSessionRenegotiationOffer {
218                session_key: session_key.clone(),
219                response,
220            },
221            |error| {
222                warn_session_command_failed(
223                    session_key,
224                    TransportCommandOp::CreateSessionRenegotiationOffer,
225                    error,
226                );
227            },
228        )
229        .await
230        .map(RtcSessionOffer::into_session_offer)
231    }
232
233    /// Applies a browser SDP answer to a pending transport offer.
234    ///
235    /// The answer can reveal transport-derived producer facts such as mapped
236    /// RTP parameters. Those facts are returned so room code can commit staged
237    /// media using values observed by the transport.
238    ///
239    /// # Errors
240    ///
241    /// Returns [`TransportAdapterError`] when the session cannot be addressed,
242    /// the answer is invalid or the active backend cannot apply it.
243    pub async fn apply_session_answer(
244        &self,
245        session_key: &TransportSessionKey,
246        answer_sdp: &str,
247    ) -> Result<AppliedSessionAnswer, TransportAdapterError> {
248        async {
249            let worker = self.require_worker_for_user(session_key)?;
250            let answer = ParsedSessionAnswer::parse(answer_sdp)?;
251            worker
252                .request_worker(|response| RtcWorkerCommand::ApplySessionAnswer {
253                    session_key: session_key.clone(),
254                    answer,
255                    response,
256                })
257                .await
258        }
259        .await
260        .inspect_err(|error| {
261            warn!(
262                ?session_key,
263                op = ?TransportCommandOp::ApplySessionAnswer,
264                answer_len = answer_sdp.len(),
265                ?error,
266                "media transport failed to apply session answer"
267            );
268        })
269    }
270
271    /// Declares a new producer on a transport session.
272    ///
273    /// `rtp_parameters` must come from router media state accepted by the core.
274    /// The returned [`TransportMediaId`] addresses the backend-local producer
275    /// for later route, activity and cleanup operations.
276    ///
277    /// # Errors
278    ///
279    /// Returns [`TransportAdapterError`] when the session cannot be addressed,
280    /// the RTP parameters are invalid or the active backend cannot create the
281    /// producer.
282    pub async fn publish_media(
283        &self,
284        session_key: &TransportSessionKey,
285        media_kind: MediaKind,
286        rtp_parameters: &RouterRtpParameters,
287    ) -> Result<TransportMediaId, TransportAdapterError> {
288        self.request_session_command(
289            session_key,
290            |response| RtcWorkerCommand::AddRecvMedia {
291                session_key: session_key.clone(),
292                media_kind: signaling_to_str0m_media_kind(media_kind),
293                rtp_parameters: rtp_parameters.clone(),
294                response,
295            },
296            |error| {
297                warn!(
298                    ?session_key,
299                    op = ?TransportCommandOp::PublishMedia,
300                    ?media_kind,
301                    mid = rtp_parameters.mid(),
302                    ?error,
303                    "media transport worker command failed"
304                );
305            },
306        )
307        .await
308    }
309
310    /// Declares a new consumer route from a source session to a consumer session.
311    ///
312    /// The source and consumer sessions must belong to the same room instance.
313    /// Cross-worker routing is an implementation detail hidden behind this
314    /// method. The initial activity controls whether the packet loop can
315    /// forward packets before later room policy updates arrive.
316    ///
317    /// # Errors
318    ///
319    /// Returns [`TransportAdapterError`] when either session cannot be
320    /// addressed, the source media id is unknown or the active backend cannot
321    /// create the consumer.
322    pub async fn consume_media(
323        &self,
324        consumer_session_key: &TransportSessionKey,
325        media_kind: MediaKind,
326        source_session_key: &TransportSessionKey,
327        source_media_id: TransportMediaId,
328        consumer_rtp_parameters: &RouterRtpParameters,
329        initial_activity: ConsumerActivity,
330    ) -> Result<TransportMediaId, TransportAdapterError> {
331        self.consume_media_with_mid(
332            consumer_session_key,
333            media_kind,
334            source_session_key,
335            source_media_id,
336            consumer_rtp_parameters,
337            initial_activity,
338        )
339        .await
340        .map(|(transport_media_id, _mid)| transport_media_id)
341    }
342
343    /// Declares consumer media and returns the MID from the same worker turn.
344    ///
345    /// # Errors
346    ///
347    /// Returns [`TransportAdapterError`] when either session is unavailable,
348    /// the source is unknown or the worker cannot create the consumer.
349    pub(in crate::engine) async fn consume_media_with_mid(
350        &self,
351        consumer_session_key: &TransportSessionKey,
352        media_kind: MediaKind,
353        source_session_key: &TransportSessionKey,
354        source_media_id: TransportMediaId,
355        consumer_rtp_parameters: &RouterRtpParameters,
356        initial_activity: ConsumerActivity,
357    ) -> Result<(TransportMediaId, String), TransportAdapterError> {
358        let result = async {
359            Self::ensure_same_room(consumer_session_key, source_session_key)?;
360            let consumer_worker = self.require_worker_for_user(consumer_session_key)?;
361            let source_worker = self.require_worker_for_user(source_session_key)?;
362            let remote_source_control = (!ptr::eq(consumer_worker, source_worker))
363                .then(|| source_worker.remote_source_control(consumer_worker));
364            consumer_worker
365                .request_worker(|response| RtcWorkerCommand::AddSendMedia {
366                    consumer_key: consumer_session_key.clone(),
367                    media_kind: signaling_to_str0m_media_kind(media_kind),
368                    source: TransportSourceKey::new(source_session_key.clone(), source_media_id),
369                    remote_source_control,
370                    consumer_rtp_parameters: consumer_rtp_parameters.clone(),
371                    active: initial_activity.is_active(),
372                    response,
373                })
374                .await
375        }
376        .await;
377        if let Err(error) = &result {
378            warn!(
379                ?consumer_session_key,
380                ?source_session_key,
381                ?source_media_id,
382                op = ?TransportCommandOp::ConsumeMedia,
383                ?media_kind,
384                mid = consumer_rtp_parameters.mid(),
385                initial_active = initial_activity.is_active(),
386                ?error,
387                "media transport worker command failed"
388            );
389        }
390        result
391    }
392
393    /// Applies a room relay route mutation.
394    ///
395    /// The transport may install, release or gate the packet-loop relay target.
396    /// Room state remains the lifecycle owner.
397    ///
398    /// # Errors
399    ///
400    /// Returns [`TransportAdapterError`] when the referenced route cannot be
401    /// addressed or the active backend cannot apply the relay mutation.
402    pub(crate) async fn apply_relay_route_effect(
403        &self,
404        effect: &TransportRelayRouteEffect,
405    ) -> Result<(), TransportAdapterError> {
406        self.execute_relay_route_effect(effect)
407            .await
408            .inspect_err(|error| {
409                warn!(
410                source = ?effect.source,
411                target_media_worker_id = effect.target_media_worker_id.as_usize(),
412                action = ?effect.action,
413                ?error,
414                "media transport failed to apply relay route effect"
415                );
416            })
417    }
418
419    /// Applies one activity revision to a remote source.
420    ///
421    /// # Errors
422    ///
423    /// Returns [`TransportAdapterError`] when the target worker is unavailable
424    /// or rejects the update.
425    pub(crate) async fn apply_remote_source_activity_effect(
426        &self,
427        effect: &TransportSourceActivityEffect,
428    ) -> Result<(), TransportAdapterError> {
429        self.execute_remote_source_activity_effect(effect)
430            .await
431            .inspect_err(|error| {
432                warn!(
433                    source = ?effect.source,
434                    target_media_worker_id = effect.target_media_worker_id.as_usize(),
435                    update = ?effect.update,
436                    ?error,
437                    "media transport failed to apply remote source activity"
438                );
439            })
440    }
441
442    /// Subscribes to source-policy invalidation signals emitted by transport
443    /// workers.
444    #[must_use]
445    pub fn source_policy_subscription(&self) -> SourcePolicyUpdateSubscription {
446        self.source_policy_signal.subscribe()
447    }
448}
449
450#[cfg(test)]
451#[expect(non_snake_case, reason = "test modules map to local TESTS directories")]
452mod TESTS;