Skip to main content

o_sfu_core/engine/room/media_graph/
mod.rs

1use std::collections::{BTreeMap, BTreeSet};
2
3use o_sfu_router::{ConsumerId, MediaKind, rtp, topology::RoutedProducerId};
4
5use crate::engine::{
6    UserId,
7    media_transport::{
8        SourceActivityRevision, TransportConsumerRoute, TransportMediaId, TransportSourceKey,
9    },
10    source_model::{ConsumerSourceSelection, PublishedSourceDescriptor, UserStreamId},
11};
12
13mod consumer_setup;
14mod producer;
15mod route_graph;
16mod source_index;
17mod subscription;
18mod topology;
19
20#[cfg(test)]
21#[expect(non_snake_case, reason = "test modules map to local TESTS directories")]
22mod TESTS;
23#[cfg(test)]
24#[path = "TESTS/route_graph.rs"]
25mod route_graph_tests;
26
27#[cfg(any(test, feature = "testing-transport"))]
28pub use self::subscription::ConsumerRouteState;
29pub(super) use self::{
30    consumer_setup::{
31        CommittedConsumerSetup, ConsumerSetupOrigin, ConsumerSetupOutcome, ConsumerSetupTarget,
32        DeclaredConsumerSetup, PendingConsumerSetup,
33    },
34    producer::{ProducerActivityCommit, PublishCommit, PublishIntentPlan, ValidatedPublish},
35    route_graph::PendingUpgrade,
36    subscription::{ReceiverRouteActivity, ReceiverRouteCommit, ReceiverRouteWork},
37    topology::{
38        CommittedTransportReceipt, RoomTopology, SessionPlacementCommit, SessionPlacementRejection,
39    },
40};
41
42#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
43pub(super) struct SubscriptionKey {
44    pub receiver: UserId,
45    pub publisher: UserId,
46    pub stream: UserStreamId,
47}
48
49impl SubscriptionKey {
50    pub fn new(receiver: &UserId, publisher: &UserId, stream: &UserStreamId) -> Self {
51        Self {
52            receiver: receiver.clone(),
53            publisher: publisher.clone(),
54            stream: stream.clone(),
55        }
56    }
57}
58
59#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
60pub(super) struct SourceKey {
61    owner_user_id: UserId,
62    stream_id: UserStreamId,
63}
64
65#[derive(Debug)]
66pub(super) struct PublishedSource {
67    pub descriptor: PublishedSourceDescriptor,
68    pub transport: TransportSourceKey,
69    pub rtp: rtp::MediaStream,
70    pub routed: RoutedProducerId,
71    pub active: bool,
72    pub activity_revision: SourceActivityRevision,
73}
74
75#[derive(Debug, Clone)]
76pub(super) struct ConsumerRouteView<'a> {
77    pub key: &'a SubscriptionKey,
78    pub route: &'a TransportConsumerRoute,
79    pub mid: &'a str,
80    /// Borrow holds so policy snapshots do not copy deadlines for every route.
81    pub pending_upgrade: Option<&'a PendingUpgrade>,
82    pub source: &'a PublishedSource,
83    pub selection: ConsumerSourceSelection,
84}
85
86impl ConsumerRouteView<'_> {
87    pub fn target(&self) -> ConsumerRouteTarget {
88        ConsumerRouteTarget::new(
89            self.route.clone(),
90            self.key.stream.clone(),
91            self.source.descriptor.media_kind(),
92        )
93    }
94}
95
96#[derive(Debug, Clone, Copy)]
97pub(super) struct PendingConsumerRouteView<'a> {
98    pub source: &'a PublishedSource,
99    pub selection: ConsumerSourceSelection,
100}
101
102#[derive(Debug, Clone, PartialEq, Eq)]
103pub struct ConsumerRouteTarget {
104    transport_route: TransportConsumerRoute,
105    stream_id: UserStreamId,
106    kind: MediaKind,
107}
108
109impl ConsumerRouteTarget {
110    fn new(
111        transport_route: TransportConsumerRoute,
112        stream_id: UserStreamId,
113        kind: MediaKind,
114    ) -> Self {
115        Self {
116            transport_route,
117            stream_id,
118            kind,
119        }
120    }
121
122    pub const fn transport_route(&self) -> &TransportConsumerRoute {
123        &self.transport_route
124    }
125
126    pub const fn consumer_media_id(&self) -> TransportMediaId {
127        self.transport_route.consumer_transport_media_id()
128    }
129
130    pub fn producer_user_id(&self) -> &UserId {
131        self.transport_route.source_session_key().user_id()
132    }
133
134    pub const fn source_media_id(&self) -> TransportMediaId {
135        self.transport_route.source_transport_media_id()
136    }
137
138    pub fn stream_id(&self) -> &UserStreamId {
139        &self.stream_id
140    }
141
142    pub fn request_keyframe_after_activity(&self, active: bool) -> bool {
143        active && self.kind == MediaKind::Video
144    }
145}
146
147impl SourceKey {
148    pub fn new(owner_user_id: &UserId, stream_id: &UserStreamId) -> Self {
149        Self {
150            owner_user_id: owner_user_id.clone(),
151            stream_id: stream_id.clone(),
152        }
153    }
154}
155
156fn remove_from_index_set<K, V>(index: &mut BTreeMap<K, BTreeSet<V>>, key: &K, value: &V)
157where
158    K: Ord,
159    V: Ord,
160{
161    let Some(values) = index.get_mut(key) else {
162        return;
163    };
164    values.remove(value);
165    if values.is_empty() {
166        index.remove(key);
167    }
168}