o_sfu_core/engine/room/media_graph/
consumer_setup.rs1use std::mem;
2
3use o_sfu_router::{
4 MediaKind as RouterMediaKind, rtp::MediaStream as RouterRtpParameters,
5 topology::RoutedProducerId,
6};
7use tracing::warn;
8
9use super::{
10 super::{
11 outbound::{OutboundSender, VersionedRemoteTrackSnapshot},
12 state::RoomState,
13 },
14 ConsumerId, ConsumerRouteTarget, PublishedSource, SubscriptionKey,
15 route_graph::{ConsumerRouteReservation, RelayRouteKey},
16};
17use crate::engine::{
18 MediaWorkerId,
19 media_transport::{
20 MediaTransport, ProducerActivity, SourceActivityUpdate, TransportConsumerRoute,
21 TransportMediaId, TransportRelayRouteEffect, TransportSessionKey,
22 TransportSourceActivityEffect, TransportSourceKey,
23 },
24 source_model::{PublishedSourceId, UserStreamId},
25};
26
27#[derive(Debug)]
28pub struct ConsumerSetupTarget {
29 pub session: TransportSessionKey,
30 pub source: TransportSourceKey,
31 pub source_id: PublishedSourceId,
32 pub stream: UserStreamId,
33 pub kind: RouterMediaKind,
34 pub routed: RoutedProducerId,
35}
36
37#[derive(Debug)]
42#[must_use = "pending consumer setups reserve route graph state and must be committed or released"]
43pub struct PendingConsumerSetup {
44 pub(super) target: ConsumerSetupTarget,
45 pub(super) consumer: ConsumerId,
46 pub(super) reservation: ConsumerRouteReservation,
47 pub(super) sender: OutboundSender,
48 pub(super) rtp: RouterRtpParameters,
49 pub(super) relays: Vec<TransportRelayRouteEffect>,
50}
51
52pub struct DeclaredConsumerSetup {
53 pub(super) pending: PendingConsumerSetup,
54 pub(super) route: TransportConsumerRoute,
55 pub(super) mid: Option<String>,
56}
57
58pub struct CommittedConsumerSetup {
59 pub(super) target: ConsumerSetupTarget,
60 pub(super) route: TransportConsumerRoute,
61 pub(super) sender: OutboundSender,
62 pub(super) transport_activity_update: Option<bool>,
63}
64
65#[expect(
66 clippy::large_enum_variant,
67 reason = "consumer setup outcomes are returned and matched immediately so boxing the committed setup would allocate on every successful consumer setup"
68)]
69#[derive(Debug)]
70pub enum ConsumerSetupOutcome {
71 Committed {
72 target: ConsumerSetupTarget,
73 route: TransportConsumerRoute,
74 sender: OutboundSender,
75 track_snapshot: VersionedRemoteTrackSnapshot,
76 remote_source_activity: Option<TransportSourceActivityEffect>,
77 transport_activity_update: Option<bool>,
78 readiness_keyframe: Option<ConsumerRouteTarget>,
79 },
80 Released(TransportConsumerRoute, Vec<TransportRelayRouteEffect>),
81}
82
83#[derive(Debug, Clone, Copy)]
84pub enum ConsumerSetupOrigin {
85 Readiness,
86 Publish,
87 Subscribe,
88}
89
90impl ConsumerSetupOrigin {
91 pub const fn as_diagnostic_str(self) -> &'static str {
92 match self {
93 Self::Readiness => "readiness",
94 Self::Publish => "publish",
95 Self::Subscribe => "subscribe",
96 }
97 }
98}
99
100impl RoomState {
101 pub fn commit_declared_consumer_setup(
106 &mut self,
107 setup: DeclaredConsumerSetup,
108 origin: ConsumerSetupOrigin,
109 ) -> ConsumerSetupOutcome {
110 let target = &setup.pending.target;
111 let session = &target.session;
112 if self
113 .user_for_connection(session.user_id(), session.connection_id())
114 .is_some_and(|user| user.parsed_client_rtp_capabilities.is_some())
115 && let Some((source_active, source_activity_revision)) = self
116 .topology
117 .published_source(target.source_id)
118 .filter(|source| target.matches_identity(source))
119 .map(|source| (source.active, source.activity_revision))
120 {
121 let selection = self.setup_selection(target, source_active);
124 let delivery_active = selection.delivery_active();
125 match self.topology.commit_consumer_setup(setup, selection) {
126 Ok(commit) => {
127 let remote_source_activity =
130 (commit.route.source().session_key().media_worker_id()
131 != commit.route.consumer_session_key().media_worker_id())
132 .then(|| TransportSourceActivityEffect {
133 source: commit.route.source().clone(),
134 target_media_worker_id: commit
135 .route
136 .consumer_session_key()
137 .media_worker_id(),
138 update: SourceActivityUpdate::new(
139 ProducerActivity::from_active(source_active),
140 source_activity_revision,
141 ),
142 });
143 let track_snapshot =
144 self.remote_track_snapshot_for_user(commit.target.session.user_id(), true);
145 let readiness_keyframe = match origin {
148 ConsumerSetupOrigin::Readiness
149 if delivery_active
150 && commit.transport_activity_update != Some(true)
151 && commit.target.kind == RouterMediaKind::Video =>
152 {
153 Some(commit.target.route_target(commit.route.clone()))
154 }
155 ConsumerSetupOrigin::Readiness
156 | ConsumerSetupOrigin::Publish
157 | ConsumerSetupOrigin::Subscribe => None,
158 };
159 ConsumerSetupOutcome::Committed {
160 target: commit.target,
161 route: commit.route,
162 sender: commit.sender,
163 track_snapshot,
164 remote_source_activity,
165 transport_activity_update: commit.transport_activity_update,
166 readiness_keyframe,
167 }
168 }
169 Err((route, relays)) => ConsumerSetupOutcome::Released(route, relays),
170 }
171 } else {
172 let DeclaredConsumerSetup { pending, route, .. } = setup;
173 ConsumerSetupOutcome::Released(route, self.topology.release_consumer_setup(pending))
174 }
175 }
176
177 pub fn release_pending_consumer_setup(
178 &mut self,
179 setup: PendingConsumerSetup,
180 ) -> Vec<TransportRelayRouteEffect> {
181 self.topology.release_consumer_setup(setup)
182 }
183}
184
185impl PendingConsumerSetup {
186 pub(in crate::engine::room) fn take_relays(&mut self) -> Vec<TransportRelayRouteEffect> {
187 mem::take(&mut self.relays)
188 }
189
190 pub(in crate::engine::room) async fn declare(
191 self,
192 media_transport: &MediaTransport,
193 origin: ConsumerSetupOrigin,
194 ) -> Result<DeclaredConsumerSetup, Self> {
195 match media_transport
196 .consume_media_with_mid(
197 &self.target.session,
198 self.target.kind,
199 self.target.source.session_key(),
200 self.target.source.transport_media_id(),
201 &self.rtp,
202 self.reservation.declared_activity(),
203 )
204 .await
205 {
206 Ok((media, mid)) => Ok(DeclaredConsumerSetup {
207 route: self.target.transport_consumer_route(media),
208 pending: self,
209 mid: Some(mid),
210 }),
211 Err(error) => {
212 warn!(
213 consumer_user_id = ?self.target.session.user_id(),
214 consumer_connection_id = ?self.target.session.connection_id(),
215 producer_user_id = ?self.target.source.session_key().user_id(),
216 producer_connection_id = ?self.target.source.session_key().connection_id(),
217 source_transport_media_id = ?self.target.source.transport_media_id(),
218 error = ?error,
219 consumer_mid = self.rtp.mid(),
220 ?origin,
221 "media transport rejected consume media declaration"
222 );
223 Err(self)
224 }
225 }
226 }
227}
228
229impl ConsumerSetupTarget {
230 pub fn new(session: TransportSessionKey, source: &PublishedSource) -> Self {
231 Self {
232 session,
233 source: source.transport.clone(),
234 source_id: source.descriptor.source_id(),
235 stream: source.descriptor.stream_id().clone(),
236 kind: source.descriptor.media_kind(),
237 routed: source.routed,
238 }
239 }
240
241 pub(super) fn subscription_key(&self) -> SubscriptionKey {
242 SubscriptionKey::new(
243 self.session.user_id(),
244 self.source.session_key().user_id(),
245 &self.stream,
246 )
247 }
248
249 pub(super) fn transport_consumer_route(
250 &self,
251 consumer_media: TransportMediaId,
252 ) -> TransportConsumerRoute {
253 TransportConsumerRoute::new(self.session.clone(), consumer_media, self.source.clone())
254 }
255
256 fn route_target(&self, route: TransportConsumerRoute) -> ConsumerRouteTarget {
257 ConsumerRouteTarget::new(route, self.stream.clone(), self.kind)
258 }
259
260 pub(super) fn relay_route_key(&self, target_worker: MediaWorkerId) -> RelayRouteKey {
261 RelayRouteKey {
262 source: self.source.clone(),
263 target_worker,
264 }
265 }
266
267 pub(super) fn matches_identity(&self, source: &PublishedSource) -> bool {
268 source.descriptor.source_id() == self.source_id
269 && source.transport == self.source
270 && source.descriptor.stream_id() == &self.stream
271 && source.descriptor.media_kind() == self.kind
272 && source.routed == self.routed
273 }
274}