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
16const 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
178pub(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 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 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 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 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 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}