1use std::{collections::BTreeSet, future::Future, io, process, sync::Arc, time::Duration};
6
7use anyhow::Result as AnyResult;
8use futures_util::{StreamExt, stream::SelectAll};
9use thiserror::Error;
10#[cfg(unix)]
11use tokio::signal::unix::{SignalKind, signal};
12use tokio::{
13 net::TcpListener,
14 runtime::Builder,
15 signal::ctrl_c,
16 task::JoinHandle,
17 time::{MissedTickBehavior, interval, sleep},
18};
19use tokio_util::{
20 sync::CancellationToken,
21 task::{AbortOnDropHandle, TaskTracker},
22};
23use tracing::{info, warn};
24
25use crate::{config::Config, core::prelude::SfuCore};
26
27pub(crate) mod auth;
28pub(crate) mod diagnostics;
29pub(crate) mod http_server;
30pub(crate) mod options;
31pub(crate) mod request_origin;
32#[cfg(test)]
33#[path = "TESTS/support.rs"]
34pub(super) mod test_support;
35pub(crate) mod websocket_server;
36
37use http_server::{serve_http, serve_http_on};
38pub(crate) use o_sfu_core::{
39 prelude::{ConnectionId, SessionBitrateLimits},
40 server::{metrics, packet_sinks, room, transport as media_transport},
41};
42pub(crate) use o_sfu_telemetry::{self as telemetry, prometheus};
43use options::{RuntimeConfig, effective_feature_flags};
44use room::{RoomAdmissionPolicy, RoomManager, RoomRuntimePolicy};
45use telemetry::{init_tracing, schema::event as telemetry_event};
46
47pub(crate) use self::{
48 media_transport::{MediaTransport, MediaTransportConfig, MediaTransportDeps},
49 metrics::RuntimeMetrics,
50 packet_sinks::RoomPacketSinkRegistry,
51};
52
53#[derive(Debug, Error)]
55pub enum ServeError {
56 #[error(transparent)]
58 Io(#[from] io::Error),
59 #[error(
61 "runtime shutdown exceeded its deadline with {remaining_sessions} WebSocket sessions remaining"
62 )]
63 ShutdownIncomplete {
64 remaining_sessions: usize,
66 },
67}
68
69#[derive(Debug)]
71pub struct Runtime {
72 config: RuntimeConfig,
73 room_manager: Arc<RoomManager>,
74 metrics: Arc<RuntimeMetrics>,
75 media_transport: MediaTransport,
76}
77
78#[derive(Debug, Clone)]
80pub(super) struct RuntimeState {
81 config: RuntimeConfig,
82 room_manager: Arc<RoomManager>,
83 media_transport: MediaTransport,
84 sfu_core: SfuCore,
85 metrics: Arc<RuntimeMetrics>,
86 pre_auth_websocket_admission: websocket_server::PreAuthWebSocketAdmission,
87 session_shutdown: CancellationToken,
88 session_tasks: TaskTracker,
89}
90
91#[derive(Default)]
92pub(super) struct RuntimeServices {
93 metrics: Arc<RuntimeMetrics>,
94 packet_sink_registry: Arc<RoomPacketSinkRegistry>,
95}
96
97impl Runtime {
98 pub fn new(config: &Config) -> AnyResult<Self> {
107 Self::from_services(config, RuntimeServices::default())
108 }
109
110 #[cfg(feature = "testing-transport")]
111 #[doc(hidden)]
112 #[must_use]
113 pub fn media_transport_for_test(&self) -> MediaTransport {
114 self.media_transport.clone()
115 }
116
117 fn from_services(config: &Config, services: RuntimeServices) -> AnyResult<Self> {
118 config.http.validate()?;
119 config.diagnostics.validate()?;
120 let runtime_config = RuntimeConfig::from_config(config)?;
121 let media_transport = build_media_transport(config, &services)?;
122 let room_runtime_policy = build_room_runtime_policy(config, &media_transport);
123 info!(
124 event = telemetry_event::RUNTIME_BOOT,
125 rtc_udp_io_backend = config.transport.rtc_udp_io_backend.wire_name(),
126 "runtime configuration loaded"
127 );
128 let room_manager = build_room_manager(
129 room_runtime_policy,
130 &services,
131 config.user.room_reservation_ttl,
132 config.user.departure_grace,
133 );
134 Ok(Self {
135 config: runtime_config,
136 room_manager,
137 metrics: services.metrics,
138 media_transport,
139 })
140 }
141
142 pub async fn serve_listener(
151 self,
152 listener: TcpListener,
153 shutdown: impl Future<Output = ()>,
154 ) -> Result<(), ServeError> {
155 self.serve(
156 |state, token| serve_http_on(listener, state, token),
157 async move {
158 shutdown.await;
159 Ok(())
160 },
161 )
162 .await
163 }
164
165 async fn serve<F, HttpServer, Shutdown>(
166 self,
167 http_server: F,
168 shutdown: Shutdown,
169 ) -> Result<(), ServeError>
170 where
171 F: FnOnce(RuntimeState, CancellationToken) -> HttpServer,
172 HttpServer: Future<Output = io::Result<()>>,
173 Shutdown: Future<Output = io::Result<()>>,
174 {
175 let timeout = self.config.http.shutdown_timeout.as_duration();
176 let tasks = RuntimeTasks::spawn(Arc::clone(&self.room_manager), self.media_transport);
177 let state = RuntimeState::from_parts(
178 self.config,
179 self.room_manager,
180 self.metrics,
181 tasks.media_transport.clone(),
182 tasks.session_shutdown.clone(),
183 tasks.session_tasks.clone(),
184 );
185 let listener_shutdown = tasks.shutdown_token.child_token();
186 let server = http_server(state, listener_shutdown.clone());
187 tokio::pin!(server);
188 tokio::pin!(shutdown);
189 let (server_done, mut failure) = tokio::select! {
190 result = &mut server => (true, result.err()),
191 result = &mut shutdown => {
192 listener_shutdown.cancel();
193 (false, result.err())
194 }
195 };
196 let deadline = sleep(timeout);
197 tokio::pin!(deadline);
198 if let Some(error) = &failure {
199 warn!(?error, "runtime serving stopped with an error");
200 }
201 let session_tasks = tasks.session_tasks.clone();
202 let teardown = async move {
203 if !server_done && let Err(error) = server.await {
204 warn!(?error, "HTTP server stopped with an error during shutdown");
205 failure.get_or_insert(error);
206 }
207 tasks.shutdown().await;
208 failure.map_or(Ok(()), |error| Err(ServeError::Io(error)))
209 };
210 tokio::select! {
211 biased;
212 () = &mut deadline => Err(ServeError::ShutdownIncomplete {
213 remaining_sessions: session_tasks.len(),
214 }),
215 result = teardown => result,
216 }
217 }
218}
219
220struct RuntimeTasks {
222 shutdown_token: CancellationToken,
223 session_shutdown: CancellationToken,
224 session_tasks: TaskTracker,
225 source_packet_policy_sync: AbortOnDropHandle<()>,
226 rooms_reservations_reaper: AbortOnDropHandle<()>,
227 media_transport: MediaTransport,
228}
229
230impl RuntimeTasks {
231 fn spawn(room_manager: Arc<RoomManager>, media_transport: MediaTransport) -> Self {
232 let shutdown_token = CancellationToken::new();
233 let source_packet_policy_sync = spawn_source_packet_policy_update_task(
234 Arc::clone(&room_manager),
235 media_transport.clone(),
236 shutdown_token.child_token(),
237 );
238 let rooms_reservations_reaper =
239 spawn_room_reservation_expiration_reaper(room_manager, shutdown_token.child_token());
240 Self {
241 session_shutdown: shutdown_token.child_token(),
242 shutdown_token,
243 session_tasks: TaskTracker::new(),
244 source_packet_policy_sync: AbortOnDropHandle::new(source_packet_policy_sync),
245 rooms_reservations_reaper: AbortOnDropHandle::new(rooms_reservations_reaper),
246 media_transport,
247 }
248 }
249
250 async fn shutdown(mut self) {
251 self.session_tasks.close();
252 self.session_shutdown.cancel();
253 self.session_tasks.wait().await;
254 self.shutdown_token.cancel();
255 if let Err(error) = (&mut self.source_packet_policy_sync).await
256 && !error.is_cancelled()
257 {
258 warn!(
259 ?error,
260 task = "source packet policy update",
261 "runtime background task stopped unexpectedly"
262 );
263 }
264 if let Err(error) = (&mut self.rooms_reservations_reaper).await
265 && !error.is_cancelled()
266 {
267 warn!(
268 ?error,
269 task = "rooms reservation reaper",
270 "runtime background task stopped unexpectedly"
271 );
272 }
273 self.media_transport.shutdown().await;
274 }
275}
276
277impl Drop for RuntimeTasks {
278 fn drop(&mut self) {
279 self.shutdown_token.cancel();
280 self.media_transport.cancel();
281 }
282}
283
284impl RuntimeState {
285 fn from_parts(
286 config: RuntimeConfig,
287 rooms: Arc<RoomManager>,
288 metrics: Arc<RuntimeMetrics>,
289 media_transport: MediaTransport,
290 session_shutdown: CancellationToken,
291 session_tasks: TaskTracker,
292 ) -> Self {
293 let sfu_core = SfuCore::new(media_transport.clone(), Arc::clone(&rooms));
294 let pre_auth_websocket_admission = websocket_server::PreAuthWebSocketAdmission::new(
295 config.auth.max_pre_auth_websocket_sessions,
296 config.auth.max_pre_auth_websocket_sessions_per_origin,
297 );
298 Self {
299 config,
300 room_manager: rooms,
301 media_transport,
302 sfu_core,
303 metrics,
304 pre_auth_websocket_admission,
305 session_shutdown,
306 session_tasks,
307 }
308 }
309}
310
311async fn shutdown_signal() -> io::Result<()> {
312 #[cfg(unix)]
313 {
314 let mut terminate = signal(SignalKind::terminate())?;
315 tokio::select! {
316 result = ctrl_c() => result,
317 _ = terminate.recv() => Ok(()),
318 }
319 }
320 #[cfg(not(unix))]
321 ctrl_c().await
322}
323
324fn spawn_source_packet_policy_update_task(
326 rooms: Arc<RoomManager>,
327 media_transport: MediaTransport,
328 shutdown_token: CancellationToken,
329) -> JoinHandle<()> {
330 info!("booted source packet policy update task");
331 let updates = media_transport.source_policy_subscription();
332 tokio::spawn(async move {
333 let mut pending = BTreeSet::new();
334 let mut running = BTreeSet::new();
335 let mut batches = SelectAll::new();
336 loop {
337 tokio::select! {
338 biased;
339 () = shutdown_token.cancelled() => break,
340 completed = batches.next(), if !batches.is_empty() => {
341 if let Some(room_id) = completed {
342 running.remove(&room_id);
343 }
344 },
345 dirty_rooms = updates.wait_for_update() => pending.extend(dirty_rooms),
346 }
347 let ready = pending
348 .difference(&running)
349 .copied()
350 .collect::<BTreeSet<_>>();
351 if ready.is_empty() {
352 continue;
353 }
354 pending.retain(|room| !ready.contains(room));
355 running.extend(ready.iter().copied());
356 batches.push(rooms.source_policy_turns(&ready, &media_transport).await);
357 }
358 batches.for_each(|_| async {}).await;
361 })
362}
363
364fn spawn_room_reservation_expiration_reaper(
365 rooms: Arc<RoomManager>,
366 shutdown_token: CancellationToken,
367) -> JoinHandle<()> {
368 info!("booted room reservation expiration reaper");
369 tokio::spawn(async move {
370 let mut interval = interval(Duration::from_secs(10));
371 interval.set_missed_tick_behavior(MissedTickBehavior::Skip);
372 loop {
373 tokio::select! {
374 biased;
375 () = shutdown_token.cancelled() => return,
376 _ = interval.tick() => {},
377 };
378 rooms.check_expired_room_reservations().await;
379 }
380 })
381}
382
383fn build_media_transport(config: &Config, services: &RuntimeServices) -> AnyResult<MediaTransport> {
384 Ok(MediaTransport::build(
385 MediaTransportConfig {
386 worker_count: config.transport.rtc_media_worker_count,
387 announced_ip: config.transport.announced_ip,
388 bitrate_limits: SessionBitrateLimits::new(
389 config.transport.max_bitrate_in,
390 config.transport.max_bitrate_out,
391 ),
392 video_bitrate_limits: config.transport.video_bitrate_limits,
393 rtc_port_range: config.transport.rtc_port_range,
394 rtc_udp_io_backend: config.transport.rtc_udp_io_backend,
395 codec_flags: config.codecs.flags,
396 codec_preferences: config.codecs.preferences,
397 media_quality_interval: config.telemetry.media_quality_interval,
398 },
399 MediaTransportDeps {
400 packet_sink_registry: Arc::clone(&services.packet_sink_registry),
401 metrics: Arc::clone(&services.metrics),
402 },
403 )?)
404}
405
406fn build_room_runtime_policy(
407 config: &Config,
408 media_transport: &MediaTransport,
409) -> RoomRuntimePolicy {
410 RoomRuntimePolicy::new(
411 RoomAdmissionPolicy::new(config.user.room_size),
412 effective_feature_flags(config.features),
413 media_transport.router_rtp_capabilities(),
414 )
415 .with_room_worker_policy(config.transport.room_worker_policy)
416 .with_media_limits(config.transport.room_media_limits)
417 .with_video_adaptation_tuning(config.transport.video_adaptation_tuning)
418}
419
420fn build_room_manager(
421 runtime_policy: RoomRuntimePolicy,
422 services: &RuntimeServices,
423 reservation_ttl: Duration,
424 departure_grace: Duration,
425) -> Arc<RoomManager> {
426 Arc::new(RoomManager::new(
427 runtime_policy,
428 Arc::clone(&services.metrics),
429 reservation_ttl,
430 departure_grace,
431 ))
432}
433
434pub fn run() -> AnyResult<()> {
439 let config = Config::from_env()?;
440 let _telemetry = init_tracing(&config.telemetry, process::id())?;
441 let runtime = Runtime::new(&config)?;
442 Ok(Builder::new_multi_thread()
443 .enable_all()
444 .build()?
445 .block_on(runtime.serve(serve_http, shutdown_signal()))?)
446}
447
448#[cfg(test)]
449#[path = "TESTS/runtime.rs"]
450mod tests;