Skip to main content

o_sfu_core/engine/room/source_policy/
turn.rs

1//! source-policy apply ownership without transport awaits under the room lock
2
3use 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/// Deferred request to recompute one room's source policy.
36///
37/// `RoomEffects` decides whether the turn runs before or after its transport
38/// work. [`Self::execute`] serializes it with publication activity.
39#[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/// Runs policy with injected bandwidth and optional prebuilt speaker observations.
140#[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        // One eligible video route can contribute at most one state update.
237        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            // Only accepted transport controls join state-only updates. Room
279            // state must not claim a selection that its worker rejected.
280            route_effects.execute(room.uuid(), media_transport).await
281        };
282        // Rejections gate only their receiver's temporal state. PR1 budget
283        // reconciliation still observes acceptance across the whole transaction.
284        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    // The allocation revision covers the captured subscriptions, sources and
440    // exact routes. Check it before selection writes advance it within this turn.
441    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            // Recovery ends continuity even when a sibling control rejects.
449            if timing.soft_pause_deadline.is_none() || accepted {
450                user.video_soft_pause_deadline = timing.soft_pause_deadline;
451            }
452            // A rejected sibling cannot cancel an already committed future hold.
453            // Newly proposed deadlines still require acceptance and due retries
454            // remain unscheduled.
455            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        // Cancellation follows eligibility, not transport acceptance. Rejected
466        // downsteps must not leave an old upgrade eligible through an interruption.
467        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    // Async route controls may accept only part of `plan`. Rebuild shared
559    // diagnostics from captured bitrates because ridless observations are
560    // unavailable after the transport await. If any participating route no longer
561    // matches its captured or planned allocation, retain the previous diagnostics
562    // rather than publish a partial receiver-wide view.
563    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}