Skip to main content

o_sfu/runtime/
mod.rs

1//! Wires process services and drains them during shutdown.
2//!
3//! Request handlers receive [`RuntimeState`] without boot or teardown control.
4
5use 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/// Failure to serve or fully drain a [`Runtime`].
54#[derive(Debug, Error)]
55pub enum ServeError {
56    /// The listener or process shutdown signal failed.
57    #[error(transparent)]
58    Io(#[from] io::Error),
59    /// The deadline elapsed before runtime drainage finished.
60    #[error(
61        "runtime shutdown exceeded its deadline with {remaining_sessions} WebSocket sessions remaining"
62    )]
63    ShutdownIncomplete {
64        /// Tracked WebSocket sessions whose finalizers had not returned.
65        remaining_sessions: usize,
66    },
67}
68
69/// Process services and lifecycle configuration.
70#[derive(Debug)]
71pub struct Runtime {
72    config: RuntimeConfig,
73    room_manager: Arc<RoomManager>,
74    metrics: Arc<RuntimeMetrics>,
75    media_transport: MediaTransport,
76}
77
78/// Cloneable request dependencies without process lifecycle control.
79#[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    /// Builds the room manager and media workers from loaded configuration.
99    ///
100    /// # Errors
101    ///
102    /// Returns [`anyhow::Error`] when the authentication key is malformed or
103    /// shorter than the HS256 minimum, the diagnostics token is invalid, proxy
104    /// mode lacks trusted proxies, listener limits are outside their supported
105    /// ranges or the media transport cannot be built.
106    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    /// Serves a caller-provided listener until `shutdown` resolves.
143    ///
144    /// Tokenless operator access follows the listener's actual local address.
145    ///
146    /// # Errors
147    ///
148    /// Returns [`ServeError::Io`] for serving failures or
149    /// [`ServeError::ShutdownIncomplete`] when the drainage deadline expires.
150    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
220/// Cancels process work when the server future is dropped.
221struct 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
324/// Recomputes room packet policy after transport observations change.
325fn 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        // A policy turn can have accepted transport effects. Finish those turns
359        // before shutdown releases their room state.
360        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
434/// # Errors
435///
436/// Returns an error when configuration, tracing, Tokio startup, serving or
437/// shutdown fails.
438pub 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;