Skip to main content

o_sfu_core/engine/room/
membership.rs

1//! Membership commit and effect boundary.
2//!
3//! [`JoinCommit`], [`ConnectionCloseCommit`] and
4//! [`DisconnectCommit`](crate::engine::room::state::DisconnectCommit) capture
5//! authoritative membership changes and the work consumed by [`RoomEffects`].
6//! Effects execute after the room-state write guard is released.
7
8#[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/// Compatibility marker that collapses every authenticated [`UserPermissions`] value.
26#[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
41/// Existing users replace their current connection without consuming another admission slot.
42pub struct JoinUserRequest {
43    pub user_id: UserId,
44    /// Ignored by room admission.
45    pub label: Option<String>,
46    /// Ignored by room admission.
47    pub permissions: UserPermissions,
48    pub sender: UserOutboundSender,
49}
50
51impl Room {
52    /// Commits the admission turn with the room router.
53    ///
54    /// # Errors
55    ///
56    /// Returns [`RoomJoinError::RoomFull`] when a new user exceeds capacity.
57    /// Returns [`RoomJoinError::NoUsableWorker`] when no configured worker is usable.
58    /// Returns [`RoomJoinError::RouterState`] when router placement cannot commit.
59    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    /// Executes context-enabled [`RoomEffects`] before returning the committed receipt
69    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    /// Returns `true` only when `connection_id` removed the current room user.
95    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    /// Returns `true` only when `connection_id` was current. A stale committed
110    /// placement may still be retired before returning `false`.
111    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    /// Captures recipients in the state snapshot that validates `connection_id`
157    /// as current for `sender_id`. Missing or stale senders are ignored.
158    ///
159    /// # Errors
160    ///
161    /// Returns [`BroadcastPayloadError`] when the payload exceeds the room
162    /// broadcast byte limit or cannot be measured as serialized JSON.
163    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    /// Ignores missing or stale connections. Publication transitions remain
188    /// authoritative for camera and screen-sharing presence.
189    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    /// Missing users are ignored.
217    ///
218    /// # Panics
219    ///
220    /// Panics if a current room user has no committed router placement.
221    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    /// Removes current sessions in one state commit and ignores missing users.
231    ///
232    /// Returns the number of sessions actually torn down.
233    ///
234    /// # Panics
235    ///
236    /// Panics if a current room user has no committed router placement.
237    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    /// Returns `None` for a missing or stale connection. The first accepted
268    /// capabilities commit realizes receiver routes waiting on negotiation.
269    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}