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::{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/// Failure to serve or fully drain a [`Runtime`].
53#[derive(Debug, Error)]
54pub enum ServeError {
55    /// The listener or process shutdown signal failed.
56    #[error(transparent)]
57    Io(#[from] io::Error),
58    /// The deadline elapsed before runtime drainage finished.
59    #[error(
60        "runtime shutdown exceeded its deadline with {remaining_sessions} WebSocket sessions remaining"
61    )]
62    ShutdownIncomplete {
63        /// Tracked WebSocket sessions whose finalizers had not returned.
64        remaining_sessions: usize,
65    },
66}
67
68/// Process services and lifecycle configuration.
69#[derive(Debug)]
70pub struct Runtime {
71    config: RuntimeConfig,
72    room_manager: Arc<RoomManager>,
73    metrics: Arc<RuntimeMetrics>,
74    media_transport: MediaTransport,
75}
76
77/// Cloneable request dependencies without process lifecycle control.
78#[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    /// Builds the room manager and media workers from loaded configuration.
98    ///
99    /// # Errors
100    ///
101    /// Returns an error when the configured media transport cannot be built.
102    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    /// Serves a caller-provided listener until `shutdown` resolves.
136    ///
137    /// Tokenless operator access follows the listener's actual local address.
138    ///
139    /// # Errors
140    ///
141    /// Returns [`ServeError::Io`] for serving failures or
142    /// [`ServeError::ShutdownIncomplete`] when the drainage deadline expires.
143    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
213/// Cancels process work when the server future is dropped.
214struct 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
317/// Recomputes room packet policy after transport observations change.
318fn 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
411/// # Errors
412///
413/// Returns an error when configuration, tracing, Tokio startup, serving or
414/// shutdown fails.
415pub 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;