Skip to main content

o_sfu_core/engine/room/manager/
mod.rs

1//! Current-room registry and lifecycle coordinator.
2//!
3//! [`RoomManager`] publishes one current room per issuer, admits WebSocket
4//! sessions and runs background room work. Lease admission holds the directory
5//! read guard to prove currentness and releases it before room work begins.
6//! Empty-room removal waits for every accepted lease to finish.
7
8#[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/// operator-facing room stats assembled from directory and transport snapshots
59#[derive(Debug, Clone, PartialEq, Eq)]
60pub struct RuntimeRoomStatsSnapshot {
61    /// room publication timestamp in RFC 3339 UTC format
62    pub create_date: String,
63    /// room uuid returned by `/v1/channel`
64    pub uuid: String,
65    /// first create request address or `unknown` when unavailable
66    pub remote_address: String,
67    /// live user, stream and bitrate stats read after the directory snapshot
68    pub users_stats: RoomUserStatsSnapshot,
69    /// room creation flag exposed by `/v1/stats`
70    pub web_rtc_enabled: bool,
71}
72
73/// committed admission result returned after join-side effects run
74#[derive(Debug, Clone)]
75pub struct RoomUserAdmission {
76    /// current room that accepted the session
77    pub room: Arc<Room>,
78    /// room-local connection id assigned to the admitted user
79    pub connection_id: ConnectionId,
80    /// transport key used by [`crate::prelude::SfuCore`] to build
81    /// [`crate::prelude::MediaSession`]
82    pub transport_session_key: TransportSessionKey,
83}
84
85/// current room directory row used by diagnostics views
86#[derive(Debug, Clone)]
87pub struct RuntimeRoomDirectorySnapshot {
88    /// current room for this directory row
89    pub room: Arc<Room>,
90    /// room publication timestamp in RFC 3339 UTC format
91    pub create_date: String,
92    /// first create request address or `unknown` when unavailable
93    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/// Coordinates current room admission and lifecycle.
120#[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    /// builds a room manager with an empty directory
171    ///
172    #[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    /// Returns the current room for `issuer` or publishes a new reservation.
192    ///
193    /// The first reservation fixes decoded key bytes, config and remote address.
194    /// The admitting runtime validates key encoding and strength before this call.
195    /// Matching requests return the same room and renew an outstanding
196    /// reservation without rearming one retired by a successful join.
197    ///
198    /// # Errors
199    ///
200    /// Returns [`RoomManagerServeError::ConflictingReservation`] when the
201    /// current room has a different `key` or `config`.
202    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    /// returns the current room for a public room uuid
238    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    /// builds `/v1/stats` rows from one directory snapshot
244    ///
245    /// the directory lock is released before transport stats are read, so the
246    /// returned rows are best-effort runtime observations rather than a global
247    /// transaction across room and media state
248    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    /// returns current directory rows for room diagnostics
261    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    /// Returns counts from rooms in one directory snapshot.
274    ///
275    /// Room states are read sequentially after the directory lock is released.
276    /// Removed rooms may contribute once. New rooms appear on the next call.
277    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    /// returns one current directory row for room diagnostics
297    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    /// Recalculates packet selection for the requested current rooms.
307    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    /// Yields each requested room after its ordered policy turn completes.
319    ///
320    /// Worker reads are shared within the batch. Missing rooms also complete.
321    /// Failed observations schedule a retry before completing their room.
322    /// Dropping the stream cancels pending turns, including accepted effects.
323    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        // Shared futures send one read per worker while each room waits only for
341        // the workers that own its published sources.
342        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    /// Admits one WebSocket connection into a current room.
409    ///
410    /// Returns after join-side room effects complete.
411    ///
412    /// # Errors
413    ///
414    /// Returns [`RoomManagerJoinError::MissingRoom`] when `room_id` is not
415    /// current. Returns [`RoomManagerJoinError::RoomFull`] when a new user
416    /// exceeds room capacity. Returns [`RoomManagerJoinError::NoUsableWorker`]
417    /// when no worker can accept placement and [`RoomManagerJoinError::RouterState`]
418    /// when router placement cannot commit.
419    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        // Retire the reservation only after membership commits. Failed admission
449        // must leave its deadline intact for the reaper.
450        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    /// closes one room connection and then re-checks empty-room removal
464    ///
465    /// returns `false` when the room is missing or the connection was not
466    /// removed by this call. the empty current room can still be removed after
467    /// stale or already-completed teardown
468    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    /// disconnects selected users from a current room and removes it if empty
492    ///
493    /// missing rooms are ignored because the caller's disconnect intent is
494    /// already satisfied
495    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    /// Claims and removes expired reservations with no active room mutations.
521    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        // Read emptiness after this mutation. A later accepted mutation can
582        // supply the final removal proof. Dropping the final lease supplies no
583        // proof and cancels the pending claim.
584        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    /// Accepts work while the directory read guard proves the row is current.
599    ///
600    /// The returned lease protects subsequent room work without retaining the
601    /// directory guard.
602    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/// lease plus room pointer accepted from the current directory row
650///
651/// dropping the lease without [`RoomManager::finish_session_mutation`] releases
652/// admission but never removes the room
653#[derive(Debug)]
654struct CurrentRoomMutation {
655    room: Arc<Room>,
656    lease: RoomLifecycleLease,
657}