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, mem};
4
5use super::{
6    action::{ConsumerPacketSelectionUpdate, FeaturedUserUpdate, RouteBudgetOutcome},
7    audio,
8    input::SourcePolicySnapshot,
9    video,
10};
11use crate::engine::{
12    media_transport::{
13        ActiveSpeakerSource, MediaTransport, ReceiverBandwidthSnapshot, ReceiverBweTargetUpdate,
14        TransportBitrateSnapshot,
15    },
16    metrics::{self, BudgetSolverOutcome},
17    room::{
18        Room, RoomEventMessage, effects::transport::RoomRouteEffects, outbound::MessageFanout,
19        state::RoomState,
20    },
21};
22
23/// Deferred request to recompute one room's source policy.
24///
25/// `RoomEffects` decides whether the turn runs before or after its transport
26/// work. [`Self::execute`] serializes it with publication activity.
27#[derive(Debug, Default)]
28pub struct SourcePolicyTurn {
29    requested: bool,
30}
31
32impl SourcePolicyTurn {
33    pub const fn packet_selection() -> Self {
34        Self { requested: true }
35    }
36
37    pub fn request(&mut self) {
38        self.requested = true;
39    }
40
41    pub async fn execute(
42        self,
43        room: &Room,
44        media_transport: Option<&MediaTransport>,
45        active_speaker_sources: Option<&[ActiveSpeakerSource]>,
46    ) {
47        if !self.requested {
48            return;
49        }
50        let _guard = room.source_policy_turn.lock().await;
51        self.execute_guarded(room, media_transport, active_speaker_sources)
52            .await;
53    }
54
55    pub(in crate::engine::room) async fn execute_guarded(
56        self,
57        room: &Room,
58        media_transport: Option<&MediaTransport>,
59        active_speaker_sources: Option<&[ActiveSpeakerSource]>,
60    ) {
61        self.execute_observed(room, media_transport, active_speaker_sources, None)
62            .await;
63    }
64
65    async fn execute_observed(
66        self,
67        room: &Room,
68        media_transport: Option<&MediaTransport>,
69        active_speaker_sources: Option<&[ActiveSpeakerSource]>,
70        bandwidth: Option<&ReceiverBandwidthSnapshot>,
71    ) -> bool {
72        if !self.requested {
73            return false;
74        }
75        let Some(media_transport) = media_transport else {
76            return false;
77        };
78        let transaction = if let Some(sources) = active_speaker_sources {
79            run_packet_selection(room, sources, media_transport, bandwidth).await
80        } else {
81            let sources = media_transport.active_speaker_source_snapshot().await;
82            run_packet_selection(room, &sources, media_transport, bandwidth).await
83        };
84        let Some(transaction) = transaction else {
85            return false;
86        };
87        transaction.commit(room, media_transport).await;
88        true
89    }
90}
91
92#[cfg(feature = "internal-benchmarks")]
93pub async fn run_source_policy_turn_for_benchmark(
94    room: &Room,
95    media_transport: &MediaTransport,
96    bandwidth: &ReceiverBandwidthSnapshot,
97) -> bool {
98    let _guard = room.source_policy_turn.lock().await;
99    SourcePolicyTurn::packet_selection()
100        .execute_observed(room, Some(media_transport), None, Some(bandwidth))
101        .await
102}
103
104async fn run_packet_selection(
105    room: &Room,
106    active_speakers: &[ActiveSpeakerSource],
107    media_transport: &MediaTransport,
108    bandwidth_override: Option<&ReceiverBandwidthSnapshot>,
109) -> Option<SourcePolicyTransaction> {
110    let sessions = {
111        let state = room.state.read().await;
112        state
113            .transport_user_entries()
114            .map(|(user_id, connection_id)| state.transport_user_key(user_id, connection_id))
115            .collect::<Vec<_>>()
116    };
117    let receiver_bandwidth = bandwidth_override.map_or_else(
118        || Cow::Owned(media_transport.receiver_bandwidth_snapshot(&sessions)),
119        Cow::Borrowed,
120    );
121    let source_bitrate = media_transport.transport_bitrate_snapshot(&sessions);
122    let state = room.state.read().await;
123    SourcePolicyTransaction::plan(
124        &state,
125        active_speakers,
126        &receiver_bandwidth,
127        &source_bitrate,
128    )
129}
130
131#[derive(Debug, Default)]
132pub(in crate::engine::room) struct SourcePolicyTransaction {
133    route_effects: RoomRouteEffects,
134    state_updates: Vec<ConsumerPacketSelectionUpdate>,
135    featured_users: Vec<FeaturedUserUpdate>,
136}
137
138impl SourcePolicyTransaction {
139    pub(in crate::engine::room) fn plan(
140        state: &RoomState,
141        active_speakers: &[ActiveSpeakerSource],
142        receiver_bandwidth: &ReceiverBandwidthSnapshot,
143        source_bitrate: &TransportBitrateSnapshot,
144    ) -> Option<Self> {
145        let mut input = SourcePolicySnapshot::from_state(
146            state,
147            active_speakers,
148            receiver_bandwidth,
149            source_bitrate,
150        );
151        let mut tx = Self::default();
152        audio::append_audio_route_activity(&mut tx, &input);
153        let receiver_bwe_targets = mem::take(&mut input.receiver_bwe_targets);
154        video::append_receiver_video_policy(&mut tx, state, &input, receiver_bwe_targets);
155        tx.featured_users = input.featured_user_updates;
156        (!tx.is_empty()).then_some(tx)
157    }
158
159    pub(super) fn push_state_update(&mut self, update: ConsumerPacketSelectionUpdate) {
160        self.state_updates.push(update);
161    }
162
163    pub(super) fn push_route_update(&mut self, update: ConsumerPacketSelectionUpdate) {
164        self.route_effects.source_policy_update(update);
165    }
166
167    pub(super) fn set_receiver_bwe_targets(&mut self, targets: Vec<ReceiverBweTargetUpdate>) {
168        self.route_effects.set_receiver_bwe_targets(targets);
169    }
170
171    async fn commit(self, room: &Room, media_transport: &MediaTransport) {
172        let mut state_updates = self.state_updates;
173        if !self.route_effects.is_empty() {
174            // Only accepted transport controls join state-only updates. Room
175            // state must not claim a selection that its worker rejected.
176            state_updates.extend(
177                self.route_effects
178                    .execute(room.uuid(), media_transport)
179                    .await,
180            );
181        }
182        // Only committed nonzero counters schedule another observation. Rejected
183        // or topology-stale work must not advance adaptation hysteresis.
184        if commit_accepted_updates(room, &state_updates, &self.featured_users).await {
185            media_transport.schedule_source_policy_follow_up(room.instance_id());
186        }
187    }
188
189    #[cfg(test)]
190    pub(in crate::engine::room) async fn execute(
191        self,
192        room: &Room,
193        media_transport: &MediaTransport,
194    ) {
195        self.commit(room, media_transport).await;
196    }
197
198    fn is_empty(&self) -> bool {
199        self.state_updates.is_empty()
200            && self.route_effects.is_empty()
201            && self.featured_users.is_empty()
202    }
203}
204
205async fn commit_accepted_updates(
206    room: &Room,
207    state_updates: &[ConsumerPacketSelectionUpdate],
208    featured_users: &[FeaturedUserUpdate],
209) -> bool {
210    if state_updates.is_empty() && featured_users.is_empty() {
211        return false;
212    }
213    record_source_selection_metrics(room, state_updates);
214    let (info_fanout, requires_follow_up) = {
215        let mut state = room.state.write().await;
216        let requires_follow_up = commit_packet_updates(&mut state, state_updates);
217        let info_fanout = commit_featured_user_updates(&mut state, featured_users);
218        drop(state);
219        (info_fanout, requires_follow_up)
220    };
221    if let Some(info_fanout) = info_fanout {
222        info_fanout.emit();
223    }
224    requires_follow_up
225}
226
227fn record_source_selection_metrics(room: &Room, updates: &[ConsumerPacketSelectionUpdate]) {
228    for update in updates {
229        if update.packet_gate.is_some() {
230            room.metrics
231                .record_source_selection_update(metrics::source_selection_kind(update.selector));
232        }
233        if let Some(outcome) = update.outcome {
234            room.metrics.record_budget_solver_outcome(match outcome {
235                RouteBudgetOutcome::Degraded => BudgetSolverOutcome::Degraded,
236                RouteBudgetOutcome::Paused => BudgetSolverOutcome::Paused,
237                RouteBudgetOutcome::Resumed => BudgetSolverOutcome::Resumed,
238            });
239        }
240    }
241}
242
243fn commit_packet_updates(state: &mut RoomState, updates: &[ConsumerPacketSelectionUpdate]) -> bool {
244    let mut requires_follow_up = false;
245    for update in updates {
246        let committed = state.topology.update_consumer_source_selection(
247            &update.key,
248            update.source_id,
249            &update.route,
250            |selection| {
251                selection.set_selector(update.selector);
252                selection.set_policy_pause_reason(update.policy_pause_reason);
253                selection.set_budget(update.budget);
254                selection.set_adaptation_observations(
255                    update.pressure_observations,
256                    update.upgrade_observations,
257                );
258            },
259        );
260        requires_follow_up |= committed && update.requires_follow_up();
261    }
262    requires_follow_up
263}
264
265fn commit_featured_user_updates(
266    state: &mut RoomState,
267    updates: &[FeaturedUserUpdate],
268) -> Option<MessageFanout> {
269    let mut changed_user_ids = Vec::new();
270    for update in updates {
271        let Some(user) = state.user_mut_for_connection(&update.user_id, update.connection_id)
272        else {
273            continue;
274        };
275        if user.featured() == update.featured {
276            continue;
277        }
278        user.set_featured(update.featured);
279        changed_user_ids.push(update.user_id.clone());
280    }
281    if changed_user_ids.is_empty() {
282        return None;
283    }
284    let snapshot = changed_user_ids
285        .into_iter()
286        .filter_map(|user_id| state.user_info_snapshot(&user_id))
287        .collect();
288    Some(state.fanout_all(&RoomEventMessage::UserInfoChanged(snapshot)))
289}