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//! Cross-worker relays are shared by subscriptions with the same source and
43//! target worker. A relay remains active while any owner has active receiver
44//! intent. An active first owner emits `Install` followed by
45//! `SetActivity(Active)`. Later aggregate activity changes emit `SetActivity`.
46//! Removing the last owner emits `Release`.
47
48use std::{
49    collections::{BTreeMap, BTreeSet, btree_map::Entry},
50    mem,
51};
52
53use o_sfu_router::topology::RoutedConsumerId;
54
55use super::{
56    ConsumerSourceSelection, SubscriptionKey, consumer_setup::ConsumerSetupTarget,
57    remove_from_index_set,
58};
59use crate::engine::{
60    ConnectionId, MediaWorkerId, UserId,
61    media_transport::{
62        RelayRouteActivity, TransportConsumerRoute, TransportMediaId, TransportRelayRouteAction,
63        TransportSessionKey, TransportSourceKey,
64    },
65    source_model::{PolicyPauseReason, PublishedSourceId, SourceSubscriptionIntent},
66};
67
68#[derive(Debug, Default)]
69pub(super) struct RouteGraph {
70    entries: BTreeMap<SubscriptionKey, Subscription>,
71    by_receiver: BTreeMap<UserId, BTreeSet<SubscriptionKey>>,
72    by_source: BTreeMap<PublishedSourceId, BTreeSet<SubscriptionKey>>,
73    relays: BTreeMap<RelayRouteKey, RelayOwners>,
74    next_reservation: RouteReservationId,
75}
76
77type RelayOwners = BTreeMap<SubscriptionKey, RelayRouteActivity>;
78
79#[derive(Debug, Default)]
80struct Subscription {
81    intent: SourceSubscriptionIntent,
82    current: Option<CurrentPublication>,
83}
84
85#[derive(Debug)]
86pub(super) struct CurrentPublication {
87    pub source_id: PublishedSourceId,
88    pub selection: ConsumerSourceSelection,
89    realization: ConsumerRealization,
90}
91
92#[derive(Debug, Default)]
93enum ConsumerRealization {
94    #[default]
95    Absent,
96    Pending(RouteReservationId, Option<RouteRelay>),
97    Committed(
98        TransportConsumerRoute,
99        String,
100        RoutedConsumerId,
101        Option<RouteRelay>,
102    ),
103}
104
105#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
106struct RouteReservationId(u64);
107
108#[derive(Debug)]
109pub struct ConsumerRouteReservation {
110    key: SubscriptionKey,
111    source_id: PublishedSourceId,
112    selection: ConsumerSourceSelection,
113    id: RouteReservationId,
114}
115
116#[derive(Debug, Clone, PartialEq, Eq)]
117struct RouteRelay {
118    route: RelayRouteKey,
119    activity: RelayRouteActivity,
120}
121
122struct TakenPending {
123    relay: Option<RouteRelay>,
124}
125
126#[derive(Debug, Clone, PartialEq, Eq)]
127pub struct RelayRouteEffect {
128    pub route: RelayRouteKey,
129    pub action: TransportRelayRouteAction,
130}
131
132#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
133pub struct RelayRouteKey {
134    pub source_user: UserId,
135    pub source_connection: ConnectionId,
136    pub source_media: TransportMediaId,
137    pub target_worker: MediaWorkerId,
138}
139
140#[derive(Debug, Default)]
141pub(super) struct RemovedRoutes {
142    pub routes: Vec<TransportConsumerRoute>,
143    pub consumers: Vec<RoutedConsumerId>,
144    pub relays: Vec<RelayRouteEffect>,
145}
146
147impl RemovedRoutes {
148    pub(super) fn extend(&mut self, mut other: Self) {
149        self.routes.append(&mut other.routes);
150        self.consumers.append(&mut other.consumers);
151        self.relays.append(&mut other.relays);
152    }
153}
154
155impl RouteGraph {
156    pub(super) fn subscription_count(&self) -> usize {
157        self.entries
158            .values()
159            .filter(|entry| entry.has_consumer_setup_or_route())
160            .count()
161    }
162
163    #[cfg(any(test, feature = "testing-transport"))]
164    pub(super) fn count(&self) -> usize {
165        self.attached()
166            .filter(|(_, current)| current.committed().is_some())
167            .count()
168    }
169
170    #[cfg(test)]
171    pub(super) fn record_count(&self) -> usize {
172        self.entries.len()
173    }
174
175    pub(super) fn merge_intent(&mut self, key: SubscriptionKey, update: SourceSubscriptionIntent) {
176        if update.is_empty() {
177            return;
178        }
179        let entry = self.entry(key);
180        entry.intent.merge(update);
181        if let (Some(active), Some(current)) = (update.active(), entry.current.as_mut()) {
182            current.selection.set_active(active);
183        }
184    }
185
186    pub(super) fn intent(&self, key: &SubscriptionKey) -> SourceSubscriptionIntent {
187        self.entries
188            .get(key)
189            .map_or_else(SourceSubscriptionIntent::default, |entry| entry.intent)
190    }
191
192    pub(super) fn attach_for_setup(
193        &mut self,
194        key: SubscriptionKey,
195        source_id: PublishedSourceId,
196    ) -> bool {
197        let entry = self.entry(key.clone());
198        if let Some(current) = &entry.current {
199            return current.source_id == source_id
200                && matches!(current.realization, ConsumerRealization::Absent);
201        }
202        entry.current = Some(CurrentPublication {
203            source_id,
204            selection: ConsumerSourceSelection::open(entry.intent.active().unwrap_or(true)),
205            realization: ConsumerRealization::Absent,
206        });
207        self.by_source.entry(source_id).or_default().insert(key);
208        true
209    }
210
211    pub(super) fn set_activity(
212        &mut self,
213        key: &SubscriptionKey,
214        source_id: PublishedSourceId,
215        connection_id: ConnectionId,
216        active: bool,
217        policy_pause_reason: Option<PolicyPauseReason>,
218    ) -> Option<Vec<RelayRouteEffect>> {
219        let relay = {
220            let current = self
221                .entries
222                .get_mut(key)
223                .and_then(|entry| entry.current.as_mut())
224                .filter(|current| current.source_id == source_id)?;
225            if let ConsumerRealization::Committed(route, ..) = &current.realization
226                && route.consumer_session_key().connection_id() != connection_id
227            {
228                return None;
229            }
230            current.selection.set_active(active);
231            if let Some(reason) = policy_pause_reason {
232                current.selection.set_policy_pause_reason(Some(reason));
233            }
234            current
235                .realization
236                .set_relay_activity(RelayRouteActivity::from_active(active))
237        };
238        Some(relay.map_or_else(Vec::new, |relay| self.set_relay_owner(key, &relay, false)))
239    }
240
241    pub(super) fn reserve_consumer_setup(
242        &mut self,
243        key: SubscriptionKey,
244        source_id: PublishedSourceId,
245        selection: ConsumerSourceSelection,
246    ) -> Option<ConsumerRouteReservation> {
247        let current = self.entries.get(&key)?.current.as_ref()?;
248        if current.source_id != source_id
249            || !matches!(current.realization, ConsumerRealization::Absent)
250        {
251            return None;
252        }
253        let id = self.next_reservation();
254        let current = self.entries.get_mut(&key)?.current.as_mut()?;
255        current.selection = selection;
256        current.realization = ConsumerRealization::Pending(id, None);
257        Some(ConsumerRouteReservation {
258            key,
259            source_id,
260            selection,
261            id,
262        })
263    }
264
265    pub(super) fn release_consumer_setup(
266        &mut self,
267        reservation: ConsumerRouteReservation,
268    ) -> Vec<RelayRouteEffect> {
269        let Some(TakenPending { relay }) = self.take_pending(&reservation) else {
270            return Vec::new();
271        };
272        let key = reservation.key;
273        relay.map_or_else(Vec::new, |relay| self.release_relay(&key, &relay))
274    }
275
276    pub(super) fn commit(
277        &mut self,
278        reservation: ConsumerRouteReservation,
279        route: TransportConsumerRoute,
280        mid: String,
281        selection: ConsumerSourceSelection,
282        accept: impl FnOnce() -> Option<RoutedConsumerId>,
283    ) -> Result<(), Vec<RelayRouteEffect>> {
284        // Consume the reservation before router acceptance so every async
285        // completion is terminal and cannot reuse its identity after failure.
286        let pending = self.take_pending(&reservation);
287        let ConsumerRouteReservation { key, source_id, .. } = reservation;
288        let Some(TakenPending { relay }) = pending else {
289            return Err(Vec::new());
290        };
291        let Some(current) = self.current_mut(&key, source_id) else {
292            return Err(relay.map_or_else(Vec::new, |relay| self.release_relay(&key, &relay)));
293        };
294        let Some(routed) = accept() else {
295            return Err(relay.map_or_else(Vec::new, |relay| self.release_relay(&key, &relay)));
296        };
297        current.selection = selection;
298        current.realization = ConsumerRealization::Committed(route, mid, routed, relay);
299        Ok(())
300    }
301
302    pub(super) fn update_selection(
303        &mut self,
304        key: &SubscriptionKey,
305        source_id: PublishedSourceId,
306        route: &TransportConsumerRoute,
307        update: impl FnOnce(&mut ConsumerSourceSelection),
308    ) -> bool {
309        let Some(current) = self
310            .entries
311            .get_mut(key)
312            .and_then(|entry| entry.current.as_mut())
313        else {
314            return false;
315        };
316        let ConsumerRealization::Committed(current_route, ..) = &current.realization else {
317            return false;
318        };
319        if current.source_id != source_id || current_route != route {
320            return false;
321        }
322        update(&mut current.selection);
323        true
324    }
325
326    pub(super) fn selection(
327        &self,
328        key: &SubscriptionKey,
329        source_id: PublishedSourceId,
330    ) -> Option<ConsumerSourceSelection> {
331        let current = self.entries.get(key)?.current.as_ref()?;
332        (current.source_id == source_id).then_some(current.selection)
333    }
334
335    pub(super) fn attached(&self) -> impl Iterator<Item = (&SubscriptionKey, &CurrentPublication)> {
336        self.entries
337            .iter()
338            .filter_map(|(key, entry)| Some((key, entry.current.as_ref()?)))
339    }
340
341    pub(super) fn current(
342        &self,
343        key: &SubscriptionKey,
344    ) -> Option<(&SubscriptionKey, &CurrentPublication)> {
345        let (key, entry) = self.entries.get_key_value(key)?;
346        Some((key, entry.current.as_ref()?))
347    }
348
349    pub(super) fn attached_for_receiver(
350        &self,
351        receiver: &UserId,
352    ) -> impl Iterator<Item = (&SubscriptionKey, &CurrentPublication)> {
353        self.by_receiver
354            .get(receiver)
355            .into_iter()
356            .flat_map(BTreeSet::iter)
357            .filter_map(|key| self.current(key))
358    }
359
360    pub(super) fn attached_for_source(
361        &self,
362        source_id: PublishedSourceId,
363    ) -> impl Iterator<Item = (&SubscriptionKey, &CurrentPublication)> {
364        self.by_source
365            .get(&source_id)
366            .into_iter()
367            .flat_map(BTreeSet::iter)
368            .filter_map(|key| self.current(key))
369    }
370
371    pub(super) fn detach_source(&mut self, source_id: PublishedSourceId) -> RemovedRoutes {
372        let keys = self.by_source.remove(&source_id).unwrap_or_default();
373        let mut removed = RemovedRoutes::default();
374        for key in keys {
375            let (realization, prune) = {
376                let Some(entry) = self.entries.get_mut(&key) else {
377                    continue;
378                };
379                let Some(current) = entry.current.take() else {
380                    continue;
381                };
382                debug_assert_eq!(current.source_id, source_id);
383                (current.realization, entry.intent.is_empty())
384            };
385            self.collect_removed_realization(&key, realization, &mut removed);
386            if prune {
387                self.entries.remove(&key);
388                remove_from_index_set(&mut self.by_receiver, &key.receiver, &key);
389            }
390        }
391        removed
392    }
393
394    pub(super) fn reset_receiver_for_replacement(
395        &mut self,
396        receiver: &UserId,
397    ) -> Vec<RelayRouteEffect> {
398        let keys = self.by_receiver.get(receiver).cloned().unwrap_or_default();
399        let mut relays = Vec::new();
400        for key in keys {
401            let relay = {
402                let Some(entry) = self.entries.get_mut(&key) else {
403                    continue;
404                };
405                let Some(current) = entry.current.as_mut() else {
406                    continue;
407                };
408                current.selection =
409                    ConsumerSourceSelection::open(entry.intent.active().unwrap_or(true));
410                match mem::take(&mut current.realization) {
411                    ConsumerRealization::Absent => None,
412                    ConsumerRealization::Pending(_, relay)
413                    | ConsumerRealization::Committed(_, _, _, relay) => relay,
414                }
415            };
416            if let Some(relay) = relay {
417                relays.extend(self.release_relay(&key, &relay));
418            }
419        }
420        relays
421    }
422
423    pub(super) fn remove_receiver(&mut self, receiver: &UserId) -> RemovedRoutes {
424        let keys = self.by_receiver.remove(receiver).unwrap_or_default();
425        let mut removed = RemovedRoutes::default();
426        for key in keys {
427            let Some(entry) = self.entries.remove(&key) else {
428                continue;
429            };
430            let Some(current) = entry.current else {
431                continue;
432            };
433            remove_from_index_set(&mut self.by_source, &current.source_id, &key);
434            self.collect_removed_realization(&key, current.realization, &mut removed);
435        }
436        removed
437    }
438
439    pub(super) fn detach_declined_consumers(
440        &mut self,
441        session: &TransportSessionKey,
442        declined: &[TransportMediaId],
443    ) -> RemovedRoutes {
444        let keys = self
445            .by_receiver
446            .get(session.user_id())
447            .cloned()
448            .unwrap_or_default();
449        let mut removed = RemovedRoutes::default();
450        for key in keys {
451            let realization = {
452                let Some(current) = self
453                    .entries
454                    .get_mut(&key)
455                    .and_then(|entry| entry.current.as_mut())
456                else {
457                    continue;
458                };
459                let ConsumerRealization::Committed(route, ..) = &current.realization else {
460                    continue;
461                };
462                if route.consumer_session_key() != session
463                    || !declined.contains(&route.consumer_transport_media_id())
464                {
465                    continue;
466                }
467                mem::take(&mut current.realization)
468            };
469            self.collect_removed_realization(&key, realization, &mut removed);
470        }
471        removed
472    }
473
474    pub(super) fn reserve_relay(
475        &mut self,
476        reservation: &ConsumerRouteReservation,
477        target: &ConsumerSetupTarget,
478        target_worker: MediaWorkerId,
479        active: bool,
480    ) -> Vec<RelayRouteEffect> {
481        let (previous, relay) = {
482            let Some(current) = self.current_mut_for(reservation) else {
483                return Vec::new();
484            };
485            let ConsumerRealization::Pending(id, relay) = &mut current.realization else {
486                return Vec::new();
487            };
488            if *id != reservation.id {
489                return Vec::new();
490            }
491            let next = RouteRelay {
492                route: target.relay_route_key(target_worker),
493                activity: RelayRouteActivity::from_active(active),
494            };
495            if relay.as_ref() == Some(&next) {
496                return Vec::new();
497            }
498            (relay.replace(next.clone()), next)
499        };
500        self.replace_relay(&reservation.key, previous, &relay)
501    }
502
503    pub(super) fn source_activity_target_workers<'a>(
504        &'a self,
505        source: &'a TransportSourceKey,
506    ) -> impl Iterator<Item = MediaWorkerId> + 'a {
507        let session = source.session_key();
508        self.relays
509            .keys()
510            .filter(move |route| {
511                route.source_user == *session.user_id()
512                    && route.source_connection == session.connection_id()
513                    && route.source_media == source.transport_media_id()
514            })
515            .map(|route| route.target_worker)
516    }
517
518    fn entry(&mut self, key: SubscriptionKey) -> &mut Subscription {
519        self.by_receiver
520            .entry(key.receiver.clone())
521            .or_default()
522            .insert(key.clone());
523        self.entries.entry(key).or_default()
524    }
525
526    fn next_reservation(&mut self) -> RouteReservationId {
527        self.next_reservation.0 += 1;
528        self.next_reservation
529    }
530
531    fn current_mut(
532        &mut self,
533        key: &SubscriptionKey,
534        source_id: PublishedSourceId,
535    ) -> Option<&mut CurrentPublication> {
536        let current = self.entries.get_mut(key)?.current.as_mut()?;
537        (current.source_id == source_id).then_some(current)
538    }
539
540    fn current_mut_for(
541        &mut self,
542        reservation: &ConsumerRouteReservation,
543    ) -> Option<&mut CurrentPublication> {
544        self.current_mut(&reservation.key, reservation.source_id)
545    }
546
547    fn take_pending(&mut self, reservation: &ConsumerRouteReservation) -> Option<TakenPending> {
548        let current = self.current_mut_for(reservation)?;
549        let pending = mem::take(&mut current.realization);
550        match pending {
551            ConsumerRealization::Pending(id, relay) if id == reservation.id => {
552                Some(TakenPending { relay })
553            }
554            other => {
555                current.realization = other;
556                None
557            }
558        }
559    }
560
561    fn collect_removed_realization(
562        &mut self,
563        key: &SubscriptionKey,
564        realization: ConsumerRealization,
565        removed: &mut RemovedRoutes,
566    ) {
567        let relay = match realization {
568            ConsumerRealization::Absent => None,
569            ConsumerRealization::Pending(_, relay) => relay,
570            ConsumerRealization::Committed(route, _, consumer, relay) => {
571                removed.routes.push(route);
572                removed.consumers.push(consumer);
573                relay
574            }
575        };
576        if let Some(relay) = relay {
577            removed.relays.extend(self.release_relay(key, &relay));
578        }
579    }
580
581    fn replace_relay(
582        &mut self,
583        key: &SubscriptionKey,
584        previous: Option<RouteRelay>,
585        relay: &RouteRelay,
586    ) -> Vec<RelayRouteEffect> {
587        match previous {
588            None => self.set_relay_owner(key, relay, true),
589            Some(previous) if previous.route == relay.route => {
590                self.set_relay_owner(key, relay, false)
591            }
592            Some(previous) => {
593                let mut effects = self.release_relay(key, &previous);
594                effects.extend(self.set_relay_owner(key, relay, true));
595                effects
596            }
597        }
598    }
599
600    fn set_relay_owner(
601        &mut self,
602        key: &SubscriptionKey,
603        relay: &RouteRelay,
604        insert_missing: bool,
605    ) -> Vec<RelayRouteEffect> {
606        let route = relay.route.clone();
607        let owners = match self.relays.entry(route.clone()) {
608            Entry::Occupied(entry) => entry.into_mut(),
609            Entry::Vacant(entry) if insert_missing => entry.insert(RelayOwners::default()),
610            Entry::Vacant(_) => return Vec::new(),
611        };
612        let before = relay_aggregate(owners);
613        owners.insert(key.clone(), relay.activity);
614        relay_effects_for(route, before, relay_aggregate(owners))
615    }
616
617    fn release_relay(
618        &mut self,
619        key: &SubscriptionKey,
620        relay: &RouteRelay,
621    ) -> Vec<RelayRouteEffect> {
622        let route = relay.route.clone();
623        let Some(owners) = self.relays.get_mut(&route) else {
624            return Vec::new();
625        };
626        let before = relay_aggregate(owners);
627        if owners.remove(key).is_none() {
628            return Vec::new();
629        }
630        let after = relay_aggregate(owners);
631        if after.is_none() {
632            self.relays.remove(&route);
633        }
634        relay_effects_for(route, before, after)
635    }
636}
637
638impl ConsumerRouteReservation {
639    pub const fn selection(&self) -> ConsumerSourceSelection {
640        self.selection
641    }
642}
643
644impl Subscription {
645    fn has_consumer_setup_or_route(&self) -> bool {
646        self.current
647            .as_ref()
648            .is_some_and(|current| !matches!(current.realization, ConsumerRealization::Absent))
649    }
650}
651
652impl CurrentPublication {
653    pub(super) const fn is_pending(&self) -> bool {
654        matches!(self.realization, ConsumerRealization::Pending(..))
655    }
656
657    pub(super) fn committed(&self) -> Option<(&TransportConsumerRoute, &str)> {
658        match &self.realization {
659            ConsumerRealization::Committed(route, mid, ..) => Some((route, mid)),
660            ConsumerRealization::Absent | ConsumerRealization::Pending(..) => None,
661        }
662    }
663}
664
665impl ConsumerRealization {
666    fn set_relay_activity(&mut self, activity: RelayRouteActivity) -> Option<RouteRelay> {
667        let relay = match self {
668            Self::Absent => return None,
669            Self::Pending(_, relay) | Self::Committed(_, _, _, relay) => relay.as_mut()?,
670        };
671        if relay.activity == activity {
672            return None;
673        }
674        relay.activity = activity;
675        Some(relay.clone())
676    }
677}
678
679fn relay_aggregate(owners: &RelayOwners) -> Option<RelayRouteActivity> {
680    (!owners.is_empty()).then(|| {
681        RelayRouteActivity::from_active(owners.values().copied().any(RelayRouteActivity::is_active))
682    })
683}
684
685fn relay_effects_for(
686    route: RelayRouteKey,
687    before: Option<RelayRouteActivity>,
688    after: Option<RelayRouteActivity>,
689) -> Vec<RelayRouteEffect> {
690    let Some(activity) = after else {
691        return vec![RelayRouteEffect {
692            route,
693            action: TransportRelayRouteAction::Release,
694        }];
695    };
696    let mut effects = Vec::new();
697    if before.is_none() {
698        effects.push(RelayRouteEffect {
699            route: route.clone(),
700            action: TransportRelayRouteAction::Install,
701        });
702    }
703    if before.unwrap_or(RelayRouteActivity::Inactive) != activity {
704        effects.push(RelayRouteEffect {
705            route,
706            action: TransportRelayRouteAction::SetActivity(activity),
707        });
708    }
709    effects
710}