Skip to main content

o_sfu_core/engine/room/state/membership/
mod.rs

1use std::{
2    collections::{BTreeMap, BTreeSet},
3    mem,
4    sync::Arc,
5};
6
7use o_sfu_router::rtp::MediaCapabilities;
8use tracing::{debug, error, warn};
9
10use super::{
11    super::{
12        BroadcastPayload, BroadcastPayloadError, RoomEventMessage, RoomJoinError, RouterPlacement,
13        UserCloseReason,
14        effects::transport::RoomTransportPlan,
15        media_graph::{
16            CommittedTransportReceipt, SessionPlacementCommit, SessionPlacementRejection,
17        },
18        outbound::{MessageFanout, OutboundSender, VersionedRemoteTrackSnapshot, fanout_all},
19    },
20    UserJoinedFanout,
21    shared::{ActiveUser, RoomState},
22};
23#[cfg(test)]
24use crate::engine::MediaWorkerId;
25use crate::engine::{ConnectionId, UserId, UserInfo, media_transport::TransportTeardown};
26
27#[cfg(test)]
28#[expect(non_snake_case, reason = "test modules map to local TESTS directories")]
29mod TESTS;
30
31#[derive(Debug, Default)]
32pub struct LifecycleEffects {
33    pub close_requests: Vec<UserCloseRequest>,
34    pub fanouts: Vec<MessageFanout>,
35    pub track_snapshots: Vec<(OutboundSender, VersionedRemoteTrackSnapshot)>,
36}
37
38impl LifecycleEffects {
39    fn push_fanout(&mut self, fanout: Option<MessageFanout>) {
40        if let Some(fanout) = fanout {
41            self.fanouts.push(fanout);
42        }
43    }
44
45    fn push_close_request(&mut self, request: Option<UserCloseRequest>) {
46        if let Some(request) = request {
47            self.close_requests.push(request);
48        }
49    }
50}
51
52#[derive(Debug)]
53pub struct UserCloseRequest {
54    pub sender: OutboundSender,
55    pub reason: UserCloseReason,
56}
57
58type RuntimeUserRemoval = (ActiveUser, RoomTransportPlan, usize);
59
60#[derive(Debug)]
61pub struct PresenceCommit {
62    pub fanout: MessageFanout,
63}
64
65#[derive(Debug)]
66pub struct JoinCommit {
67    pub effects: LifecycleEffects,
68    pub receipt: CommittedTransportReceipt,
69    pub transport_plan: RoomTransportPlan,
70}
71
72#[expect(
73    clippy::large_enum_variant,
74    reason = "connection close is cold and boxing the room transport plan would add allocation without simplifying ownership"
75)]
76#[derive(Debug)]
77pub enum ConnectionCloseCommit {
78    Current {
79        session_teardown: Option<TransportTeardown>,
80        effects: LifecycleEffects,
81        transport_plan: RoomTransportPlan,
82        evictions: usize,
83    },
84    StalePlacement {
85        session_teardown: TransportTeardown,
86    },
87}
88
89#[derive(Debug)]
90pub struct DisconnectCommit {
91    pub session_teardowns: Vec<TransportTeardown>,
92    pub effects: LifecycleEffects,
93    pub transport_plan: RoomTransportPlan,
94    pub evictions: usize,
95}
96
97impl RoomState {
98    pub fn fanout_all(&self, message: RoomEventMessage) -> MessageFanout {
99        fanout_all(self.users.values().map(|user| user.sender.clone()), message)
100    }
101
102    pub fn fanout_all_except(
103        &self,
104        message: RoomEventMessage,
105        excluded_user_id: &UserId,
106    ) -> MessageFanout {
107        fanout_all(
108            self.users
109                .iter()
110                .filter(|(user_id, _session)| excluded_user_id != *user_id)
111                .map(|(_user_id, user)| user.sender.clone()),
112            message,
113        )
114    }
115
116    fn apply_join_routing(
117        &mut self,
118        user_id: &UserId,
119        connection_id: ConnectionId,
120        previous_connection: Option<ConnectionId>,
121        home_placement: RouterPlacement,
122    ) -> Result<SessionPlacementCommit, RoomJoinError> {
123        self.topology
124            .commit_session_placement(user_id, connection_id, previous_connection, home_placement)
125            .map_err(|rejection| {
126                match rejection {
127                    SessionPlacementRejection::MissingPreviousSession {
128                        previous_connection,
129                    } => {
130                        error!(
131                            ?user_id,
132                            connection_id = ?previous_connection,
133                            "missing committed routing session for replacement join"
134                        );
135                    }
136                    SessionPlacementRejection::Router(error) => {
137                        error!(
138                            ?user_id,
139                            ?error,
140                            "failed to mirror user join into room router"
141                        );
142                    }
143                }
144                RoomJoinError::RouterState
145            })
146    }
147
148    #[cfg(test)]
149    fn fallback_join_placement(&self) -> RouterPlacement {
150        RouterPlacement {
151            router: self.topology.router().placement_snapshot().primary(),
152            media_worker: MediaWorkerId::from_raw(0),
153        }
154    }
155
156    fn install_joined_session(
157        &mut self,
158        user_id: &UserId,
159        sender: OutboundSender,
160        connection_id: ConnectionId,
161    ) -> Option<OutboundSender> {
162        if let Some(user) = self.users.get_mut(user_id) {
163            let old_sender = mem::replace(&mut user.sender, sender);
164            user.reset_presentation();
165            user.parsed_client_rtp_capabilities = None;
166            user.connection_id = connection_id;
167            user.video_soft_pause_deadline = None;
168            return Some(old_sender);
169        }
170        self.users.insert(
171            user_id.clone(),
172            ActiveUser {
173                user_id: Arc::new(user_id.clone()),
174                info: UserInfo::default(),
175                server_featured: None,
176                parsed_client_rtp_capabilities: None,
177                connection_id,
178                video_soft_pause_deadline: None,
179                sender,
180            },
181        );
182        None
183    }
184
185    #[cfg(test)]
186    pub fn apply_join(
187        &mut self,
188        user_id: &UserId,
189        sender: OutboundSender,
190    ) -> Result<JoinCommit, RoomJoinError> {
191        self.apply_join_on_placement(
192            user_id,
193            sender,
194            UserJoinedFanout::Suppress,
195            self.fallback_join_placement(),
196        )
197    }
198
199    pub fn apply_join_on_placement(
200        &mut self,
201        user_id: &UserId,
202        sender: OutboundSender,
203        joined_fanout: UserJoinedFanout,
204        home_placement: RouterPlacement,
205    ) -> Result<JoinCommit, RoomJoinError> {
206        let previous_connection = self.users.get(user_id).map(|user| user.connection_id);
207        let is_new = previous_connection.is_none();
208        if is_new && self.users.len() >= self.admission_policy.max_sessions {
209            return Err(RoomJoinError::RoomFull);
210        }
211        let connection_id = ConnectionId::allocate(&mut self.next_connection_id);
212        let mut source_recipients = if previous_connection.is_some() {
213            self.topology
214                .committed_consumer_user_ids_for_owner_sources(user_id)
215        } else {
216            BTreeSet::new()
217        };
218        source_recipients.remove(user_id);
219        let placement =
220            self.apply_join_routing(user_id, connection_id, previous_connection, home_placement)?;
221        let receipt = placement.receipt;
222        let mut transport_plan = placement.replacement_transport_plan;
223        if let Some(previous_connection) = previous_connection {
224            transport_plan.extend_teardown(
225                self.staged_publishes
226                    .take_teardowns_for_connection(user_id, previous_connection),
227            );
228        }
229
230        let previous_sender = self.install_joined_session(user_id, sender, connection_id);
231        let had_previous_sender = previous_sender.is_some();
232
233        let mut effects = LifecycleEffects::default();
234        effects.push_close_request(previous_sender.map(|sender| UserCloseRequest {
235            sender,
236            reason: UserCloseReason::Replaced,
237        }));
238        effects
239            .track_snapshots
240            .extend(self.remote_track_snapshots_for_users(source_recipients, true));
241        effects.push_fanout(had_previous_sender.then(|| {
242            self.fanout_all_except(
243                RoomEventMessage::UserDeparted {
244                    user_id: user_id.clone(),
245                },
246                user_id,
247            )
248        }));
249        effects.push_fanout(if joined_fanout == UserJoinedFanout::Emit {
250            self.user_info_snapshot(user_id)
251                .map(|(joined_user_id, info)| {
252                    self.fanout_all_except(
253                        RoomEventMessage::UserJoined {
254                            user_id: joined_user_id,
255                            info,
256                        },
257                        user_id,
258                    )
259                })
260        } else {
261            None
262        });
263        Ok(JoinCommit {
264            effects,
265            receipt,
266            transport_plan,
267        })
268    }
269
270    fn remove_runtime_user(&mut self, user_id: &UserId) -> Option<RuntimeUserRemoval> {
271        let user = self.users.remove(user_id)?;
272        let (mut transport_plan, evictions) = self.topology.remove_session(user_id);
273        transport_plan.extend_teardown(
274            self.staged_publishes
275                .take_teardowns_for_connection(user_id, user.connection_id),
276        );
277        Some((user, transport_plan, evictions))
278    }
279
280    pub fn close_connection(
281        &mut self,
282        user_id: &UserId,
283        connection_id: ConnectionId,
284    ) -> Option<ConnectionCloseCommit> {
285        if self
286            .users
287            .get(user_id)
288            .is_none_or(|user| user.connection_id != connection_id)
289        {
290            let session_key = self
291                .topology
292                .retire_committed_placement(user_id, connection_id)?;
293            return Some(ConnectionCloseCommit::StalePlacement {
294                session_teardown: TransportTeardown::CloseSession { session_key },
295            });
296        }
297        let session_teardown = self
298            .committed_transport_user_key(user_id, connection_id)
299            .map(|session_key| TransportTeardown::CloseSession { session_key });
300        let mut source_recipients = self
301            .topology
302            .committed_consumer_user_ids_for_owner_sources(user_id);
303        source_recipients.remove(user_id);
304        let (user, transport_plan, evictions) = self.remove_runtime_user(user_id)?;
305        Some(ConnectionCloseCommit::Current {
306            session_teardown,
307            effects: LifecycleEffects {
308                close_requests: vec![UserCloseRequest {
309                    sender: user.sender,
310                    reason: UserCloseReason::RemovedByRuntime,
311                }],
312                fanouts: vec![self.fanout_all(RoomEventMessage::UserDeparted {
313                    user_id: user_id.clone(),
314                })],
315                track_snapshots: self.remote_track_snapshots_for_users(source_recipients, true),
316            },
317            transport_plan,
318            evictions,
319        })
320    }
321
322    pub fn apply_presence_update(
323        &mut self,
324        user_id: &UserId,
325        connection_id: ConnectionId,
326        info: &UserInfo,
327    ) -> Option<PresenceCommit> {
328        let Some(current_user) = self.users.get_mut(user_id) else {
329            warn!(
330                ?user_id,
331                connection_id = ?connection_id,
332                ?info,
333                "discarding user presence update because the user is missing"
334            );
335            return None;
336        };
337        if current_user.connection_id != connection_id {
338            warn!(
339                ?user_id,
340                connection_id = ?connection_id,
341                current_connection_id = ?current_user.connection_id,
342                ?info,
343                "discarding user presence update because the connection is stale"
344            );
345            return None;
346        }
347        current_user.apply_info_update(info);
348        let snapshot = BTreeMap::from([self.user_info_snapshot(user_id)?]);
349        debug!(
350            ?user_id,
351            connection_id = ?connection_id,
352            ?info,
353            snapshot_len = snapshot.len(),
354            "applied user presence update and staged user info fanout"
355        );
356        Some(PresenceCommit {
357            fanout: self.fanout_all(RoomEventMessage::UserInfoChanged(snapshot)),
358        })
359    }
360
361    pub fn set_user_negotiated(
362        &mut self,
363        user_id: &UserId,
364        connection_id: ConnectionId,
365        capabilities: MediaCapabilities,
366    ) -> Option<bool> {
367        let user = self.user_mut_for_connection(user_id, connection_id)?;
368        let became_ready = user.parsed_client_rtp_capabilities.is_none();
369        user.parsed_client_rtp_capabilities = Some(capabilities);
370        Some(became_ready)
371    }
372
373    pub fn apply_disconnect_users(&mut self, user_ids: &[UserId]) -> DisconnectCommit {
374        let mut source_recipients = BTreeSet::new();
375        for user_id in user_ids {
376            source_recipients.extend(
377                self.topology
378                    .committed_consumer_user_ids_for_owner_sources(user_id),
379            );
380        }
381        for user_id in user_ids {
382            source_recipients.remove(user_id);
383        }
384        let mut close_requests = Vec::new();
385        let mut session_teardowns = Vec::new();
386        let mut fanouts = Vec::new();
387        let mut transport_plan = RoomTransportPlan::default();
388        let mut evictions = 0;
389        for user_id in user_ids {
390            let Some(connection_id) = self.users.get(user_id).map(|user| user.connection_id) else {
391                continue;
392            };
393            let session_teardown = TransportTeardown::CloseSession {
394                session_key: self.transport_user_key(user_id, connection_id),
395            };
396            let Some((user, user_transport_plan, user_evictions)) =
397                self.remove_runtime_user(user_id)
398            else {
399                continue;
400            };
401            transport_plan.extend(user_transport_plan);
402            evictions += user_evictions;
403            session_teardowns.push(session_teardown);
404            close_requests.push(UserCloseRequest {
405                sender: user.sender,
406                reason: UserCloseReason::RemovedByRuntime,
407            });
408            fanouts.push(self.fanout_all(RoomEventMessage::UserDeparted {
409                user_id: user_id.clone(),
410            }));
411        }
412        DisconnectCommit {
413            session_teardowns,
414            effects: LifecycleEffects {
415                close_requests,
416                fanouts,
417                track_snapshots: self.remote_track_snapshots_for_users(source_recipients, true),
418            },
419            transport_plan,
420            evictions,
421        }
422    }
423
424    /// Plans a broadcast to every current user except `user_id`.
425    ///
426    /// Returns `Ok(None)` when `connection_id` is missing or stale for
427    /// `user_id`.
428    ///
429    /// # Errors
430    ///
431    /// Returns [`BroadcastPayloadError::TooLarge`] when the serialized payload
432    /// exceeds the room broadcast limit. Returns
433    /// [`BroadcastPayloadError::JsonSerialization`] when JSON serialization
434    /// fails.
435    pub fn broadcast_fanout(
436        &self,
437        user_id: &UserId,
438        connection_id: ConnectionId,
439        message: serde_json::Value,
440    ) -> Result<Option<MessageFanout>, BroadcastPayloadError> {
441        if self.user_for_connection(user_id, connection_id).is_none() {
442            return Ok(None);
443        }
444        let message = BroadcastPayload::try_new(message)?;
445        Ok(Some(self.fanout_all_except(
446            RoomEventMessage::Broadcast {
447                sender_id: user_id.clone(),
448                message,
449            },
450            user_id,
451        )))
452    }
453}