1use o_sfu_router::MediaKind;
2use o_sfu_telemetry::schema::event as telemetry_event;
3use tracing::{info, warn};
4
5use super::receiver_route::ReceiverSetupTurn;
6use crate::engine::{
7 media_transport::{
8 ConsumerActivity, ConsumerRouteControl, ConsumerRouteControlOutcome, MediaControlPlan,
9 MediaTransport, ProducerActivity, ReceiverBweTargetUpdate, SourceActivityUpdate,
10 TransportConsumerRoute, TransportRelayRouteAction, TransportRelayRouteEffect,
11 TransportSourceActivityEffect, TransportSourceKey, TransportTeardown,
12 },
13 room::{
14 Room,
15 media_graph::{
16 ConsumerRouteTarget, ConsumerSetupOrigin, ReceiverRouteActivity, ReceiverRouteWork,
17 },
18 source_policy::ConsumerPacketSelectionUpdate,
19 },
20 source_model::UserStreamId,
21};
22
23#[derive(Debug, Default)]
24pub(in crate::engine::room) struct RoomTransportPlan {
25 relays: Vec<TransportRelayRouteEffect>,
26 remote_source_activity: Vec<TransportSourceActivityEffect>,
27 teardown: Vec<TransportTeardown>,
28 route_control: RoomRouteEffects,
29 setup_turns: Vec<ReceiverSetupTurn>,
30}
31
32impl RoomTransportPlan {
33 pub(in crate::engine::room) fn from_relays_and_teardown(
34 mut relays: Vec<TransportRelayRouteEffect>,
35 additional_teardown: impl IntoIterator<Item = TransportTeardown>,
36 ) -> Self {
37 let mut teardown = Vec::new();
38 extract_relay_teardown(&mut relays, &mut teardown);
39 teardown.extend(additional_teardown);
40 Self {
41 relays,
42 teardown,
43 ..Self::default()
44 }
45 }
46
47 pub(in crate::engine::room) fn extend(&mut self, other: Self) {
48 self.relays.extend(other.relays);
49 self.remote_source_activity
50 .extend(other.remote_source_activity);
51 self.teardown.extend(other.teardown);
52 self.route_control.append(other.route_control);
53 self.setup_turns.extend(other.setup_turns);
54 }
55
56 pub(in crate::engine::room) fn extend_teardown(
57 &mut self,
58 teardown: impl IntoIterator<Item = TransportTeardown>,
59 ) {
60 self.teardown.extend(teardown);
61 }
62
63 pub(super) fn push_producer(
64 &mut self,
65 source: TransportSourceKey,
66 stream_id: UserStreamId,
67 update: SourceActivityUpdate,
68 ) {
69 self.route_control
70 .producer_activity(source, stream_id, update);
71 }
72
73 pub(super) fn extend_remote_source_activity(
74 &mut self,
75 effects: impl IntoIterator<Item = TransportSourceActivityEffect>,
76 ) {
77 self.remote_source_activity.extend(effects);
78 }
79
80 pub(super) fn push_receiver_work(
81 &mut self,
82 mut work: ReceiverRouteWork,
83 origin: ConsumerSetupOrigin,
84 ) {
85 extract_relay_teardown(&mut work.relays, &mut self.teardown);
86 self.relays.extend(work.relays);
87 self.teardown.extend(work.teardown);
88 for activity in work.activities {
89 self.route_control.receiver_activity(activity);
90 }
91 for setup in work.setups {
92 self.setup_turns.push(ReceiverSetupTurn::new(setup, origin));
93 }
94 }
95
96 pub(super) async fn execute(self, room: &Room, media_transport: Option<&MediaTransport>) {
97 let Some(media_transport) = media_transport else {
98 return;
99 };
100 execute_relay_route_effects(media_transport, self.relays).await;
101 execute_remote_source_activity_effects(media_transport, self.remote_source_activity).await;
102 self.route_control
103 .execute(room.uuid(), media_transport)
104 .await;
105 media_transport.teardown(self.teardown).await;
106 for turn in self.setup_turns {
107 turn.execute(room, media_transport).await;
108 }
109 }
110
111 #[cfg(test)]
112 pub(in crate::engine::room) fn relays_and_teardown(
113 &self,
114 ) -> (&[TransportRelayRouteEffect], &[TransportTeardown]) {
115 (&self.relays, &self.teardown)
116 }
117}
118
119#[derive(Debug, Default)]
120#[must_use = "room route effects must be executed or intentionally dropped"]
121pub(in crate::engine::room) struct RoomRouteEffects(
122 MediaControlPlan<ProducerRouteFinish, ConsumerRouteFinish>,
123);
124
125impl RoomRouteEffects {
126 fn producer_activity(
127 &mut self,
128 source: TransportSourceKey,
129 stream_id: UserStreamId,
130 update: SourceActivityUpdate,
131 ) {
132 self.0.push_producer(
133 source.clone(),
134 update,
135 ProducerRouteFinish {
136 source,
137 stream_id,
138 activity: update.activity(),
139 },
140 );
141 }
142
143 fn receiver_activity(&mut self, activity: ReceiverRouteActivity) {
144 let target = activity.target();
145 let active = activity.active();
146 self.0.push_consumer(
147 ConsumerRouteControl::new(target.transport_route().clone())
148 .activity(ConsumerActivity::from_active(active))
149 .request_keyframe(target.request_keyframe_after_activity(active)),
150 ConsumerRouteFinish::Activity(activity),
151 );
152 }
153
154 pub(super) fn setup_activity(
155 &mut self,
156 route: TransportConsumerRoute,
157 kind: MediaKind,
158 active: bool,
159 ) {
160 self.0.push_consumer(
161 ConsumerRouteControl::new(route.clone())
162 .activity(ConsumerActivity::from_active(active))
163 .request_keyframe(active && kind == MediaKind::Video),
164 ConsumerRouteFinish::SetupActivity(route, active),
165 );
166 }
167
168 pub(super) fn keyframe(&mut self, target: ConsumerRouteTarget) {
169 self.0.push_consumer(
170 ConsumerRouteControl::new(target.transport_route().clone()).request_keyframe(true),
171 ConsumerRouteFinish::Keyframe(target),
172 );
173 }
174
175 pub(in crate::engine::room) fn source_policy_update(
176 &mut self,
177 update: ConsumerPacketSelectionUpdate,
178 ) {
179 self.0.push_consumer(
180 update.route_control(),
181 ConsumerRouteFinish::SourcePolicy(update),
182 );
183 }
184
185 pub(in crate::engine::room) fn set_receiver_bwe_targets(
186 &mut self,
187 targets: Vec<ReceiverBweTargetUpdate>,
188 ) {
189 self.0.set_receiver_bwe_targets(targets);
190 }
191
192 fn append(&mut self, other: Self) {
193 self.0.append(other.0);
194 }
195
196 pub(in crate::engine::room) fn is_empty(&self) -> bool {
197 self.0.is_empty()
198 }
199
200 pub(in crate::engine::room) async fn execute(
201 self,
202 room_id: &str,
203 media_transport: &MediaTransport,
204 ) -> Vec<ConsumerPacketSelectionUpdate> {
205 let transport_outcome = media_transport.apply_media_control(self.0).await;
206
207 let mut accepted_policy_updates = Vec::new();
208 for (completion, result) in transport_outcome.producers {
209 if let Err(error) = result {
210 warn!(
211 ?error,
212 source = ?completion.source,
213 active = completion.activity.is_active(),
214 "media transport failed to update producer route activity"
215 );
216 }
217 completion.emit_activity_event(room_id);
218 }
219 for (completion, result) in transport_outcome.consumers {
220 completion.finish(room_id, result, &mut accepted_policy_updates);
221 }
222 accepted_policy_updates
223 }
224}
225
226#[derive(Debug)]
227struct ProducerRouteFinish {
228 source: TransportSourceKey,
229 stream_id: UserStreamId,
230 activity: ProducerActivity,
231}
232
233impl ProducerRouteFinish {
234 fn emit_activity_event(&self, room_id: &str) {
235 let session = self.source.session_key();
236 info!(
237 event = telemetry_event::PUBLICATION_ACTIVITY_CHANGED,
238 room_id,
239 user_id = %session.user_id().path_segment(),
240 connection_id = session.connection_id().as_u64(),
241 media_worker_id = session.media_worker_id().as_usize(),
242 transport_media_id = self.source.transport_media_id().as_u64(),
243 active = self.activity.is_active(),
244 stream_id = %self.stream_id,
245 "publication activity changed"
246 );
247 }
248}
249
250#[derive(Debug)]
251enum ConsumerRouteFinish {
252 Activity(ReceiverRouteActivity),
253 Keyframe(ConsumerRouteTarget),
254 SetupActivity(TransportConsumerRoute, bool),
255 SourcePolicy(ConsumerPacketSelectionUpdate),
256}
257
258impl ConsumerRouteFinish {
259 #[allow(
260 clippy::cognitive_complexity,
261 reason = "closed route completion policy is clearer than one use finish helpers"
262 )]
263 fn finish(
264 self,
265 room_id: &str,
266 result: ConsumerRouteControlOutcome,
267 accepted_policy_updates: &mut Vec<ConsumerPacketSelectionUpdate>,
268 ) {
269 match self {
270 Self::Activity(activity) => {
271 let target = activity.target();
272 if result.activity_failed() {
273 warn!(
274 error = ?result.error(),
275 route = ?target.transport_route(),
276 stream_id = %target.stream_id(),
277 active = activity.active(),
278 "media transport failed to update consumer route activity"
279 );
280 } else if result.keyframe_failed() {
281 warn!(
282 error = ?result.error(),
283 route = ?target.transport_route(),
284 stream_id = %target.stream_id(),
285 "media transport failed to request a consumer keyframe refresh"
286 );
287 }
288 let session = target.transport_route().consumer_session_key();
289 info!(
290 event = telemetry_event::SUBSCRIPTION_ACTIVITY_CHANGED,
291 room_id,
292 user_id = %session.user_id().path_segment(),
293 connection_id = session.connection_id().as_u64(),
294 media_worker_id = session.media_worker_id().as_usize(),
295 transport_media_id = target.consumer_media_id().as_u64(),
296 active = activity.active(),
297 producer_user_id = %target.producer_user_id().path_segment(),
298 source_transport_media_id = target.source_media_id().as_u64(),
299 stream_id = %target.stream_id(),
300 "subscription activity changed"
301 );
302 }
303 Self::Keyframe(target) => {
304 if result.keyframe_failed() {
305 warn!(
306 error = ?result.error(),
307 route = ?target.transport_route(),
308 "media transport failed to request a refreshed consumer keyframe"
309 );
310 }
311 }
312 Self::SetupActivity(route, active) => {
313 if result.activity_failed() {
314 warn!(
315 error = ?result.error(),
316 ?route,
317 active,
318 "media transport failed to correct in-flight consumer setup activity"
319 );
320 } else if result.keyframe_failed() {
321 warn!(
322 error = ?result.error(),
323 ?route,
324 "media transport failed to request keyframe after consumer setup activity correction"
325 );
326 }
327 }
328 Self::SourcePolicy(update) => {
329 if result.packet_gate_failed() || result.activity_failed() {
330 warn!(
331 route = ?update.route,
332 error = ?result.error(),
333 route_active = update.route_active(),
334 "media transport rejected the receiver-driven packet selection update"
335 );
336 return;
337 }
338 if result.keyframe_failed() {
339 warn!(
340 error = ?result.error(),
341 route = ?update.route,
342 "media transport failed to request an adaptation keyframe refresh"
343 );
344 }
345 accepted_policy_updates.push(update);
346 }
347 }
348 }
349}
350
351pub(super) async fn execute_relays_and_teardown(
352 media_transport: &MediaTransport,
353 mut relays: Vec<TransportRelayRouteEffect>,
354 additional_teardown: impl IntoIterator<Item = TransportTeardown>,
355) -> bool {
356 let mut teardown = Vec::new();
357 extract_relay_teardown(&mut relays, &mut teardown);
358 teardown.extend(additional_teardown);
359 let relays_applied = execute_relay_route_effects(media_transport, relays).await;
360 media_transport.teardown(teardown).await;
361 relays_applied
362}
363
364fn extract_relay_teardown(
365 relays: &mut Vec<TransportRelayRouteEffect>,
366 teardown: &mut Vec<TransportTeardown>,
367) {
368 teardown.extend(
369 relays
370 .extract_if(.., |effect| {
371 effect.action == TransportRelayRouteAction::Release
372 })
373 .map(|effect| TransportTeardown::ReleaseRelayRoute {
374 source: effect.source,
375 target_media_worker_id: effect.target_media_worker_id,
376 }),
377 );
378}
379
380async fn execute_relay_route_effects(
381 media_transport: &MediaTransport,
382 effects: impl IntoIterator<Item = TransportRelayRouteEffect>,
383) -> bool {
384 let mut applied = true;
385 for effect in effects {
386 if let Err(error) = media_transport.apply_relay_route_effect(&effect).await {
387 applied = false;
388 warn!(
389 ?effect,
390 ?error,
391 "media transport failed to apply relay route effect"
392 );
393 }
394 }
395 applied
396}
397
398pub(super) async fn execute_remote_source_activity_effects(
399 media_transport: &MediaTransport,
400 effects: impl IntoIterator<Item = TransportSourceActivityEffect>,
401) {
402 for effect in effects {
403 let _ = media_transport
404 .apply_remote_source_activity_effect(&effect)
405 .await;
406 }
407}
408
409#[cfg(test)]
410#[path = "TESTS/route.rs"]
411mod route_tests;