Skip to main content

o_sfu_core/engine/room/effects/
transport.rs

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;