o_sfu_core/engine/room/state/
shared.rs1use 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 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 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 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 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}