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 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    /// Selects the worker that owns a transport session.
31    ///
32    /// The mapping is deterministic and depends only on the runtime-assigned
33    /// media-worker id in the session key. Room and signaling code must not
34    /// infer topology from user identity.
35    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    /// Returns the source and consumer workers needed for cross-worker relay.
40    ///
41    /// `None` means both sessions are on the same worker and local routing is
42    /// enough. A returned pair means the source worker must activate relay
43    /// forwarding toward the consumer worker before the consumer route can be
44    /// fully installed.
45    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    /// Returns the latest bitrate estimates for the requested sessions.
59    ///
60    /// Missing sessions are omitted from the snapshot. Estimates are suitable
61    /// for diagnostics and policy input, not for accounting.
62    #[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    /// Returns receiver-side bandwidth estimates for the requested sessions.
77    ///
78    /// Room policy may use these estimates as source-selection input. They are
79    /// best-effort observations from the transport backend.
80    #[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    /// Returns sampled transport-quality observations for the requested sessions.
94    #[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    /// Returns transport health for the requested sessions with one lock per worker.
107    ///
108    /// Missing sessions and unavailable worker snapshots contribute no facts
109    #[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    /// Returns source activity and active-speaker facts with one command per worker.
122    ///
123    /// Missing sources and worker dispatch failures contribute no facts
124    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    /// Returns transport pressure for every media worker.
158    #[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    /// Returns each media source's newest active-speaker observation in recency
177    /// order, preferring its strongest audio level on timestamp ties.
178    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        // A relayed source is observed on its owner and consumer workers. Keep one
184        // policy fact per media ID, preferring recency then the strongest level as
185        // the deterministic same-timestamp tie-break.
186        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    /// Applies one cross-worker relay mutation on the source worker.
204    ///
205    /// The target worker contributes its relay identity and mailbox. The source
206    /// packet loop owns registration and activity because it decides fanout before
207    /// packets cross workers.
208    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    /// Returns the MID stored by the current transport media handle.
252    ///
253    /// `None` means `session_key` selects no worker, the worker request failed or
254    /// `transport_media_id` has no registered handle.
255    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    /// Returns the latest known transport health for one session.
271    ///
272    /// Health is connectivity evidence only. It should not be used as the
273    /// source of truth for whether a participant belongs to a room.
274    #[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    /// Enforces room isolation before worker or media lookup.
333    ///
334    /// `MediaWorkerId` selects execution ownership only. It never authorizes a
335    /// route between room instances.
336    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}