Skip to main content

o_sfu_core/engine/room/
factory.rs

1//! Room construction for rooms that are new to the runtime directory.
2//!
3//! `RoomManager` owns idempotent lookup, directory publication and creation
4//! diagnostics. This module contains the cold-path allocation step used
5//! after lookup misses, before the new room is visible to other runtime
6//! entrypoints.
7//!
8//! A factory-created room receives fresh process-local placement, the immutable
9//! runtime policy selected at boot and the shared runtime metrics catalog. It
10//! does not register the room or emit creation events.
11//!
12//! Same-room worker placement is not decided here. The factory
13//! gives a room its stable instance id and primary router id, while
14//! `RoomManager::join_user` assigns workers from packet-loop heartbeat delays
15//! when sessions arrive.
16
17use std::sync::{Arc, Mutex};
18
19use o_sfu_router::{RouterId, rtp::MediaCapabilities};
20use secrecy::SecretSlice;
21
22use super::{Room, RoomRuntimeContext};
23use crate::{
24    RoomMediaLimits, RoomWorkerPolicy, RuntimeFeatureFlags, VideoAdaptationTuning,
25    engine::{RoomInstanceId, metrics::RuntimeMetrics, sync::lock_unpoisoned},
26};
27
28/// admission limits that stay fixed for one room lifetime
29///
30/// this is kept separate from the wider runtime policy because admission is a
31/// narrow concern with its own tests and state checks
32///
33/// the policy is passed into `RoomState` at construction time and then treated
34/// as immutable room configuration
35#[derive(Debug, Clone, Copy, PartialEq, Eq)]
36pub struct RoomAdmissionPolicy {
37    /// maximum number of live users the room accepts at once
38    ///
39    /// replaced connections still consume this budget until the room transition
40    /// finishes and the old live user has been removed
41    ///
42    /// Also bounds absent publisher targets retained per receiver. Present
43    /// members do not consume this pending-intent allowance.
44    pub max_sessions: usize,
45}
46
47impl RoomAdmissionPolicy {
48    #[must_use]
49    pub const fn new(max_sessions: usize) -> Self {
50        Self { max_sessions }
51    }
52}
53
54/// stable runtime policy bundle shared by the room and its state model
55///
56/// this groups the room rules that are fixed for the room lifetime and read by
57/// more than one boundary during join, negotiation and observability work
58#[derive(Debug, Clone, PartialEq, Eq)]
59pub struct RoomRuntimePolicy {
60    /// room-level admission limits enforced by room state
61    pub admission_policy: RoomAdmissionPolicy,
62    /// feature surface the room advertises to clients
63    pub feature_flags: RuntimeFeatureFlags,
64    /// router-native capability baseline used for negotiation and bootstrap
65    pub router_rtp_capabilities: MediaCapabilities,
66    /// same-room local worker-placement policy selected at runtime boot
67    pub room_worker_policy: RoomWorkerPolicy,
68    /// room media activation caps applied by source policy
69    pub media_limits: RoomMediaLimits,
70    /// receiver video adaptation knobs applied by source policy
71    pub video_adaptation_tuning: VideoAdaptationTuning,
72}
73
74impl RoomRuntimePolicy {
75    #[must_use]
76    pub fn new(
77        admission_policy: RoomAdmissionPolicy,
78        feature_flags: RuntimeFeatureFlags,
79        router_rtp_capabilities: MediaCapabilities,
80    ) -> Self {
81        Self {
82            admission_policy,
83            feature_flags,
84            router_rtp_capabilities,
85            room_worker_policy: RoomWorkerPolicy::strict_single_router(),
86            media_limits: RoomMediaLimits::default(),
87            video_adaptation_tuning: VideoAdaptationTuning::default(),
88        }
89    }
90
91    /// return a room policy that uses the provided same-room worker policy
92    #[must_use]
93    pub fn with_room_worker_policy(mut self, room_worker_policy: RoomWorkerPolicy) -> Self {
94        self.room_worker_policy = room_worker_policy;
95        self
96    }
97
98    /// return a room policy that uses the provided media activation limits
99    #[must_use]
100    pub fn with_media_limits(mut self, media_limits: RoomMediaLimits) -> Self {
101        self.media_limits = media_limits;
102        self
103    }
104
105    #[must_use]
106    pub fn with_video_adaptation_tuning(
107        mut self,
108        video_adaptation_tuning: VideoAdaptationTuning,
109    ) -> Self {
110        self.video_adaptation_tuning = video_adaptation_tuning;
111        self
112    }
113}
114
115/// external room config passed in from the http or runtime edge
116///
117/// this type keeps room identity separate from operator-facing knobs and
118/// compatibility toggles that may be chosen per room at creation time
119#[derive(Debug, Clone, PartialEq, Eq)]
120pub struct RoomConfig {
121    /// whether this room should expose WebRTC to clients at all
122    pub web_rtc_enabled: bool,
123    /// compatibility recording address from `/v1/channel`
124    pub recording_address: Option<String>,
125}
126
127impl Default for RoomConfig {
128    fn default() -> Self {
129        Self {
130            web_rtc_enabled: true,
131            recording_address: None,
132        }
133    }
134}
135
136pub(super) struct RoomInit {
137    /// runtime-local instance and primary placement for the room
138    pub(super) runtime_context: RoomRuntimeContext,
139    /// validated room policy copied from runtime startup
140    pub(super) runtime_policy: RoomRuntimePolicy,
141    /// compatibility-facing issuer captured at room creation
142    pub(super) issuer: String,
143    /// room key captured from the first create request
144    pub(super) key: SecretSlice<u8>,
145    /// room-level compatibility configuration
146    pub(super) config: RoomConfig,
147    /// process metric catalog used by room observers
148    pub(super) metrics: Arc<RuntimeMetrics>,
149}
150
151/// Monotonic placement counters assigned by the current process.
152///
153/// Room instance ids and router ids are allocated under one lock so every
154/// new room receives one coherent runtime placement. The counters are not a
155/// distributed identity source and must not leak into the Odoo-facing room
156/// contract.
157#[derive(Debug)]
158struct RoomRuntimeAllocator {
159    next_room_instance_id: u64,
160    /// Next router id to allocate for room-local topology.
161    ///
162    /// Room creation consumes one primary router id. Dynamic spillover
163    /// placement consumes additional router ids when sessions join.
164    next_router_id: u64,
165}
166
167/// Cold-path constructor for rooms that are new to the directory.
168///
169/// `RoomFactory` keeps runtime-wide creation dependencies behind the manager
170/// so `RoomManager::serve_room` can focus on idempotent lookup and
171/// publication. Each call to [`Self::create`] returns an unpublished
172/// [`Room`] with fresh process-local placement. The caller must insert it in
173/// the directory before exposing it to other runtime entrypoints.
174#[derive(Debug)]
175pub(crate) struct RoomFactory {
176    /// Runtime-wide room rules cloned into each room.
177    ///
178    /// Keeping the policy here makes every room start from the validated
179    /// boot-time policy while still letting the room own its copy.
180    runtime_policy: RoomRuntimePolicy,
181    /// Process metric catalog cloned into each room.
182    metrics: Arc<RuntimeMetrics>,
183    /// Serialized allocator for process-local placement ids.
184    ///
185    /// This keeps concurrent create requests from receiving the same runtime
186    /// placement.
187    allocator: Mutex<RoomRuntimeAllocator>,
188}
189
190impl RoomFactory {
191    /// Builds the factory for one [`RoomManager`](super::RoomManager) lifetime.
192    #[must_use]
193    pub(crate) fn new(runtime_policy: RoomRuntimePolicy, metrics: Arc<RuntimeMetrics>) -> Self {
194        Self {
195            runtime_policy,
196            metrics,
197            allocator: Mutex::new(RoomRuntimeAllocator {
198                next_room_instance_id: 0,
199                next_router_id: 0,
200            }),
201        }
202    }
203
204    /// Creates an unpublished room from one manager lookup miss
205    ///
206    /// The room emits no creation diagnostics. `RoomManager` publishes it
207    /// before emitting its creation event
208    #[must_use]
209    pub(crate) fn create(
210        &self,
211        issuer: &str,
212        key: SecretSlice<u8>,
213        config: &RoomConfig,
214    ) -> Arc<Room> {
215        Arc::new(Room::new(RoomInit {
216            runtime_context: self.allocate_runtime_context(),
217            runtime_policy: self.runtime_policy.clone(),
218            issuer: issuer.to_owned(),
219            key,
220            config: config.clone(),
221            metrics: Arc::clone(&self.metrics),
222        }))
223    }
224
225    /// Reserves runtime-local placement for one new room.
226    ///
227    /// The primary router id is allocated here, but worker placement remains
228    /// unset until the first session join assigns the room from live load data.
229    ///
230    /// The mutex is poisoned-tolerant because placement allocation has no
231    /// partial side effect beyond the counters themselves. Recovering the inner
232    /// value keeps later room creation possible after an unrelated panic.
233    fn allocate_runtime_context(&self) -> RoomRuntimeContext {
234        let (room_instance_id, primary_router_id) = {
235            let mut allocator = lock_unpoisoned(&self.allocator);
236            let room_instance_id = RoomInstanceId::allocate(&mut allocator.next_room_instance_id);
237            let primary_router_id = RouterId(allocator.next_router_id);
238            allocator.next_router_id = allocator.next_router_id.saturating_add(1);
239            drop(allocator);
240            (room_instance_id, primary_router_id)
241        };
242        RoomRuntimeContext::new_unassigned(room_instance_id, primary_router_id)
243    }
244
245    /// reserve a new process-unique identifier for a spillover router
246    ///
247    /// this provides the room engine with a thread-safe way to allocate new router
248    /// identities on the fly when media load exceeds the primary worker capacity.
249    /// the allocator lock is held only long enough to increment the counter, keeping
250    /// the cold-path creation from blocking active request loops
251    pub(super) fn allocate_spillover_router(&self) -> RouterId {
252        let mut allocator = lock_unpoisoned(&self.allocator);
253        let router_id = RouterId(allocator.next_router_id);
254        allocator.next_router_id = allocator.next_router_id.saturating_add(1);
255        router_id
256    }
257}