o_sfu/runtime/diagnostics/
queries.rs1use 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}