Skip to main content

o_sfu_core/engine/media_transport/
workers.rs

1//! RTC worker ownership below the media transport boundary.
2//!
3//! `MediaTransport` owns the process-local RTC worker topology, maps each
4//! transport session to the worker selected by its media-worker id and
5//! coordinates cross-worker relay cleanup. Packet-loop hot paths live inside
6//! `engine::media_transport::rtc`.
7
8#[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/// Worker availability is independent of missing or delayed loop samples.
31#[derive(Debug, Clone, Copy, PartialEq, Eq)]
32pub(crate) enum WorkerPlacementState {
33    Running(Option<u64>),
34    Unavailable,
35}
36
37impl MediaTransport {
38    /// Selects the worker that owns a transport session.
39    ///
40    /// The mapping is deterministic and depends only on the runtime-assigned
41    /// media-worker id in the session key. Room and signaling code must not
42    /// infer topology from user identity.
43    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    /// Returns the latest bitrate estimates for the requested sessions.
48    ///
49    /// Missing sessions are omitted from the snapshot. Estimates are suitable
50    /// for diagnostics and policy input, not for accounting.
51    #[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    /// Returns receiver-side bandwidth estimates for the requested sessions.
66    ///
67    /// Room policy may use these estimates as source-selection input. They are
68    /// best-effort observations from the transport backend.
69    #[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    /// Returns sampled transport-quality observations for the requested sessions.
83    #[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    /// Returns transport health for the requested sessions with one lock per worker.
96    ///
97    /// Missing sessions and unavailable worker snapshots contribute no facts
98    #[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    /// Returns source activity and active-speaker facts with one command per worker.
111    ///
112    /// # Errors
113    ///
114    /// Returns [`TransportAdapterError::TransportUnavailable`] if any source
115    /// selects an unavailable worker or its observation does not complete.
116    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    /// Returns transport pressure for every media worker.
159    #[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    /// Reports whether the worker still owns usable transport sessions.
179    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    /// Returns one worker's active-speaker inventory.
198    ///
199    /// # Errors
200    ///
201    /// Returns [`TransportAdapterError::TransportUnavailable`] when the worker
202    /// is missing or cannot answer before the observation deadline.
203    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    /// Returns one observation from each requested worker before merging relayed sources.
214    ///
215    /// # Errors
216    ///
217    /// Returns [`TransportAdapterError::TransportUnavailable`] if a requested
218    /// worker is unavailable. Callers must retry rather than use partial facts.
219    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        // A relayed source can be observed on its owner and consumer workers.
231        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    /// Returns active-speaker observations from all workers.
249    ///
250    /// # Errors
251    ///
252    /// Returns [`TransportAdapterError::TransportUnavailable`] if any worker
253    /// observation fails. No partial observation is returned to policy.
254    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    /// Applies one cross-worker relay mutation on the source worker.
263    ///
264    /// The target worker contributes its relay identity and mailbox. The source
265    /// packet loop owns registration and activity because it decides fanout before
266    /// packets cross workers.
267    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    /// Returns the MID stored by the current transport media handle.
311    ///
312    /// `None` means `transport_media_id` has no registered handle.
313    ///
314    /// # Errors
315    ///
316    /// Returns [`TransportAdapterError::TransportUnavailable`] if the selected
317    /// worker is unavailable or cannot answer before the observation deadline.
318    #[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    /// Returns the latest known transport health for one session.
333    ///
334    /// Health is connectivity evidence only. It should not be used as the
335    /// source of truth for whether a participant belongs to a room.
336    #[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    /// Enforces room isolation before worker or media lookup.
397    ///
398    /// `MediaWorkerId` selects execution ownership only. It never authorizes a
399    /// route between room instances.
400    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}