Skip to main content

o_sfu_core/engine/media_transport/rtc/
commands.rs

1//! mailbox command contract for the RTC worker engine
2//!
3//! `MediaTransport` production paths and cfg-gated worker harnesses translate
4//! transport intent into these values before the packet-loop task dispatches
5//! them while it owns mutable rtc state
6//! request commands carry a oneshot response
7//! fire-and-forget route controls are best-effort because they may target a
8//! worker that has already torn down the corresponding relay or session
9
10use std::sync::Arc;
11
12use o_sfu_rfc::webrtc::sdp;
13use o_sfu_router::rtp::MediaStream as RouterRtpParameters;
14use str0m::{
15    change::{SdpAnswer, SdpOffer},
16    media::{KeyframeRequestKind, MediaKind, Rid},
17};
18use tokio::sync::{mpsc, oneshot};
19
20use super::{
21    codec::{ParsedAnswerRids, RepairSummary, validate_answer_sdp},
22    state::{
23        relay_registry::{RelayPacketMailbox, RelayTargetId},
24        route_control::PacketLayerGate,
25    },
26};
27use crate::engine::{
28    media_transport::{
29        ActiveSpeakerSource, AppliedSessionAnswer, ConsumerRouteControl,
30        ConsumerRouteControlOutcome, ProducerRouteControl, ReceiverBweTargetUpdate, SessionOffer,
31        SessionUploadSlot, SourceActivityUpdate, TransportAdapterError, TransportConsumerRoute,
32        TransportMediaId, TransportResult, TransportSessionKey, TransportSourceDiagnosticsSnapshot,
33        TransportSourceKey,
34    },
35    metrics::{
36        RtcMetricsRecorder, RtcRemoteControlDropKind, RtcRemotePacketGateConvergence,
37        RtcWorkerObservationKind,
38    },
39};
40
41/// command handle used by remote consumers to push control back to a source worker
42///
43/// a route that consumes media from another worker keeps this handle beside the
44/// remote-source registration
45/// later keyframe or layer-gate requests can then reach the worker that owns
46/// the producer without exposing the source worker internals
47///
48/// sends are deliberately best-effort
49/// stale remote routes, closed workers and full mailboxes are normal during
50/// teardown or topology churn
51#[derive(Debug, Clone)]
52pub struct RemoteSourceControl {
53    tx: mpsc::Sender<RtcWorkerCommand>,
54    target_id: RelayTargetId,
55    rtc_metrics: Arc<RtcMetricsRecorder>,
56}
57
58#[derive(Debug, Clone, Copy, PartialEq, Eq)]
59pub(super) enum RemoteControlSendOutcome {
60    Forwarded,
61    Full,
62    Closed,
63}
64
65impl RemoteSourceControl {
66    /// creates a source-control handle for one relay target on a worker mailbox
67    pub(super) fn new(
68        tx: mpsc::Sender<RtcWorkerCommand>,
69        target_id: RelayTargetId,
70        rtc_metrics: Arc<RtcMetricsRecorder>,
71    ) -> Self {
72        Self {
73            tx,
74            target_id,
75            rtc_metrics,
76        }
77    }
78
79    /// asks the source worker to request a keyframe for a remote consumer
80    ///
81    /// this never waits for the source worker
82    /// `Full` retains caller-side retry authority while `Closed` ends it
83    pub(super) fn request_kf(
84        &self,
85        source: &TransportSourceKey,
86        rid: Option<Rid>,
87        kind: KeyframeRequestKind,
88    ) -> RemoteControlSendOutcome {
89        self.send_command(
90            RouteControlRequest::RequestRemoteKeyframe {
91                source: source.clone(),
92                target_id: self.target_id,
93                rid,
94                kind,
95            },
96            RtcRemoteControlDropKind::Keyframe,
97        )
98    }
99
100    /// publishes the effective remote-source packet gate to the source worker
101    ///
102    /// `Full` requires a later retry. `Closed` ends retry authority.
103    pub(super) fn set_pkt_gate(
104        &self,
105        source: &TransportSourceKey,
106        packet_gate: PacketLayerGate,
107    ) -> RemoteControlSendOutcome {
108        self.send_command(
109            RouteControlRequest::SetRemoteSourcePacketGate {
110                source: source.clone(),
111                target_id: self.target_id,
112                packet_gate,
113            },
114            RtcRemoteControlDropKind::PacketGate,
115        )
116    }
117
118    pub(super) fn record_pkt_gate_retry(&self) {
119        self.rtc_metrics
120            .record_rtc_remote_packet_gate_convergence(RtcRemotePacketGateConvergence::Retry);
121    }
122
123    pub(super) fn record_pkt_gate_flushed(&self) {
124        self.rtc_metrics
125            .record_rtc_remote_packet_gate_convergence(RtcRemotePacketGateConvergence::Flushed);
126    }
127
128    #[inline]
129    fn send_command(
130        &self,
131        request: RouteControlRequest,
132        drop_kind: RtcRemoteControlDropKind,
133    ) -> RemoteControlSendOutcome {
134        match self.tx.try_send(RtcWorkerCommand::RouteControl {
135            request,
136            response: None,
137        }) {
138            Ok(()) => RemoteControlSendOutcome::Forwarded,
139            Err(mpsc::error::TrySendError::Full(_)) => {
140                self.rtc_metrics.record_rtc_remote_control_drop(drop_kind);
141                RemoteControlSendOutcome::Full
142            }
143            Err(mpsc::error::TrySendError::Closed(_)) => {
144                self.rtc_metrics.record_rtc_remote_control_drop(drop_kind);
145                RemoteControlSendOutcome::Closed
146            }
147        }
148    }
149}
150
151/// Response channel for a command executed by the packet loop.
152///
153/// Dropping the receiver cancels only the API wait and does not retract an
154/// enqueued command. Terminal worker cancellation may still win before dispatch.
155pub type RtcWorkerResponse<T> = oneshot::Sender<TransportResult<T>>;
156
157/// SDP answer parsed before packet-loop command delivery.
158pub struct ParsedSessionAnswer {
159    pub(super) answer: SdpAnswer,
160    pub(super) rids: ParsedAnswerRids,
161    pub(super) repair: RepairSummary,
162}
163
164impl ParsedSessionAnswer {
165    /// Validates and parses a remote SDP answer.
166    ///
167    /// # Errors
168    ///
169    /// Returns [`TransportAdapterError::InvalidInput`] for malformed repair
170    /// topology or SDP that str0m cannot parse.
171    pub(in crate::engine::media_transport) fn parse(answer_sdp: &str) -> TransportResult<Self> {
172        let repair = validate_answer_sdp(answer_sdp)?;
173        let answer = SdpAnswer::from_sdp_string(answer_sdp)
174            .map_err(|_error| TransportAdapterError::InvalidInput)?;
175        let rids = ParsedAnswerRids::parse(answer_sdp, &answer);
176        Ok(Self {
177            answer,
178            rids,
179            repair,
180        })
181    }
182}
183
184pub struct RtcSessionOffer {
185    offer: SdpOffer,
186    upload_slots: Vec<SessionUploadSlot>,
187}
188
189impl RtcSessionOffer {
190    pub(super) fn new(offer: SdpOffer, upload_slots: Vec<SessionUploadSlot>) -> Self {
191        Self {
192            offer,
193            upload_slots,
194        }
195    }
196
197    pub(in crate::engine::media_transport) fn into_session_offer(self) -> SessionOffer {
198        let mut sdp = self.offer.to_sdp_string();
199        let media_line_start = sdp.find("\r\nm=").map_or(sdp.len(), |index| index + 2);
200        sdp.insert_str(media_line_start, sdp::EOC_LINE);
201        SessionOffer::new(sdp).with_upload_slots(self.upload_slots)
202    }
203}
204
205#[derive(Debug)]
206pub enum WorkerMediaControlBatch {
207    ReceiverBwe(Vec<(usize, ReceiverBweTargetUpdate)>),
208    ProducerActivity(Vec<(usize, ProducerRouteControl)>),
209    ConsumerGates {
210        source: TransportSourceKey,
211        updates: Vec<(usize, TransportConsumerRoute, PacketLayerGate)>,
212    },
213    ConsumerFollowUp(Vec<(usize, ConsumerRouteControl)>),
214}
215
216#[derive(Debug)]
217pub enum WorkerMediaControlBatchOutcome {
218    Applied(Vec<TransportResult<()>>),
219    Consumers(Vec<ConsumerRouteControlOutcome>),
220}
221
222pub enum RouteControlRequest {
223    AddRelayTarget {
224        source: TransportSourceKey,
225        target_id: RelayTargetId,
226        target: RelayPacketMailbox,
227    },
228    RemoveRelayTarget {
229        source: TransportSourceKey,
230        target_id: RelayTargetId,
231    },
232    SetRelayTargetActive {
233        source: TransportSourceKey,
234        target_id: RelayTargetId,
235        active: bool,
236    },
237    SetRemoteSourceActivity {
238        source: TransportSourceKey,
239        update: SourceActivityUpdate,
240    },
241    RequestRemoteKeyframe {
242        source: TransportSourceKey,
243        target_id: RelayTargetId,
244        rid: Option<Rid>,
245        kind: KeyframeRequestKind,
246    },
247    SetRemoteSourcePacketGate {
248        source: TransportSourceKey,
249        target_id: RelayTargetId,
250        packet_gate: PacketLayerGate,
251    },
252}
253
254/// production command handled by the rtc packet-loop task
255///
256/// variants are grouped by ownership boundary: negotiation mutates str0m SDP
257/// state, media commands mutate producer or consumer registrations, relay
258/// commands mutate cross-worker fanout and observability commands read
259/// worker-local snapshots
260pub enum RtcWorkerCommand {
261    /// create the initial capability-probe offer before media registration
262    ///
263    /// worker startup binds the shared UDP socket, then this command lazily creates
264    /// session RTC state from its candidate address
265    /// it rejects a pending offer, a committed initial answer or registered media
266    CreateInitialSessionOffer {
267        room_id: Arc<str>,
268        session_key: TransportSessionKey,
269        response: RtcWorkerResponse<RtcSessionOffer>,
270    },
271    /// drain a staged follow-up offer after media topology changed
272    ///
273    /// media add and remove commands stage the SDP work before this command
274    /// runs
275    /// this command hands the staged offer to the worker API and preserves the
276    /// one-outstanding-offer rule owned by the worker
277    CreateSessionRenegotiationOffer {
278        session_key: TransportSessionKey,
279        response: RtcWorkerResponse<RtcSessionOffer>,
280    },
281    /// read active-speaker sources from worker-local route-control state
282    ///
283    /// the result is a cold-path observation for room policy
284    /// it does not mutate route state or packet-loop scheduling
285    ActiveSpeakerSourceSnapshot {
286        response: RtcWorkerResponse<Vec<ActiveSpeakerSource>>,
287    },
288    /// Read packet activity and active-speaker facts for selected sources.
289    SourceDiagnosticsSnapshot {
290        transport_media_ids: Vec<TransportMediaId>,
291        response: RtcWorkerResponse<TransportSourceDiagnosticsSnapshot>,
292    },
293    /// accept the answer for the current pending local offer
294    ///
295    /// this commits str0m SDP state, marks the session dirty, refreshes
296    /// negotiated producer parameters, registers remote candidate recovery hints
297    /// and returns the producer details that became usable after the answer
298    ApplySessionAnswer {
299        session_key: TransportSessionKey,
300        answer: ParsedSessionAnswer,
301        response: RtcWorkerResponse<AppliedSessionAnswer>,
302    },
303    /// remove a session from worker state
304    ///
305    /// teardown removes rtc state, media handles, route destinations, demux
306    /// indexes, bitrate counters and snapshot entries owned by the session
307    CloseSession {
308        session_key: TransportSessionKey,
309        response: RtcWorkerResponse<()>,
310    },
311    /// remove one producer or consumer media registration owned by a session
312    ///
313    /// producer removal drops incoming bitrate tracking and the source route
314    /// consumer removal drops local rewrite state and the destination route
315    /// negotiated media removal may stage the next SDP offer before the handle
316    /// leaves the public registry
317    RemoveMedia {
318        session_key: TransportSessionKey,
319        transport_media_id: TransportMediaId,
320        response: RtcWorkerResponse<()>,
321    },
322    /// resolve negotiated producer parameters for adapter tests
323    ///
324    /// this test-only command reads the answer-derived producer state after
325    /// negotiation so adapter tests can assert the transport boundary without
326    /// reaching into private registries
327    #[cfg(test)]
328    ResolveNegotiatedProducerParameters {
329        session_key: TransportSessionKey,
330        transport_media_id: TransportMediaId,
331        response: RtcWorkerResponse<RouterRtpParameters>,
332    },
333    /// resolve the MID stored by a registered media handle
334    ///
335    /// returns `None` when this worker has no current handle for the id, including
336    /// after media removal
337    #[cfg(test)]
338    ResolveMediaMid {
339        transport_media_id: TransportMediaId,
340        response: RtcWorkerResponse<Option<String>>,
341    },
342    /// register one browser upload as worker-local producer media
343    ///
344    /// before the initial answer this can declare receive state directly in
345    /// str0m
346    /// after negotiation it stages a recv-only m-section plus pending
347    /// receive identities, then registers bitrate counters and the producer
348    /// media handle
349    AddRecvMedia {
350        session_key: TransportSessionKey,
351        media_kind: MediaKind,
352        rtp_parameters: RouterRtpParameters,
353        response: RtcWorkerResponse<TransportMediaId>,
354    },
355    /// register one browser download as consumer media for a source
356    ///
357    /// the worker validates local or remote source ownership, stages or declares
358    /// send-only media, registers the consumer handle and creates the packet-loop
359    /// route destination
360    /// remote sources install rollback-protected control so failed consumer
361    /// setup does not leave stale relay state behind
362    AddSendMedia {
363        consumer_key: TransportSessionKey,
364        media_kind: MediaKind,
365        source: TransportSourceKey,
366        remote_source_control: Option<RemoteSourceControl>,
367        consumer_rtp_parameters: RouterRtpParameters,
368        active: bool,
369        response: RtcWorkerResponse<(TransportMediaId, String)>,
370    },
371    ApplyMediaControlBatch {
372        batch: WorkerMediaControlBatch,
373        response: RtcWorkerResponse<WorkerMediaControlBatchOutcome>,
374    },
375    RouteControl {
376        request: RouteControlRequest,
377        response: Option<RtcWorkerResponse<()>>,
378    },
379}
380
381impl RtcWorkerCommand {
382    /// Classifies mailbox reads without assuming a command's response type is read-only.
383    pub(super) fn observation_kind(&self) -> Option<RtcWorkerObservationKind> {
384        match self {
385            Self::ActiveSpeakerSourceSnapshot { .. } => {
386                Some(RtcWorkerObservationKind::ActiveSpeakerSources)
387            }
388            Self::SourceDiagnosticsSnapshot { .. } => {
389                Some(RtcWorkerObservationKind::SourceDiagnostics)
390            }
391            #[cfg(test)]
392            Self::ResolveMediaMid { .. } => Some(RtcWorkerObservationKind::ResolveMediaMid),
393            #[cfg(test)]
394            Self::ResolveNegotiatedProducerParameters { .. } => {
395                Some(RtcWorkerObservationKind::NegotiatedProducerParameters)
396            }
397            Self::CreateInitialSessionOffer { .. }
398            | Self::CreateSessionRenegotiationOffer { .. }
399            | Self::ApplySessionAnswer { .. }
400            | Self::CloseSession { .. }
401            | Self::RemoveMedia { .. }
402            | Self::AddRecvMedia { .. }
403            | Self::AddSendMedia { .. }
404            | Self::ApplyMediaControlBatch { .. }
405            | Self::RouteControl { .. } => None,
406        }
407    }
408
409    pub(super) fn may_change_demux_topology(&self) -> bool {
410        match self {
411            Self::CreateInitialSessionOffer { .. }
412            | Self::ApplySessionAnswer { .. }
413            | Self::CloseSession { .. } => true,
414            Self::CreateSessionRenegotiationOffer { .. }
415            | Self::ActiveSpeakerSourceSnapshot { .. }
416            | Self::SourceDiagnosticsSnapshot { .. }
417            | Self::RemoveMedia { .. }
418            | Self::AddRecvMedia { .. }
419            | Self::AddSendMedia { .. }
420            | Self::ApplyMediaControlBatch { .. }
421            | Self::RouteControl { .. } => false,
422            #[cfg(test)]
423            Self::ResolveMediaMid { .. } => false,
424            #[cfg(test)]
425            Self::ResolveNegotiatedProducerParameters { .. } => false,
426        }
427    }
428}