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 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}