Skip to main content

o_sfu_core/engine/room/state/
shared.rs

1use std::{
2    collections::{BTreeMap, BTreeSet},
3    sync::Arc,
4    time::Instant,
5};
6
7use o_sfu_router::rtp::{MediaCapabilities, MediaCapabilities as RouterRtpCapabilities};
8
9use super::super::{
10    RoomAdmissionPolicy, RoomMediaCounts,
11    media_graph::{ConsumerRouteView, RoomTopology},
12    outbound::{
13        OutboundSender, RemoteTrackProjection, RemoteTrackSnapshot, VersionedRemoteTrackSnapshot,
14    },
15    transition::StagedPublishes,
16};
17use crate::{
18    RoomMediaLimits, VideoAdaptationTuning,
19    engine::{
20        ConnectionId, MediaWorkerId, PeerSnapshot, RecordingState, UserId, UserInfo,
21        media_transport::TransportSessionKey, room::placement::PlacementSnapshot,
22        source_model::UserStreamId,
23    },
24};
25
26#[derive(Debug)]
27pub struct RoomState {
28    pub(super) admission_policy: RoomAdmissionPolicy,
29    pub media_limits: RoomMediaLimits,
30    pub video_adaptation_tuning: VideoAdaptationTuning,
31    pub users: BTreeMap<UserId, ActiveUser>,
32    /// rejects stale async callbacks from previous connections
33    pub(super) next_connection_id: u64,
34    pub next_consumer_id: u64,
35    next_track_snapshot_revision: u64,
36    pub(super) recording_state: RecordingState,
37    pub(in crate::engine::room) staged_publishes: StagedPublishes,
38    pub(in crate::engine::room) topology: RoomTopology,
39}
40
41#[derive(Debug)]
42pub struct ActiveUser {
43    pub(super) user_id: Arc<UserId>,
44    pub(super) info: UserInfo,
45    pub(super) server_featured: Option<bool>,
46    pub parsed_client_rtp_capabilities: Option<RouterRtpCapabilities>,
47    pub connection_id: ConnectionId,
48    /// Continuous receiver pressure survives victim changes within this connection.
49    pub(in crate::engine::room) video_soft_pause_deadline: Option<Instant>,
50    pub sender: OutboundSender,
51}
52
53impl ActiveUser {
54    pub(super) fn reset_presentation(&mut self) {
55        self.info = UserInfo::default();
56        self.server_featured = None;
57    }
58
59    pub(super) fn apply_info_update(&mut self, info: &UserInfo) {
60        self.info.apply_partial_update(info);
61    }
62
63    pub(in crate::engine::room) const fn featured(&self) -> Option<bool> {
64        self.server_featured
65    }
66
67    pub(in crate::engine::room) const fn is_deaf(&self) -> bool {
68        matches!(self.info.is_deaf, Some(true))
69    }
70
71    pub(in crate::engine::room) const fn is_screensharing(&self) -> bool {
72        matches!(self.info.is_screen_sharing_on, Some(true))
73    }
74
75    pub(in crate::engine::room) fn set_featured(&mut self, featured: Option<bool>) {
76        self.server_featured = featured;
77    }
78
79    pub(super) fn project_info(&self) -> UserInfo {
80        self.info
81            .clone()
82            .with_featured(self.server_featured)
83            .snapshot_complete()
84    }
85}
86
87impl RoomState {
88    pub fn new(
89        runtime_context: &super::super::RoomRuntimeContext,
90        admission_policy: RoomAdmissionPolicy,
91        media_limits: RoomMediaLimits,
92        video_adaptation_tuning: VideoAdaptationTuning,
93        router_rtp_capabilities: MediaCapabilities,
94    ) -> Self {
95        Self {
96            admission_policy,
97            media_limits,
98            video_adaptation_tuning,
99            users: BTreeMap::new(),
100            next_connection_id: 0,
101            next_consumer_id: 1,
102            next_track_snapshot_revision: 1,
103            recording_state: RecordingState {
104                recording: Some(false),
105                audio: Some(false),
106                transcription: Some(false),
107                video: Some(false),
108            },
109            staged_publishes: StagedPublishes::default(),
110            topology: RoomTopology::new(
111                runtime_context,
112                router_rtp_capabilities,
113                admission_policy.max_sessions,
114            ),
115        }
116    }
117
118    pub fn user_for_connection(
119        &self,
120        user_id: &UserId,
121        connection_id: ConnectionId,
122    ) -> Option<&ActiveUser> {
123        let user = self.users.get(user_id)?;
124        if user.connection_id != connection_id {
125            return None;
126        }
127        Some(user)
128    }
129
130    pub fn user_mut_for_connection(
131        &mut self,
132        user_id: &UserId,
133        connection_id: ConnectionId,
134    ) -> Option<&mut ActiveUser> {
135        let user = self.users.get_mut(user_id)?;
136        if user.connection_id != connection_id {
137            return None;
138        }
139        Some(user)
140    }
141
142    pub fn recording_state(&self) -> RecordingState {
143        self.recording_state.clone()
144    }
145
146    #[cfg(any(test, feature = "testing-transport"))]
147    pub fn router_rtp_capabilities(&self) -> MediaCapabilities {
148        self.topology.router().rtp_capabilities().clone()
149    }
150
151    pub fn transport_user_entries(&self) -> impl Iterator<Item = (&UserId, ConnectionId)> {
152        self.users
153            .iter()
154            .map(|(user_id, user)| (user_id, user.connection_id))
155    }
156
157    /// Returns the transport key for an exact committed router placement.
158    ///
159    /// Use [`Self::committed_transport_user_key`] when placement may be stale.
160    ///
161    /// # Panics
162    ///
163    /// Panics when `user_id` and `connection_id` have no committed router
164    /// placement.
165    pub fn transport_user_key(
166        &self,
167        user_id: &UserId,
168        connection_id: ConnectionId,
169    ) -> TransportSessionKey {
170        if let Some(user) = self.user_for_connection(user_id, connection_id) {
171            return self
172                .topology
173                .transport_user_key(Arc::clone(&user.user_id), connection_id);
174        }
175        self.topology
176            .transport_user_key(user_id.clone(), connection_id)
177    }
178
179    /// Returns `None` unless the exact router placement remains committed.
180    pub fn committed_transport_user_key(
181        &self,
182        user_id: &UserId,
183        connection_id: ConnectionId,
184    ) -> Option<TransportSessionKey> {
185        if let Some(user) = self.user_for_connection(user_id, connection_id) {
186            return self
187                .topology
188                .committed_transport_user_key(Arc::clone(&user.user_id), connection_id);
189        }
190        self.topology
191            .committed_transport_user_key(user_id.clone(), connection_id)
192    }
193
194    pub fn placement_usage_snapshot(&self) -> PlacementSnapshot {
195        self.topology.router().placement_snapshot()
196    }
197
198    pub fn assigned_primary_media_worker_id(&self) -> Option<MediaWorkerId> {
199        self.topology.router().primary_worker()
200    }
201
202    pub(in crate::engine::room) fn committed_consumer_routes(
203        &self,
204    ) -> impl Iterator<Item = ConsumerRouteView<'_>> {
205        self.topology.committed_consumer_routes().filter(|route| {
206            self.user_connection_id(&route.key.receiver)
207                == Some(route.route.consumer_session_key().connection_id())
208        })
209    }
210
211    pub fn user_connection_id(&self, user_id: &UserId) -> Option<ConnectionId> {
212        self.users.get(user_id).map(|user| user.connection_id)
213    }
214
215    pub fn user_count(&self) -> usize {
216        self.users.len()
217    }
218
219    pub fn user_snapshots_except(&self, excluded_user_id: &UserId) -> Vec<PeerSnapshot> {
220        self.users
221            .iter()
222            .filter(|(user_id, _session)| *user_id != excluded_user_id)
223            .map(|(user_id, user)| PeerSnapshot {
224                user_id: user_id.clone(),
225                info: user.project_info(),
226            })
227            .collect()
228    }
229
230    pub fn user_info_snapshot(&self, user_id: &UserId) -> Option<(UserId, UserInfo)> {
231        let user = self.users.get(user_id)?;
232        Some((user_id.clone(), user.project_info()))
233    }
234
235    pub(in crate::engine::room) fn remote_track_snapshot_for_user(
236        &mut self,
237        user_id: &UserId,
238        requires_negotiation: bool,
239    ) -> VersionedRemoteTrackSnapshot {
240        let revision = self.next_track_snapshot_revision;
241        self.next_track_snapshot_revision = revision.saturating_add(1);
242        let connection_id = self.user_connection_id(user_id);
243        let snapshot = RemoteTrackSnapshot {
244            tracks: self
245                .topology
246                .committed_consumer_routes_for_user(user_id)
247                .filter(|route| {
248                    connection_id == Some(route.route.consumer_session_key().connection_id())
249                })
250                .map(|route| {
251                    let source = &route.source.descriptor;
252                    RemoteTrackProjection {
253                        consumer_mid: route.mid.to_owned(),
254                        user_id: source.owner().user_id().clone(),
255                        stream_id: source.stream_id().clone(),
256                        producer_active: route.source.active,
257                    }
258                })
259                .collect(),
260            requires_negotiation,
261        };
262        VersionedRemoteTrackSnapshot { snapshot, revision }
263    }
264
265    pub(in crate::engine::room) fn remote_track_snapshots_for_users(
266        &mut self,
267        user_ids: BTreeSet<UserId>,
268        requires_negotiation: bool,
269    ) -> Vec<(OutboundSender, VersionedRemoteTrackSnapshot)> {
270        user_ids
271            .into_iter()
272            .filter_map(|user_id| {
273                Some((
274                    self.users.get(&user_id)?.sender.clone(),
275                    self.remote_track_snapshot_for_user(&user_id, requires_negotiation),
276                ))
277            })
278            .collect()
279    }
280
281    pub fn user_stats_counts(&self) -> (u64, BTreeMap<UserStreamId, u64>) {
282        (
283            u64::try_from(self.users.len()).unwrap_or(u64::MAX),
284            self.topology.active_stream_user_counts(),
285        )
286    }
287
288    pub fn media_counts(&self) -> RoomMediaCounts {
289        self.topology.media_counts()
290    }
291
292    pub fn is_empty(&self) -> bool {
293        self.users.is_empty()
294    }
295}