Skip to main content

o_sfu_core/engine/room/media_graph/
topology.rs

1//! Room publications and subscriptions share transport-placement authority.
2
3use 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/// Keeps source records and route realization beside [`Router`] so transport
44/// effects resolve from the same committed placement.
45///
46/// Transport-facing mutations return resolved work for execution after the room
47/// state lock is released.
48#[derive(Debug)]
49pub struct RoomTopology {
50    instance: RoomInstanceId,
51    sources: PublishedSources,
52    route_graph: RouteGraph,
53    router: Router,
54    next_producer_id: u64,
55    // Revision of committed video allocation inputs across transport awaits.
56    video_allocation_revision: u64,
57}
58
59/// Receipt acknowledging that [`RoomTopology`] committed a session placement.
60///
61/// It snapshots the committed connection identity and worker-resolved transport
62/// key across the room-state lock boundary. Placement lifetime remains controlled by
63/// [`RoomTopology::commit_session_placement`], [`RoomTopology::remove_session`]
64/// or [`RoomTopology::retire_committed_placement`].
65///
66/// # Admission handoff
67///
68/// Membership keeps the receipt while
69/// [`RoomEffects`](crate::engine::room::effects::batch::RoomEffects) consumes the
70/// join effects:
71///
72/// ```rust,ignore
73/// let JoinCommit { receipt, effects, transport_plan } = admission.commit(self, joined_fanout).await?;
74/// RoomEffects::from_join(effects, transport_plan).execute(self, context).await;
75/// Ok(receipt)
76/// ```
77#[derive(Debug, Clone, PartialEq, Eq)]
78pub struct CommittedTransportReceipt {
79    /// Transport identity resolved from the committed media worker placement.
80    pub transport_session_key: TransportSessionKey,
81}
82
83/// The new placement is authoritative before displaced-session cleanup is returned.
84///
85/// Bundling the receipt with resolved cleanup lets membership release `room.state`
86/// without looking up the displaced placement again.
87#[derive(Debug)]
88pub struct SessionPlacementCommit {
89    pub receipt: CommittedTransportReceipt,
90    /// Empty for a first placement.
91    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    /// Returns `None` unless the exact user connection remains committed.
142    #[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    /// Requires the exact user connection to remain committed.
156    ///
157    /// Use [`Self::committed_transport_user_key`] for stale callbacks or teardown
158    /// races where the placement may already be retired.
159    ///
160    /// # Panics
161    ///
162    /// Panics when no committed router placement exists.
163    #[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    /// Returns `None` when `connection_id` is not the user's committed placement.
193    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    /// Counts each logical subscription once while pending or committed.
205    #[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    /// Counts committed consumers, including inactive selections.
214    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    /// Returns the source ID only when owner, connection and stream remain current.
264    #[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    /// Iterates committed sources in source-ID order.
277    pub(in crate::engine::room) fn published_sources(
278        &self,
279    ) -> impl Iterator<Item = &PublishedSource> {
280        self.sources.iter()
281    }
282
283    /// Counts distinct users with active publications for each logical stream.
284    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    /// Returns a detector owner with an active promotable source in the same group.
299    #[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    /// Returns `None` when the attached source differs from `source_id`.
374    #[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    /// Returns merged receiver intent or the default for a missing subscription.
384    #[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    /// Attaches `source_id` to each eligible receiver lacking a route realization.
434    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    /// Skips self-consumption and attaches only an absent realization.
457    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    /// Attaches each matching source with no realization to the exact committed receiver.
472    ///
473    /// Returns an empty list when `connection_id` is stale.
474    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    /// Updates selection while subscription, source and exact route still match.
496    ///
497    /// Returns `false` when async transport work refers to a displaced route.
498    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    /// Commits a hold only for the active source, subscription and exact route.
515    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    /// Commits a negotiated source only after its router producer succeeds.
577    ///
578    /// # Errors
579    ///
580    /// Returns [`PublicationCommitError::Source`] when descriptor allocation or
581    /// validation fails. Returns [`PublicationCommitError::Router`] when the
582    /// router rejects the producer dependency.
583    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    /// Advances the revision only when the exact source connection changes activity.
628    ///
629    /// Returns `None` when the source is missing, the connection is stale or the
630    /// requested activity already matches.
631    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    /// Commits the new placement before returning cleanup for `previous_connection`.
669    ///
670    /// Replacement joins must pass the currently committed `previous_connection`.
671    ///
672    /// # Errors
673    ///
674    /// Returns [`SessionPlacementRejection::MissingPreviousSession`] when the
675    /// expected replacement target is no longer committed. Returns
676    /// [`SessionPlacementRejection::Router`] when the router rejects the new
677    /// connection or placement.
678    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        // Preserve the displaced key before the router replaces its user mapping.
686        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                // Session close removes source media while subscriber routes need
714                // explicit teardown.
715                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    /// Returns the cleanup plan even if router removal fails.
735    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    /// Selects MID from the transport declaration, negotiated RTP MID then the
751    /// consumer identity.
752    ///
753    /// # Errors
754    ///
755    /// Returns the declared route and relay release effects when the reservation
756    /// is stale or the router rejects the consumer dependency.
757    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        // Remove pending state first so router rejection cannot strand a reservation.
784        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    /// Returns `None` unless the route realization remains absent.
812    ///
813    /// A cross-worker reservation also claims relay ownership.
814    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        // Relay activity follows receiver intent rather than policy-gated delivery.
824        // A temporary policy pause must remain resumable without rebuilding the
825        // shared cross-worker source path.
826        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    /// A stale reservation releases no relay ownership.
849    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    /// Merges sparse intent before projecting activity to an attached route.
857    ///
858    /// Empty updates have no effect. Intent survives absent publications and
859    /// unready receivers. Returns no activity work for layout-only updates,
860    /// missing sources, changed attachments or displaced consumer connections.
861    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        // Deafening pauses audio delivery without replacing explicit receiver intent.
888        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}