1use std::{borrow::Cow, collections::BTreeMap, time::Instant};
4
5use o_sfu_router::MediaKind;
6use o_sfu_telemetry::schema::event as telemetry_event;
7use tracing::info;
8
9use super::{
10 action::{
11 ConsumerPacketSelectionUpdate, FeaturedUserUpdate, ReceiverPolicyTiming,
12 ReceiverVideoBudgetPlan, UpgradeChange, VideoRouteTransition,
13 },
14 audio,
15 input::SourcePolicySnapshot,
16 video,
17};
18use crate::engine::{
19 ConnectionId,
20 media_transport::{
21 ActiveSpeakerSource, MediaTransport, ReceiverBandwidthSnapshot, ReceiverBweTargetUpdate,
22 TransportBitrateSnapshot,
23 },
24 metrics::{self, BudgetSolverOutcome},
25 room::{
26 Room, RoomEventMessage, SourcePolicyGuard, effects::transport::RoomRouteEffects,
27 outbound::MessageFanout, state::RoomState,
28 },
29 source_model::{
30 PolicyPauseReason, ReceiverVideoBudgetDiagnostics, SourceAdaptationPolicy,
31 SourceEncodingId, SourceSelector,
32 },
33};
34
35#[derive(Debug, Default)]
40pub struct SourcePolicyTurn {
41 requested: bool,
42}
43
44impl SourcePolicyTurn {
45 pub const fn packet_selection() -> Self {
46 Self { requested: true }
47 }
48
49 pub fn request(&mut self) {
50 self.requested = true;
51 }
52
53 pub async fn execute(
54 self,
55 room: &Room,
56 media_transport: Option<&MediaTransport>,
57 active_speaker_sources: Option<&[ActiveSpeakerSource]>,
58 ) {
59 if !self.requested {
60 return;
61 }
62 let guard = room.lock_source_policy().await;
63 self.execute_guarded(&guard, media_transport, active_speaker_sources)
64 .await;
65 }
66
67 pub(in crate::engine::room) async fn execute_guarded(
68 self,
69 guard: &SourcePolicyGuard<'_>,
70 media_transport: Option<&MediaTransport>,
71 active_speaker_sources: Option<&[ActiveSpeakerSource]>,
72 ) {
73 self.execute_observed(
74 guard,
75 media_transport,
76 active_speaker_sources,
77 None,
78 Instant::now(),
79 )
80 .await;
81 }
82
83 async fn execute_observed(
84 self,
85 guard: &SourcePolicyGuard<'_>,
86 media_transport: Option<&MediaTransport>,
87 active_speaker_sources: Option<&[ActiveSpeakerSource]>,
88 bandwidth: Option<&ReceiverBandwidthSnapshot>,
89 now: Instant,
90 ) -> bool {
91 if !self.requested {
92 return false;
93 }
94 let Some(media_transport) = media_transport else {
95 return false;
96 };
97 let room = guard.room();
98 let transaction = if let Some(sources) = active_speaker_sources {
99 run_packet_selection(room, sources, media_transport, bandwidth, now).await
100 } else {
101 let worker_indices = room.source_policy_worker_indices(media_transport).await;
102 let Ok(sources) = media_transport
103 .active_speaker_source_snapshot_for_workers(&worker_indices)
104 .await
105 else {
106 media_transport.schedule_source_policy_retry(room.instance_id());
107 return false;
108 };
109 run_packet_selection(room, &sources, media_transport, bandwidth, now).await
110 };
111 let Some(transaction) = transaction else {
112 media_transport.set_source_policy_deadline(room.instance_id(), None);
113 return false;
114 };
115 transaction.commit(room, media_transport).await;
116 true
117 }
118}
119
120impl Room {
121 pub(in crate::engine::room) async fn source_policy_worker_indices(
122 &self,
123 media_transport: &MediaTransport,
124 ) -> Vec<usize> {
125 let state = self.state.read().await;
126 let mut indices = state
127 .topology
128 .published_sources()
129 .map(|source| source.transport.session_key().media_worker_id().as_usize())
130 .filter(|&worker| media_transport.worker_is_usable(worker))
131 .collect::<Vec<_>>();
132 drop(state);
133 indices.sort_unstable();
134 indices.dedup();
135 indices
136 }
137}
138
139#[cfg(feature = "internal-benchmarks")]
141pub async fn run_source_policy_turn_for_benchmark(
142 room: &Room,
143 media_transport: &MediaTransport,
144 bandwidth: &ReceiverBandwidthSnapshot,
145 active_speakers: Option<&[ActiveSpeakerSource]>,
146 now: Instant,
147) -> bool {
148 let guard = room.lock_source_policy().await;
149 SourcePolicyTurn::packet_selection()
150 .execute_observed(
151 &guard,
152 Some(media_transport),
153 active_speakers,
154 Some(bandwidth),
155 now,
156 )
157 .await
158}
159
160async fn run_packet_selection(
161 room: &Room,
162 active_speakers: &[ActiveSpeakerSource],
163 media_transport: &MediaTransport,
164 bandwidth_override: Option<&ReceiverBandwidthSnapshot>,
165 now: Instant,
166) -> Option<SourcePolicyTransaction> {
167 let sessions = {
168 let state = room.state.read().await;
169 state
170 .transport_user_entries()
171 .map(|(user_id, connection_id)| state.transport_user_key(user_id, connection_id))
172 .collect::<Vec<_>>()
173 };
174 let receiver_bandwidth = bandwidth_override.map_or_else(
175 || Cow::Owned(media_transport.receiver_bandwidth_snapshot(&sessions)),
176 Cow::Borrowed,
177 );
178 let source_bitrate = media_transport.transport_bitrate_snapshot(&sessions);
179 let state = room.state.read().await;
180 SourcePolicyTransaction::plan(
181 &state,
182 active_speakers,
183 &receiver_bandwidth,
184 &source_bitrate,
185 now,
186 )
187}
188
189#[derive(Debug)]
190pub(in crate::engine::room) struct SourcePolicyTransaction {
191 route_effects: RoomRouteEffects,
192 state_updates: Vec<ConsumerPacketSelectionUpdate>,
193 receiver_video_budget_plans: Vec<ReceiverVideoBudgetPlan>,
194 featured_users: Vec<FeaturedUserUpdate>,
195 video_allocation_revision: u64,
196 outstanding_controls: BTreeMap<ConnectionId, usize>,
197 receiver_timing: Vec<ReceiverPolicyTiming>,
198 planned_at: Instant,
199}
200
201impl SourcePolicyTransaction {
202 pub(in crate::engine::room) fn plan(
203 state: &RoomState,
204 active_speakers: &[ActiveSpeakerSource],
205 receiver_bandwidth: &ReceiverBandwidthSnapshot,
206 source_bitrate: &TransportBitrateSnapshot,
207 now: Instant,
208 ) -> Option<Self> {
209 let input = SourcePolicySnapshot::from_state(
210 state,
211 active_speakers,
212 receiver_bandwidth,
213 source_bitrate,
214 );
215 let mut tx = Self {
216 video_allocation_revision: state.topology.video_allocation_revision(),
217 route_effects: RoomRouteEffects::default(),
218 state_updates: Vec::new(),
219 receiver_video_budget_plans: Vec::new(),
220 featured_users: Vec::new(),
221 outstanding_controls: BTreeMap::new(),
222 receiver_timing: Vec::new(),
223 planned_at: now,
224 };
225 audio::append_audio_route_activity(&mut tx, &input);
226 video::append_receiver_video_policy(&mut tx, state, &input, now);
227 tx.featured_users = input.featured_user_updates;
228 (!tx.is_empty()).then_some(tx)
229 }
230
231 pub(super) fn push_state_update(
232 &mut self,
233 update: ConsumerPacketSelectionUpdate,
234 max_video_updates: usize,
235 ) {
236 if self.state_updates.is_empty() {
238 self.state_updates.reserve(max_video_updates);
239 }
240 self.state_updates.push(update);
241 }
242
243 pub(super) fn push_route_update(&mut self, update: ConsumerPacketSelectionUpdate) {
244 self.expect_route_update(update.route.consumer_session_key().connection_id());
245 self.route_effects.source_policy_update(update);
246 }
247
248 pub(super) fn expect_route_update(&mut self, connection: ConnectionId) {
249 *self.outstanding_controls.entry(connection).or_default() += 1;
250 }
251
252 pub(super) fn push_receiver_timing(&mut self, timing: ReceiverPolicyTiming) {
253 self.receiver_timing.push(timing);
254 }
255
256 pub(super) fn push_receiver_video_budget_plan(&mut self, plan: ReceiverVideoBudgetPlan) {
257 self.receiver_video_budget_plans.push(plan);
258 }
259
260 pub(super) fn set_receiver_bwe_targets(&mut self, targets: Vec<ReceiverBweTargetUpdate>) {
261 self.route_effects.set_receiver_bwe_targets(targets);
262 }
263
264 async fn commit(self, room: &Room, media_transport: &MediaTransport) {
265 let Self {
266 route_effects,
267 state_updates,
268 receiver_video_budget_plans,
269 featured_users,
270 video_allocation_revision,
271 mut outstanding_controls,
272 receiver_timing,
273 planned_at,
274 } = self;
275 let accepted_route_updates = if route_effects.is_empty() {
276 Vec::new()
277 } else {
278 route_effects.execute(room.uuid(), media_transport).await
281 };
282 for update in &accepted_route_updates {
285 if let Some(remaining) =
286 outstanding_controls.get_mut(&update.route.consumer_session_key().connection_id())
287 {
288 *remaining -= 1;
289 }
290 }
291 let mut updates = [state_updates, accepted_route_updates];
292 if updates.iter().all(Vec::is_empty)
293 && receiver_video_budget_plans.is_empty()
294 && featured_users.is_empty()
295 && receiver_timing.is_empty()
296 {
297 media_transport.set_source_policy_deadline(room.instance_id(), None);
298 return;
299 }
300 let mut state = room.state.write().await;
301 let deadline = commit_packet_updates(
302 &mut state,
303 &mut updates,
304 &receiver_video_budget_plans,
305 video_allocation_revision,
306 &outstanding_controls,
307 &receiver_timing,
308 planned_at,
309 );
310 let info_fanout = commit_featured_user_updates(&mut state, &featured_users);
311 drop(state);
312 for batch in &updates {
313 record_committed_selection_updates(room, batch);
314 }
315 if let Some(info_fanout) = info_fanout {
316 info_fanout.emit();
317 }
318 media_transport.set_source_policy_deadline(room.instance_id(), deadline);
319 }
320
321 #[cfg(test)]
322 pub(in crate::engine::room) async fn execute(
323 self,
324 room: &Room,
325 media_transport: &MediaTransport,
326 ) {
327 self.commit(room, media_transport).await;
328 }
329
330 fn is_empty(&self) -> bool {
331 self.state_updates.is_empty()
332 && self.route_effects.is_empty()
333 && self.receiver_video_budget_plans.is_empty()
334 && self.featured_users.is_empty()
335 && self.receiver_timing.is_empty()
336 }
337}
338
339fn record_committed_selection_updates(room: &Room, updates: &[ConsumerPacketSelectionUpdate]) {
340 for update in updates {
341 if update.packet_gate.is_some() {
342 room.metrics
343 .record_source_selection_update(metrics::source_selection_kind(update.selector));
344 }
345 let Some(transition) = update.transition else {
346 continue;
347 };
348 let (metric_outcome, outcome, reason) = match transition {
349 VideoRouteTransition::Degraded => (BudgetSolverOutcome::Degraded, "degraded", None),
350 VideoRouteTransition::Paused { reason } => {
351 (BudgetSolverOutcome::Paused, "paused", Some(reason))
352 }
353 VideoRouteTransition::Resumed { cleared_reason } => (
354 BudgetSolverOutcome::Resumed,
355 "resumed",
356 Some(cleared_reason),
357 ),
358 };
359 room.metrics.record_budget_solver_outcome(metric_outcome);
360 let consumer = update.route.consumer_session_key();
361 let source = update.route.source_session_key();
362 info!(
363 event = telemetry_event::SOURCE_POLICY_ROUTE_CHANGED,
364 room_id = room.uuid(),
365 user_id = %consumer.user_id().path_segment(),
366 connection_id = consumer.connection_id().as_u64(),
367 media_worker_id = consumer.media_worker_id().as_usize(),
368 transport_media_id = update.route.consumer_transport_media_id().as_u64(),
369 producer_user_id = %source.user_id().path_segment(),
370 source_transport_media_id = update.route.source_transport_media_id().as_u64(),
371 stream_id = %update.key.stream,
372 outcome,
373 reason = reason.map(policy_pause_reason_name),
374 latest_receiver_bandwidth_estimate_bps = update
375 .planned_budget
376 .latest_receiver_bandwidth()
377 .map(crate::Bitrate::as_bps),
378 selected_video_budget_bps = update
379 .planned_budget
380 .selected_video_budget()
381 .map(crate::Bitrate::as_bps),
382 planned_active_video_route_count = update.planned_budget.active_video_route_count(),
383 planned_selected_video_bitrate_bps = update
384 .planned_budget
385 .selected_video_bitrate()
386 .as_bps(),
387 selector = source_selector_name(update.selector),
388 selected_encoding_id = update
389 .selector
390 .selected_encoding()
391 .map(SourceEncodingId::as_u64),
392 selected_estimated_bitrate_bps = update
393 .selected_estimated_bitrate
394 .map(crate::Bitrate::as_bps),
395 "source policy route changed"
396 );
397 }
398}
399
400const fn source_selector_name(selector: SourceSelector) -> &'static str {
401 match selector {
402 SourceSelector::Open => "open",
403 SourceSelector::Encoding(_) => "encoding",
404 }
405}
406
407const fn policy_pause_reason_name(reason: PolicyPauseReason) -> &'static str {
408 match reason {
409 PolicyPauseReason::BudgetPressure => "budget_pressure",
410 PolicyPauseReason::HiddenTile => "hidden_tile",
411 PolicyPauseReason::OverflowTile => "overflow_tile",
412 PolicyPauseReason::MissingUsableLayer => "missing_usable_layer",
413 PolicyPauseReason::AudioSpeakerLimit => "audio_speaker_limit",
414 PolicyPauseReason::ReceiverDeafened => "receiver_deafened",
415 PolicyPauseReason::VideoDownloadLimit => "video_download_limit",
416 PolicyPauseReason::SourceBitrateLimit => "source_bitrate_limit",
417 }
418}
419
420fn commit_packet_updates(
421 state: &mut RoomState,
422 updates: &mut [Vec<ConsumerPacketSelectionUpdate>],
423 receiver_video_budget_plans: &[ReceiverVideoBudgetPlan],
424 video_allocation_revision: u64,
425 outstanding_controls: &BTreeMap<ConnectionId, usize>,
426 receiver_timing: &[ReceiverPolicyTiming],
427 now: Instant,
428) -> Option<Instant> {
429 let allocation_plan_is_current =
430 state.topology.video_allocation_revision() == video_allocation_revision;
431 let all_route_controls_accepted = outstanding_controls.values().all(|count| *count == 0);
432 let reconcile_planned_budgets = !allocation_plan_is_current || !all_route_controls_accepted;
433 let receiver_controls_accepted = |connection| {
434 outstanding_controls
435 .get(&connection)
436 .is_none_or(|count| *count == 0)
437 };
438 let mut next_deadline = None;
439 if allocation_plan_is_current {
442 for timing in receiver_timing {
443 let accepted = receiver_controls_accepted(timing.connection_id);
444 let Some(user) = state.user_mut_for_connection(&timing.receiver, timing.connection_id)
445 else {
446 continue;
447 };
448 if timing.soft_pause_deadline.is_none() || accepted {
450 user.video_soft_pause_deadline = timing.soft_pause_deadline;
451 }
452 let deadline = if accepted {
456 timing.next_deadline
457 } else {
458 timing.retained_deadline
459 };
460 if let Some(deadline) = deadline.filter(|deadline| *deadline > now) {
461 next_deadline =
462 Some(next_deadline.map_or(deadline, |next: Instant| next.min(deadline)));
463 }
464 }
465 for route in receiver_video_budget_plans
468 .iter()
469 .flat_map(|plan| &plan.routes)
470 .filter(|route| route.interrupts_upgrade)
471 {
472 state
473 .topology
474 .update_consumer_upgrade(&route.key, route.source_id, &route.route, None);
475 }
476 for update in updates.iter().flatten() {
477 if update.interrupts_upgrade {
478 state.topology.update_consumer_upgrade(
479 &update.key,
480 update.source_id,
481 &update.route,
482 None,
483 );
484 }
485 }
486 }
487 for batch in updates {
488 batch.retain_mut(|update| {
489 let commit_planned_budget = allocation_plan_is_current
490 && (!reconcile_planned_budgets
491 || !receiver_video_budget_plans
492 .iter()
493 .any(|plan| plan.receiver == update.key.receiver));
494 let committed = state.topology.update_consumer_source_selection(
495 &update.key,
496 update.source_id,
497 &update.route,
498 |selection| {
499 selection.set_selector(update.selector);
500 selection.set_policy_pause_reason(update.policy_pause_reason);
501 if commit_planned_budget {
502 selection.set_budget(update.planned_budget);
503 }
504 },
505 );
506 if committed
507 && update.transition.is_some()
508 && !commit_planned_budget
509 && !route_transition_remains_observable(state, update)
510 {
511 update.transition = None;
512 }
513 if let UpgradeChange::Set(pending_upgrade) = update.upgrade
514 && committed
515 && allocation_plan_is_current
516 && receiver_controls_accepted(update.route.consumer_session_key().connection_id())
517 && state.user_connection_id(&update.key.receiver)
518 == Some(update.route.consumer_session_key().connection_id())
519 {
520 state.topology.update_consumer_upgrade(
521 &update.key,
522 update.source_id,
523 &update.route,
524 pending_upgrade,
525 );
526 }
527 committed
528 });
529 }
530 if reconcile_planned_budgets {
531 for plan in receiver_video_budget_plans {
532 reconcile_receiver_video_budget(state, plan);
533 }
534 }
535 next_deadline
536}
537
538fn route_transition_remains_observable(
539 state: &RoomState,
540 update: &ConsumerPacketSelectionUpdate,
541) -> bool {
542 state.user_connection_id(&update.key.receiver)
543 == Some(update.route.consumer_session_key().connection_id())
544 && state
545 .topology
546 .committed_consumer_route_for_key(&update.key)
547 .is_some_and(|route| {
548 route.source.descriptor.source_id() == update.source_id
549 && route.route == &update.route
550 && route.source.active
551 && route.selection.active()
552 })
553}
554
555fn reconcile_receiver_video_budget(state: &mut RoomState, plan: &ReceiverVideoBudgetPlan) {
556 let receiver = &plan.receiver;
557 let receiver_connection_id = state.user_connection_id(receiver);
558 let mut active_route_count = 0;
564 let mut selected_video_bitrate = crate::Bitrate::zero();
565 let mut update_targets = Vec::with_capacity(plan.routes.len());
566 let mut plan_index = 0;
567 for route in state
568 .topology
569 .committed_consumer_routes_for_user(receiver)
570 .filter(|route| {
571 receiver_connection_id == Some(route.route.consumer_session_key().connection_id())
572 })
573 {
574 let policy = route.source.descriptor.policy();
575 if route.source.descriptor.media_kind() != MediaKind::Video
576 || (policy.adaptation() == SourceAdaptationPolicy::None
577 && policy.video_bitrate_cap().is_none())
578 {
579 continue;
580 }
581 while plan
582 .routes
583 .get(plan_index)
584 .is_some_and(|planned| &planned.key < route.key)
585 {
586 plan_index += 1;
587 }
588 let participating = route.source.active && route.selection.active();
589 let Some(planned) = plan
590 .routes
591 .get(plan_index)
592 .filter(|planned| &planned.key == route.key)
593 else {
594 if participating {
595 return;
596 }
597 continue;
598 };
599 plan_index += 1;
600 if planned.source_id != route.source.descriptor.source_id() || planned.route != *route.route
601 {
602 if participating {
603 return;
604 }
605 continue;
606 }
607 let planned_state = planned.planned;
608 let selected_bitrate = if planned_state.matches_selection(route.selection) {
609 planned_state.selected_bitrate
610 } else if planned.captured.matches_selection(route.selection) {
611 planned.captured.selected_bitrate
612 } else {
613 if participating {
614 return;
615 }
616 continue;
617 };
618 update_targets.push(planned);
619 if participating && route.selection.policy_pause_reason().is_none() {
620 active_route_count += 1;
621 selected_video_bitrate = selected_video_bitrate.saturating_add(selected_bitrate);
622 }
623 }
624 let budget = ReceiverVideoBudgetDiagnostics::new(
625 plan.planned_budget.latest_receiver_bandwidth(),
626 plan.planned_budget.selected_video_budget(),
627 active_route_count,
628 selected_video_bitrate,
629 );
630 for target in update_targets {
631 let updated = state.topology.update_consumer_source_selection(
632 &target.key,
633 target.source_id,
634 &target.route,
635 |selection| selection.set_budget(budget),
636 );
637 debug_assert!(
638 updated,
639 "validated route should remain committed under the lock"
640 );
641 }
642}
643
644fn commit_featured_user_updates(
645 state: &mut RoomState,
646 updates: &[FeaturedUserUpdate],
647) -> Option<MessageFanout> {
648 let mut changed_user_ids = Vec::new();
649 for update in updates {
650 let Some(user) = state.user_mut_for_connection(&update.user_id, update.connection_id)
651 else {
652 continue;
653 };
654 if user.featured() == update.featured {
655 continue;
656 }
657 user.set_featured(update.featured);
658 changed_user_ids.push(update.user_id.clone());
659 }
660 if changed_user_ids.is_empty() {
661 return None;
662 }
663 let snapshot = changed_user_ids
664 .into_iter()
665 .filter_map(|user_id| state.user_info_snapshot(&user_id))
666 .collect();
667 Some(state.fanout_all(RoomEventMessage::UserInfoChanged(snapshot)))
668}