o_sfu_core/engine/room/media_graph/
mod.rs1use 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 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}