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, RelayRouteEffect, 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}
56
57#[derive(Debug, Clone, PartialEq, Eq)]
78pub struct CommittedTransportReceipt {
79 pub connection_id: ConnectionId,
81 pub transport_session_key: TransportSessionKey,
83}
84
85#[derive(Debug)]
90pub struct SessionPlacementCommit {
91 pub receipt: CommittedTransportReceipt,
92 pub replacement_transport_plan: RoomTransportPlan,
94}
95
96#[derive(Debug)]
97pub(super) struct ConsumerActivityCommit {
98 pub(super) update: Option<ReceiverRouteActivity>,
99 pub(super) relay_effects: Vec<TransportRelayRouteEffect>,
100}
101
102#[derive(Debug)]
103pub enum SessionPlacementRejection {
104 MissingPreviousSession { previous_connection: ConnectionId },
105 Router(RouterError),
106}
107
108impl RoomTopology {
109 pub fn new(
110 runtime_context: &RoomRuntimeContext,
111 router_rtp_capabilities: MediaCapabilities,
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::default(),
123 router,
124 next_producer_id: 1,
125 }
126 }
127
128 pub(in crate::engine::room) fn router(&self) -> &Router {
129 &self.router
130 }
131
132 #[must_use]
134 pub fn committed_transport_user_key(
135 &self,
136 user_id: impl Into<Arc<UserId>>,
137 connection_id: ConnectionId,
138 ) -> Option<TransportSessionKey> {
139 let user_id = user_id.into();
140 let worker = self
141 .router
142 .committed_media_worker_id(user_id.as_ref(), connection_id)?;
143 Some(self.transport_session_key(user_id, connection_id, worker))
144 }
145
146 #[must_use]
171 #[expect(
172 clippy::unreachable,
173 reason = "current room operations require committed connection placement and must not synthesize a transport worker"
174 )]
175 pub fn transport_user_key(
176 &self,
177 user_id: impl Into<Arc<UserId>>,
178 connection_id: ConnectionId,
179 ) -> TransportSessionKey {
180 let user_id = user_id.into();
181 let Some(worker) = self
182 .router
183 .committed_media_worker_id(user_id.as_ref(), connection_id)
184 else {
185 unreachable!("transport session key lookup requires committed connection placement");
186 };
187 self.transport_session_key(user_id, connection_id, worker)
188 }
189
190 fn transport_session_key(
191 &self,
192 user_id: Arc<UserId>,
193 connection_id: ConnectionId,
194 media_worker_id: MediaWorkerId,
195 ) -> TransportSessionKey {
196 TransportSessionKey::new(self.instance, media_worker_id, connection_id, user_id)
197 }
198
199 pub fn retire_committed_placement(
201 &mut self,
202 user_id: &UserId,
203 connection_id: ConnectionId,
204 ) -> Option<TransportSessionKey> {
205 let media_worker = self
206 .router
207 .retire_committed_placement(user_id, connection_id)?;
208 Some(self.transport_session_key(user_id.clone().into(), connection_id, media_worker))
209 }
210
211 #[must_use]
213 pub fn media_counts(&self) -> RoomMediaCounts {
214 RoomMediaCounts {
215 publications: self.sources.publication_count(),
216 subscriptions: self.route_graph.subscription_count(),
217 }
218 }
219
220 #[cfg(any(test, feature = "testing-transport"))]
221 pub(in crate::engine::room) fn consumer_count(&self) -> usize {
222 self.route_graph.count()
223 }
224
225 #[cfg(any(test, feature = "testing-transport"))]
226 pub(in crate::engine::room) fn first_published_transport_media_id(
227 &self,
228 ) -> Option<TransportMediaId> {
229 self.sources.first_transport_media_id()
230 }
231
232 #[cfg(any(test, feature = "testing-transport"))]
233 pub(in crate::engine::room) fn producer_transport_media_id(
234 &self,
235 user_id: &UserId,
236 connection_id: ConnectionId,
237 stream_id: &UserStreamId,
238 ) -> Option<TransportMediaId> {
239 self.sources
240 .transport_media_id(user_id, connection_id, stream_id)
241 }
242
243 #[must_use]
244 pub(in crate::engine::room) fn source_descriptor(
245 &self,
246 source_id: PublishedSourceId,
247 ) -> Option<&PublishedSourceDescriptor> {
248 self.sources
249 .source(source_id)
250 .map(|source| &source.descriptor)
251 }
252
253 #[must_use]
254 pub(in crate::engine::room) fn source_for_transport_media(
255 &self,
256 transport_media_id: TransportMediaId,
257 ) -> Option<&PublishedSource> {
258 self.sources.source_for_transport(transport_media_id)
259 }
260
261 #[must_use]
262 pub(in crate::engine::room) fn source_id_for_owner_stream(
263 &self,
264 owner_user_id: &UserId,
265 stream_id: &UserStreamId,
266 ) -> Option<PublishedSourceId> {
267 self.sources.id_for_owner_stream(owner_user_id, stream_id)
268 }
269
270 #[must_use]
272 pub(in crate::engine::room) fn published_source_id(
273 &self,
274 owner: &UserId,
275 connection: ConnectionId,
276 stream_id: &UserStreamId,
277 ) -> Option<PublishedSourceId> {
278 let id = self.sources.id_for_owner_stream(owner, stream_id)?;
279 let source = self.sources.source(id)?;
280 (source.transport.session_key().connection_id() == connection).then_some(id)
281 }
282
283 pub(in crate::engine::room) fn published_sources(
285 &self,
286 ) -> impl Iterator<Item = &PublishedSource> {
287 self.sources.iter()
288 }
289
290 pub(in crate::engine::room) fn active_stream_user_counts(&self) -> BTreeMap<UserStreamId, u64> {
292 let mut users_by_stream: BTreeMap<UserStreamId, BTreeSet<UserId>> = BTreeMap::new();
293 for source in self.sources.iter().filter(|source| source.active) {
294 users_by_stream
295 .entry(source.descriptor.stream_id().clone())
296 .or_default()
297 .insert(source.descriptor.owner().user_id().clone());
298 }
299 users_by_stream
300 .into_iter()
301 .map(|(stream_id, users)| (stream_id, u64::try_from(users.len()).unwrap_or(u64::MAX)))
302 .collect()
303 }
304
305 #[must_use]
307 pub(in crate::engine::room) fn active_speaker_detector_owner(
308 &self,
309 transport_media_id: TransportMediaId,
310 ) -> Option<UserId> {
311 let source = self.source_for_transport_media(transport_media_id)?;
312 let detector_policy = source.descriptor.policy().active_speaker()?;
313 if detector_policy.role() != ActiveSpeakerSourceRole::Detector {
314 return None;
315 }
316 let owner = source.descriptor.owner().user_id();
317 self.sources
318 .owner_has_promotable_source_in_group(owner, detector_policy.group())
319 .then(|| owner.clone())
320 }
321
322 pub(in crate::engine::room) fn committed_consumer_routes(
323 &self,
324 ) -> impl Iterator<Item = ConsumerRouteView<'_>> {
325 self.route_graph
326 .attached()
327 .filter_map(|(key, current)| self.consumer_route(key, current))
328 }
329
330 pub(in crate::engine::room) fn committed_consumer_routes_for_user(
331 &self,
332 user_id: &UserId,
333 ) -> impl Iterator<Item = ConsumerRouteView<'_>> {
334 self.route_graph
335 .attached_for_receiver(user_id)
336 .filter_map(|(key, current)| self.consumer_route(key, current))
337 }
338
339 pub(super) fn detach_declined_consumers(
340 &mut self,
341 session: &TransportSessionKey,
342 declined: &[TransportMediaId],
343 ) -> (Vec<TransportRelayRouteEffect>, Vec<TransportTeardown>, bool) {
344 let RemovedRoutes {
345 routes,
346 consumers,
347 relays,
348 } = self
349 .route_graph
350 .detach_declined_consumers(session, declined);
351 let detached = !routes.is_empty();
352 let relays = self.resolve_relay_effects(relays);
353 for consumer in consumers {
354 if let Some(error) = self.router.remove_consumer(consumer).err() {
355 error!(?consumer, ?error, "failed to remove declined room consumer");
356 }
357 }
358 let teardown = Self::media_teardowns([], routes).collect();
359 (relays, teardown, detached)
360 }
361
362 pub(in crate::engine::room) fn pending_consumer_routes_for_user(
363 &self,
364 user_id: &UserId,
365 ) -> impl Iterator<Item = PendingConsumerRouteView<'_>> {
366 self.route_graph
367 .attached_for_receiver(user_id)
368 .filter(|(_, current)| current.is_pending())
369 .filter_map(|(_, current)| {
370 let source = self.sources.source(current.source_id)?;
371 Some(PendingConsumerRouteView {
372 source,
373 selection: current.selection,
374 })
375 })
376 }
377
378 #[must_use]
380 pub(in crate::engine::room) fn consumer_source_selection(
381 &self,
382 key: &SubscriptionKey,
383 source_id: PublishedSourceId,
384 ) -> Option<ConsumerSourceSelection> {
385 self.route_graph.selection(key, source_id)
386 }
387
388 pub(in crate::engine::room) fn merge_subscription_intent(
390 &mut self,
391 key: SubscriptionKey,
392 intent: SourceSubscriptionIntent,
393 ) {
394 self.route_graph.merge_intent(key, intent);
395 }
396
397 #[must_use]
399 pub(in crate::engine::room) fn subscription_intent(
400 &self,
401 key: &SubscriptionKey,
402 ) -> SourceSubscriptionIntent {
403 self.route_graph.intent(key)
404 }
405
406 pub(in crate::engine::room) fn committed_consumer_user_ids_for_source(
407 &self,
408 source_id: PublishedSourceId,
409 ) -> BTreeSet<UserId> {
410 self.route_graph
411 .attached_for_source(source_id)
412 .filter(|(_, current)| current.committed().is_some())
413 .map(|(key, _)| key.receiver.clone())
414 .collect()
415 }
416
417 pub(in crate::engine::room) fn committed_consumer_user_ids_for_owner_sources(
418 &self,
419 user_id: &UserId,
420 ) -> BTreeSet<UserId> {
421 let route_graph = &self.route_graph;
422 self.sources
423 .ids_for_owner(user_id)
424 .flat_map(|source_id| route_graph.attached_for_source(source_id))
425 .filter(|(_, current)| current.committed().is_some())
426 .map(|(key, _)| key.receiver.clone())
427 .collect()
428 }
429
430 #[must_use]
431 pub(in crate::engine::room) fn committed_consumer_route_for_key(
432 &self,
433 key: &SubscriptionKey,
434 ) -> Option<ConsumerRouteView<'_>> {
435 let (key, current) = self.route_graph.current(key)?;
436 self.consumer_route(key, current)
437 }
438
439 #[must_use]
440 pub(in crate::engine::room) fn published_source(
441 &self,
442 source_id: PublishedSourceId,
443 ) -> Option<&PublishedSource> {
444 self.sources.source(source_id)
445 }
446
447 pub(super) fn missing_consumer_targets_for_source<'a>(
449 &mut self,
450 source_id: PublishedSourceId,
451 receivers: impl IntoIterator<Item = (&'a UserId, ConnectionId)>,
452 ) -> Vec<ConsumerSetupTarget> {
453 receivers
454 .into_iter()
455 .filter_map(|(user, connection)| self.consumer_target(user, connection, source_id))
456 .collect()
457 }
458
459 fn consumer_target(
460 &mut self,
461 user_id: &UserId,
462 connection_id: ConnectionId,
463 source_id: PublishedSourceId,
464 ) -> Option<ConsumerSetupTarget> {
465 let consumer_session = self.committed_transport_user_key(user_id.clone(), connection_id)?;
466 let (sources, route_graph) = (&self.sources, &mut self.route_graph);
467 Self::attach_consumer_target(route_graph, &consumer_session, sources.source(source_id)?)
468 }
469
470 fn attach_consumer_target(
472 route_graph: &mut RouteGraph,
473 consumer_session: &TransportSessionKey,
474 source: &PublishedSource,
475 ) -> Option<ConsumerSetupTarget> {
476 if source.descriptor.owner().user_id() == consumer_session.user_id() {
477 return None;
478 }
479 let target = ConsumerSetupTarget::new(consumer_session.clone(), source);
480 route_graph
481 .attach_for_setup(target.subscription_key(), target.source_id)
482 .then_some(target)
483 }
484
485 pub(super) fn missing_consumer_targets(
489 &mut self,
490 user_id: &UserId,
491 connection_id: ConnectionId,
492 include_source: impl Fn(&PublishedSource) -> bool,
493 ) -> Vec<ConsumerSetupTarget> {
494 let Some(consumer_session) =
495 self.committed_transport_user_key(user_id.clone(), connection_id)
496 else {
497 return Vec::new();
498 };
499 let (sources, route_graph) = (&self.sources, &mut self.route_graph);
500 sources
501 .iter()
502 .filter(|source| include_source(source))
503 .filter_map(|source| {
504 Self::attach_consumer_target(route_graph, &consumer_session, source)
505 })
506 .collect()
507 }
508
509 pub fn update_consumer_source_selection(
513 &mut self,
514 key: &SubscriptionKey,
515 source_id: PublishedSourceId,
516 route: &TransportConsumerRoute,
517 update: impl FnOnce(&mut ConsumerSourceSelection),
518 ) -> bool {
519 self.route_graph
520 .update_selection(key, source_id, route, update)
521 }
522
523 fn detach_user_sources(
524 &mut self,
525 user_id: &UserId,
526 ) -> (Vec<TransportSourceKey>, RemovedRoutes) {
527 let mut sources = Vec::new();
528 let mut removed = RemovedRoutes::default();
529 let source_ids = self.sources.ids_for_owner(user_id).collect::<Vec<_>>();
530 for source_id in source_ids {
531 if let Some((source, routes)) = self.remove_source(source_id) {
532 sources.push(source.transport);
533 removed.extend(routes);
534 }
535 }
536 (sources, removed)
537 }
538
539 fn remove_source(
540 &mut self,
541 source_id: PublishedSourceId,
542 ) -> Option<(PublishedSource, RemovedRoutes)> {
543 let source = self.sources.remove(source_id)?;
544 let routes = self.route_graph.detach_source(source_id);
545 Some((source, routes))
546 }
547
548 fn consumer_route<'a>(
549 &'a self,
550 key: &'a SubscriptionKey,
551 current: &'a CurrentPublication,
552 ) -> Option<ConsumerRouteView<'a>> {
553 let (route, mid) = current.committed()?;
554 let source = self.sources.source(current.source_id)?;
555 Some(ConsumerRouteView {
556 key,
557 route,
558 mid,
559 source,
560 selection: current.selection,
561 })
562 }
563
564 pub(in crate::engine::room) fn commit_publication(
572 &mut self,
573 publish: ValidatedPublish,
574 rtp: RouterRtpParameters,
575 encodings: &[SessionUploadEncoding],
576 media: TransportMediaId,
577 ) -> Result<PublishedSourceId, PublicationCommitError> {
578 let descriptor = allocate_source_descriptor(&mut self.sources, &publish, &rtp, encodings)?;
579 let source_id = descriptor.source_id();
580 let producer_id = ProducerId::allocate(&mut self.next_producer_id);
581 let routed = self
582 .router
583 .add_producer(publish.session_key.user_id(), producer_id)?;
584 self.sources.insert(PublishedSource {
585 descriptor,
586 transport: TransportSourceKey::new(publish.session_key, media),
587 rtp,
588 routed,
589 active: true,
590 activity_revision: SourceActivityRevision::default(),
591 });
592 Ok(source_id)
593 }
594
595 fn media_teardowns(
596 sources: impl IntoIterator<Item = TransportSourceKey>,
597 routes: impl IntoIterator<Item = TransportConsumerRoute>,
598 ) -> impl Iterator<Item = TransportTeardown> {
599 sources
600 .into_iter()
601 .map(|source| TransportTeardown::RemoveMedia {
602 session_key: source.session_key().clone(),
603 transport_media_id: source.transport_media_id(),
604 })
605 .chain(
606 routes
607 .into_iter()
608 .map(|route| TransportTeardown::RemoveMedia {
609 session_key: route.consumer_session_key().clone(),
610 transport_media_id: route.consumer_transport_media_id(),
611 }),
612 )
613 }
614
615 fn resolve_relay_effects(
620 &self,
621 effects: impl IntoIterator<Item = RelayRouteEffect>,
622 ) -> Vec<TransportRelayRouteEffect> {
623 effects
624 .into_iter()
625 .map(|effect| {
626 let route = effect.route;
627 TransportRelayRouteEffect {
628 source: TransportSourceKey::new(
629 self.transport_user_key(route.source_user, route.source_connection),
630 route.source_media,
631 ),
632 target_media_worker_id: route.target_worker,
633 action: effect.action,
634 }
635 })
636 .collect()
637 }
638
639 fn resolve_relay_effects_with_displaced(
645 &self,
646 effects: impl IntoIterator<Item = RelayRouteEffect>,
647 user_id: &UserId,
648 session_key: &TransportSessionKey,
649 ) -> Vec<TransportRelayRouteEffect> {
650 effects
651 .into_iter()
652 .map(|effect| {
653 let route = effect.route;
654 let source_session_key = if route.source_user == *user_id
655 && route.source_connection == session_key.connection_id()
656 {
657 session_key.clone()
660 } else {
661 self.transport_user_key(route.source_user, route.source_connection)
662 };
663 TransportRelayRouteEffect {
664 source: TransportSourceKey::new(source_session_key, route.source_media),
665 target_media_worker_id: route.target_worker,
666 action: effect.action,
667 }
668 })
669 .collect()
670 }
671
672 pub fn set_published_source_activity(
677 &mut self,
678 source_id: PublishedSourceId,
679 connection_id: ConnectionId,
680 active: bool,
681 ) -> Option<SourceActivityRevision> {
682 let source = self.sources.source_mut(source_id)?;
683 if source.transport.session_key().connection_id() != connection_id
684 || source.active == active
685 {
686 return None;
687 }
688 source.active = active;
689 source.activity_revision = source.activity_revision.next();
690 Some(source.activity_revision)
691 }
692
693 pub(super) fn source_activity_effects(
694 &self,
695 source: &TransportSourceKey,
696 update: SourceActivityUpdate,
697 ) -> Vec<TransportSourceActivityEffect> {
698 self.route_graph
699 .source_activity_target_workers(source)
700 .map(|target_media_worker_id| TransportSourceActivityEffect {
701 source: source.clone(),
702 target_media_worker_id,
703 update,
704 })
705 .collect()
706 }
707
708 pub fn commit_session_placement(
724 &mut self,
725 user_id: &UserId,
726 connection_id: ConnectionId,
727 previous_connection: Option<ConnectionId>,
728 home_placement: RouterPlacement,
729 ) -> Result<SessionPlacementCommit, SessionPlacementRejection> {
730 let previous_session_key = if let Some(previous_connection) = previous_connection {
732 let Some(key) = self.committed_transport_user_key(user_id.clone(), previous_connection)
733 else {
734 return Err(SessionPlacementRejection::MissingPreviousSession {
735 previous_connection,
736 });
737 };
738 Some(key)
739 } else {
740 None
741 };
742 let media_worker = self
743 .router
744 .commit_session_placement(user_id, connection_id, home_placement)
745 .map_err(SessionPlacementRejection::Router)?;
746 let session_key =
747 self.transport_session_key(user_id.clone().into(), connection_id, media_worker);
748 let receipt = CommittedTransportReceipt {
749 connection_id,
750 transport_session_key: session_key,
751 };
752 let replacement_transport_plan = previous_session_key.as_ref().map_or_else(
753 RoomTransportPlan::default,
754 |replaced_session_key| {
755 let close_session = TransportTeardown::CloseSession {
756 session_key: replaced_session_key.clone(),
757 };
758 let (_, removed_sources) = self.detach_user_sources(user_id);
761 let receiver_relays = self.route_graph.reset_receiver_for_replacement(user_id);
762 let RemovedRoutes {
763 routes, mut relays, ..
764 } = removed_sources;
765 relays.extend(receiver_relays);
766 let relay_effects = self.resolve_relay_effects_with_displaced(
767 relays,
768 user_id,
769 replaced_session_key,
770 );
771 let teardown = Self::media_teardowns([], routes).chain([close_session]);
772 RoomTransportPlan::from_relays_and_teardown(relay_effects, teardown)
773 },
774 );
775 Ok(SessionPlacementCommit {
776 receipt,
777 replacement_transport_plan,
778 })
779 }
780
781 pub fn remove_session(&mut self, user_id: &UserId) -> RoomTransportPlan {
787 let (sources, mut removed) = self.detach_user_sources(user_id);
788 removed.extend(self.route_graph.remove_receiver(user_id));
789 let teardown = Self::media_teardowns(sources, removed.routes);
790 let relay_effects = self.resolve_relay_effects(removed.relays);
792 if let Some(error) = self.router.remove_session(user_id).err() {
793 error!(?user_id, ?error, "failed to remove user from room router");
794 }
795 RoomTransportPlan::from_relays_and_teardown(relay_effects, teardown)
796 }
797
798 pub(super) fn commit_consumer_setup(
806 &mut self,
807 setup: DeclaredConsumerSetup,
808 selection: ConsumerSourceSelection,
809 ) -> Result<CommittedConsumerSetup, (TransportConsumerRoute, Vec<TransportRelayRouteEffect>)>
810 {
811 let DeclaredConsumerSetup {
812 pending:
813 PendingConsumerSetup {
814 target,
815 consumer,
816 reservation,
817 sender,
818 rtp,
819 relays: _,
820 },
821 route,
822 mid,
823 } = setup;
824 let active = selection.delivery_active();
825 let declared_active = reservation.selection().delivery_active();
826 let committed_mid = mid.unwrap_or_else(|| {
827 rtp.mid()
828 .map_or_else(|| consumer.to_string(), ToOwned::to_owned)
829 });
830 let (route_graph, room_router) = (&mut self.route_graph, &mut self.router);
831 let result =
833 route_graph.commit(reservation, route.clone(), committed_mid, selection, || {
834 match room_router.add_consumer(target.session.user_id(), consumer, target.routed) {
835 Ok(routed_consumer) => Some(routed_consumer),
836 Err(error) => {
837 warn!(
838 consumer_user_id = ?target.session.user_id(),
839 source_id = ?target.source_id,
840 ?error,
841 "router rejected consumer creation"
842 );
843 None
844 }
845 }
846 });
847 if let Err(relays) = result {
848 return Err((route, self.resolve_relay_effects(relays)));
849 }
850 Ok(CommittedConsumerSetup {
851 target,
852 route,
853 sender,
854 transport_activity_update: (active != declared_active).then_some(active),
855 })
856 }
857
858 pub(super) fn reserve_consumer_setup(
862 &mut self,
863 target: ConsumerSetupTarget,
864 consumer: ConsumerId,
865 selection: ConsumerSourceSelection,
866 sender: OutboundSender,
867 rtp: RouterRtpParameters,
868 ) -> Option<PendingConsumerSetup> {
869 let key = target.subscription_key();
870 let relay_active = selection.active();
874 let reservation =
875 self.route_graph
876 .reserve_consumer_setup(key, target.source_id, selection)?;
877 let source_worker = target.source.session_key().media_worker_id();
878 let target_worker = target.session.media_worker_id();
879 let relays = if source_worker == target_worker {
880 Vec::new()
881 } else {
882 let relays =
883 self.route_graph
884 .reserve_relay(&reservation, &target, target_worker, relay_active);
885 self.resolve_relay_effects(relays)
886 };
887 Some(PendingConsumerSetup {
888 target,
889 consumer,
890 reservation,
891 sender,
892 rtp,
893 relays,
894 })
895 }
896
897 pub(super) fn release_consumer_setup(
899 &mut self,
900 setup: PendingConsumerSetup,
901 ) -> Vec<TransportRelayRouteEffect> {
902 let relays = self.route_graph.release_consumer_setup(setup.reservation);
903 self.resolve_relay_effects(relays)
904 }
905
906 pub(super) fn set_consumer_activity(
909 &mut self,
910 user_id: &UserId,
911 connection_id: ConnectionId,
912 target_user_id: &UserId,
913 stream_id: &UserStreamId,
914 active: bool,
915 receiver_deafened: bool,
916 ) -> Option<ConsumerActivityCommit> {
917 let key = SubscriptionKey::new(user_id, target_user_id, stream_id);
918 let source_id = self.source_id_for_owner_stream(target_user_id, stream_id)?;
919 let policy_pause_reason = (receiver_deafened
921 && self
922 .source_descriptor(source_id)
923 .is_some_and(|source| source.media_kind() == MediaKind::Audio))
924 .then_some(PolicyPauseReason::ReceiverDeafened);
925 let relay_effects = self.route_graph.set_activity(
926 &key,
927 source_id,
928 connection_id,
929 active,
930 policy_pause_reason,
931 )?;
932 let relay_effects = self.resolve_relay_effects(relay_effects);
933 let update = self
934 .committed_consumer_route_for_key(&key)
935 .filter(|route| route.route.consumer_session_key().connection_id() == connection_id)
936 .map(|route| {
937 ReceiverRouteActivity::new(route.target(), route.selection.delivery_active())
938 });
939 Some(ConsumerActivityCommit {
940 update,
941 relay_effects,
942 })
943 }
944}