1#[cfg(any(test, feature = "testing-transport"))]
9use std::sync::Mutex;
10use std::{
11 cmp::Reverse,
12 collections::{BTreeMap, BTreeSet},
13 future::Future,
14 sync::Arc,
15 time::Duration,
16};
17
18use futures_util::{
19 FutureExt, StreamExt,
20 future::try_join_all,
21 stream::{self, BoxStream, FuturesUnordered},
22};
23use o_sfu_telemetry::schema::event as telemetry_event;
24use secrecy::SecretSlice;
25use tokio::sync::{RwLock, Semaphore};
26use tracing::{info, warn};
27
28#[cfg(any(test, feature = "testing-transport"))]
29pub use super::placement::JoinPlacementTestGate;
30use super::{
31 Room, RoomConfig, RoomJoinError, RoomManagerJoinError, RoomRuntimePolicy,
32 RoomUserStatsSnapshot,
33 directory::{RoomDirectory, RoomDirectoryEntry, RoomLifecycleLease},
34 effects::batch::RoomEffectContext,
35 factory::RoomFactory,
36 membership::JoinUserRequest,
37 placement::JoinAdmissionTurn,
38 source_policy::SourcePolicyTurn,
39};
40use crate::engine::{
41 ConnectionId, RoomInstanceId, UserId,
42 media_transport::{
43 ActiveSpeakerSource, MediaTransport, TransportAdapterError, TransportSessionKey,
44 },
45 metrics::{RoomGaugeValues, RuntimeMetrics},
46 room::{
47 directory::{ExpiryReason, RoomRemovalPolicy},
48 instance::RoomManagerServeError,
49 },
50};
51
52#[cfg(any(test, feature = "testing-transport"))]
53#[path = "TESTS/support.rs"]
54mod test_support;
55
56const SOURCE_POLICY_CONCURRENCY: usize = 8;
57
58#[derive(Debug, Clone, PartialEq, Eq)]
60pub struct RuntimeRoomStatsSnapshot {
61 pub create_date: String,
63 pub uuid: String,
65 pub remote_address: String,
67 pub users_stats: RoomUserStatsSnapshot,
69 pub web_rtc_enabled: bool,
71}
72
73#[derive(Debug, Clone)]
75pub struct RoomUserAdmission {
76 pub room: Arc<Room>,
78 pub connection_id: ConnectionId,
80 pub transport_session_key: TransportSessionKey,
83}
84
85#[derive(Debug, Clone)]
87pub struct RuntimeRoomDirectorySnapshot {
88 pub room: Arc<Room>,
90 pub create_date: String,
92 pub remote_address: String,
94}
95
96fn retrieve_room_reservation(
97 directory: &RoomDirectory,
98 issuer: &str,
99 key: &SecretSlice<u8>,
100 config: &RoomConfig,
101) -> Result<Option<Arc<Room>>, RoomManagerServeError> {
102 let Some(entry) = directory.entry_by_issuer(issuer) else {
103 return Ok(None);
104 };
105 if !entry.room.definition.matches_reservation(key, config) {
106 warn!(
107 event = telemetry_event::ROOM_RESERVATION_CONFLICT,
108 room_id = entry.room.uuid(),
109 issuer,
110 config = ?config,
111 "conflicting reservation request"
112 );
113 return Err(RoomManagerServeError::ConflictingReservation);
114 }
115 entry.lifecycle.renew_reservation();
116 Ok(Some(Arc::clone(&entry.room)))
117}
118
119#[derive(Debug)]
121pub struct RoomManager {
122 directory: RwLock<RoomDirectory>,
123 factory: RoomFactory,
124 reservation_ttl: Duration,
125 departure_grace: Duration,
126 policy_permits: Semaphore,
127 #[cfg(any(test, feature = "testing-transport"))]
128 join_placement_gate: Mutex<Option<Arc<JoinPlacementTestGate>>>,
129}
130
131fn room_active_speaker_sources(
132 snapshots: &[Arc<Vec<ActiveSpeakerSource>>],
133 source_ids: &[u64],
134) -> Vec<ActiveSpeakerSource> {
135 let mut by_media = BTreeMap::<u64, ActiveSpeakerSource>::new();
136 for &media_id in source_ids {
137 for snapshot in snapshots {
138 let candidate = snapshot
139 .binary_search_by_key(&media_id, |source| source.transport_media_id().as_u64())
140 .ok()
141 .and_then(|index| snapshot.get(index))
142 .copied();
143 if let Some(candidate) = candidate {
144 let candidate_rank = (
145 candidate.observed_at(),
146 candidate.last_audio_level_dbov().unwrap_or(i8::MIN),
147 );
148 let entry = by_media.entry(media_id).or_insert(candidate);
149 let current_rank = (
150 entry.observed_at(),
151 entry.last_audio_level_dbov().unwrap_or(i8::MIN),
152 );
153 if candidate_rank > current_rank {
154 *entry = candidate;
155 }
156 }
157 }
158 }
159 let mut sources = by_media.into_values().collect::<Vec<_>>();
160 sources.sort_unstable_by_key(|source| {
161 (
162 Reverse(source.observed_at()),
163 source.transport_media_id().as_u64(),
164 )
165 });
166 sources
167}
168
169impl RoomManager {
170 #[must_use]
173 pub fn new(
174 runtime_policy: RoomRuntimePolicy,
175 metrics: Arc<RuntimeMetrics>,
176 reservation_ttl: Duration,
177 departure_grace: Duration,
178 ) -> Self {
179 let factory = RoomFactory::new(runtime_policy, metrics);
180 Self {
181 directory: RwLock::new(RoomDirectory::default()),
182 factory,
183 reservation_ttl,
184 departure_grace,
185 policy_permits: Semaphore::new(SOURCE_POLICY_CONCURRENCY),
186 #[cfg(any(test, feature = "testing-transport"))]
187 join_placement_gate: Mutex::new(None),
188 }
189 }
190
191 pub async fn serve_room(
203 &self,
204 issuer: &str,
205 key: SecretSlice<u8>,
206 config: &RoomConfig,
207 remote_address: Option<&str>,
208 ) -> Result<Arc<Room>, RoomManagerServeError> {
209 {
210 let directory = self.directory.read().await;
211 if let Some(room) = retrieve_room_reservation(&directory, issuer, &key, config)? {
212 return Ok(room);
213 }
214 }
215 let mut directory = self.directory.write().await;
216 if let Some(room) = retrieve_room_reservation(&directory, issuer, &key, config)? {
217 return Ok(room);
218 }
219 let room = self.factory.create(issuer, key, config);
220 directory.insert(
221 Arc::clone(&room),
222 remote_address,
223 self.reservation_ttl,
224 self.departure_grace,
225 );
226 drop(directory);
227 info!(
228 event = telemetry_event::ROOM_CREATED,
229 room_id = room.uuid(),
230 remote_address = remote_address.unwrap_or("unknown"),
231 web_rtc_enabled = config.web_rtc_enabled,
232 "room created"
233 );
234 Ok(room)
235 }
236
237 pub async fn get_by_uuid(&self, uuid: &str) -> Option<Arc<Room>> {
239 let directory = self.directory.read().await;
240 directory.get_by_uuid(uuid)
241 }
242
243 pub async fn stats_snapshots(
249 &self,
250 media_transport: &MediaTransport,
251 ) -> Vec<RuntimeRoomStatsSnapshot> {
252 let entries = self.directory_entries().await;
253 let mut snapshots = Vec::with_capacity(entries.len());
254 for entry in entries {
255 snapshots.push(self.entry_stats_snapshot(entry, media_transport).await);
256 }
257 snapshots
258 }
259
260 pub async fn directory_snapshots(&self) -> Vec<RuntimeRoomDirectorySnapshot> {
262 self.directory_entries()
263 .await
264 .into_iter()
265 .map(|entry| RuntimeRoomDirectorySnapshot {
266 room: entry.room,
267 create_date: entry.create_date,
268 remote_address: entry.remote_address,
269 })
270 .collect()
271 }
272
273 pub async fn room_gauges(&self) -> RoomGaugeValues {
278 let rooms = self.directory.read().await.rooms();
279 let mut gauges = RoomGaugeValues {
280 rooms: rooms.len(),
281 ..RoomGaugeValues::default()
282 };
283 for room in rooms {
284 let state = room.state.read().await;
285 let media = state.media_counts();
286 gauges.users = gauges.users.saturating_add(state.user_count());
287 gauges.publications = gauges.publications.saturating_add(media.publications);
288 gauges.subscriptions = gauges.subscriptions.saturating_add(media.subscriptions);
289 gauges.recording_rooms = gauges
290 .recording_rooms
291 .saturating_add(usize::from(state.recording_state().recording == Some(true)));
292 }
293 gauges
294 }
295
296 pub async fn directory_snapshot(&self, room_id: &str) -> Option<RuntimeRoomDirectorySnapshot> {
298 let entry = self.entry(room_id).await?;
299 Some(RuntimeRoomDirectorySnapshot {
300 room: entry.room,
301 create_date: entry.create_date,
302 remote_address: entry.remote_address,
303 })
304 }
305
306 pub async fn sync_source_packet_selection_policies_for_runtime_ids(
308 &self,
309 room_instance_ids: &BTreeSet<RoomInstanceId>,
310 media_transport: &MediaTransport,
311 ) {
312 self.source_policy_turns(room_instance_ids, media_transport)
313 .await
314 .for_each(|_| async {})
315 .await;
316 }
317
318 pub async fn source_policy_turns<'a>(
324 &'a self,
325 room_instance_ids: &BTreeSet<RoomInstanceId>,
326 media_transport: &'a MediaTransport,
327 ) -> BoxStream<'a, RoomInstanceId> {
328 let rooms = self
329 .directory_entries_for_instance_ids(room_instance_ids)
330 .await;
331 let mut absent = room_instance_ids.clone();
332 let mut jobs = Vec::with_capacity(rooms.len());
333 let mut needed_workers = BTreeSet::new();
334 for room in rooms {
335 absent.remove(&room.instance_id());
336 let worker_indices = room.source_policy_worker_indices(media_transport).await;
337 needed_workers.extend(worker_indices.iter().copied());
338 jobs.push((room, worker_indices));
339 }
340 let observations = needed_workers
343 .into_iter()
344 .map(|worker_index| {
345 let future = async move {
346 let mut sources = media_transport
347 .active_speaker_source_snapshot_for_worker(worker_index)
348 .await?;
349 sources.sort_unstable_by_key(|source| {
350 (
351 source.transport_media_id().as_u64(),
352 Reverse(source.observed_at()),
353 Reverse(source.last_audio_level_dbov().unwrap_or(i8::MIN)),
354 )
355 });
356 sources.dedup_by_key(|source| source.transport_media_id());
357 Ok::<_, TransportAdapterError>(Arc::new(sources))
358 }
359 .boxed()
360 .shared();
361 (worker_index, future)
362 })
363 .collect::<BTreeMap<_, _>>();
364 let turns = jobs
365 .into_iter()
366 .map(|(room, worker_indices)| {
367 let worker_reads = worker_indices
368 .iter()
369 .map(|worker_index| observations.get(worker_index).cloned())
370 .collect::<Option<Vec<_>>>();
371 async move {
372 let room_id = room.instance_id();
373 let Some(worker_reads) = worker_reads else {
374 media_transport.schedule_source_policy_retry(room_id);
375 return room_id;
376 };
377 let Ok(snapshots) = try_join_all(worker_reads).await else {
378 media_transport.schedule_source_policy_retry(room_id);
379 return room_id;
380 };
381 let guard = room.lock_source_policy().await;
382 let Ok(_permit) = self.policy_permits.acquire().await else {
383 return room_id;
384 };
385 if room.source_policy_worker_indices(media_transport).await != worker_indices {
386 media_transport.schedule_source_policy_retry(room_id);
387 return room_id;
388 }
389 let source_ids = {
390 let state = room.state.read().await;
391 state
392 .topology
393 .published_sources()
394 .map(|source| source.transport.transport_media_id().as_u64())
395 .collect::<Vec<_>>()
396 };
397 let sources = room_active_speaker_sources(&snapshots, &source_ids);
398 SourcePolicyTurn::packet_selection()
399 .execute_guarded(&guard, Some(media_transport), Some(&sources))
400 .await;
401 room_id
402 }
403 })
404 .collect::<FuturesUnordered<_>>();
405 stream::iter(absent).chain(turns).boxed()
406 }
407
408 pub async fn join_user(
420 &self,
421 room_id: &str,
422 request: JoinUserRequest,
423 media_transport: &MediaTransport,
424 ) -> Result<RoomUserAdmission, RoomManagerJoinError> {
425 let mutation = self
426 .begin_current_room_mutation(room_id)
427 .await
428 .ok_or(RoomManagerJoinError::MissingRoom)?;
429 let room = Arc::clone(&mutation.room);
430
431 let admission = JoinAdmissionTurn::from_factory(request, media_transport, &self.factory);
432 #[cfg(any(test, feature = "testing-transport"))]
433 let admission = admission.with_gate(self.join_placement_gate_for_test());
434 let join_commit = match room
435 .commit_admission(admission, RoomEffectContext::runtime(media_transport))
436 .await
437 {
438 Ok(commit) => commit,
439 Err(err) => {
440 self.finish_session_mutation(room_id, mutation, None).await;
441 return Err(match err {
442 RoomJoinError::RoomFull => RoomManagerJoinError::RoomFull,
443 RoomJoinError::RouterState => RoomManagerJoinError::RouterState,
444 RoomJoinError::NoUsableWorker => RoomManagerJoinError::NoUsableWorker,
445 });
446 }
447 };
448 mutation.lease.clear_expiration();
451 let receipt = room
452 .finalize_admission(join_commit, RoomEffectContext::runtime(media_transport))
453 .await;
454
455 self.finish_session_mutation(room_id, mutation, None).await;
456 Ok(RoomUserAdmission {
457 room,
458 connection_id: receipt.transport_session_key.connection_id(),
459 transport_session_key: receipt.transport_session_key,
460 })
461 }
462
463 pub async fn close_session(
469 &self,
470 room_id: &str,
471 user_id: &UserId,
472 connection_id: ConnectionId,
473 media_transport: &MediaTransport,
474 ) -> bool {
475 let Some(did_remove_active_session) = self
476 .run_current_room_mutation(
477 room_id,
478 |room| async move {
479 room.remove_user(user_id, connection_id, media_transport)
480 .await
481 },
482 Some(RoomRemovalPolicy::AfterGrace),
483 )
484 .await
485 else {
486 return false;
487 };
488 did_remove_active_session
489 }
490
491 pub async fn disconnect_users(
496 &self,
497 room_id: &str,
498 user_ids: &[UserId],
499 media_transport: &MediaTransport,
500 ) {
501 let result = self
502 .run_current_room_mutation(
503 room_id,
504 |room| async move { room.disconnect_users(user_ids, media_transport).await },
505 Some(RoomRemovalPolicy::Immediately),
506 )
507 .await;
508 let (disconnected_session_count, outcome) =
509 result.map_or((0, "room_missing"), |count| (count, "success"));
510 info!(
511 event = telemetry_event::USERS_BULK_DISCONNECTED,
512 disconnected_session_count,
513 outcome,
514 requested_user_count = user_ids.len(),
515 room_id,
516 "bulk disconnect handled"
517 );
518 }
519
520 pub async fn check_expired_room_reservations(&self) {
522 let mut directory = self.directory.write().await;
523 for entry in directory.entries() {
524 if let Some(expiry_reason) = entry.lifecycle.claim_expired_room()
525 && directory.remove_if_current(entry.room.uuid(), &entry.room)
526 {
527 match expiry_reason {
528 ExpiryReason::ReservationLapsed => {
529 info!(
530 event = telemetry_event::ROOM_RESERVATION_EXPIRED,
531 room_id = entry.room.uuid(),
532 "room reservation expired"
533 );
534 }
535 ExpiryReason::GraceElapsed => {
536 info!(
537 event = telemetry_event::ROOM_DESTROYED,
538 room_id = entry.room.uuid(),
539 "room destroyed"
540 );
541 }
542 }
543 }
544 }
545 drop(directory);
546 }
547
548 #[cfg(test)]
549 pub(super) async fn with_current_room<T, F, Fut>(&self, room_id: &str, action: F) -> Option<T>
550 where
551 F: FnOnce(Arc<Room>) -> Fut,
552 Fut: Future<Output = T>,
553 {
554 self.run_current_room_mutation(room_id, action, None).await
555 }
556
557 async fn run_current_room_mutation<T, F, Fut>(
558 &self,
559 room_id: &str,
560 action: F,
561 on_empty: Option<RoomRemovalPolicy>,
562 ) -> Option<T>
563 where
564 F: FnOnce(Arc<Room>) -> Fut,
565 Fut: Future<Output = T>,
566 {
567 let mutation = self.begin_current_room_mutation(room_id).await?;
568 let output = action(Arc::clone(&mutation.room)).await;
569 self.finish_session_mutation(room_id, mutation, on_empty)
570 .await;
571 Some(output)
572 }
573
574 async fn finish_session_mutation(
575 &self,
576 room_id: &str,
577 mutation: CurrentRoomMutation,
578 on_empty: Option<RoomRemovalPolicy>,
579 ) {
580 let CurrentRoomMutation { room, lease } = mutation;
581 if lease.finish(on_empty, room.is_empty().await)
585 && self
586 .directory
587 .write()
588 .await
589 .remove_if_current(room_id, &room)
590 {
591 info!(
592 event = telemetry_event::ROOM_DESTROYED,
593 room_id, "room destroyed"
594 );
595 }
596 }
597
598 async fn begin_current_room_mutation(&self, room_id: &str) -> Option<CurrentRoomMutation> {
603 let directory = self.directory.read().await;
604 let entry = directory.entry(room_id)?;
605 let lease = entry.lifecycle.begin()?;
606 let room = Arc::clone(&entry.room);
607 drop(directory);
608 Some(CurrentRoomMutation { room, lease })
609 }
610
611 async fn entry(&self, room_id: &str) -> Option<RoomDirectoryEntry> {
612 let directory = self.directory.read().await;
613 directory.entry(room_id).cloned()
614 }
615
616 async fn directory_entries(&self) -> Vec<RoomDirectoryEntry> {
617 let directory = self.directory.read().await;
618 directory.entries()
619 }
620
621 async fn directory_entries_for_instance_ids(
622 &self,
623 room_instance_ids: &BTreeSet<RoomInstanceId>,
624 ) -> Vec<Arc<Room>> {
625 let directory = self.directory.read().await;
626 room_instance_ids
627 .iter()
628 .filter_map(|room_instance_id| directory.get_by_instance_id(*room_instance_id))
629 .collect()
630 }
631
632 async fn entry_stats_snapshot(
633 &self,
634 entry: RoomDirectoryEntry,
635 media_transport: &MediaTransport,
636 ) -> RuntimeRoomStatsSnapshot {
637 let room = entry.room;
638 let users_stats = room.session_stats_snapshot(media_transport).await;
639 RuntimeRoomStatsSnapshot {
640 create_date: entry.create_date,
641 uuid: room.uuid().to_owned(),
642 remote_address: entry.remote_address,
643 users_stats,
644 web_rtc_enabled: room.web_rtc_enabled(),
645 }
646 }
647}
648
649#[derive(Debug)]
654struct CurrentRoomMutation {
655 room: Arc<Room>,
656 lease: RoomLifecycleLease,
657}