Skip to main content

o_sfu_core/engine/room/effects/
receiver_route.rs

1use o_sfu_telemetry::schema::event as telemetry_event;
2use tracing::info;
3
4use super::transport::{
5    RoomRouteEffects, execute_relays_and_teardown, execute_remote_source_activity_effects,
6};
7use crate::engine::{
8    media_transport::{MediaTransport, TransportTeardown},
9    room::{
10        Room,
11        media_graph::{ConsumerSetupOrigin, ConsumerSetupOutcome, PendingConsumerSetup},
12    },
13};
14
15#[derive(Debug)]
16pub(super) struct ReceiverSetupTurn {
17    setup: PendingConsumerSetup,
18    origin: ConsumerSetupOrigin,
19}
20
21impl ReceiverSetupTurn {
22    pub(super) const fn new(setup: PendingConsumerSetup, origin: ConsumerSetupOrigin) -> Self {
23        Self { setup, origin }
24    }
25
26    pub(super) async fn execute(self, room: &Room, media_transport: &MediaTransport) {
27        let Self {
28            setup: mut pending,
29            origin,
30        } = self;
31        if !execute_relays_and_teardown(media_transport, pending.take_relays(), []).await {
32            Self::release_pending_setup(room, pending, media_transport).await;
33            return;
34        }
35        let setup = match pending.declare(media_transport, origin).await {
36            Ok(setup) => setup,
37            Err(pending) => {
38                Self::release_pending_setup(room, pending, media_transport).await;
39                return;
40            }
41        };
42        let setup_outcome = {
43            let mut state = room.state.write().await;
44            state.commit_declared_consumer_setup(setup, origin)
45        };
46        match setup_outcome {
47            ConsumerSetupOutcome::Committed {
48                target,
49                route,
50                sender,
51                track_snapshot,
52                remote_source_activity,
53                transport_activity_update,
54                readiness_keyframe,
55            } => {
56                let consumer = route.consumer_session_key();
57                let source = route.source();
58                info!(
59                    event = telemetry_event::SUBSCRIBE_SUCCEEDED,
60                    room_id = room.uuid(),
61                    user_id = %consumer.user_id().path_segment(),
62                    connection_id = consumer.connection_id().as_u64(),
63                    media_worker_id = consumer.media_worker_id().as_usize(),
64                    transport_media_id = route.consumer_transport_media_id().as_u64(),
65                    producer_user_id = %source.session_key().user_id().path_segment(),
66                    source_transport_media_id = source.transport_media_id().as_u64(),
67                    stream_id = %target.stream,
68                    origin = origin.as_diagnostic_str(),
69                    "subscription committed"
70                );
71                if let Some(effect) = remote_source_activity {
72                    execute_remote_source_activity_effects(media_transport, [effect]).await;
73                }
74                let mut route_effects = RoomRouteEffects::default();
75                if let Some(active) = transport_activity_update {
76                    route_effects.setup_activity(route, target.kind, active);
77                }
78                if let Some(kf_target) = readiness_keyframe {
79                    route_effects.keyframe(kf_target);
80                }
81                if !route_effects.is_empty() {
82                    route_effects.execute(room.uuid(), media_transport).await;
83                }
84                let _ = sender.send_remote_tracks(track_snapshot);
85            }
86            ConsumerSetupOutcome::Released(route, relays) => {
87                let teardown = [TransportTeardown::RemoveMedia {
88                    session_key: route.consumer_session_key().clone(),
89                    transport_media_id: route.consumer_transport_media_id(),
90                }];
91                execute_relays_and_teardown(media_transport, relays, teardown).await;
92            }
93        }
94    }
95
96    async fn release_pending_setup(
97        room: &Room,
98        setup: PendingConsumerSetup,
99        media_transport: &MediaTransport,
100    ) {
101        let relays = {
102            let mut state = room.state.write().await;
103            state.release_pending_consumer_setup(setup)
104        };
105        execute_relays_and_teardown(media_transport, relays, []).await;
106    }
107}