o_sfu_core/engine/room/source_policy/
turn.rs1use 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#[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 state_updates.extend(
177 self.route_effects
178 .execute(room.uuid(), media_transport)
179 .await,
180 );
181 }
182 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}