o_sfu_core/engine/room/effects/
receiver_route.rs1use 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}