o_sfu_core/engine/media_transport/
workers.rs1#[cfg(any(test, feature = "testing-transport"))]
9use std::sync::atomic::Ordering;
10use std::{cmp::Reverse, collections::BTreeMap, ptr};
11
12use str0m::media::MediaKind as Str0mMediaKind;
13
14use super::rtc::{RtcWorker, RtcWorkerCommand};
15use crate::engine::{
16 MediaWorkerId,
17 media_transport::{
18 ActiveSpeakerSource, MediaTransport, ReceiverBandwidthSnapshot, TransportAdapterError,
19 TransportBitrateSnapshot, TransportHealthSnapshot, TransportMediaId,
20 TransportQualitySnapshot, TransportRelayRouteAction, TransportRelayRouteEffect,
21 TransportSessionHealth, TransportSessionKey, TransportSourceActivityEffect,
22 TransportSourceDiagnosticsSnapshot, TransportSourceKey, TransportTeardown,
23 TransportWorkerPressureSnapshot,
24 },
25};
26
27type RelayRegistrationWorkers<'a> = Option<(&'a RtcWorker, &'a RtcWorker)>;
28
29impl MediaTransport {
30 pub(super) fn worker_for_user(&self, session_key: &TransportSessionKey) -> Option<&RtcWorker> {
36 self.worker_for_index(session_key.media_worker_id().as_usize())
37 }
38
39 pub(super) fn relay_registration_workers(
46 &self,
47 consumer_session_key: &TransportSessionKey,
48 source_session_key: &TransportSessionKey,
49 ) -> Result<RelayRegistrationWorkers<'_>, TransportAdapterError> {
50 let consumer_worker = self.require_worker_for_user(consumer_session_key)?;
51 let source_worker = self.require_worker_for_user(source_session_key)?;
52 if ptr::eq(consumer_worker, source_worker) {
53 return Ok(None);
54 }
55 Ok(Some((source_worker, consumer_worker)))
56 }
57
58 #[must_use]
63 pub fn transport_bitrate_snapshot(
64 &self,
65 session_keys: &[TransportSessionKey],
66 ) -> TransportBitrateSnapshot {
67 let mut snapshot = TransportBitrateSnapshot::default();
68 self.for_session_workers(session_keys, |worker, worker_session_keys| {
69 let worker_snapshot = worker.transport_bitrate_snapshot(worker_session_keys);
70 snapshot.total = snapshot.total.saturating_add(worker_snapshot.total);
71 snapshot.per_media.extend(worker_snapshot.per_media);
72 });
73 snapshot
74 }
75
76 #[must_use]
81 pub fn receiver_bandwidth_snapshot(
82 &self,
83 session_keys: &[TransportSessionKey],
84 ) -> ReceiverBandwidthSnapshot {
85 let mut snapshot = ReceiverBandwidthSnapshot::default();
86 self.for_session_workers(session_keys, |worker, worker_session_keys| {
87 let worker_snapshot = worker.receiver_bandwidth_snapshot(worker_session_keys);
88 snapshot.per_session.extend(worker_snapshot.per_session);
89 });
90 snapshot
91 }
92
93 #[must_use]
95 pub fn transport_quality_snapshot(
96 &self,
97 session_keys: &[TransportSessionKey],
98 ) -> TransportQualitySnapshot {
99 let mut snapshot = TransportQualitySnapshot::default();
100 self.for_session_workers(session_keys, |worker, worker_session_keys| {
101 snapshot.extend(worker.transport_quality_snapshot(worker_session_keys));
102 });
103 snapshot
104 }
105
106 #[must_use]
110 pub fn transport_health_snapshot(
111 &self,
112 session_keys: &[TransportSessionKey],
113 ) -> TransportHealthSnapshot {
114 let mut snapshot = TransportHealthSnapshot::default();
115 self.for_session_workers(session_keys, |worker, worker_session_keys| {
116 snapshot.extend(worker.transport_health_snapshot(worker_session_keys));
117 });
118 snapshot
119 }
120
121 pub async fn source_diagnostics_snapshot(
125 &self,
126 sources: &[TransportSourceKey],
127 ) -> TransportSourceDiagnosticsSnapshot {
128 let mut snapshot = TransportSourceDiagnosticsSnapshot::default();
129 let mut media_ids_by_worker = BTreeMap::<usize, Vec<TransportMediaId>>::new();
130 for source in sources {
131 let Some(worker_index) = self.worker_index_for_user(source.session_key()) else {
132 continue;
133 };
134 media_ids_by_worker
135 .entry(worker_index)
136 .or_default()
137 .push(source.transport_media_id());
138 }
139 for (worker_index, transport_media_ids) in media_ids_by_worker {
140 let Some(worker) = self.worker_for_index(worker_index) else {
141 continue;
142 };
143 #[cfg(any(test, feature = "testing-transport"))]
144 self.source_diagnostics_requests
145 .fetch_add(1, Ordering::Relaxed);
146 let worker_snapshot = worker
147 .source_diagnostics_snapshot(&transport_media_ids)
148 .await;
149 snapshot.activity.extend(worker_snapshot.activity);
150 snapshot
151 .active_speaker_diagnostics
152 .extend(worker_snapshot.active_speaker_diagnostics);
153 }
154 snapshot
155 }
156
157 #[must_use]
159 pub fn worker_pressure_snapshots(&self) -> Vec<TransportWorkerPressureSnapshot> {
160 self.workers
161 .iter()
162 .enumerate()
163 .map(|(worker_index, worker)| {
164 worker.worker_pressure_snapshot(MediaWorkerId::from_raw(worker_index))
165 })
166 .collect()
167 }
168
169 pub(crate) fn packet_loop_delays_ms(&self) -> Vec<Option<u64>> {
170 self.workers
171 .iter()
172 .map(RtcWorker::packet_loop_delay_ms)
173 .collect()
174 }
175
176 pub async fn active_speaker_source_snapshot(&self) -> Vec<ActiveSpeakerSource> {
179 let mut snapshot = Vec::new();
180 for worker in self.workers.iter() {
181 snapshot.extend(worker.active_speaker_source_snapshot().await);
182 }
183 snapshot.sort_unstable_by_key(|source| {
187 (
188 source.transport_media_id().as_u64(),
189 Reverse(source.observed_at()),
190 Reverse(source.last_audio_level_dbov().unwrap_or(i8::MIN)),
191 )
192 });
193 snapshot.dedup_by_key(|source| source.transport_media_id());
194 snapshot.sort_unstable_by_key(|source| {
195 (
196 Reverse(source.observed_at()),
197 source.transport_media_id().as_u64(),
198 )
199 });
200 snapshot
201 }
202
203 pub(super) async fn execute_relay_route_effect(
209 &self,
210 effect: &TransportRelayRouteEffect,
211 ) -> Result<(), TransportAdapterError> {
212 if effect.action == TransportRelayRouteAction::Release {
213 self.teardown([TransportTeardown::ReleaseRelayRoute {
214 source: effect.source.clone(),
215 target_media_worker_id: effect.target_media_worker_id,
216 }])
217 .await;
218 return Ok(());
219 }
220 let source_worker = self.require_worker_for_user(effect.source.session_key())?;
221 let target_worker =
222 self.require_worker_for_media_worker_id(effect.target_media_worker_id)?;
223 if ptr::eq(source_worker, target_worker) {
224 return Ok(());
225 }
226 let request = target_worker.relay_route_request(effect.source.clone(), effect.action);
227 source_worker
228 .request_worker(|response| RtcWorkerCommand::RouteControl {
229 request,
230 response: Some(response),
231 })
232 .await
233 }
234
235 pub(super) async fn execute_remote_source_activity_effect(
236 &self,
237 effect: &TransportSourceActivityEffect,
238 ) -> Result<(), TransportAdapterError> {
239 let target_worker =
240 self.require_worker_for_media_worker_id(effect.target_media_worker_id)?;
241 let request =
242 RtcWorker::remote_source_activity_request(effect.source.clone(), effect.update);
243 target_worker
244 .request_worker(|response| RtcWorkerCommand::RouteControl {
245 request,
246 response: Some(response),
247 })
248 .await
249 }
250
251 pub(crate) async fn transport_media_mid(
256 &self,
257 session_key: &TransportSessionKey,
258 transport_media_id: TransportMediaId,
259 ) -> Option<String> {
260 self.worker_for_user(session_key)?
261 .request_worker(|response| RtcWorkerCommand::ResolveMediaMid {
262 transport_media_id,
263 response,
264 })
265 .await
266 .ok()
267 .flatten()
268 }
269
270 #[must_use]
275 pub fn session_transport_health(
276 &self,
277 session_key: &TransportSessionKey,
278 ) -> Option<TransportSessionHealth> {
279 self.worker_for_user(session_key)?
280 .session_transport_health(session_key)
281 }
282
283 fn worker_index_for_user(&self, session_key: &TransportSessionKey) -> Option<usize> {
284 let worker_index = session_key.media_worker_id().as_usize();
285 (worker_index < self.workers.len()).then_some(worker_index)
286 }
287
288 pub(super) fn require_worker_for_user(
289 &self,
290 session_key: &TransportSessionKey,
291 ) -> Result<&RtcWorker, TransportAdapterError> {
292 self.worker_for_user(session_key)
293 .ok_or(TransportAdapterError::TransportUnavailable)
294 }
295
296 pub(super) fn require_worker_for_media_worker_id(
297 &self,
298 media_worker_id: MediaWorkerId,
299 ) -> Result<&RtcWorker, TransportAdapterError> {
300 self.worker_for_index(media_worker_id.as_usize())
301 .ok_or(TransportAdapterError::TransportUnavailable)
302 }
303
304 fn session_keys_by_worker(
305 &self,
306 session_keys: &[TransportSessionKey],
307 ) -> BTreeMap<usize, Vec<TransportSessionKey>> {
308 let mut keys_by_worker = BTreeMap::<usize, Vec<TransportSessionKey>>::new();
309 for session_key in session_keys {
310 if let Some(worker_index) = self.worker_index_for_user(session_key) {
311 keys_by_worker
312 .entry(worker_index)
313 .or_default()
314 .push(session_key.clone());
315 }
316 }
317 keys_by_worker
318 }
319
320 fn for_session_workers(
321 &self,
322 session_keys: &[TransportSessionKey],
323 mut visit: impl FnMut(&RtcWorker, &[TransportSessionKey]),
324 ) {
325 for (worker_index, worker_session_keys) in self.session_keys_by_worker(session_keys) {
326 if let Some(worker) = self.worker_for_index(worker_index) {
327 visit(worker, &worker_session_keys);
328 }
329 }
330 }
331
332 pub(super) fn ensure_same_room(
337 consumer_session_key: &TransportSessionKey,
338 source_session_key: &TransportSessionKey,
339 ) -> Result<(), TransportAdapterError> {
340 if consumer_session_key.room_instance_id() == source_session_key.room_instance_id() {
341 return Ok(());
342 }
343 Err(TransportAdapterError::InvalidInput)
344 }
345
346 pub(super) fn worker_for_index(&self, worker_index: usize) -> Option<&RtcWorker> {
347 self.workers.get(worker_index)
348 }
349
350 #[cfg(any(test, feature = "testing-transport"))]
351 pub(super) fn all_workers(&self) -> impl Iterator<Item = &RtcWorker> {
352 self.workers.iter()
353 }
354}
355
356pub(super) fn signaling_to_str0m_media_kind(kind: o_sfu_router::MediaKind) -> Str0mMediaKind {
357 match kind {
358 o_sfu_router::MediaKind::Audio => Str0mMediaKind::Audio,
359 o_sfu_router::MediaKind::Video => Str0mMediaKind::Video,
360 }
361}