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 futures_util::{StreamExt, TryStreamExt, future::try_join_all, stream};
13use str0m::media::MediaKind as Str0mMediaKind;
14
15use super::rtc::{RtcWorker, RtcWorkerCommand};
16use crate::engine::{
17 MediaWorkerId,
18 media_transport::{
19 ActiveSpeakerSource, MediaTransport, ReceiverBandwidthSnapshot, TransportAdapterError,
20 TransportBitrateSnapshot, TransportHealthSnapshot, TransportMediaId,
21 TransportQualitySnapshot, TransportRelayRouteAction, TransportRelayRouteEffect,
22 TransportSessionHealth, TransportSessionKey, TransportSourceActivityEffect,
23 TransportSourceDiagnosticsSnapshot, TransportSourceKey, TransportTeardown,
24 TransportWorkerPressureSnapshot,
25 },
26};
27
28const WORKER_OBSERVATION_CONCURRENCY: usize = 8;
29
30#[derive(Debug, Clone, Copy, PartialEq, Eq)]
32pub(crate) enum WorkerPlacementState {
33 Running(Option<u64>),
34 Unavailable,
35}
36
37impl MediaTransport {
38 pub(super) fn worker_for_user(&self, session_key: &TransportSessionKey) -> Option<&RtcWorker> {
44 self.worker_for_index(session_key.media_worker_id().as_usize())
45 }
46
47 #[must_use]
52 pub fn transport_bitrate_snapshot(
53 &self,
54 session_keys: &[TransportSessionKey],
55 ) -> TransportBitrateSnapshot {
56 let mut snapshot = TransportBitrateSnapshot::default();
57 self.for_session_workers(session_keys, |worker, worker_session_keys| {
58 let worker_snapshot = worker.transport_bitrate_snapshot(worker_session_keys);
59 snapshot.total = snapshot.total.saturating_add(worker_snapshot.total);
60 snapshot.per_media.extend(worker_snapshot.per_media);
61 });
62 snapshot
63 }
64
65 #[must_use]
70 pub fn receiver_bandwidth_snapshot(
71 &self,
72 session_keys: &[TransportSessionKey],
73 ) -> ReceiverBandwidthSnapshot {
74 let mut snapshot = ReceiverBandwidthSnapshot::default();
75 self.for_session_workers(session_keys, |worker, worker_session_keys| {
76 let worker_snapshot = worker.receiver_bandwidth_snapshot(worker_session_keys);
77 snapshot.per_session.extend(worker_snapshot.per_session);
78 });
79 snapshot
80 }
81
82 #[must_use]
84 pub fn transport_quality_snapshot(
85 &self,
86 session_keys: &[TransportSessionKey],
87 ) -> TransportQualitySnapshot {
88 let mut snapshot = TransportQualitySnapshot::default();
89 self.for_session_workers(session_keys, |worker, worker_session_keys| {
90 snapshot.extend(worker.transport_quality_snapshot(worker_session_keys));
91 });
92 snapshot
93 }
94
95 #[must_use]
99 pub fn transport_health_snapshot(
100 &self,
101 session_keys: &[TransportSessionKey],
102 ) -> TransportHealthSnapshot {
103 let mut snapshot = TransportHealthSnapshot::default();
104 self.for_session_workers(session_keys, |worker, worker_session_keys| {
105 snapshot.extend(worker.transport_health_snapshot(worker_session_keys));
106 });
107 snapshot
108 }
109
110 pub async fn source_diagnostics_snapshot(
117 &self,
118 sources: &[TransportSourceKey],
119 ) -> Result<TransportSourceDiagnosticsSnapshot, TransportAdapterError> {
120 let mut media_ids_by_worker = BTreeMap::<usize, Vec<TransportMediaId>>::new();
121 for source in sources {
122 let worker_index = self
123 .worker_index_for_user(source.session_key())
124 .ok_or(TransportAdapterError::TransportUnavailable)?;
125 media_ids_by_worker
126 .entry(worker_index)
127 .or_default()
128 .push(source.transport_media_id());
129 }
130 let worker_snapshots = stream::iter(media_ids_by_worker)
131 .map(|(worker_index, transport_media_ids)| async move {
132 let worker = self
133 .worker_for_index(worker_index)
134 .ok_or(TransportAdapterError::TransportUnavailable)?;
135 #[cfg(any(test, feature = "testing-transport"))]
136 self.source_diagnostics_requests
137 .fetch_add(1, Ordering::Relaxed);
138 let snapshot = worker
139 .source_diagnostics_snapshot(&transport_media_ids)
140 .await?;
141 Ok::<_, TransportAdapterError>((worker_index, snapshot))
142 })
143 .buffer_unordered(WORKER_OBSERVATION_CONCURRENCY)
144 .try_collect::<Vec<_>>()
145 .await?;
146 let mut worker_snapshots = worker_snapshots;
147 worker_snapshots.sort_unstable_by_key(|(worker_index, _)| *worker_index);
148 let mut snapshot = TransportSourceDiagnosticsSnapshot::default();
149 for (_, worker_snapshot) in worker_snapshots {
150 snapshot.activity.extend(worker_snapshot.activity);
151 snapshot
152 .active_speaker_diagnostics
153 .extend(worker_snapshot.active_speaker_diagnostics);
154 }
155 Ok(snapshot)
156 }
157
158 #[must_use]
160 pub fn worker_pressure_snapshots(&self) -> Vec<TransportWorkerPressureSnapshot> {
161 self.workers
162 .iter()
163 .enumerate()
164 .map(|(worker_index, worker)| {
165 worker.worker_pressure_snapshot(MediaWorkerId::from_raw(worker_index))
166 })
167 .collect()
168 }
169
170 #[cfg(any(test, feature = "testing-transport"))]
171 pub(crate) fn packet_loop_delays_ms(&self) -> Vec<Option<u64>> {
172 self.workers
173 .iter()
174 .map(RtcWorker::packet_loop_delay_ms)
175 .collect()
176 }
177
178 pub(crate) fn worker_is_usable(&self, worker_index: usize) -> bool {
180 self.worker_for_index(worker_index)
181 .is_some_and(RtcWorker::is_usable)
182 }
183
184 pub(crate) fn worker_placement_states(&self) -> Vec<WorkerPlacementState> {
185 self.workers
186 .iter()
187 .map(|worker| {
188 if worker.is_usable() {
189 WorkerPlacementState::Running(worker.packet_loop_delay_ms())
190 } else {
191 WorkerPlacementState::Unavailable
192 }
193 })
194 .collect()
195 }
196
197 pub(crate) async fn active_speaker_source_snapshot_for_worker(
204 &self,
205 worker_index: usize,
206 ) -> Result<Vec<ActiveSpeakerSource>, TransportAdapterError> {
207 let worker = self
208 .worker_for_index(worker_index)
209 .ok_or(TransportAdapterError::TransportUnavailable)?;
210 worker.active_speaker_source_snapshot().await
211 }
212
213 pub(crate) async fn active_speaker_source_snapshot_for_workers(
220 &self,
221 worker_indices: &[usize],
222 ) -> Result<Vec<ActiveSpeakerSource>, TransportAdapterError> {
223 let snapshots = try_join_all(
224 worker_indices
225 .iter()
226 .map(|&worker_index| self.active_speaker_source_snapshot_for_worker(worker_index)),
227 )
228 .await?;
229 let mut snapshot = snapshots.into_iter().flatten().collect::<Vec<_>>();
230 snapshot.sort_unstable_by_key(|source| {
232 (
233 source.transport_media_id().as_u64(),
234 Reverse(source.observed_at()),
235 Reverse(source.last_audio_level_dbov().unwrap_or(i8::MIN)),
236 )
237 });
238 snapshot.dedup_by_key(|source| source.transport_media_id());
239 snapshot.sort_unstable_by_key(|source| {
240 (
241 Reverse(source.observed_at()),
242 source.transport_media_id().as_u64(),
243 )
244 });
245 Ok(snapshot)
246 }
247
248 pub async fn active_speaker_source_snapshot(
255 &self,
256 ) -> Result<Vec<ActiveSpeakerSource>, TransportAdapterError> {
257 let workers = (0..self.workers.len()).collect::<Vec<_>>();
258 self.active_speaker_source_snapshot_for_workers(&workers)
259 .await
260 }
261
262 pub(super) async fn execute_relay_route_effect(
268 &self,
269 effect: &TransportRelayRouteEffect,
270 ) -> Result<(), TransportAdapterError> {
271 if effect.action == TransportRelayRouteAction::Release {
272 self.teardown([TransportTeardown::ReleaseRelayRoute {
273 source: effect.source.clone(),
274 target_media_worker_id: effect.target_media_worker_id,
275 }])
276 .await;
277 return Ok(());
278 }
279 let source_worker = self.require_worker_for_user(effect.source.session_key())?;
280 let target_worker =
281 self.require_worker_for_media_worker_id(effect.target_media_worker_id)?;
282 if ptr::eq(source_worker, target_worker) {
283 return Ok(());
284 }
285 let request = target_worker.relay_route_request(effect.source.clone(), effect.action);
286 source_worker
287 .request_worker(|response| RtcWorkerCommand::RouteControl {
288 request,
289 response: Some(response),
290 })
291 .await
292 }
293
294 pub(super) async fn execute_remote_source_activity_effect(
295 &self,
296 effect: &TransportSourceActivityEffect,
297 ) -> Result<(), TransportAdapterError> {
298 let target_worker =
299 self.require_worker_for_media_worker_id(effect.target_media_worker_id)?;
300 let request =
301 RtcWorker::remote_source_activity_request(effect.source.clone(), effect.update);
302 target_worker
303 .request_worker(|response| RtcWorkerCommand::RouteControl {
304 request,
305 response: Some(response),
306 })
307 .await
308 }
309
310 #[cfg(test)]
319 pub(crate) async fn transport_media_mid(
320 &self,
321 session_key: &TransportSessionKey,
322 transport_media_id: TransportMediaId,
323 ) -> Result<Option<String>, TransportAdapterError> {
324 self.require_worker_for_user(session_key)?
325 .request_worker(|response| RtcWorkerCommand::ResolveMediaMid {
326 transport_media_id,
327 response,
328 })
329 .await
330 }
331
332 #[must_use]
337 pub fn session_transport_health(
338 &self,
339 session_key: &TransportSessionKey,
340 ) -> Option<TransportSessionHealth> {
341 self.worker_for_user(session_key)?
342 .session_transport_health(session_key)
343 }
344
345 fn worker_index_for_user(&self, session_key: &TransportSessionKey) -> Option<usize> {
346 let worker_index = session_key.media_worker_id().as_usize();
347 (worker_index < self.workers.len()).then_some(worker_index)
348 }
349
350 pub(super) fn require_worker_for_user(
351 &self,
352 session_key: &TransportSessionKey,
353 ) -> Result<&RtcWorker, TransportAdapterError> {
354 self.worker_for_user(session_key)
355 .filter(|worker| worker.is_usable())
356 .ok_or(TransportAdapterError::TransportUnavailable)
357 }
358
359 pub(super) fn require_worker_for_media_worker_id(
360 &self,
361 media_worker_id: MediaWorkerId,
362 ) -> Result<&RtcWorker, TransportAdapterError> {
363 self.worker_for_index(media_worker_id.as_usize())
364 .filter(|worker| worker.is_usable())
365 .ok_or(TransportAdapterError::TransportUnavailable)
366 }
367
368 fn session_keys_by_worker(
369 &self,
370 session_keys: &[TransportSessionKey],
371 ) -> BTreeMap<usize, Vec<TransportSessionKey>> {
372 let mut keys_by_worker = BTreeMap::<usize, Vec<TransportSessionKey>>::new();
373 for session_key in session_keys {
374 if let Some(worker_index) = self.worker_index_for_user(session_key) {
375 keys_by_worker
376 .entry(worker_index)
377 .or_default()
378 .push(session_key.clone());
379 }
380 }
381 keys_by_worker
382 }
383
384 fn for_session_workers(
385 &self,
386 session_keys: &[TransportSessionKey],
387 mut visit: impl FnMut(&RtcWorker, &[TransportSessionKey]),
388 ) {
389 for (worker_index, worker_session_keys) in self.session_keys_by_worker(session_keys) {
390 if let Some(worker) = self.worker_for_index(worker_index) {
391 visit(worker, &worker_session_keys);
392 }
393 }
394 }
395
396 pub(super) fn ensure_same_room(
401 consumer_session_key: &TransportSessionKey,
402 source_session_key: &TransportSessionKey,
403 ) -> Result<(), TransportAdapterError> {
404 if consumer_session_key.room_instance_id() == source_session_key.room_instance_id() {
405 return Ok(());
406 }
407 Err(TransportAdapterError::InvalidInput)
408 }
409
410 pub(super) fn worker_for_index(&self, worker_index: usize) -> Option<&RtcWorker> {
411 self.workers.get(worker_index)
412 }
413
414 #[cfg(any(test, feature = "testing-transport"))]
415 pub(super) fn all_workers(&self) -> impl Iterator<Item = &RtcWorker> {
416 self.workers.iter()
417 }
418}
419
420pub(super) fn signaling_to_str0m_media_kind(kind: o_sfu_router::MediaKind) -> Str0mMediaKind {
421 match kind {
422 o_sfu_router::MediaKind::Audio => Str0mMediaKind::Audio,
423 o_sfu_router::MediaKind::Video => Str0mMediaKind::Video,
424 }
425}