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