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, 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/// 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}
56
57/// Receipt acknowledging that [`RoomTopology`] committed a session placement.
58///
59/// It snapshots the committed connection identity and worker-resolved transport
60/// key across the room-state lock boundary. Placement lifetime remains controlled by
61/// [`RoomTopology::commit_session_placement`], [`RoomTopology::remove_session`]
62/// or [`RoomTopology::retire_committed_placement`].
63///
64/// # Admission handoff
65///
66/// Membership keeps the receipt while
67/// [`RoomEffects`](crate::engine::room::effects::batch::RoomEffects) consumes the
68/// join commit:
69///
70/// ```rust,ignore
71/// let commit = admission.commit(self, joined_fanout).await?;
72/// let receipt = commit.receipt.clone();
73///
74/// RoomEffects::from_join(commit).execute(self, context).await;
75/// Ok(receipt)
76/// ```
77#[derive(Debug, Clone, PartialEq, Eq)]
78pub struct CommittedTransportReceipt {
79    /// Room-local connection identity used to reject stale operations.
80    pub connection_id: ConnectionId,
81    /// Transport identity resolved from the committed media worker placement.
82    pub transport_session_key: TransportSessionKey,
83}
84
85/// The new placement is authoritative before displaced-session cleanup is returned.
86///
87/// Bundling the receipt with resolved cleanup lets membership release `room.state`
88/// without looking up the displaced placement again.
89#[derive(Debug)]
90pub struct SessionPlacementCommit {
91    pub receipt: CommittedTransportReceipt,
92    /// Empty for a first placement.
93    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    /// Returns `None` unless the exact user connection remains committed.
133    #[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    /// Requires the exact user connection to remain committed.
147    ///
148    /// Use [`Self::committed_transport_user_key`] for stale callbacks or teardown
149    /// races where the placement may already be retired.
150    ///
151    /// # Lookup choice
152    ///
153    /// Receiver work may carry a stale connection while relay effects from
154    /// committed graph state use the strict lookup:
155    ///
156    /// ```rust,ignore
157    /// let Some(consumer_session) =
158    ///     topology.committed_transport_user_key(user_id.clone(), connection_id)
159    /// else {
160    ///     return Vec::new();
161    /// };
162    ///
163    /// let source_session =
164    ///     topology.transport_user_key(route.source_user, route.source_connection);
165    /// ```
166    ///
167    /// # Panics
168    ///
169    /// Panics when no committed router placement exists.
170    #[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    /// Returns `None` when `connection_id` is not the user's committed placement.
200    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    /// Counts each logical subscription once while pending or committed.
212    #[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    /// Returns the source ID only when owner, connection and stream remain current.
271    #[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    /// Iterates committed sources in source-ID order.
284    pub(in crate::engine::room) fn published_sources(
285        &self,
286    ) -> impl Iterator<Item = &PublishedSource> {
287        self.sources.iter()
288    }
289
290    /// Counts distinct users with active publications for each logical stream.
291    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    /// Returns a detector owner with an active promotable source in the same group.
306    #[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    /// Returns `None` when the attached source differs from `source_id`.
379    #[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    /// Ignores empty updates and applies `active` to an attached selection.
389    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    /// Returns merged receiver intent or the default for a missing subscription.
398    #[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    /// Attaches `source_id` to each eligible receiver lacking a route realization.
448    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    /// Skips self-consumption and attaches only an absent realization.
471    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    /// Attaches each matching source with no realization to the exact committed receiver.
486    ///
487    /// Returns an empty list when `connection_id` is stale.
488    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    /// Updates selection while subscription, source and exact route still match.
510    ///
511    /// Returns `false` when async transport work refers to a displaced route.
512    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    /// Commits a negotiated source only after its router producer succeeds.
565    ///
566    /// # Errors
567    ///
568    /// Returns [`PublicationCommitError::Source`] when descriptor allocation or
569    /// validation fails. Returns [`PublicationCommitError::Router`] when the
570    /// router rejects the producer dependency.
571    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    /// # Panics
616    ///
617    /// Panics if `effects` violates the topology invariant that every relay source
618    /// has a committed router placement.
619    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    /// Resolves relay effects while retaining a displaced source's transport key.
640    ///
641    /// # Panics
642    ///
643    /// Panics if any non-displaced relay source lacks a committed placement.
644    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                    // Current lookup now resolves the replacement. Use the key
658                    // captured before displacement for this source's cleanup.
659                    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    /// Advances the revision only when the exact source connection changes activity.
673    ///
674    /// Returns `None` when the source is missing, the connection is stale or the
675    /// requested activity already matches.
676    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    /// Commits the new placement before returning cleanup for `previous_connection`.
709    ///
710    /// Replacement joins must pass the currently committed `previous_connection`.
711    ///
712    /// # Errors
713    ///
714    /// Returns [`SessionPlacementRejection::MissingPreviousSession`] when the
715    /// expected replacement target is no longer committed. Returns
716    /// [`SessionPlacementRejection::Router`] when the router rejects the new
717    /// connection or placement.
718    ///
719    /// # Panics
720    ///
721    /// Panics when existing relay state refers to another uncommitted source
722    /// placement.
723    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        // Preserve the displaced key before the router replaces its user mapping.
731        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                // Session close removes source media while subscriber routes need
759                // explicit teardown.
760                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    /// Returns the cleanup plan even if router removal fails.
782    ///
783    /// # Panics
784    ///
785    /// Panics when detached relay state refers to an uncommitted source placement.
786    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        // Resolve relay keys before router removal makes source placement unavailable.
791        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    /// Selects MID from the transport declaration, negotiated RTP MID then the
799    /// consumer identity.
800    ///
801    /// # Errors
802    ///
803    /// Returns the declared route and relay release effects when the reservation
804    /// is stale or the router rejects the consumer dependency.
805    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        // Remove pending state first so router rejection cannot strand a reservation.
832        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    /// Returns `None` unless the route realization remains absent.
859    ///
860    /// A cross-worker reservation also claims relay ownership.
861    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        // Relay activity follows receiver intent rather than policy-gated delivery.
871        // A temporary policy pause must remain resumable without rebuilding the
872        // shared cross-worker source path.
873        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    /// A stale reservation releases no relay ownership.
898    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    /// Returns `None` when the source is missing, its attachment changed or the
907    /// committed consumer belongs to another receiver connection.
908    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        // Deafening pauses audio delivery without replacing explicit receiver intent.
920        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}