Skip to main content

o_sfu/runtime/diagnostics/
queries.rs

1//! Endpoint-specific diagnostics queries over passive room captures and RTC observations.
2
3use std::{collections::BTreeMap, slice};
4
5use o_sfu_core::{
6    MediaWorkerId,
7    server::{
8        room::{RoomOverviewCapture, RuntimeRoomDirectorySnapshot},
9        transport::{
10            MediaTransport, TransportAdapterError, TransportHealthSnapshot, TransportSessionKey,
11        },
12    },
13};
14use o_sfu_telemetry::diagnostics::{
15    DiagnosticsRoomDetail, DiagnosticsRoomSummary, DiagnosticsSummaryResponse,
16    DiagnosticsUserDetail, DiagnosticsUserSummary, DiagnosticsWorkerPressure,
17    DiagnosticsWorkerSummary,
18};
19
20use crate::{application::stream_catalog::DiscussStream, runtime::room::RoomManager};
21
22pub(crate) async fn summary_response(
23    rooms: &RoomManager,
24    transport: &MediaTransport,
25) -> DiagnosticsSummaryResponse {
26    let captures = overview_captures(rooms).await;
27    let health = transport.transport_health_snapshot(&overview_session_keys(&captures));
28    let mut response = DiagnosticsSummaryResponse {
29        rooms_active: captures.len(),
30        ..Default::default()
31    };
32    for (_, capture) in captures {
33        let counts = capture.media_counts;
34        let transport_counts = capture.transport_counts(&health);
35        response.users_active = response
36            .users_active
37            .saturating_add(capture.session_keys.len());
38        response.publications_active = response
39            .publications_active
40            .saturating_add(counts.publications);
41        response.subscriptions_active = response
42            .subscriptions_active
43            .saturating_add(counts.subscriptions);
44        response.recording_rooms_active = response
45            .recording_rooms_active
46            .saturating_add(usize::from(capture.recording_state.recording == Some(true)));
47        response.transport.connected = response
48            .transport
49            .connected
50            .saturating_add(transport_counts.connected);
51        response.transport.disconnected = response
52            .transport
53            .disconnected
54            .saturating_add(transport_counts.disconnected);
55        response.transport.unknown = response
56            .transport
57            .unknown
58            .saturating_add(transport_counts.unknown);
59    }
60    response.transport.total = response.users_active;
61    response
62}
63
64pub(crate) async fn rooms_response(
65    rooms: &RoomManager,
66    transport: &MediaTransport,
67) -> Vec<DiagnosticsRoomSummary> {
68    let captures = overview_captures(rooms).await;
69    let health = transport.transport_health_snapshot(&overview_session_keys(&captures));
70    captures
71        .into_iter()
72        .map(|captured| room_summary(captured, &health))
73        .collect()
74}
75
76/// Captures a room detail without presenting unavailable worker facts as empty.
77///
78/// # Errors
79///
80/// Returns [`TransportAdapterError::TransportUnavailable`] when source
81/// diagnostics cannot be observed from every assigned worker.
82pub(crate) async fn room_detail_response(
83    rooms: &RoomManager,
84    transport: &MediaTransport,
85    room_id: &str,
86) -> Result<Option<DiagnosticsRoomDetail>, TransportAdapterError> {
87    let Some(entry) = rooms.directory_snapshot(room_id).await else {
88        return Ok(None);
89    };
90    let capture = entry.room.diagnostics_detail_capture().await;
91    let session_keys = capture.session_keys();
92    let source_keys = capture.source_keys().cloned().collect::<Vec<_>>();
93    let bitrate = transport.transport_bitrate_snapshot(session_keys);
94    let quality = transport.transport_quality_snapshot(session_keys);
95    let health = transport.transport_health_snapshot(session_keys);
96    let source_diagnostics = transport.source_diagnostics_snapshot(&source_keys).await?;
97    let (overview, users, sources) =
98        capture.into_views(&bitrate, &quality, &health, &source_diagnostics);
99    Ok(Some(DiagnosticsRoomDetail {
100        summary: room_summary((entry, overview), &health),
101        sources,
102        users,
103    }))
104}
105
106pub(crate) async fn room_users_response(
107    rooms: &RoomManager,
108    transport: &MediaTransport,
109    room_id: &str,
110) -> Option<Vec<DiagnosticsUserSummary>> {
111    let room = rooms.get_by_uuid(room_id).await?;
112    let capture = room.diagnostics_users_capture().await;
113    let session_keys = capture.session_keys().cloned().collect::<Vec<_>>();
114    let bitrate = transport.transport_bitrate_snapshot(&session_keys);
115    let health = transport.transport_health_snapshot(&session_keys);
116    Some(capture.into_user_summaries(
117        room_id,
118        &bitrate,
119        &health,
120        DiscussStream::all().map(DiscussStream::label),
121    ))
122}
123
124pub(crate) async fn workers_response(
125    rooms: &RoomManager,
126    transport: &MediaTransport,
127) -> Vec<DiagnosticsWorkerSummary> {
128    let mut captures = Vec::new();
129    for entry in rooms.directory_snapshots().await {
130        captures.push(entry.room.diagnostics_users_capture().await);
131    }
132    let session_keys = captures
133        .iter()
134        .flat_map(|capture| capture.session_keys().cloned())
135        .collect::<Vec<_>>();
136    let health = transport.transport_health_snapshot(&session_keys);
137    let mut workers = BTreeMap::new();
138    for snapshot in transport.worker_pressure_snapshots() {
139        let id = snapshot.media_worker_id.as_usize();
140        workers.insert(
141            id,
142            DiagnosticsWorkerSummary {
143                media_worker_id: id,
144                pressure: DiagnosticsWorkerPressure {
145                    command_backlog_depth: snapshot.command_backlog_depth,
146                    egress_bitrate_bps: snapshot.egress_bitrate.as_bps(),
147                    packet_loop_delay_ms: snapshot.packet_loop_delay_ms,
148                    relay_mailbox_depth: snapshot.relay_mailbox_depth,
149                    worker_pressure_score: snapshot.worker_pressure_score,
150                },
151                ..Default::default()
152            },
153        );
154    }
155    for capture in captures {
156        capture.add_to_worker_summaries(&health, &mut workers);
157    }
158    workers.into_values().collect()
159}
160
161pub(crate) async fn user_detail_response(
162    rooms: &RoomManager,
163    transport: &MediaTransport,
164    room_id: &str,
165    user_key: &str,
166) -> Option<DiagnosticsUserDetail> {
167    let room = rooms.get_by_uuid(room_id).await?;
168    let capture = room.diagnostics_user_capture(user_key).await?;
169    let session_keys = slice::from_ref(capture.session_key());
170    let bitrate = transport.transport_bitrate_snapshot(session_keys);
171    let quality = transport.transport_quality_snapshot(session_keys);
172    let health = transport.transport_health_snapshot(session_keys);
173    let (recording_state, user) = capture.into_view(&bitrate, &quality, &health);
174    Some(DiagnosticsUserDetail {
175        room_id: room.uuid().to_owned(),
176        recording_state,
177        user,
178    })
179}
180
181async fn overview_captures(
182    rooms: &RoomManager,
183) -> Vec<(RuntimeRoomDirectorySnapshot, RoomOverviewCapture)> {
184    let entries = rooms.directory_snapshots().await;
185    let mut captures = Vec::with_capacity(entries.len());
186    for entry in entries {
187        let capture = entry.room.diagnostics_overview_capture().await;
188        captures.push((entry, capture));
189    }
190    captures
191}
192
193fn overview_session_keys(
194    captures: &[(RuntimeRoomDirectorySnapshot, RoomOverviewCapture)],
195) -> Vec<TransportSessionKey> {
196    captures
197        .iter()
198        .flat_map(|(_, capture)| capture.session_keys.iter().cloned())
199        .collect()
200}
201
202fn room_summary(
203    (entry, capture): (RuntimeRoomDirectorySnapshot, RoomOverviewCapture),
204    health: &TransportHealthSnapshot,
205) -> DiagnosticsRoomSummary {
206    let counts = capture.media_counts;
207    let transport = capture.transport_counts(health);
208    DiagnosticsRoomSummary {
209        create_date: entry.create_date,
210        media_worker_id: capture
211            .primary_media_worker_id
212            .map_or(0, MediaWorkerId::as_usize),
213        publication_count: counts.publications,
214        recording_state: capture.recording_state,
215        remote_address: entry.remote_address,
216        source_count: counts.publications,
217        user_count: capture.session_keys.len(),
218        subscription_count: counts.subscriptions,
219        transport,
220        uuid: entry.room.uuid().to_owned(),
221        web_rtc_enabled: entry.room.web_rtc_enabled(),
222    }
223}