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