Skip to main content

o_sfu_core/engine/media_transport/
route_control.rs

1use std::{collections::BTreeMap, mem::take};
2
3use tracing::{debug, warn};
4
5use super::{
6    ConsumerActivity, MediaTransport, ReceiverBweTargetUpdate, SourceActivityUpdate,
7    SourcePacketGate, TransportAdapterError, TransportConsumerRoute, TransportResult,
8    TransportSourceKey,
9    rtc::{
10        PacketLayerGate, RtcWorkerCommand, WorkerMediaControlBatch as WorkerBatch,
11        WorkerMediaControlBatchOutcome as WorkerBatchOutcome,
12    },
13};
14use crate::MediaWorkerId;
15
16/// Maximum mutations in one synchronous worker command.
17///
18/// Larger plans are split so one room effect cannot monopolize control
19/// dispatch. The packet loop runs its pump phase between chunks.
20const MAX_MEDIA_CONTROL_BATCH_ITEMS: usize = 64;
21const UNAVAILABLE: TransportAdapterError = TransportAdapterError::TransportUnavailable;
22
23#[derive(Debug)]
24pub(crate) struct TransportSourceActivityEffect {
25    pub(crate) source: TransportSourceKey,
26    pub(crate) target_media_worker_id: MediaWorkerId,
27    pub(crate) update: SourceActivityUpdate,
28}
29
30#[derive(Debug)]
31#[must_use = "media-control plans must be executed or intentionally dropped"]
32pub(crate) struct MediaControlPlan<P = (), C = ()> {
33    producers: Vec<(ProducerRouteControl, P)>,
34    consumers: Vec<(ConsumerRouteControl, C)>,
35    receiver_bwe_targets: Vec<ReceiverBweTargetUpdate>,
36}
37
38impl<P, C> MediaControlPlan<P, C> {
39    pub(crate) fn push_producer(
40        &mut self,
41        source: TransportSourceKey,
42        update: SourceActivityUpdate,
43        finish: P,
44    ) {
45        self.producers
46            .push((ProducerRouteControl { source, update }, finish));
47    }
48
49    pub(crate) fn push_consumer(&mut self, control: ConsumerRouteControl, finish: C) {
50        assert!(
51            control.packet_gate.is_some() || control.activity.is_some() || control.request_keyframe,
52            "consumer media control must contain an operation"
53        );
54        self.consumers.push((control, finish));
55    }
56
57    pub(crate) fn set_receiver_bwe_targets(&mut self, updates: Vec<ReceiverBweTargetUpdate>) {
58        self.receiver_bwe_targets = updates;
59    }
60
61    pub(crate) fn append(&mut self, other: Self) {
62        self.producers.extend(other.producers);
63        self.consumers.extend(other.consumers);
64        self.receiver_bwe_targets.extend(other.receiver_bwe_targets);
65    }
66
67    pub(crate) fn is_empty(&self) -> bool {
68        self.producers.is_empty()
69            && self.consumers.is_empty()
70            && self.receiver_bwe_targets.is_empty()
71    }
72}
73
74impl<P, C> Default for MediaControlPlan<P, C> {
75    fn default() -> Self {
76        Self {
77            producers: Vec::new(),
78            consumers: Vec::new(),
79            receiver_bwe_targets: Vec::new(),
80        }
81    }
82}
83
84#[derive(Debug)]
85pub(in crate::engine::media_transport) struct ProducerRouteControl {
86    pub(in crate::engine::media_transport) source: TransportSourceKey,
87    pub(in crate::engine::media_transport) update: SourceActivityUpdate,
88}
89
90#[derive(Debug)]
91pub(crate) struct ConsumerRouteControl {
92    pub(in crate::engine::media_transport) route: TransportConsumerRoute,
93    packet_gate: Option<SourcePacketGate>,
94    pub(in crate::engine::media_transport) activity: Option<ConsumerActivity>,
95    pub(in crate::engine::media_transport) request_keyframe: bool,
96}
97
98impl ConsumerRouteControl {
99    pub(crate) const fn new(route: TransportConsumerRoute) -> Self {
100        Self {
101            route,
102            packet_gate: None,
103            activity: None,
104            request_keyframe: false,
105        }
106    }
107
108    pub(crate) fn packet_gate(mut self, packet_gate: SourcePacketGate) -> Self {
109        self.packet_gate = Some(packet_gate);
110        self
111    }
112
113    pub(crate) const fn activity(mut self, activity: ConsumerActivity) -> Self {
114        self.activity = Some(activity);
115        self
116    }
117
118    pub(crate) const fn request_keyframe(mut self, request: bool) -> Self {
119        self.request_keyframe = request;
120        self
121    }
122
123    fn failure(&self, error: TransportAdapterError) -> ConsumerRouteControlOutcome {
124        let failure = match (self.packet_gate.is_some(), self.activity.is_some()) {
125            (true, _) => ConsumerRouteControlFailure::PacketGate(error),
126            (false, true) => ConsumerRouteControlFailure::Activity(error),
127            (false, false) => ConsumerRouteControlFailure::Keyframe(error),
128        };
129        ConsumerRouteControlOutcome(Some(failure))
130    }
131}
132
133#[derive(Debug)]
134pub(crate) struct MediaControlOutcome<P, C> {
135    pub(crate) producers: Vec<(P, TransportResult<()>)>,
136    pub(crate) consumers: Vec<(C, ConsumerRouteControlOutcome)>,
137}
138
139#[derive(Debug, Clone, Copy, PartialEq, Eq)]
140pub(in crate::engine::media_transport) enum ConsumerRouteControlFailure {
141    PacketGate(TransportAdapterError),
142    Activity(TransportAdapterError),
143    Keyframe(TransportAdapterError),
144}
145
146#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
147pub(crate) struct ConsumerRouteControlOutcome(
148    pub(in crate::engine::media_transport) Option<ConsumerRouteControlFailure>,
149);
150
151impl ConsumerRouteControlOutcome {
152    pub(crate) const fn packet_gate_failed(self) -> bool {
153        matches!(self.0, Some(ConsumerRouteControlFailure::PacketGate(_)))
154    }
155
156    pub(crate) const fn activity_failed(self) -> bool {
157        matches!(self.0, Some(ConsumerRouteControlFailure::Activity(_)))
158    }
159
160    pub(crate) const fn keyframe_failed(self) -> bool {
161        matches!(self.0, Some(ConsumerRouteControlFailure::Keyframe(_)))
162    }
163
164    pub(crate) fn error(self) -> Option<TransportAdapterError> {
165        self.0.map(|failure| match failure {
166            ConsumerRouteControlFailure::PacketGate(error)
167            | ConsumerRouteControlFailure::Activity(error)
168            | ConsumerRouteControlFailure::Keyframe(error) => error,
169        })
170    }
171
172    #[cfg(test)]
173    pub(crate) const fn keyframe_error(error: TransportAdapterError) -> Self {
174        Self(Some(ConsumerRouteControlFailure::Keyframe(error)))
175    }
176}
177
178/// Returns one result per planned item or fails the whole batch.
179///
180/// A mismatched response shape has no safe positional mapping back to room
181/// completions.
182pub(super) fn reconcile_applied(
183    response: TransportResult<WorkerBatchOutcome>,
184    expected: usize,
185) -> Vec<TransportResult<()>> {
186    match response {
187        Ok(WorkerBatchOutcome::Applied(results)) if results.len() == expected => results,
188        Err(error) => vec![Err(error); expected],
189        _ => vec![Err(UNAVAILABLE); expected],
190    }
191}
192
193type WorkerBatches<T> = BTreeMap<usize, Vec<(usize, T)>>;
194type GateUpdate = (usize, TransportConsumerRoute, PacketLayerGate);
195
196fn push<K: Ord, T>(batches: &mut BTreeMap<K, Vec<T>>, key: K, item: T) {
197    batches.entry(key).or_default().push(item);
198}
199
200fn into_batches<T>(mut items: Vec<T>) -> impl Iterator<Item = Vec<T>> {
201    let mut batches = Vec::with_capacity(items.len().div_ceil(MAX_MEDIA_CONTROL_BATCH_ITEMS));
202    // Split off the tail so each split_off moves at most one batch. rev() below
203    // restores source order.
204    while items.len() > MAX_MEDIA_CONTROL_BATCH_ITEMS {
205        let tail_len = (items.len() - 1) % MAX_MEDIA_CONTROL_BATCH_ITEMS + 1;
206        batches.push(items.split_off(items.len() - tail_len));
207    }
208    batches.push(items);
209    batches.into_iter().rev().filter(|batch| !batch.is_empty())
210}
211
212struct MediaControlExecution<P, C> {
213    receiver_bwe: WorkerBatches<ReceiverBweTargetUpdate>,
214    producers: WorkerBatches<ProducerRouteControl>,
215    producer_results: Vec<(P, TransportResult<()>)>,
216    gates: BTreeMap<(usize, TransportSourceKey), Vec<GateUpdate>>,
217    consumers: WorkerBatches<ConsumerRouteControl>,
218    consumer_results: Vec<(C, ConsumerRouteControlOutcome)>,
219}
220
221#[allow(
222    clippy::indexing_slicing,
223    reason = "plan indexes are created by enumerate alongside their result slots"
224)]
225impl<P, C> MediaControlExecution<P, C> {
226    fn new(plan: MediaControlPlan<P, C>) -> Self {
227        let MediaControlPlan {
228            producers,
229            consumers,
230            receiver_bwe_targets,
231        } = plan;
232        let mut execution = Self {
233            receiver_bwe: BTreeMap::new(),
234            producers: BTreeMap::new(),
235            producer_results: Vec::with_capacity(producers.len()),
236            gates: BTreeMap::new(),
237            consumers: BTreeMap::new(),
238            consumer_results: Vec::with_capacity(consumers.len()),
239        };
240        for (index, update) in receiver_bwe_targets.into_iter().enumerate() {
241            let worker = update.session_key().media_worker_id().as_usize();
242            push(&mut execution.receiver_bwe, worker, (index, update));
243        }
244        for (index, (control, finish)) in producers.into_iter().enumerate() {
245            let worker = control.source.session_key().media_worker_id().as_usize();
246            execution.producer_results.push((finish, Err(UNAVAILABLE)));
247            push(&mut execution.producers, worker, (index, control));
248        }
249        for (index, (mut control, finish)) in consumers.into_iter().enumerate() {
250            // Worker placement is not room authority. Reject cross-room controls
251            // before translating them into worker-local mutations.
252            if !control.route.is_single_room() {
253                execution
254                    .consumer_results
255                    .push((finish, control.failure(TransportAdapterError::InvalidInput)));
256                continue;
257            }
258            execution
259                .consumer_results
260                .push((finish, ConsumerRouteControlOutcome::default()));
261            let worker = control
262                .route
263                .consumer_session_key()
264                .media_worker_id()
265                .as_usize();
266            if let Some(gate) = control.packet_gate.take() {
267                let gate = match gate {
268                    SourcePacketGate::Open => PacketLayerGate::Open,
269                    SourcePacketGate::Rid(rid) => PacketLayerGate::Rid(rid.as_str().into()),
270                };
271                push(
272                    &mut execution.gates,
273                    (worker, control.route.source().clone()),
274                    (index, control.route.clone(), gate),
275                );
276            }
277            if control.activity.is_some() || control.request_keyframe {
278                push(&mut execution.consumers, worker, (index, control));
279            }
280        }
281        execution
282    }
283
284    async fn apply_receiver_bwe(&mut self, transport: &MediaTransport) {
285        for (worker, updates) in take(&mut self.receiver_bwe) {
286            for updates in into_batches(updates) {
287                let details: Vec<_> = updates
288                    .iter()
289                    .map(|(_, update)| (update.session_key().clone(), update.target()))
290                    .collect();
291                let response = transport
292                    .execute_batch(worker, WorkerBatch::ReceiverBwe(updates))
293                    .await;
294                let results = reconcile_applied(response, details.len());
295                for ((session_key, target), result) in details.into_iter().zip(results) {
296                    match result {
297                        Ok(()) => {}
298                        Err(TransportAdapterError::InvalidInput) => debug!(
299                            ?session_key,
300                            target = target.as_bps(),
301                            "media transport skipped a stale receiver BWE target"
302                        ),
303                        Err(error) => warn!(
304                            ?error,
305                            ?session_key,
306                            target = target.as_bps(),
307                            "media transport failed to update a receiver BWE target"
308                        ),
309                    }
310                }
311            }
312        }
313    }
314
315    async fn apply_producers(&mut self, transport: &MediaTransport) {
316        for (worker, controls) in take(&mut self.producers) {
317            for controls in into_batches(controls) {
318                let indexes: Vec<_> = controls.iter().map(|(index, _)| *index).collect();
319                let response = transport
320                    .execute_batch(worker, WorkerBatch::ProducerActivity(controls))
321                    .await;
322                let results = reconcile_applied(response, indexes.len());
323                for (index, result) in indexes.into_iter().zip(results) {
324                    self.producer_results[index].1 = result;
325                }
326            }
327        }
328    }
329
330    async fn apply_gates(&mut self, transport: &MediaTransport) {
331        for ((worker, source), updates) in take(&mut self.gates) {
332            let indexes: Vec<_> = updates.iter().map(|(index, _, _)| *index).collect();
333            let response = transport
334                .execute_batch(worker, WorkerBatch::ConsumerGates { source, updates })
335                .await;
336            let results = reconcile_applied(response, indexes.len());
337            for (index, result) in indexes.into_iter().zip(results) {
338                if let Err(error) = result {
339                    let failure = ConsumerRouteControlFailure::PacketGate(error);
340                    self.consumer_results[index].1 = ConsumerRouteControlOutcome(Some(failure));
341                }
342            }
343        }
344    }
345
346    async fn apply_consumers(&mut self, transport: &MediaTransport) {
347        let results = &self.consumer_results;
348        // Room policy cannot commit a selection whose gate was rejected. Skip its
349        // dependent activity and keyframe phases.
350        for controls in self.consumers.values_mut() {
351            controls.retain(
352                |(index, _)| matches!(results.get(*index), Some((_, result)) if result.error().is_none()),
353            );
354        }
355        for (worker, controls) in take(&mut self.consumers) {
356            for controls in into_batches(controls) {
357                let indexes: Vec<_> = controls
358                    .iter()
359                    .map(|(index, control)| (*index, control.activity.is_some()))
360                    .collect();
361                let response = transport
362                    .execute_batch(worker, WorkerBatch::ConsumerFollowUp(controls))
363                    .await;
364                let results = match response {
365                    Ok(WorkerBatchOutcome::Consumers(results))
366                        if results.len() == indexes.len() =>
367                    {
368                        results
369                    }
370                    response => {
371                        let error = response.err().unwrap_or(UNAVAILABLE);
372                        indexes
373                            .iter()
374                            .map(|(_, has_activity)| {
375                                let failure = if *has_activity {
376                                    ConsumerRouteControlFailure::Activity(error)
377                                } else {
378                                    ConsumerRouteControlFailure::Keyframe(error)
379                                };
380                                ConsumerRouteControlOutcome(Some(failure))
381                            })
382                            .collect()
383                    }
384                };
385                for ((index, _), result) in indexes.into_iter().zip(results) {
386                    self.consumer_results[index].1 = result;
387                }
388            }
389        }
390    }
391}
392
393impl MediaTransport {
394    async fn execute_batch(
395        &self,
396        worker_index: usize,
397        batch: WorkerBatch,
398    ) -> TransportResult<WorkerBatchOutcome> {
399        let worker = self.worker_for_index(worker_index).ok_or(UNAVAILABLE)?;
400        #[cfg(test)]
401        self.observe_media_control_batch(worker_index, &batch);
402        worker
403            .request_worker(|response| RtcWorkerCommand::ApplyMediaControlBatch { batch, response })
404            .await
405    }
406
407    /// Applies a media-control plan and returns completions in insertion order.
408    pub(crate) async fn apply_media_control<P, C>(
409        &self,
410        plan: MediaControlPlan<P, C>,
411    ) -> MediaControlOutcome<P, C> {
412        let mut execution = MediaControlExecution::new(plan);
413        execution.apply_receiver_bwe(self).await;
414        execution.apply_producers(self).await;
415        // Room policy may commit only a selection whose gate and activity reached
416        // worker state. Apply gates before their dependent activity and keyframe work.
417        execution.apply_gates(self).await;
418        execution.apply_consumers(self).await;
419        MediaControlOutcome {
420            producers: execution.producer_results,
421            consumers: execution.consumer_results,
422        }
423    }
424}