1#[cfg(any(test, feature = "testing-transport"))]
9use o_sfu_router::rtp::MediaCapabilities;
10use o_sfu_telemetry::schema::event as telemetry_event;
11use tracing::{info, warn};
12
13use super::{
14 BroadcastPayloadError, Room, RoomJoinError, UserOutboundSender,
15 effects::batch::{RoomEffectContext, RoomEffects},
16 media_graph::CommittedTransportReceipt,
17 placement::JoinAdmissionTurn,
18 state::ConnectionCloseCommit,
19};
20use crate::engine::{
21 ConnectionId, MediaWorkerId, UserId, UserInfo, UserPermissions,
22 media_transport::MediaTransport, room::state::JoinCommit,
23};
24
25#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
27pub struct RoomUserPermissions;
28
29impl From<UserPermissions> for RoomUserPermissions {
30 fn from(_value: UserPermissions) -> Self {
31 Self
32 }
33}
34
35#[derive(Debug, Clone, Copy, PartialEq, Eq)]
36pub enum UserCloseReason {
37 Replaced,
38 RemovedByRuntime,
39}
40
41pub struct JoinUserRequest {
43 pub user_id: UserId,
44 pub label: Option<String>,
46 pub permissions: UserPermissions,
48 pub sender: UserOutboundSender,
49}
50
51impl Room {
52 pub(super) async fn commit_admission(
60 &self,
61 admission: JoinAdmissionTurn<'_, impl FnOnce() -> o_sfu_router::RouterId>,
62 context: RoomEffectContext<'_>,
63 ) -> Result<JoinCommit, RoomJoinError> {
64 let joined_fanout = context.user_joined_fanout();
65 admission.commit(self, joined_fanout).await
66 }
67
68 pub(super) async fn finalize_admission(
70 &self,
71 commit: JoinCommit,
72 context: RoomEffectContext<'_>,
73 ) -> CommittedTransportReceipt {
74 let JoinCommit {
75 receipt,
76 effects,
77 transport_plan,
78 } = commit;
79 RoomEffects::from_join(effects, transport_plan)
80 .execute(self, context)
81 .await;
82 let session = &receipt.transport_session_key;
83 info!(
84 event = telemetry_event::USER_JOINED,
85 room_id = self.uuid(),
86 user_id = %session.user_id().path_segment(),
87 connection_id = session.connection_id().as_u64(),
88 media_worker_id = session.media_worker_id().as_usize(),
89 "user joined room"
90 );
91 receipt
92 }
93
94 pub(crate) async fn remove_user(
96 &self,
97 user_id: &UserId,
98 connection_id: ConnectionId,
99 media_transport: &MediaTransport,
100 ) -> bool {
101 self.remove_user_with_teardown(
102 user_id,
103 connection_id,
104 RoomEffectContext::runtime(media_transport),
105 )
106 .await
107 }
108
109 pub async fn remove_user_with_teardown(
112 &self,
113 user_id: &UserId,
114 connection_id: ConnectionId,
115 context: RoomEffectContext<'_>,
116 ) -> bool {
117 let Some(commit) = self
118 .state
119 .write()
120 .await
121 .close_connection(user_id, connection_id)
122 else {
123 return false;
124 };
125 if let ConnectionCloseCommit::Current { evictions, .. } = &commit {
126 self.metrics
127 .record_subscription_intent_evictions(*evictions);
128 }
129 let (removed_current_user, media_worker_id) = match &commit {
130 ConnectionCloseCommit::Current {
131 session_teardown, ..
132 } => (
133 true,
134 session_teardown
135 .as_ref()
136 .map(|teardown| teardown.session_key().media_worker_id()),
137 ),
138 ConnectionCloseCommit::StalePlacement { .. } => (false, None),
139 };
140 RoomEffects::from_connection_close(commit)
141 .execute(self, context)
142 .await;
143 if removed_current_user {
144 info!(
145 event = telemetry_event::USER_CLOSED,
146 room_id = self.uuid(),
147 user_id = %user_id.path_segment(),
148 connection_id = connection_id.as_u64(),
149 media_worker_id = media_worker_id.map(MediaWorkerId::as_usize),
150 "user closed"
151 );
152 }
153 removed_current_user
154 }
155
156 pub(crate) async fn broadcast(
164 &self,
165 sender_id: &UserId,
166 connection_id: ConnectionId,
167 message: serde_json::Value,
168 ) -> Result<(), BroadcastPayloadError> {
169 let fanout = {
170 let state = self.state.read().await;
171 state.broadcast_fanout(sender_id, connection_id, message)
172 }?;
173 if let Some(fanout) = fanout {
174 fanout.emit();
175 }
176 Ok(())
177 }
178
179 pub(crate) async fn has_connection(
180 &self,
181 user_id: &UserId,
182 connection_id: ConnectionId,
183 ) -> bool {
184 self.state.read().await.user_connection_id(user_id) == Some(connection_id)
185 }
186
187 pub(crate) async fn update_user_info(
190 &self,
191 user_id: &UserId,
192 connection_id: ConnectionId,
193 media_transport: &MediaTransport,
194 mut info: UserInfo,
195 ) {
196 info.is_camera_on = None;
197 info.is_screen_sharing_on = None;
198 let commit = {
199 let mut state = self.state.write().await;
200 state.apply_presence_update(user_id, connection_id, &info)
201 };
202 if let Some(commit) = commit {
203 RoomEffects::from_presence(commit)
204 .execute(self, RoomEffectContext::runtime(media_transport))
205 .await;
206 } else {
207 warn!(
208 ?user_id,
209 connection_id = ?connection_id,
210 ?info,
211 "user info update was rejected by room state"
212 );
213 }
214 }
215
216 pub(crate) async fn disconnect_users(
222 &self,
223 user_ids: &[UserId],
224 media_transport: &MediaTransport,
225 ) -> usize {
226 self.disconnect_users_with_teardown(user_ids, RoomEffectContext::runtime(media_transport))
227 .await
228 }
229
230 pub async fn disconnect_users_with_teardown(
238 &self,
239 user_ids: &[UserId],
240 context: RoomEffectContext<'_>,
241 ) -> usize {
242 let commit = self.state.write().await.apply_disconnect_users(user_ids);
243 self.metrics
244 .record_subscription_intent_evictions(commit.evictions);
245 let sessions = commit
246 .session_teardowns
247 .iter()
248 .map(|teardown| teardown.session_key().clone())
249 .collect::<Vec<_>>();
250 RoomEffects::from_disconnect(commit)
251 .execute(self, context)
252 .await;
253 for session in &sessions {
254 info!(
255 event = telemetry_event::USER_DISCONNECTED,
256 room_id = self.uuid(),
257 user_id = %session.user_id().path_segment(),
258 connection_id = session.connection_id().as_u64(),
259 media_worker_id = session.media_worker_id().as_usize(),
260 "user disconnected"
261 );
262 }
263 sessions.len()
264 }
265
266 #[cfg(any(test, feature = "testing-transport"))]
267 pub async fn apply_session_negotiated(
270 &self,
271 user_id: &UserId,
272 connection_id: ConnectionId,
273 capabilities: MediaCapabilities,
274 media_port: &MediaTransport,
275 ) -> Option<()> {
276 self.user_operation(user_id, connection_id, media_port)
277 .apply_session_negotiated(capabilities, &[])
278 .await
279 }
280
281 #[cfg(test)]
282 pub(super) async fn user_count(&self) -> usize {
283 self.state.read().await.user_count()
284 }
285
286 pub(super) async fn is_empty(&self) -> bool {
287 self.state.read().await.is_empty()
288 }
289}