o_sfu/lib.rs
1//! Odoo's Selective Forwarding Unit (SFU) for audio/video calls.
2//!
3//! `o-sfu` is a multi-tenant SFU designed to be compatible with the Odoo saas/.sh
4//! architecture, which means that a single o-sfu server can serve
5//! independent (individually authenticated) "rooms" which can be used by many
6//! independent tenants (like odoo saas "databases").
7//!
8//! # Reading Map
9//!
10//! - **Start or embed the server**: [`run`], [`Runtime`] and [`config`].
11//! - **Understand admission**: [`auth`], [`http`] and [`websocket`].
12//! - **Change room and media behavior**: [`core::server::room::Room`],
13//! [`o_sfu_router::Router`] and [`core::server::transport::MediaTransport`].
14//! - **Integrate a browser client**: `SfuClient` exposes the Odoo-facing API.
15//! [`o_sfu_protocol::host::ProtocolCore`] supplies its signaling state machine.
16//!
17//!
18//! # Core Concepts
19//!
20//! `o-sfu` separates the **control plane** for admission and room policy, the
21//! **routing plane** for user-to-connection placement and the **packet plane**
22//! for RTP forwarding on worker loops.
23//!
24//! - **[`Runtime`]**: Owns the process lifecycle, the HTTP/WebSocket servers and graceful shutdown.
25//! - **[`core::server::room::Room`]**: The control plane boundary for a set of participants. It commits membership and media relationships.
26//! - **[`o_sfu_router::Router`]**: The routing plane, a pure, sans-I/O engine owning the placement graph that maps users to connections.
27//! - **[`core::prelude::MediaSession`]**: Orchestrates a user's connection, bridging room intent to transport effects.
28//! - **[`core::server::transport::MediaTransport`]**: Owns the media workers and abstracts their threading model. Each worker holds a packet loop and applies projected routes to incoming datagrams.
29//!
30//! # Architecture
31//!
32//! Control and routing decisions happen above the packet loops, which apply the
33//! resulting transport state to UDP datagrams.
34//!
35//! ```text
36//! HTTP control API ------------------------------> RoomManager
37//! |
38//! WebSocket session -> SfuCore -> MediaSession |
39//! \ /
40//! +----> Room <----+
41//! |
42//! +<----> Router
43//! |
44//! v
45//! MediaTransport
46//! | | |
47//! v v v
48//! RTC workers / packet loops
49//! (one thread per worker)
50//! | | |
51//! UDP socket UDP socket UDP socket
52//! | | |
53//! v v v
54//! fanout fanout fanout
55//! UDP OUT UDP OUT UDP OUT
56//! ```
57//!
58//! # Admission Edge
59//!
60//! Applications provision rooms through
61//! [`RoomManager::serve_room`](core::server::room::RoomManager::serve_room).
62//! WebSocket clients join through
63//! [`SfuCore::admit_user`](core::prelude::SfuCore::admit_user). Both paths require
64//! JWT authentication before admission.
65//!
66//! Upgraded sockets retain their connection permits. Before authentication,
67//! they have a first-frame deadline and global plus per-origin caps
68//! (one bucket per IPv4 address or IPv6 /64).
69//!
70//! ```text
71//! Incoming TCP connection
72//! |
73//! v
74//! connection cap + HTTP header deadline
75//! |
76//! v
77//! Axum Router
78//! |
79//! +-> GET /v1/noop
80//! | public liveness -> noop
81//! |
82//! +-> GET /v1/channel
83//! | Authorization JWT -> VerifiedRoomRequest -> room
84//! |
85//! +-> POST /v1/disconnect
86//! | request-body JWT -> VerifiedDisconnectClaims -> disconnect
87//! |
88//! +-> WebSocket /
89//! | upgrade -> first-frame JWT -> admit_user -> session
90//! |
91//! +-> GET /v1/stats, /metrics and /internal/diagnostics/...
92//! authorize_operator -> stats, metrics or diagnostics
93//! ```
94//!
95//! - **HTTP**: Parses server-to-server requests using [`http::CreateRoomQuery`]. Verifies [`auth::HttpRoomClaims`]. The request that creates the current room fixes its signing key.
96//! - **WebSocket**: Client connection frames are decoded by [`websocket::decode_auth_payload_text`]. A hint selects a candidate key. Claims are verified, normalized into [`auth::WebSocketConnectClaims`] and trusted for access.
97//!
98//! # Security Model
99//!
100//! `o-sfu` secures two planes independently. Application-layer JWTs gate room
101//! admission on the control plane. `str0m` encrypts media on the packet plane
102//! with DTLS-SRTP. An external reverse proxy provides TLS for HTTP and
103//! WebSocket signaling.
104//!
105//! ```text
106//! control plane JWT HS256 admission trust
107//! packet plane DTLS-SRTP (str0m) media confidentiality
108//! signaling wire TLS (via proxy) transport confidentiality
109//! ```
110//!
111//! ## JWT Admission
112//!
113//! Tokens are `HS256` only. [`auth::verify`] rejects any other `alg`, checks the
114//! HMAC in constant time and requires an unexpired `exp`. It validates `nbf`
115//! and the `iat` future-skew bound when present. Current Odoo does not issue
116//! `iat`. Tokens without `exp` fail with [`auth::AuthenticationError::MissingExpiry`].
117//! It caps token size at
118//! [`auth::MAX_JWT_TOKEN_BYTES`].
119//!
120//! There are two different keys:
121//!
122//! - **Server-to-server key**: `AUTH_KEY` (base64, at least 32 decoded bytes) verifies
123//! the HTTP [`http::CreateRoomQuery`] path through [`auth::HttpRoomClaims`] and
124//! [`auth::HttpDisconnectClaims`]. See [`config`].
125//! - **Per-room key**: the request that creates the current room pins the signing
126//! key from the `key` or `keySeed` claim in [`auth::HttpRoomClaims`]. A direct
127//! key must decode to at least 32 bytes or room creation returns `400`.
128//! `keySeed` instead derives the room's signing bytes from `AUTH_KEY`:
129//! ```text
130//! room_key = HMAC-SHA256(
131//! key = Base64Decode(AUTH_KEY),
132//! message = Base64Decode(keySeed)
133//! )
134//! ```
135//! WebSocket [`auth::WebSocketConnectClaims`] verify against that room key,
136//! never against `AUTH_KEY`. Runtime construction decodes the global key once.
137//! Room creation retains decoded secret bytes, so verification does not decode
138//! stored keys and equivalent base64 encodings match the same reservation.
139//!
140//! HTTP room creation uses the
141//! `Authorization` header, HTTP disconnect uses the request body and the
142//! WebSocket client sends a first-frame auth envelope decoded by
143//! [`websocket::decode_auth_payload_text`]. An unverified room id selects only a
144//! candidate key, then the same token is verified once against it. Modern
145//! [`auth::WebSocketConnectClaims`] must name the selected room. Legacy Odoo
146//! tokens select it through the auth envelope's `channel` and are normalized
147//! only after verification with that room's key. Odoo may send both
148//! `session_id` and account `user_id`. The session takes precedence as the
149//! participant identity. Tokens without either identity are rejected.
150//!
151//! Admission establishes identity and room scope. It does not enforce the
152//! per-user `permissions` claim provided by each tenant.
153//!
154//! ## Media Transport
155//!
156//! `str0m` handles ICE, DTLS and SRTP over UDP. `o-sfu` builds and drives the
157//! `str0m` sessions to forward RTP between participants.
158//!
159//! - **Keying**: the DTLS handshake derives SRTP keys per RFC 5764 DTLS-SRTP.
160//! - **Certificate**: `str0m` generates a self-signed certificate when each RTC
161//! session is built and advertises its SHA-256 fingerprint in the SDP offer.
162//! Accepting the answer stores the expected remote fingerprint. The DTLS
163//! handshake verifies it against the peer certificate.
164//! - **ICE**: `o-sfu` runs ICE-lite with `a=setup:actpass` and advertises
165//! `ANNOUNCED_IP`, so media UDP must reach the host directly.
166//!
167//! ## Signaling Transport
168//!
169//! HTTPS and WSS are expected to be terminated by an external reverse proxy.
170//! Setting `PROXY=true` in [`config`] requires `TRUSTED_PROXIES` CIDRs.
171//! Forwarded headers are honored only when the TCP peer belongs to that set.
172//! The public edge must overwrite client-supplied forwarded host and protocol
173//! headers. See [`http::resolve_request_origin`] for origin
174//! resolution and [`http`] for operator route access.
175//!
176//! # Room and Router Ownership
177//!
178//! Room transitions produce typed commits while holding short exclusive state
179//! locks. `RoomEffects` consumes deferred transport, source-policy and WebSocket
180//! output work after the lock is released.
181//!
182//! ```text
183//! room state lock held lock released
184//! +--------------------------------+ +------------------------------+
185//! | validate user and connection | | MediaTransport commands |
186//! | commit room topology | | source-policy turn |
187//! | capture transition commit |--->| websocket output |
188//! | | | idempotent teardown |
189//! +--------------------------------+ +------------------------------+
190//! ```
191//!
192//! [`o_sfu_router::Router`] owns exact user-to-connection placement. Receiver shadows are foreign local sessions derived from active consumer dependencies, disappearing with their final consumer.
193//!
194//! Absent subscription targets are capped per receiver. Eviction drops the
195//! oldest absent target's intent. Present members do not consume that allowance.
196//!
197//! # Signaling and Client Bundle
198//!
199//! Browsers use `SfuClient` for connection, publication, subscription and room
200//! control. Signaling state stays in [`o_sfu_protocol::host::ProtocolCore`] and
201//! yields ordered [`o_sfu_protocol::host::Command`] values. The WASM bridge
202//! serializes commands for `BrowserRuntime`, which executes effects through
203//! browser `WebSocket`, `RTCPeerConnection` and timer APIs. Protocol events are
204//! mapped to Odoo bundle updates during command serialization. `BrowserRuntime`
205//! applies those updates to client state and notifies the application.
206//!
207//! Outbound overflow closes with `4110` (`Overloaded`), allowing reconnect and
208//! intent replay. Replacement or explicit removal uses terminal `4108` (`Kicked`).
209//!
210//! ```text
211//! SfuClient (public API)
212//! |
213//! v
214//! BrowserRuntime
215//! |
216//! v
217//! ProtocolCore (sans-I/O) -> Vec<Command>
218//! |
219//! v
220//! WASM serialization -> TypeScript command union
221//! |
222//! v
223//! BrowserRuntime -> WebSocket, RTCPeerConnection, timers
224//! ^ |
225//! +--------------- browser events ---------+
226//! ```
227//!
228//! # Packet Path
229//!
230//! [`core::server::transport::MediaTransport`] owns the media workers, which hold the packet loops. These loops receive UDP datagrams, drive WebRTC state (`str0m`), apply route tables and forward RTP.
231//!
232//! ```text
233//! UDP datagram
234//! |
235//! v
236//! worker ingress and `str0m` drain
237//! |
238//! v
239//! packet facts (source, RID and codec)
240//! |
241//! +-> origin packet sinks
242//! |
243//! +-> source and destination packet gates
244//! |
245//! +-> relay fanout
246//! +-> local RTC -> RTP identity and codec rewrite
247//! ```
248//!
249//! Registered origin packet sinks observe publisher packets before route gates,
250//! including publishers without active receivers. Source, relay and receiver
251//! gates then narrow routed fanout. Same-process relays share payload data with
252//! another worker for local delivery.
253//!
254//! Worker BWE and audio observations feed into room source policy, which updates route gates for later packets.
255//!
256//! ```text
257//! worker BWE and audio observations
258//! |
259//! v
260//! room source policy
261//! |
262//! v
263//! route gates for later packets
264//! ```
265//!
266//! # Observability
267//!
268//! Monitored through the [`o_sfu_telemetry`] sub-crate. See [`http::telemetry`] for the HTTP contracts.
269//!
270//! - **Metrics**: [`http::telemetry::metrics`] exposes Prometheus text exposition.
271//! - **Diagnostics**: [`http::telemetry::diagnostics`] exposes JSON state summaries.
272//!
273//! # Scaling
274//!
275//! Rooms use one [`o_sfu_router::Router`] facade and default to one local router.
276//! [`config::RoomWorkerPolicy`] can enable additional same-process local routers.
277//! Only running workers are eligible. Joins prefer assigned workers with a known
278//! delay below the configured threshold. When none qualifies, a join may attach
279//! an unused healthy worker within the router cap and worker count. Otherwise,
280//! joins reuse a running assigned worker even when it exceeds the delay threshold.
281//! Placements on failed workers do not count toward the cap. If no assigned worker
282//! is running, a later join can attach a fresh placement on another running worker
283//! even under the single-router policy. Admission fails when no worker is running.
284//!
285//! # Feature Flags
286//!
287//! Core media behavior is configured at runtime through [`config`], not Cargo features.
288//! The default feature `otel-tracing` enables OpenTelemetry tracing support through [`o_sfu_telemetry::TraceExportConfig`].
289//! Other features are only used for tests and benchmarking.
290//!
291//! # Sub-crates
292//!
293//! | Crate | Role |
294//! | --- | --- |
295//! | [`o_sfu_rfc`] | RFC-backed JWT, RTP, RTCP, SDP and WebRTC consts/types |
296//! | `o-sfu-model` | Shared call data ([`o_sfu_protocol::wire::UserId`], etc.) |
297//! | [`o_sfu_router`] | Sans-I/O [`o_sfu_router::Router`] facade for room placement and routed media lifetimes |
298//! | [`o_sfu_core`] | Room engine, [`core::prelude::SourcePolicy`], recording taps and [`core::server::transport::MediaTransport`] projection |
299//! | [`o_sfu_protocol`] | Sans-I/O [`o_sfu_protocol::host::ProtocolCore`] and typed commands |
300//! | [`o_sfu_telemetry`] | Tracing setup, metrics and diagnostics response types |
301pub mod config;
302pub mod core {
303 pub use o_sfu_core::{prelude, server};
304}
305pub(crate) mod application;
306mod runtime;
307
308pub mod auth {
309 pub use crate::runtime::auth::{
310 AuthenticationError, HttpDisconnectClaims, HttpRoomClaims, MAX_JWT_TOKEN_BYTES,
311 RegisteredJwtClaims, WebSocketConnectClaims, sign, verify,
312 };
313}
314
315#[cfg(any(test, feature = "testing"))]
316pub mod test_support {
317 pub use crate::runtime::auth::test_support::TestHttpRoomClaims;
318}
319
320/// HTTP route and payload contracts.
321///
322/// `/v1/stats`, `/metrics` and diagnostics require the configured
323/// [`crate::config::DiagnosticsConfig::auth_token`] on every listener. Without
324/// one, the actual listener must be loopback. Missing or invalid tokens return
325/// `401 Unauthorized` with `WWW-Authenticate: Bearer realm="o-sfu"`.
326/// Tokenless non-loopback access returns `403 Forbidden`.
327pub mod http {
328 pub use crate::runtime::{
329 http_server::contract::{
330 CreateRoomQuery, IncomingBitRateStatsResponse, NoopResponse, RoomResponse,
331 StatsResponse, route,
332 },
333 request_origin::{RequestOrigin, resolve_request_origin},
334 };
335
336 /// Operator-facing metrics and diagnostics contracts.
337 pub mod telemetry {
338 /// Prometheus metric scrape contract.
339 ///
340 /// `GET` [`metrics::PATH`] returns `200 OK` Prometheus text exposition
341 /// with [`metrics::CONTENT_TYPE`].
342 ///
343 /// This is a scrape endpoint.
344 /// Configure Prometheus to scrape [`metrics::PATH`], then issue `PromQL`
345 /// queries to Prometheus.
346 /// Histogram families render `<name>_bucket` with the additional `le`
347 /// label plus `<name>_sum` and `<name>_count`.
348 ///
349 /// # Scrape and Query
350 ///
351 /// ```yaml
352 /// scrape_configs:
353 /// - job_name: o-sfu
354 /// scheme: https
355 /// metrics_path: /metrics
356 /// authorization:
357 /// type: Bearer
358 /// credentials_file: /run/secrets/o_sfu_diagnostics_token
359 /// tls_config:
360 /// ca_file: /run/secrets/o_sfu_observability_ca
361 /// server_name: o-sfu-observability.internal
362 /// static_configs:
363 /// - targets: ["o-sfu-observability.internal:443"]
364 /// ```
365 ///
366 /// The endpoint returns Prometheus text exposition.
367 ///
368 /// ```text
369 /// # HELP osfu_rooms_active Current number of active rooms owned by this runtime.
370 /// # TYPE osfu_rooms_active gauge
371 /// osfu_rooms_active 3
372 /// # TYPE osfu_worker_rtp_packets_total counter
373 /// osfu_worker_rtp_packets_total{media_worker_id="0",direction="ingress"} 1240
374 /// ```
375 ///
376 /// Query clients send `PromQL` to the Prometheus-compatible backend
377 /// rather than to o-sfu.
378 /// These examples cover a gauge, counter and histogram.
379 ///
380 /// ```promql
381 /// sum(osfu_users_active)
382 /// sum by (stage) (rate(osfu_ws_connections_total[5m]))
383 /// histogram_quantile(
384 /// 0.95,
385 /// sum by (le, route) (rate(osfu_http_request_duration_seconds_bucket[10m]))
386 /// )
387 /// ```
388 ///
389 /// Read the [`metrics::MetricName`] variants below to find every
390 /// exported name and its meaning.
391 /// Query [`metrics::PATH`] to see each family's `HELP`, `TYPE`, label
392 /// keys and current label values before building selectors.
393 pub mod metrics {
394 pub use o_sfu_telemetry::{
395 metrics::MetricName, prometheus::PROMETHEUS_CONTENT_TYPE as CONTENT_TYPE,
396 };
397
398 pub use crate::http::route::METRICS as PATH;
399 }
400
401 /// JSON diagnostics contract.
402 ///
403 /// Every constant in [`diagnostics::route`] is a `GET` endpoint returning
404 /// `200 OK` JSON on success.
405 ///
406 /// # Routes and Parameters
407 ///
408 /// | request | JSON response | parameter source |
409 /// | --- | --- | --- |
410 /// | `GET /internal/diagnostics/summary` | one [`diagnostics::DiagnosticsSummaryResponse`] | none |
411 /// | `GET /internal/diagnostics/rooms` | array of [`diagnostics::DiagnosticsRoomSummary`] | none |
412 /// | `GET /internal/diagnostics/workers` | array of [`diagnostics::DiagnosticsWorkerSummary`] | none |
413 /// | `GET /internal/diagnostics/rooms/{uuid}` | one [`diagnostics::DiagnosticsRoomDetail`] | `uuid` from the rooms response |
414 /// | `GET /internal/diagnostics/rooms/{uuid}/users` | array of [`diagnostics::DiagnosticsUserSummary`] | `uuid` from the rooms response |
415 /// | `GET /internal/diagnostics/rooms/{uuid}/users/{id}` | one [`diagnostics::DiagnosticsUserDetail`] | `uuid` from rooms and `userKey` from room users |
416 ///
417 /// `userId` may be a JSON number or string.
418 /// `userKey` is always the string to put into `{id}`.
419 /// URL-encode both path values before substitution.
420 ///
421 /// # Summary Request and Response over HTTPS
422 ///
423 /// ```text
424 /// GET /internal/diagnostics/summary HTTP/1.1
425 /// Host: o-sfu-observability.internal
426 /// Authorization: Bearer <diagnostics-token>
427 /// Accept: application/json
428 ///
429 /// HTTP/1.1 200 OK
430 /// Content-Type: application/json
431 ///
432 /// {
433 /// "roomsActive": 1,
434 /// "publicationsActive": 1,
435 /// "recordingRoomsActive": 0,
436 /// "usersActive": 2,
437 /// "subscriptionsActive": 1,
438 /// "transport": {
439 /// "connectedUsers": 2,
440 /// "disconnectedUsers": 0,
441 /// "totalUsers": 2,
442 /// "unknownUsers": 0
443 /// }
444 /// }
445 /// ```
446 ///
447 /// # JavaScript Fetch Example
448 ///
449 /// ```javascript
450 /// const origin = "https://o-sfu-observability.internal";
451 /// const headers = {
452 /// Authorization: `Bearer ${process.env.DIAGNOSTICS_AUTH_TOKEN}`,
453 /// };
454 ///
455 /// async function getJson(path) {
456 /// const response = await fetch(`${origin}${path}`, { headers });
457 /// if (!response.ok) {
458 /// throw new Error(`${response.status} ${await response.text()}`);
459 /// }
460 /// return response.json();
461 /// }
462 ///
463 /// async function main() {
464 /// const rooms = await getJson("/internal/diagnostics/rooms");
465 /// const roomUuid = encodeURIComponent(rooms[0].uuid);
466 /// const room = await getJson(`/internal/diagnostics/rooms/${roomUuid}`);
467 /// const users = await getJson(`/internal/diagnostics/rooms/${roomUuid}/users`);
468 /// const userKey = encodeURIComponent(users[0].userKey);
469 /// const user = await getJson(
470 /// `/internal/diagnostics/rooms/${roomUuid}/users/${userKey}`,
471 /// );
472 ///
473 /// console.log(room.summary, room.users, room.sources);
474 /// console.log(user.user, user.recordingState);
475 /// }
476 ///
477 /// main().catch((error) => {
478 /// console.error(error);
479 /// process.exitCode = 1;
480 /// });
481 /// ```
482 ///
483 /// The rooms response has this shape.
484 ///
485 /// ```json
486 /// [
487 /// {
488 /// "createDate": "2026-07-15T10:20:30.000Z",
489 /// "mediaWorkerId": 0,
490 /// "publicationCount": 1,
491 /// "recordingState": {
492 /// "recording": false,
493 /// "audio": false,
494 /// "transcription": false,
495 /// "video": false
496 /// },
497 /// "remoteAddress": "203.0.113.10",
498 /// "sourceCount": 1,
499 /// "userCount": 2,
500 /// "subscriptionCount": 1,
501 /// "transport": {
502 /// "connectedUsers": 2,
503 /// "disconnectedUsers": 0,
504 /// "totalUsers": 2,
505 /// "unknownUsers": 0
506 /// },
507 /// "uuid": "550e8400-e29b-41d4-a716-446655440000",
508 /// "webRtcEnabled": true
509 /// }
510 /// ]
511 /// ```
512 ///
513 /// The room users response has this shape.
514 ///
515 /// ```json
516 /// [
517 /// {
518 /// "audioIncomingBitrateBps": 32000,
519 /// "cameraIncomingBitrateBps": 600000,
520 /// "connectionId": 91,
521 /// "health": "connected",
522 /// "incomingBitrateBps": 632000,
523 /// "mediaWorkerId": 0,
524 /// "publicationCount": 2,
525 /// "roomId": "550e8400-e29b-41d4-a716-446655440000",
526 /// "screenIncomingBitrateBps": 0,
527 /// "subscriptionCount": 1,
528 /// "userId": 42,
529 /// "userKey": "42"
530 /// }
531 /// ]
532 /// ```
533 ///
534 /// The response structs below list every field in each payload.
535 /// Wire names are `camelCase` unless a field documents an exception.
536 ///
537 /// User detail is room-scoped because the same user key can be active
538 /// in several rooms.
539 pub mod diagnostics {
540 pub use o_sfu_telemetry::diagnostics::{
541 DiagnosticsRoomDetail, DiagnosticsRoomSummary, DiagnosticsSummaryResponse,
542 DiagnosticsUserDetail, DiagnosticsUserSummary, DiagnosticsWorkerSummary,
543 };
544
545 pub use crate::http::route::diagnostics as route;
546 }
547 }
548}
549
550pub mod websocket {
551 pub use crate::runtime::websocket_server::{
552 ClientBatchDecodeError, ClientBatchDecodeFailureKind, MAX_CLIENT_BATCH_ENVELOPES,
553 MAX_CLIENT_FRAME_BYTES, decode_auth_payload_text, decode_client_batch,
554 };
555}
556
557pub use self::runtime::{Runtime, ServeError, run};