Skip to main content

o_sfu_core/engine/room/media_graph/
route_graph.rs

1//! Logical subscriptions and their transport realization.
2//!
3//! Explicit receiver intent survives while no publication is attached. A
4//! current publication realizes as `Absent -> Pending -> Committed`. Current
5//! setup failures, router rejection and receiver replacement return realization
6//! to `Absent`. Reservation IDs reject async completion after source detach or
7//! replacement.
8//!
9//! ```text
10//! logical subscription realization:
11//!      publication attached
12//!              |
13//!              v
14//!   +--------------------+  reserve setup
15//!   | ConsumerRealization| ------------------> +-------------------------------------+
16//!   |      Absent        |                     | ConsumerRealization::Pending        |
17//!   +--------------------+ <------------------ | (RouteReservationId, Option<Relay>) |
18//!              ^       setup fails/rejected    +-------------------------------------+
19//!              |                                                  |
20//!              |              decline/replace                     v setup accepted
21//!              +---------------------------------- +---------------------------------+
22//!                                                  | ConsumerRealization::Committed  |
23//!                                                  +---------------------------------+
24//!
25//! cross-worker relay sharing:
26//!   receiver 1 (worker B) \
27//!   receiver 2 (worker B)  --> [ RelayRouteKey (source worker A -> worker B) ]
28//!   receiver 3 (worker B) /               |
29//!                                         v
30//!                         shared relay channel (worker A -> B)
31//!                                         |
32//!                             +-----------+-----------+
33//!                             v           v           v
34//!                          recv 1      recv 2      recv 3 (local fanout on worker B)
35//! ```
36//!
37//! Intent without a current publication keeps `Subscription.current` at `None`.
38//! Accepted setup requires a current identity, a successful transport declaration
39//! and router dependency acceptance. Source detach returns to `None` while
40//! retaining intent. `remove_receiver` deletes the record and its intent.
41//!
42//! Each receiver retains at most the room session limit of absent publisher
43//! targets. Sparse updates keep the target's original age. A new target evicts
44//! every stream preference of the oldest absent target. Present members do not
45//! consume this allowance, including members with no publications. Joining
46//! promotes existing intent and departure enters the allowance as a new target.
47//!
48//! Cross-worker relays are shared by subscriptions with the same source and
49//! target worker. A relay remains active while any owner has active receiver
50//! intent. An active first owner emits `Install` followed by
51//! `SetActivity(Active)`. Later aggregate activity changes emit `SetActivity`.
52//! Removing the last owner emits `Release`.
53
54use std::{
55    collections::{BTreeMap, BTreeSet, VecDeque, btree_map::Entry},
56    mem,
57    time::Instant,
58};
59
60use o_sfu_router::topology::RoutedConsumerId;
61
62use super::{
63    ConsumerSourceSelection, SubscriptionKey, consumer_setup::ConsumerSetupTarget,
64    remove_from_index_set,
65};
66use crate::engine::{
67    ConnectionId, MediaWorkerId, UserId,
68    media_transport::{
69        ConsumerActivity, RelayRouteActivity, TransportConsumerRoute, TransportMediaId,
70        TransportRelayRouteAction, TransportRelayRouteEffect, TransportSessionKey,
71        TransportSourceKey,
72    },
73    source_model::{
74        PolicyPauseReason, PublishedSourceId, SourceSelector, SourceSubscriptionIntent,
75    },
76};
77
78#[derive(Debug)]
79pub(super) struct RouteGraph {
80    entries: BTreeMap<SubscriptionKey, Subscription>,
81    by_receiver: BTreeMap<UserId, BTreeSet<SubscriptionKey>>,
82    by_source: BTreeMap<PublishedSourceId, BTreeSet<SubscriptionKey>>,
83    relays: BTreeMap<RelayRouteKey, RelayOwners>,
84    next_reservation: RouteReservationId,
85    pending_targets: BTreeMap<UserId, VecDeque<UserId>>,
86    pending_target_limit: usize,
87}
88
89type RelayOwners = BTreeMap<SubscriptionKey, RelayRouteActivity>;
90
91#[derive(Debug, Default)]
92struct Subscription {
93    intent: SourceSubscriptionIntent,
94    current: Option<CurrentPublication>,
95}
96
97#[derive(Debug)]
98pub(super) struct CurrentPublication {
99    pub source_id: PublishedSourceId,
100    pub selection: ConsumerSourceSelection,
101    realization: ConsumerRealization,
102}
103
104#[derive(Debug, Default)]
105enum ConsumerRealization {
106    #[default]
107    Absent,
108    Pending(RouteReservationId, Option<RouteRelay>),
109    Committed(CommittedConsumerRoute),
110}
111
112/// Transport realization and adaptation hold share the exact route lifetime.
113#[derive(Debug)]
114pub(super) struct CommittedConsumerRoute {
115    pub route: TransportConsumerRoute,
116    pub mid: String,
117    routed: RoutedConsumerId,
118    relay: Option<RouteRelay>,
119    pub pending_upgrade: Option<PendingUpgrade>,
120}
121
122/// Eligibility must remain continuous for this exact deliverable post-fit target.
123#[derive(Debug, Clone, Copy, PartialEq, Eq)]
124pub(in crate::engine::room) struct PendingUpgrade {
125    pub selector: SourceSelector,
126    pub deadline: Instant,
127}
128
129#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
130struct RouteReservationId(u64);
131
132#[derive(Debug)]
133pub struct ConsumerRouteReservation {
134    key: SubscriptionKey,
135    source_id: PublishedSourceId,
136    declared_activity: ConsumerActivity,
137    id: RouteReservationId,
138}
139
140#[derive(Debug, Clone, PartialEq, Eq)]
141struct RouteRelay {
142    route: RelayRouteKey,
143    activity: RelayRouteActivity,
144}
145
146struct TakenPending {
147    relay: Option<RouteRelay>,
148}
149
150/// Captured source placement survives replacement until its relay owners release it.
151#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
152pub struct RelayRouteKey {
153    pub source: TransportSourceKey,
154    pub target_worker: MediaWorkerId,
155}
156
157#[derive(Debug, Default)]
158pub(super) struct RemovedRoutes {
159    pub routes: Vec<TransportConsumerRoute>,
160    pub consumers: Vec<RoutedConsumerId>,
161    pub relays: Vec<TransportRelayRouteEffect>,
162}
163
164impl RemovedRoutes {
165    pub(super) fn extend(&mut self, mut other: Self) {
166        self.routes.append(&mut other.routes);
167        self.consumers.append(&mut other.consumers);
168        self.relays.append(&mut other.relays);
169    }
170}
171
172impl RouteGraph {
173    pub(super) fn new(pending_target_limit: usize) -> Self {
174        Self {
175            entries: BTreeMap::new(),
176            by_receiver: BTreeMap::new(),
177            by_source: BTreeMap::new(),
178            relays: BTreeMap::new(),
179            next_reservation: RouteReservationId::default(),
180            pending_targets: BTreeMap::new(),
181            pending_target_limit,
182        }
183    }
184
185    /// Membership promotion preserves every stream preference and frees its
186    /// absent-target allowance without duplicating the logical subscriptions.
187    pub(super) fn publisher_joined(&mut self, publisher: &UserId) {
188        self.pending_targets.retain(|_, pending| {
189            pending.retain(|target| target != publisher);
190            !pending.is_empty()
191        });
192    }
193
194    /// Call after detaching every source and removing this user's receiver state.
195    /// A departed member becomes the newest absent target for each receiver
196    /// retaining explicit intent, even when it never published media.
197    pub(super) fn publisher_left(&mut self, publisher: &UserId) -> usize {
198        let receivers: Vec<_> = self
199            .by_receiver
200            .iter()
201            .filter(|(_, keys)| keys.iter().any(|key| &key.publisher == publisher))
202            .map(|(receiver, _)| receiver.clone())
203            .collect();
204        receivers
205            .iter()
206            .map(|receiver| self.retain_pending_target(receiver, publisher))
207            .sum()
208    }
209
210    fn retain_pending_target(&mut self, receiver: &UserId, publisher: &UserId) -> usize {
211        let pending = self.pending_targets.entry(receiver.clone()).or_default();
212        if pending.contains(publisher) {
213            return 0;
214        }
215        pending.push_back(publisher.clone());
216        if pending.len() <= self.pending_target_limit {
217            return 0;
218        }
219        let Some(evicted) = pending.pop_front() else {
220            return 0;
221        };
222        if pending.is_empty() {
223            self.pending_targets.remove(receiver);
224        }
225        // Only absent publishers enter this index. Their detached records have
226        // no transport or relay ownership, so eviction removes intent alone.
227        if let Some(keys) = self.by_receiver.get_mut(receiver) {
228            keys.retain(|key| {
229                if key.publisher != evicted {
230                    return true;
231                }
232                let removed = self.entries.remove(key);
233                debug_assert!(removed.is_none_or(|entry| entry.current.is_none()));
234                false
235            });
236            if keys.is_empty() {
237                self.by_receiver.remove(receiver);
238            }
239        }
240        1
241    }
242
243    pub(super) fn subscription_count(&self) -> usize {
244        self.entries
245            .values()
246            .filter(|entry| entry.has_consumer_setup_or_route())
247            .count()
248    }
249
250    pub(super) fn count(&self) -> usize {
251        self.attached()
252            .filter(|(_, current)| current.committed().is_some())
253            .count()
254    }
255
256    #[cfg(test)]
257    pub(super) fn record_count(&self) -> usize {
258        self.entries.len()
259    }
260
261    pub(super) fn merge_intent(
262        &mut self,
263        key: &SubscriptionKey,
264        update: SourceSubscriptionIntent,
265        publisher_present: bool,
266    ) -> usize {
267        if update.is_empty() {
268            return 0;
269        }
270        let entry = self.entry(key.clone());
271        entry.intent.merge(update);
272        if let (Some(active), Some(current)) = (update.active(), entry.current.as_mut()) {
273            current.set_active(active);
274        }
275        if publisher_present {
276            0
277        } else {
278            self.retain_pending_target(&key.receiver, &key.publisher)
279        }
280    }
281
282    pub(super) fn intent(&self, key: &SubscriptionKey) -> SourceSubscriptionIntent {
283        self.entries
284            .get(key)
285            .map_or_else(SourceSubscriptionIntent::default, |entry| entry.intent)
286    }
287
288    pub(super) fn attach_for_setup(
289        &mut self,
290        key: SubscriptionKey,
291        source_id: PublishedSourceId,
292    ) -> bool {
293        let entry = self.entry(key.clone());
294        if let Some(current) = &entry.current {
295            return current.source_id == source_id
296                && matches!(current.realization, ConsumerRealization::Absent);
297        }
298        entry.current = Some(CurrentPublication {
299            source_id,
300            selection: ConsumerSourceSelection::open(entry.intent.active().unwrap_or(true)),
301            realization: ConsumerRealization::Absent,
302        });
303        self.by_source.entry(source_id).or_default().insert(key);
304        true
305    }
306
307    pub(super) fn set_activity(
308        &mut self,
309        key: &SubscriptionKey,
310        source_id: PublishedSourceId,
311        connection_id: ConnectionId,
312        active: bool,
313        policy_pause_reason: Option<PolicyPauseReason>,
314    ) -> Option<Vec<TransportRelayRouteEffect>> {
315        let relay = {
316            let current = self
317                .entries
318                .get_mut(key)
319                .and_then(|entry| entry.current.as_mut())
320                .filter(|current| current.source_id == source_id)?;
321            if let ConsumerRealization::Committed(committed) = &current.realization
322                && committed.route.consumer_session_key().connection_id() != connection_id
323            {
324                return None;
325            }
326            current.set_active(active);
327            if let Some(reason) = policy_pause_reason {
328                current.selection.set_policy_pause_reason(Some(reason));
329            }
330            current
331                .realization
332                .set_relay_activity(RelayRouteActivity::from_active(active))
333        };
334        Some(relay.map_or_else(Vec::new, |relay| self.set_relay_owner(key, &relay, false)))
335    }
336
337    pub(super) fn reserve_consumer_setup(
338        &mut self,
339        key: SubscriptionKey,
340        source_id: PublishedSourceId,
341        selection: ConsumerSourceSelection,
342    ) -> Option<ConsumerRouteReservation> {
343        let current = self.entries.get_mut(&key)?.current.as_mut()?;
344        if current.source_id != source_id
345            || !matches!(current.realization, ConsumerRealization::Absent)
346        {
347            return None;
348        }
349        self.next_reservation.0 += 1;
350        let id = self.next_reservation;
351        current.selection = selection;
352        current.realization = ConsumerRealization::Pending(id, None);
353        Some(ConsumerRouteReservation {
354            key,
355            source_id,
356            declared_activity: ConsumerActivity::from_active(selection.delivery_active()),
357            id,
358        })
359    }
360
361    pub(super) fn release_consumer_setup(
362        &mut self,
363        reservation: ConsumerRouteReservation,
364    ) -> Vec<TransportRelayRouteEffect> {
365        let Some(TakenPending { relay }) = self.take_pending(&reservation) else {
366            return Vec::new();
367        };
368        let key = reservation.key;
369        relay.map_or_else(Vec::new, |relay| self.release_relay(&key, &relay))
370    }
371
372    pub(super) fn commit(
373        &mut self,
374        reservation: ConsumerRouteReservation,
375        route: TransportConsumerRoute,
376        mid: String,
377        selection: ConsumerSourceSelection,
378        accept: impl FnOnce() -> Option<RoutedConsumerId>,
379    ) -> Result<(), Vec<TransportRelayRouteEffect>> {
380        // Consume the reservation before router acceptance so every async
381        // completion is terminal and cannot reuse its identity after failure.
382        let pending = self.take_pending(&reservation);
383        let ConsumerRouteReservation { key, source_id, .. } = reservation;
384        let Some(TakenPending { relay }) = pending else {
385            return Err(Vec::new());
386        };
387        let Some(current) = self.current_mut(&key, source_id) else {
388            return Err(relay.map_or_else(Vec::new, |relay| self.release_relay(&key, &relay)));
389        };
390        let Some(routed) = accept() else {
391            return Err(relay.map_or_else(Vec::new, |relay| self.release_relay(&key, &relay)));
392        };
393        current.selection = selection;
394        current.realization = ConsumerRealization::Committed(CommittedConsumerRoute {
395            route,
396            mid,
397            routed,
398            relay,
399            pending_upgrade: None,
400        });
401        Ok(())
402    }
403
404    pub(super) fn update_selection(
405        &mut self,
406        key: &SubscriptionKey,
407        source_id: PublishedSourceId,
408        route: &TransportConsumerRoute,
409        update: impl FnOnce(&mut ConsumerSourceSelection),
410    ) -> bool {
411        let Some(current) = self.current_mut(key, source_id) else {
412            return false;
413        };
414        let ConsumerRealization::Committed(committed) = &current.realization else {
415            return false;
416        };
417        if &committed.route != route {
418            return false;
419        }
420        update(&mut current.selection);
421        if !current.selection.active() {
422            current.clear_upgrade();
423        }
424        true
425    }
426
427    pub(super) fn update_upgrade(
428        &mut self,
429        key: &SubscriptionKey,
430        source_id: PublishedSourceId,
431        route: &TransportConsumerRoute,
432        pending_upgrade: Option<PendingUpgrade>,
433    ) -> bool {
434        let Some(current) = self.current_mut(key, source_id) else {
435            return false;
436        };
437        let ConsumerRealization::Committed(committed) = &mut current.realization else {
438            return false;
439        };
440        if &committed.route != route || !current.selection.active() {
441            return false;
442        }
443        committed.pending_upgrade = pending_upgrade;
444        true
445    }
446
447    pub(super) fn clear_source_upgrades(&mut self, source_id: PublishedSourceId) {
448        let (entries, by_source) = (&mut self.entries, &self.by_source);
449        for key in by_source.get(&source_id).into_iter().flatten() {
450            if let Some(current) = entries
451                .get_mut(key)
452                .and_then(|entry| entry.current.as_mut())
453            {
454                current.clear_upgrade();
455            }
456        }
457    }
458
459    pub(super) fn selection(
460        &self,
461        key: &SubscriptionKey,
462        source_id: PublishedSourceId,
463    ) -> Option<ConsumerSourceSelection> {
464        let current = self.entries.get(key)?.current.as_ref()?;
465        (current.source_id == source_id).then_some(current.selection)
466    }
467
468    pub(super) fn attached(&self) -> impl Iterator<Item = (&SubscriptionKey, &CurrentPublication)> {
469        self.entries
470            .iter()
471            .filter_map(|(key, entry)| Some((key, entry.current.as_ref()?)))
472    }
473
474    pub(super) fn current(
475        &self,
476        key: &SubscriptionKey,
477    ) -> Option<(&SubscriptionKey, &CurrentPublication)> {
478        let (key, entry) = self.entries.get_key_value(key)?;
479        Some((key, entry.current.as_ref()?))
480    }
481
482    pub(super) fn attached_for_receiver(
483        &self,
484        receiver: &UserId,
485    ) -> impl Iterator<Item = (&SubscriptionKey, &CurrentPublication)> {
486        self.by_receiver
487            .get(receiver)
488            .into_iter()
489            .flat_map(BTreeSet::iter)
490            .filter_map(|key| self.current(key))
491    }
492
493    pub(super) fn attached_for_source(
494        &self,
495        source_id: PublishedSourceId,
496    ) -> impl Iterator<Item = (&SubscriptionKey, &CurrentPublication)> {
497        self.by_source
498            .get(&source_id)
499            .into_iter()
500            .flat_map(BTreeSet::iter)
501            .filter_map(|key| self.current(key))
502    }
503
504    pub(super) fn detach_source(&mut self, source_id: PublishedSourceId) -> RemovedRoutes {
505        let keys = self.by_source.remove(&source_id).unwrap_or_default();
506        let mut removed = RemovedRoutes::default();
507        for key in keys {
508            let (realization, prune) = {
509                let Some(entry) = self.entries.get_mut(&key) else {
510                    continue;
511                };
512                let Some(current) = entry.current.take() else {
513                    continue;
514                };
515                debug_assert_eq!(current.source_id, source_id);
516                (current.realization, entry.intent.is_empty())
517            };
518            self.collect_removed_realization(&key, realization, &mut removed);
519            if prune {
520                self.entries.remove(&key);
521                remove_from_index_set(&mut self.by_receiver, &key.receiver, &key);
522            }
523        }
524        removed
525    }
526
527    pub(super) fn reset_receiver_for_replacement(
528        &mut self,
529        receiver: &UserId,
530    ) -> Vec<TransportRelayRouteEffect> {
531        let keys = self.by_receiver.get(receiver).cloned().unwrap_or_default();
532        let mut relays = Vec::new();
533        for key in keys {
534            let relay = {
535                let Some(entry) = self.entries.get_mut(&key) else {
536                    continue;
537                };
538                let Some(current) = entry.current.as_mut() else {
539                    continue;
540                };
541                current.selection =
542                    ConsumerSourceSelection::open(entry.intent.active().unwrap_or(true));
543                match mem::take(&mut current.realization) {
544                    ConsumerRealization::Absent => None,
545                    ConsumerRealization::Pending(_, relay) => relay,
546                    ConsumerRealization::Committed(committed) => committed.relay,
547                }
548            };
549            if let Some(relay) = relay {
550                relays.extend(self.release_relay(&key, &relay));
551            }
552        }
553        relays
554    }
555
556    pub(super) fn remove_receiver(&mut self, receiver: &UserId) -> RemovedRoutes {
557        self.pending_targets.remove(receiver);
558        let keys = self.by_receiver.remove(receiver).unwrap_or_default();
559        let mut removed = RemovedRoutes::default();
560        for key in keys {
561            let Some(entry) = self.entries.remove(&key) else {
562                continue;
563            };
564            let Some(current) = entry.current else {
565                continue;
566            };
567            remove_from_index_set(&mut self.by_source, &current.source_id, &key);
568            self.collect_removed_realization(&key, current.realization, &mut removed);
569        }
570        removed
571    }
572
573    pub(super) fn detach_declined_consumers(
574        &mut self,
575        session: &TransportSessionKey,
576        declined: &[TransportMediaId],
577    ) -> RemovedRoutes {
578        let keys = self
579            .by_receiver
580            .get(session.user_id())
581            .cloned()
582            .unwrap_or_default();
583        let mut removed = RemovedRoutes::default();
584        for key in keys {
585            let realization = {
586                let Some(current) = self
587                    .entries
588                    .get_mut(&key)
589                    .and_then(|entry| entry.current.as_mut())
590                else {
591                    continue;
592                };
593                let ConsumerRealization::Committed(committed) = &current.realization else {
594                    continue;
595                };
596                if committed.route.consumer_session_key() != session
597                    || !declined.contains(&committed.route.consumer_transport_media_id())
598                {
599                    continue;
600                }
601                mem::take(&mut current.realization)
602            };
603            self.collect_removed_realization(&key, realization, &mut removed);
604        }
605        removed
606    }
607
608    pub(super) fn reserve_relay(
609        &mut self,
610        reservation: &ConsumerRouteReservation,
611        target: &ConsumerSetupTarget,
612        target_worker: MediaWorkerId,
613        active: bool,
614    ) -> Vec<TransportRelayRouteEffect> {
615        let (previous, relay) = {
616            let Some(current) = self.current_mut_for(reservation) else {
617                return Vec::new();
618            };
619            let ConsumerRealization::Pending(id, relay) = &mut current.realization else {
620                return Vec::new();
621            };
622            if *id != reservation.id {
623                return Vec::new();
624            }
625            let next = RouteRelay {
626                route: target.relay_route_key(target_worker),
627                activity: RelayRouteActivity::from_active(active),
628            };
629            if relay.as_ref() == Some(&next) {
630                return Vec::new();
631            }
632            (relay.replace(next.clone()), next)
633        };
634        self.replace_relay(&reservation.key, previous, &relay)
635    }
636
637    pub(super) fn source_activity_target_workers<'a>(
638        &'a self,
639        source: &'a TransportSourceKey,
640    ) -> impl Iterator<Item = MediaWorkerId> + 'a {
641        self.relays
642            .keys()
643            .filter(move |route| &route.source == source)
644            .map(|route| route.target_worker)
645    }
646
647    fn entry(&mut self, key: SubscriptionKey) -> &mut Subscription {
648        self.by_receiver
649            .entry(key.receiver.clone())
650            .or_default()
651            .insert(key.clone());
652        self.entries.entry(key).or_default()
653    }
654
655    fn current_mut(
656        &mut self,
657        key: &SubscriptionKey,
658        source_id: PublishedSourceId,
659    ) -> Option<&mut CurrentPublication> {
660        let current = self.entries.get_mut(key)?.current.as_mut()?;
661        (current.source_id == source_id).then_some(current)
662    }
663
664    fn current_mut_for(
665        &mut self,
666        reservation: &ConsumerRouteReservation,
667    ) -> Option<&mut CurrentPublication> {
668        self.current_mut(&reservation.key, reservation.source_id)
669    }
670
671    fn take_pending(&mut self, reservation: &ConsumerRouteReservation) -> Option<TakenPending> {
672        let current = self.current_mut_for(reservation)?;
673        let pending = mem::take(&mut current.realization);
674        match pending {
675            ConsumerRealization::Pending(id, relay) if id == reservation.id => {
676                Some(TakenPending { relay })
677            }
678            other => {
679                current.realization = other;
680                None
681            }
682        }
683    }
684
685    fn collect_removed_realization(
686        &mut self,
687        key: &SubscriptionKey,
688        realization: ConsumerRealization,
689        removed: &mut RemovedRoutes,
690    ) {
691        let relay = match realization {
692            ConsumerRealization::Absent => None,
693            ConsumerRealization::Pending(_, relay) => relay,
694            ConsumerRealization::Committed(committed) => {
695                removed.routes.push(committed.route);
696                removed.consumers.push(committed.routed);
697                committed.relay
698            }
699        };
700        if let Some(relay) = relay {
701            removed.relays.extend(self.release_relay(key, &relay));
702        }
703    }
704
705    fn replace_relay(
706        &mut self,
707        key: &SubscriptionKey,
708        previous: Option<RouteRelay>,
709        relay: &RouteRelay,
710    ) -> Vec<TransportRelayRouteEffect> {
711        match previous {
712            None => self.set_relay_owner(key, relay, true),
713            Some(previous) if previous.route == relay.route => {
714                self.set_relay_owner(key, relay, false)
715            }
716            Some(previous) => {
717                let mut effects = self.release_relay(key, &previous);
718                effects.extend(self.set_relay_owner(key, relay, true));
719                effects
720            }
721        }
722    }
723
724    fn set_relay_owner(
725        &mut self,
726        key: &SubscriptionKey,
727        relay: &RouteRelay,
728        insert_missing: bool,
729    ) -> Vec<TransportRelayRouteEffect> {
730        let route = relay.route.clone();
731        let owners = match self.relays.entry(route.clone()) {
732            Entry::Occupied(entry) => entry.into_mut(),
733            Entry::Vacant(entry) if insert_missing => entry.insert(RelayOwners::default()),
734            Entry::Vacant(_) => return Vec::new(),
735        };
736        let before = relay_aggregate(owners);
737        owners.insert(key.clone(), relay.activity);
738        relay_effects_for(route, before, relay_aggregate(owners))
739    }
740
741    fn release_relay(
742        &mut self,
743        key: &SubscriptionKey,
744        relay: &RouteRelay,
745    ) -> Vec<TransportRelayRouteEffect> {
746        let route = relay.route.clone();
747        let Some(owners) = self.relays.get_mut(&route) else {
748            return Vec::new();
749        };
750        let before = relay_aggregate(owners);
751        if owners.remove(key).is_none() {
752            return Vec::new();
753        }
754        let after = relay_aggregate(owners);
755        if after.is_none() {
756            self.relays.remove(&route);
757        }
758        relay_effects_for(route, before, after)
759    }
760}
761
762impl ConsumerRouteReservation {
763    pub const fn declared_activity(&self) -> ConsumerActivity {
764        self.declared_activity
765    }
766}
767
768impl Subscription {
769    fn has_consumer_setup_or_route(&self) -> bool {
770        self.current
771            .as_ref()
772            .is_some_and(|current| !matches!(current.realization, ConsumerRealization::Absent))
773    }
774}
775
776impl CurrentPublication {
777    fn set_active(&mut self, active: bool) {
778        self.selection.set_active(active);
779        if !active {
780            self.clear_upgrade();
781        }
782    }
783
784    fn clear_upgrade(&mut self) {
785        if let ConsumerRealization::Committed(committed) = &mut self.realization {
786            committed.pending_upgrade = None;
787        }
788    }
789
790    pub(super) const fn is_pending(&self) -> bool {
791        matches!(self.realization, ConsumerRealization::Pending(..))
792    }
793
794    pub(super) fn committed(&self) -> Option<&CommittedConsumerRoute> {
795        match &self.realization {
796            ConsumerRealization::Committed(committed) => Some(committed),
797            ConsumerRealization::Absent | ConsumerRealization::Pending(..) => None,
798        }
799    }
800}
801
802impl ConsumerRealization {
803    fn set_relay_activity(&mut self, activity: RelayRouteActivity) -> Option<RouteRelay> {
804        let relay = match self {
805            Self::Absent => return None,
806            Self::Pending(_, relay) => relay.as_mut()?,
807            Self::Committed(committed) => committed.relay.as_mut()?,
808        };
809        if relay.activity == activity {
810            return None;
811        }
812        relay.activity = activity;
813        Some(relay.clone())
814    }
815}
816
817fn relay_aggregate(owners: &RelayOwners) -> Option<RelayRouteActivity> {
818    (!owners.is_empty()).then(|| {
819        RelayRouteActivity::from_active(owners.values().copied().any(RelayRouteActivity::is_active))
820    })
821}
822
823fn relay_effects_for(
824    route: RelayRouteKey,
825    before: Option<RelayRouteActivity>,
826    after: Option<RelayRouteActivity>,
827) -> Vec<TransportRelayRouteEffect> {
828    let Some(activity) = after else {
829        return vec![TransportRelayRouteEffect {
830            source: route.source,
831            target_media_worker_id: route.target_worker,
832            action: TransportRelayRouteAction::Release,
833        }];
834    };
835    let mut effects = Vec::new();
836    if before.is_none() {
837        effects.push(TransportRelayRouteEffect {
838            source: route.source.clone(),
839            target_media_worker_id: route.target_worker,
840            action: TransportRelayRouteAction::Install,
841        });
842    }
843    if before.unwrap_or(RelayRouteActivity::Inactive) != activity {
844        effects.push(TransportRelayRouteEffect {
845            source: route.source,
846            target_media_worker_id: route.target_worker,
847            action: TransportRelayRouteAction::SetActivity(activity),
848        });
849    }
850    effects
851}