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