1use std::{
4 collections::{BTreeMap, BTreeSet},
5 sync::Arc,
6};
7
8use o_sfu_router::{
9 MediaKind, ProducerId, Router, RouterError,
10 rtp::{MediaCapabilities, MediaStream as RouterRtpParameters},
11};
12use tracing::{error, warn};
13
14use super::{
15 CommittedConsumerSetup, ConsumerId, ConsumerRouteView, ConsumerSetupTarget,
16 DeclaredConsumerSetup, PendingConsumerRouteView, PendingConsumerSetup, PublishedSource,
17 ReceiverRouteActivity, SubscriptionKey, ValidatedPublish,
18 producer::{PublicationCommitError, allocate_source_descriptor},
19 route_graph::{CurrentPublication, PendingUpgrade, RemovedRoutes, RouteGraph},
20 source_index::PublishedSources,
21};
22use crate::engine::{
23 ConnectionId, MediaWorkerId, RoomInstanceId, UserId,
24 media_transport::{
25 SessionUploadEncoding, SourceActivityRevision, SourceActivityUpdate,
26 TransportConsumerRoute, TransportMediaId, TransportRelayRouteEffect, TransportSessionKey,
27 TransportSourceActivityEffect, TransportSourceKey, TransportTeardown,
28 },
29 room::{
30 RoomMediaCounts, RoomRuntimeContext, RouterPlacement,
31 effects::transport::RoomTransportPlan, outbound::OutboundSender,
32 },
33 source_model::{
34 ActiveSpeakerSourceRole, ConsumerSourceSelection, PolicyPauseReason,
35 PublishedSourceDescriptor, PublishedSourceId, SourceSubscriptionIntent, UserStreamId,
36 },
37};
38
39#[cfg(test)]
40#[path = "TESTS/topology_support.rs"]
41mod topology_support;
42
43#[derive(Debug)]
49pub struct RoomTopology {
50 instance: RoomInstanceId,
51 sources: PublishedSources,
52 route_graph: RouteGraph,
53 router: Router,
54 next_producer_id: u64,
55 video_allocation_revision: u64,
57}
58
59#[derive(Debug, Clone, PartialEq, Eq)]
78pub struct CommittedTransportReceipt {
79 pub transport_session_key: TransportSessionKey,
81}
82
83#[derive(Debug)]
88pub struct SessionPlacementCommit {
89 pub receipt: CommittedTransportReceipt,
90 pub replacement_transport_plan: RoomTransportPlan,
92}
93
94#[derive(Debug)]
95pub(in crate::engine::room) struct ConsumerActivityCommit {
96 pub(super) update: Option<ReceiverRouteActivity>,
97 pub(super) relay_effects: Vec<TransportRelayRouteEffect>,
98 pub(super) evictions: usize,
99}
100
101#[derive(Debug)]
102pub enum SessionPlacementRejection {
103 MissingPreviousSession { previous_connection: ConnectionId },
104 Router(RouterError),
105}
106
107impl RoomTopology {
108 pub fn new(
109 runtime_context: &RoomRuntimeContext,
110 router_rtp_capabilities: MediaCapabilities,
111 pending_target_limit: usize,
112 ) -> Self {
113 let router = match runtime_context.initial_router_placements() {
114 Some(placements) => {
115 Router::with_placements(placements.clone(), router_rtp_capabilities)
116 }
117 None => Router::new(runtime_context.primary_router(), router_rtp_capabilities),
118 };
119 Self {
120 instance: runtime_context.instance(),
121 sources: PublishedSources::default(),
122 route_graph: RouteGraph::new(pending_target_limit),
123 router,
124 next_producer_id: 1,
125 video_allocation_revision: 0,
126 }
127 }
128
129 pub(in crate::engine::room) const fn video_allocation_revision(&self) -> u64 {
130 self.video_allocation_revision
131 }
132
133 fn invalidate_video_allocation(&mut self) {
134 self.video_allocation_revision = self.video_allocation_revision.wrapping_add(1);
135 }
136
137 pub(in crate::engine::room) fn router(&self) -> &Router {
138 &self.router
139 }
140
141 #[must_use]
143 pub fn committed_transport_user_key(
144 &self,
145 user_id: impl Into<Arc<UserId>>,
146 connection_id: ConnectionId,
147 ) -> Option<TransportSessionKey> {
148 let user_id = user_id.into();
149 let worker = self
150 .router
151 .committed_media_worker_id(user_id.as_ref(), connection_id)?;
152 Some(self.transport_session_key(user_id, connection_id, worker))
153 }
154
155 #[must_use]
164 #[expect(
165 clippy::unreachable,
166 reason = "current room operations require committed connection placement and must not synthesize a transport worker"
167 )]
168 pub fn transport_user_key(
169 &self,
170 user_id: impl Into<Arc<UserId>>,
171 connection_id: ConnectionId,
172 ) -> TransportSessionKey {
173 let user_id = user_id.into();
174 let Some(worker) = self
175 .router
176 .committed_media_worker_id(user_id.as_ref(), connection_id)
177 else {
178 unreachable!("transport session key lookup requires committed connection placement");
179 };
180 self.transport_session_key(user_id, connection_id, worker)
181 }
182
183 fn transport_session_key(
184 &self,
185 user_id: Arc<UserId>,
186 connection_id: ConnectionId,
187 media_worker_id: MediaWorkerId,
188 ) -> TransportSessionKey {
189 TransportSessionKey::new(self.instance, media_worker_id, connection_id, user_id)
190 }
191
192 pub fn retire_committed_placement(
194 &mut self,
195 user_id: &UserId,
196 connection_id: ConnectionId,
197 ) -> Option<TransportSessionKey> {
198 let media_worker = self
199 .router
200 .retire_committed_placement(user_id, connection_id)?;
201 Some(self.transport_session_key(user_id.clone().into(), connection_id, media_worker))
202 }
203
204 #[must_use]
206 pub fn media_counts(&self) -> RoomMediaCounts {
207 RoomMediaCounts {
208 publications: self.sources.publication_count(),
209 subscriptions: self.route_graph.subscription_count(),
210 }
211 }
212
213 pub(in crate::engine::room) fn consumer_count(&self) -> usize {
215 self.route_graph.count()
216 }
217
218 #[cfg(any(test, feature = "testing-transport"))]
219 pub(in crate::engine::room) fn first_published_transport_media_id(
220 &self,
221 ) -> Option<TransportMediaId> {
222 self.sources.first_transport_media_id()
223 }
224
225 #[cfg(any(test, feature = "testing-transport"))]
226 pub(in crate::engine::room) fn producer_transport_media_id(
227 &self,
228 user_id: &UserId,
229 connection_id: ConnectionId,
230 stream_id: &UserStreamId,
231 ) -> Option<TransportMediaId> {
232 self.sources
233 .transport_media_id(user_id, connection_id, stream_id)
234 }
235
236 #[must_use]
237 pub(in crate::engine::room) fn source_descriptor(
238 &self,
239 source_id: PublishedSourceId,
240 ) -> Option<&PublishedSourceDescriptor> {
241 self.sources
242 .source(source_id)
243 .map(|source| &source.descriptor)
244 }
245
246 #[must_use]
247 pub(in crate::engine::room) fn source_for_transport_media(
248 &self,
249 transport_media_id: TransportMediaId,
250 ) -> Option<&PublishedSource> {
251 self.sources.source_for_transport(transport_media_id)
252 }
253
254 #[must_use]
255 pub(in crate::engine::room) fn source_id_for_owner_stream(
256 &self,
257 owner_user_id: &UserId,
258 stream_id: &UserStreamId,
259 ) -> Option<PublishedSourceId> {
260 self.sources.id_for_owner_stream(owner_user_id, stream_id)
261 }
262
263 #[must_use]
265 pub(in crate::engine::room) fn published_source_id(
266 &self,
267 owner: &UserId,
268 connection: ConnectionId,
269 stream_id: &UserStreamId,
270 ) -> Option<PublishedSourceId> {
271 let id = self.sources.id_for_owner_stream(owner, stream_id)?;
272 let source = self.sources.source(id)?;
273 (source.transport.session_key().connection_id() == connection).then_some(id)
274 }
275
276 pub(in crate::engine::room) fn published_sources(
278 &self,
279 ) -> impl Iterator<Item = &PublishedSource> {
280 self.sources.iter()
281 }
282
283 pub(in crate::engine::room) fn active_stream_user_counts(&self) -> BTreeMap<UserStreamId, u64> {
285 let mut users_by_stream: BTreeMap<UserStreamId, BTreeSet<UserId>> = BTreeMap::new();
286 for source in self.sources.iter().filter(|source| source.active) {
287 users_by_stream
288 .entry(source.descriptor.stream_id().clone())
289 .or_default()
290 .insert(source.descriptor.owner().user_id().clone());
291 }
292 users_by_stream
293 .into_iter()
294 .map(|(stream_id, users)| (stream_id, u64::try_from(users.len()).unwrap_or(u64::MAX)))
295 .collect()
296 }
297
298 #[must_use]
300 pub(in crate::engine::room) fn active_speaker_detector_owner(
301 &self,
302 transport_media_id: TransportMediaId,
303 ) -> Option<UserId> {
304 let source = self.source_for_transport_media(transport_media_id)?;
305 let detector_policy = source.descriptor.policy().active_speaker()?;
306 if detector_policy.role() != ActiveSpeakerSourceRole::Detector {
307 return None;
308 }
309 let owner = source.descriptor.owner().user_id();
310 self.sources
311 .owner_has_promotable_source_in_group(owner, detector_policy.group())
312 .then(|| owner.clone())
313 }
314
315 pub(in crate::engine::room) fn committed_consumer_routes(
316 &self,
317 ) -> impl Iterator<Item = ConsumerRouteView<'_>> {
318 self.route_graph
319 .attached()
320 .filter_map(|(key, current)| self.consumer_route(key, current))
321 }
322
323 pub(in crate::engine::room) fn committed_consumer_routes_for_user(
324 &self,
325 user_id: &UserId,
326 ) -> impl Iterator<Item = ConsumerRouteView<'_>> {
327 self.route_graph
328 .attached_for_receiver(user_id)
329 .filter_map(|(key, current)| self.consumer_route(key, current))
330 }
331
332 pub(super) fn detach_declined_consumers(
333 &mut self,
334 session: &TransportSessionKey,
335 declined: &[TransportMediaId],
336 ) -> (Vec<TransportRelayRouteEffect>, Vec<TransportTeardown>, bool) {
337 let RemovedRoutes {
338 routes,
339 consumers,
340 relays,
341 } = self
342 .route_graph
343 .detach_declined_consumers(session, declined);
344 let detached = !routes.is_empty();
345 if detached {
346 self.invalidate_video_allocation();
347 }
348 for consumer in consumers {
349 if let Some(error) = self.router.remove_consumer(consumer).err() {
350 error!(?consumer, ?error, "failed to remove declined room consumer");
351 }
352 }
353 let teardown = Self::media_teardowns([], routes).collect();
354 (relays, teardown, detached)
355 }
356
357 pub(in crate::engine::room) fn pending_consumer_routes_for_user(
358 &self,
359 user_id: &UserId,
360 ) -> impl Iterator<Item = PendingConsumerRouteView<'_>> {
361 self.route_graph
362 .attached_for_receiver(user_id)
363 .filter(|(_, current)| current.is_pending())
364 .filter_map(|(_, current)| {
365 let source = self.sources.source(current.source_id)?;
366 Some(PendingConsumerRouteView {
367 source,
368 selection: current.selection,
369 })
370 })
371 }
372
373 #[must_use]
375 pub(in crate::engine::room) fn consumer_source_selection(
376 &self,
377 key: &SubscriptionKey,
378 source_id: PublishedSourceId,
379 ) -> Option<ConsumerSourceSelection> {
380 self.route_graph.selection(key, source_id)
381 }
382
383 #[must_use]
385 pub(in crate::engine::room) fn subscription_intent(
386 &self,
387 key: &SubscriptionKey,
388 ) -> SourceSubscriptionIntent {
389 self.route_graph.intent(key)
390 }
391
392 pub(in crate::engine::room) fn committed_consumer_user_ids_for_source(
393 &self,
394 source_id: PublishedSourceId,
395 ) -> BTreeSet<UserId> {
396 self.route_graph
397 .attached_for_source(source_id)
398 .filter(|(_, current)| current.committed().is_some())
399 .map(|(key, _)| key.receiver.clone())
400 .collect()
401 }
402
403 pub(in crate::engine::room) fn committed_consumer_user_ids_for_owner_sources(
404 &self,
405 user_id: &UserId,
406 ) -> BTreeSet<UserId> {
407 let route_graph = &self.route_graph;
408 self.sources
409 .ids_for_owner(user_id)
410 .flat_map(|source_id| route_graph.attached_for_source(source_id))
411 .filter(|(_, current)| current.committed().is_some())
412 .map(|(key, _)| key.receiver.clone())
413 .collect()
414 }
415
416 #[must_use]
417 pub(in crate::engine::room) fn committed_consumer_route_for_key(
418 &self,
419 key: &SubscriptionKey,
420 ) -> Option<ConsumerRouteView<'_>> {
421 let (key, current) = self.route_graph.current(key)?;
422 self.consumer_route(key, current)
423 }
424
425 #[must_use]
426 pub(in crate::engine::room) fn published_source(
427 &self,
428 source_id: PublishedSourceId,
429 ) -> Option<&PublishedSource> {
430 self.sources.source(source_id)
431 }
432
433 pub(super) fn missing_consumer_targets_for_source<'a>(
435 &mut self,
436 source_id: PublishedSourceId,
437 receivers: impl IntoIterator<Item = (&'a UserId, ConnectionId)>,
438 ) -> Vec<ConsumerSetupTarget> {
439 receivers
440 .into_iter()
441 .filter_map(|(user, connection)| self.consumer_target(user, connection, source_id))
442 .collect()
443 }
444
445 fn consumer_target(
446 &mut self,
447 user_id: &UserId,
448 connection_id: ConnectionId,
449 source_id: PublishedSourceId,
450 ) -> Option<ConsumerSetupTarget> {
451 let consumer_session = self.committed_transport_user_key(user_id.clone(), connection_id)?;
452 let (sources, route_graph) = (&self.sources, &mut self.route_graph);
453 Self::attach_consumer_target(route_graph, &consumer_session, sources.source(source_id)?)
454 }
455
456 fn attach_consumer_target(
458 route_graph: &mut RouteGraph,
459 consumer_session: &TransportSessionKey,
460 source: &PublishedSource,
461 ) -> Option<ConsumerSetupTarget> {
462 if source.descriptor.owner().user_id() == consumer_session.user_id() {
463 return None;
464 }
465 let target = ConsumerSetupTarget::new(consumer_session.clone(), source);
466 route_graph
467 .attach_for_setup(target.subscription_key(), target.source_id)
468 .then_some(target)
469 }
470
471 pub(super) fn missing_consumer_targets(
475 &mut self,
476 user_id: &UserId,
477 connection_id: ConnectionId,
478 include_source: impl Fn(&PublishedSource) -> bool,
479 ) -> Vec<ConsumerSetupTarget> {
480 let Some(consumer_session) =
481 self.committed_transport_user_key(user_id.clone(), connection_id)
482 else {
483 return Vec::new();
484 };
485 let (sources, route_graph) = (&self.sources, &mut self.route_graph);
486 sources
487 .iter()
488 .filter(|source| include_source(source))
489 .filter_map(|source| {
490 Self::attach_consumer_target(route_graph, &consumer_session, source)
491 })
492 .collect()
493 }
494
495 pub fn update_consumer_source_selection(
499 &mut self,
500 key: &SubscriptionKey,
501 source_id: PublishedSourceId,
502 route: &TransportConsumerRoute,
503 update: impl FnOnce(&mut ConsumerSourceSelection),
504 ) -> bool {
505 let updated = self
506 .route_graph
507 .update_selection(key, source_id, route, update);
508 if updated {
509 self.invalidate_video_allocation();
510 }
511 updated
512 }
513
514 pub(in crate::engine::room) fn update_consumer_upgrade(
516 &mut self,
517 key: &SubscriptionKey,
518 source_id: PublishedSourceId,
519 route: &TransportConsumerRoute,
520 pending_upgrade: Option<PendingUpgrade>,
521 ) -> bool {
522 if !self
523 .sources
524 .source(source_id)
525 .is_some_and(|source| source.active)
526 {
527 return false;
528 }
529 self.route_graph
530 .update_upgrade(key, source_id, route, pending_upgrade)
531 }
532
533 fn detach_user_sources(
534 &mut self,
535 user_id: &UserId,
536 ) -> (Vec<TransportSourceKey>, RemovedRoutes) {
537 let mut sources = Vec::new();
538 let mut removed = RemovedRoutes::default();
539 let source_ids = self.sources.ids_for_owner(user_id).collect::<Vec<_>>();
540 for source_id in source_ids {
541 if let Some((source, routes)) = self.remove_source(source_id) {
542 sources.push(source.transport);
543 removed.extend(routes);
544 }
545 }
546 (sources, removed)
547 }
548
549 fn remove_source(
550 &mut self,
551 source_id: PublishedSourceId,
552 ) -> Option<(PublishedSource, RemovedRoutes)> {
553 let source = self.sources.remove(source_id)?;
554 let routes = self.route_graph.detach_source(source_id);
555 self.invalidate_video_allocation();
556 Some((source, routes))
557 }
558
559 fn consumer_route<'a>(
560 &'a self,
561 key: &'a SubscriptionKey,
562 current: &'a CurrentPublication,
563 ) -> Option<ConsumerRouteView<'a>> {
564 let committed = current.committed()?;
565 let source = self.sources.source(current.source_id)?;
566 Some(ConsumerRouteView {
567 key,
568 route: &committed.route,
569 mid: &committed.mid,
570 pending_upgrade: committed.pending_upgrade.as_ref(),
571 source,
572 selection: current.selection,
573 })
574 }
575
576 pub(in crate::engine::room) fn commit_publication(
584 &mut self,
585 publish: ValidatedPublish,
586 rtp: RouterRtpParameters,
587 encodings: &[SessionUploadEncoding],
588 media: TransportMediaId,
589 ) -> Result<PublishedSourceId, PublicationCommitError> {
590 let descriptor = allocate_source_descriptor(&mut self.sources, &publish, &rtp, encodings)?;
591 let source_id = descriptor.source_id();
592 let producer_id = ProducerId::allocate(&mut self.next_producer_id);
593 let routed = self
594 .router
595 .add_producer(publish.session_key.user_id(), producer_id)?;
596 self.sources.insert(PublishedSource {
597 descriptor,
598 transport: TransportSourceKey::new(publish.session_key, media),
599 rtp,
600 routed,
601 active: true,
602 activity_revision: SourceActivityRevision::default(),
603 });
604 Ok(source_id)
605 }
606
607 fn media_teardowns(
608 sources: impl IntoIterator<Item = TransportSourceKey>,
609 routes: impl IntoIterator<Item = TransportConsumerRoute>,
610 ) -> impl Iterator<Item = TransportTeardown> {
611 sources
612 .into_iter()
613 .map(|source| TransportTeardown::RemoveMedia {
614 session_key: source.session_key().clone(),
615 transport_media_id: source.transport_media_id(),
616 })
617 .chain(
618 routes
619 .into_iter()
620 .map(|route| TransportTeardown::RemoveMedia {
621 session_key: route.consumer_session_key().clone(),
622 transport_media_id: route.consumer_transport_media_id(),
623 }),
624 )
625 }
626
627 pub fn set_published_source_activity(
632 &mut self,
633 source_id: PublishedSourceId,
634 connection_id: ConnectionId,
635 active: bool,
636 ) -> Option<SourceActivityRevision> {
637 let source = self.sources.source_mut(source_id)?;
638 if source.transport.session_key().connection_id() != connection_id
639 || source.active == active
640 {
641 return None;
642 }
643 source.active = active;
644 source.activity_revision = source.activity_revision.next();
645 let revision = source.activity_revision;
646 if !active {
647 self.route_graph.clear_source_upgrades(source_id);
648 }
649 self.invalidate_video_allocation();
650 Some(revision)
651 }
652
653 pub(super) fn source_activity_effects(
654 &self,
655 source: &TransportSourceKey,
656 update: SourceActivityUpdate,
657 ) -> Vec<TransportSourceActivityEffect> {
658 self.route_graph
659 .source_activity_target_workers(source)
660 .map(|target_media_worker_id| TransportSourceActivityEffect {
661 source: source.clone(),
662 target_media_worker_id,
663 update,
664 })
665 .collect()
666 }
667
668 pub fn commit_session_placement(
679 &mut self,
680 user_id: &UserId,
681 connection_id: ConnectionId,
682 previous_connection: Option<ConnectionId>,
683 home_placement: RouterPlacement,
684 ) -> Result<SessionPlacementCommit, SessionPlacementRejection> {
685 let previous_session_key = if let Some(previous_connection) = previous_connection {
687 let Some(key) = self.committed_transport_user_key(user_id.clone(), previous_connection)
688 else {
689 return Err(SessionPlacementRejection::MissingPreviousSession {
690 previous_connection,
691 });
692 };
693 Some(key)
694 } else {
695 None
696 };
697 let media_worker = self
698 .router
699 .commit_session_placement(user_id, connection_id, home_placement)
700 .map_err(SessionPlacementRejection::Router)?;
701 self.route_graph.publisher_joined(user_id);
702 let session_key =
703 self.transport_session_key(user_id.clone().into(), connection_id, media_worker);
704 let receipt = CommittedTransportReceipt {
705 transport_session_key: session_key,
706 };
707 let replacement_transport_plan = previous_session_key.as_ref().map_or_else(
708 RoomTransportPlan::default,
709 |replaced_session_key| {
710 let close_session = TransportTeardown::CloseSession {
711 session_key: replaced_session_key.clone(),
712 };
713 let (_, removed_sources) = self.detach_user_sources(user_id);
716 let receiver_relays = self.route_graph.reset_receiver_for_replacement(user_id);
717 let RemovedRoutes {
718 routes, mut relays, ..
719 } = removed_sources;
720 relays.extend(receiver_relays);
721 let teardown = Self::media_teardowns([], routes).chain([close_session]);
722 RoomTransportPlan::from_relays_and_teardown(relays, teardown)
723 },
724 );
725 if previous_session_key.is_some() {
726 self.invalidate_video_allocation();
727 }
728 Ok(SessionPlacementCommit {
729 receipt,
730 replacement_transport_plan,
731 })
732 }
733
734 pub fn remove_session(&mut self, user_id: &UserId) -> (RoomTransportPlan, usize) {
736 let (sources, mut removed) = self.detach_user_sources(user_id);
737 removed.extend(self.route_graph.remove_receiver(user_id));
738 let evictions = self.route_graph.publisher_left(user_id);
739 self.invalidate_video_allocation();
740 let teardown = Self::media_teardowns(sources, removed.routes);
741 if let Some(error) = self.router.remove_session(user_id).err() {
742 error!(?user_id, ?error, "failed to remove user from room router");
743 }
744 (
745 RoomTransportPlan::from_relays_and_teardown(removed.relays, teardown),
746 evictions,
747 )
748 }
749
750 pub(super) fn commit_consumer_setup(
758 &mut self,
759 setup: DeclaredConsumerSetup,
760 selection: ConsumerSourceSelection,
761 ) -> Result<CommittedConsumerSetup, (TransportConsumerRoute, Vec<TransportRelayRouteEffect>)>
762 {
763 let DeclaredConsumerSetup {
764 pending:
765 PendingConsumerSetup {
766 target,
767 consumer,
768 reservation,
769 sender,
770 rtp,
771 relays: _,
772 },
773 route,
774 mid,
775 } = setup;
776 let active = selection.delivery_active();
777 let declared_active = reservation.declared_activity().is_active();
778 let committed_mid = mid.unwrap_or_else(|| {
779 rtp.mid()
780 .map_or_else(|| consumer.to_string(), ToOwned::to_owned)
781 });
782 let (route_graph, room_router) = (&mut self.route_graph, &mut self.router);
783 let result =
785 route_graph.commit(reservation, route.clone(), committed_mid, selection, || {
786 match room_router.add_consumer(target.session.user_id(), consumer, target.routed) {
787 Ok(routed_consumer) => Some(routed_consumer),
788 Err(error) => {
789 warn!(
790 consumer_user_id = ?target.session.user_id(),
791 source_id = ?target.source_id,
792 ?error,
793 "router rejected consumer creation"
794 );
795 None
796 }
797 }
798 });
799 if let Err(relays) = result {
800 return Err((route, relays));
801 }
802 self.invalidate_video_allocation();
803 Ok(CommittedConsumerSetup {
804 target,
805 route,
806 sender,
807 transport_activity_update: (active != declared_active).then_some(active),
808 })
809 }
810
811 pub(super) fn reserve_consumer_setup(
815 &mut self,
816 target: ConsumerSetupTarget,
817 consumer: ConsumerId,
818 selection: ConsumerSourceSelection,
819 sender: OutboundSender,
820 rtp: RouterRtpParameters,
821 ) -> Option<PendingConsumerSetup> {
822 let key = target.subscription_key();
823 let relay_active = selection.active();
827 let reservation =
828 self.route_graph
829 .reserve_consumer_setup(key, target.source_id, selection)?;
830 let source_worker = target.source.session_key().media_worker_id();
831 let target_worker = target.session.media_worker_id();
832 let relays = if source_worker == target_worker {
833 Vec::new()
834 } else {
835 self.route_graph
836 .reserve_relay(&reservation, &target, target_worker, relay_active)
837 };
838 Some(PendingConsumerSetup {
839 target,
840 consumer,
841 reservation,
842 sender,
843 rtp,
844 relays,
845 })
846 }
847
848 pub(super) fn release_consumer_setup(
850 &mut self,
851 setup: PendingConsumerSetup,
852 ) -> Vec<TransportRelayRouteEffect> {
853 self.route_graph.release_consumer_setup(setup.reservation)
854 }
855
856 pub(in crate::engine::room) fn apply_subscription_intent(
862 &mut self,
863 key: &SubscriptionKey,
864 connection_id: ConnectionId,
865 intent: SourceSubscriptionIntent,
866 receiver_deafened: bool,
867 publisher_present: bool,
868 ) -> ConsumerActivityCommit {
869 let mut commit = ConsumerActivityCommit {
870 update: None,
871 relay_effects: Vec::new(),
872 evictions: 0,
873 };
874 if intent.is_empty() {
875 return commit;
876 }
877 commit.evictions = self
878 .route_graph
879 .merge_intent(key, intent, publisher_present);
880 self.invalidate_video_allocation();
881 let Some(active) = intent.active() else {
882 return commit;
883 };
884 let Some(source_id) = self.source_id_for_owner_stream(&key.publisher, &key.stream) else {
885 return commit;
886 };
887 let policy_pause_reason = (active
889 && receiver_deafened
890 && self
891 .source_descriptor(source_id)
892 .is_some_and(|source| source.media_kind() == MediaKind::Audio))
893 .then_some(PolicyPauseReason::ReceiverDeafened);
894 let Some(relay_effects) = self.route_graph.set_activity(
895 key,
896 source_id,
897 connection_id,
898 active,
899 policy_pause_reason,
900 ) else {
901 return commit;
902 };
903 let update = self
904 .committed_consumer_route_for_key(key)
905 .filter(|route| route.route.consumer_session_key().connection_id() == connection_id)
906 .map(|route| {
907 ReceiverRouteActivity::new(route.target(), route.selection.delivery_active())
908 });
909 commit.update = update;
910 commit.relay_effects = relay_effects;
911 commit
912 }
913}