Skip to main content

o_sfu/runtime/http_server/rooms/
disconnect.rs

1//! Bulk disconnect with bounded body JWTs and normalized runtime user IDs.
2
3use std::{str, sync::Arc};
4
5use axum::{
6    body::Bytes,
7    extract::{DefaultBodyLimit, FromRef, FromRequest, Request, State},
8    http::StatusCode,
9    response::{IntoResponse, Response},
10    routing::{MethodRouter, post},
11};
12use secrecy::SecretString;
13use tracing::Instrument;
14
15use crate::runtime::{
16    MediaTransport, RuntimeMetrics, RuntimeState,
17    auth::{self, HttpDisconnectClaims},
18    room::RoomManager,
19    telemetry,
20};
21
22const MAX_BODY_BYTES: usize = 16 * 1024;
23
24pub(super) fn route() -> MethodRouter<RuntimeState> {
25    post(disconnect).layer(DefaultBodyLimit::max(MAX_BODY_BYTES))
26}
27
28#[derive(Debug, Clone)]
29struct Services {
30    room_manager: Arc<RoomManager>,
31    media_transport: MediaTransport,
32    metrics: Arc<RuntimeMetrics>,
33}
34
35impl FromRef<RuntimeState> for Services {
36    fn from_ref(state: &RuntimeState) -> Self {
37        Self {
38            room_manager: Arc::clone(&state.room_manager),
39            media_transport: state.media_transport.clone(),
40            metrics: Arc::clone(&state.metrics),
41        }
42    }
43}
44
45/// Body JWT verified with the global auth key and runtime-normalized user IDs.
46///
47/// # Errors
48///
49/// Extraction returns a [`Response`]: `400 Bad Request` for non-UTF-8 bodies
50/// or `422 Unprocessable Entity` when JWT verification fails. Both record the
51/// corresponding disconnect rejection counter. Body buffering failures retain
52/// Axum's rejection response, including `413 Payload Too Large` for the route's
53/// body limit, without incrementing those counters.
54#[derive(Debug, Clone, PartialEq, Eq)]
55struct VerifiedDisconnectClaims(HttpDisconnectClaims);
56
57/// Bulk-disconnect endpoint used by Odoo to remove users from active rooms.
58///
59/// [`VerifiedDisconnectClaims`] owns JWT verification and request-body decoding.
60async fn disconnect(
61    State(services): State<Services>,
62    VerifiedDisconnectClaims(claims): VerifiedDisconnectClaims,
63) -> Response {
64    async {
65        for (room_id, user_ids) in &claims.user_ids_by_room {
66            services
67                .room_manager
68                .disconnect_users(room_id, user_ids, &services.media_transport)
69                .await;
70        }
71        services.metrics.record_http_disconnect_success();
72        StatusCode::OK.into_response()
73    }
74    .instrument(telemetry::http_request_span("disconnect"))
75    .await
76}
77
78impl FromRequest<RuntimeState> for VerifiedDisconnectClaims {
79    type Rejection = Response;
80
81    async fn from_request(req: Request, state: &RuntimeState) -> Result<Self, Self::Rejection> {
82        let body = Bytes::from_request(req, state)
83            .await
84            .map_err(IntoResponse::into_response)?;
85        let token = str::from_utf8(&body)
86            .map_err(|_error| record_rejection(state, StatusCode::BAD_REQUEST))?;
87        let token = SecretString::from(token);
88        let mut claims =
89            auth::verify_claims::<HttpDisconnectClaims>(&token, &state.config.auth.key)
90                .map_err(|_error| record_rejection(state, StatusCode::UNPROCESSABLE_ENTITY))?;
91        claims.normalize_runtime_user_ids();
92        Ok(Self(claims))
93    }
94}
95
96fn record_rejection(state: &RuntimeState, status: StatusCode) -> Response {
97    match status {
98        StatusCode::BAD_REQUEST => state.metrics.record_http_disconnect_bad_request(),
99        StatusCode::UNPROCESSABLE_ENTITY => {
100            state.metrics.record_http_disconnect_unprocessable_entity();
101        }
102        _ => {}
103    }
104    status.into_response()
105}