1mod 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#[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 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 pub fn cancel(&self) {
131 for worker in self.workers.iter() {
132 worker.cancel();
133 }
134 }
135
136 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 #[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 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 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 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 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 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 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 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 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 #[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;