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 ConsumerActivity, MediaTransport, ProducerActivity, SourceActivityUpdate,
21 TransportConsumerRoute, 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#[allow(
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 let activity =
196 ConsumerActivity::from_active(self.reservation.selection().delivery_active());
197 match media_transport
198 .consume_media(
199 &self.target.session,
200 self.target.kind,
201 self.target.source.session_key(),
202 self.target.source.transport_media_id(),
203 &self.rtp,
204 activity,
205 )
206 .await
207 {
208 Ok(media) => {
209 let mid = media_transport
210 .transport_media_mid(&self.target.session, media)
211 .await;
212 Ok(DeclaredConsumerSetup {
213 route: self.target.transport_consumer_route(media),
214 pending: self,
215 mid,
216 })
217 }
218 Err(error) => {
219 warn!(
220 consumer_user_id = ?self.target.session.user_id(),
221 consumer_connection_id = ?self.target.session.connection_id(),
222 producer_user_id = ?self.target.source.session_key().user_id(),
223 producer_connection_id = ?self.target.source.session_key().connection_id(),
224 source_transport_media_id = ?self.target.source.transport_media_id(),
225 error = ?error,
226 consumer_mid = self.rtp.mid(),
227 ?origin,
228 "media transport rejected consume media declaration"
229 );
230 Err(self)
231 }
232 }
233 }
234}
235
236impl ConsumerSetupTarget {
237 pub fn new(session: TransportSessionKey, source: &PublishedSource) -> Self {
238 Self {
239 session,
240 source: source.transport.clone(),
241 source_id: source.descriptor.source_id(),
242 stream: source.descriptor.stream_id().clone(),
243 kind: source.descriptor.media_kind(),
244 routed: source.routed,
245 }
246 }
247
248 pub(super) fn subscription_key(&self) -> SubscriptionKey {
249 SubscriptionKey::new(
250 self.session.user_id(),
251 self.source.session_key().user_id(),
252 &self.stream,
253 )
254 }
255
256 pub(super) fn transport_consumer_route(
257 &self,
258 consumer_media: TransportMediaId,
259 ) -> TransportConsumerRoute {
260 TransportConsumerRoute::new(self.session.clone(), consumer_media, self.source.clone())
261 }
262
263 fn route_target(&self, route: TransportConsumerRoute) -> ConsumerRouteTarget {
264 ConsumerRouteTarget::new(route, self.stream.clone(), self.kind)
265 }
266
267 pub(super) fn relay_route_key(&self, target_worker: MediaWorkerId) -> RelayRouteKey {
268 RelayRouteKey {
269 source_user: self.source.session_key().user_id().clone(),
270 source_connection: self.source.session_key().connection_id(),
271 source_media: self.source.transport_media_id(),
272 target_worker,
273 }
274 }
275
276 pub(super) fn matches_identity(&self, source: &PublishedSource) -> bool {
277 source.descriptor.source_id() == self.source_id
278 && source.transport == self.source
279 && source.descriptor.stream_id() == &self.stream
280 && source.descriptor.media_kind() == self.kind
281 && source.routed == self.routed
282 }
283}