Skip to main content

o_sfu_telemetry/graph/
user.rs

1//! User-centered Grafana node graph formatting.
2//!
3//! The graph is centered on one selected user while still showing
4//! the packet path through local media-worker ownership:
5//!
6//! ```text
7//! Outbound from selected user:
8//!
9//!   selected user
10//!        |
11//!        v
12//!   selected user's media worker
13//!        |
14//!        v
15//!   published source
16//!        |
17//!        v
18//!   receiver media worker
19//!        |
20//!        v
21//!   peer user
22//!
23//! Inbound to selected user:
24//!
25//!   peer user
26//!        |
27//!        v
28//!   peer media worker
29//!        |
30//!        v
31//!   published source
32//!        |
33//!        v
34//!   selected user's media worker
35//!        |
36//!        v
37//!   selected user
38//! ```
39
40use std::collections::HashSet;
41
42use serde_json::{Value, json};
43
44use super::common::{
45    download_main_stat, push_unique_edge, push_unique_node, route_state_color, route_state_label,
46    source_by_id, stream_id_color, stream_id_label, transport_health_label, user_by_id,
47};
48use crate::diagnostics::types::{
49    DiagnosticsRoomDetail, DiagnosticsSource, DiagnosticsSubscription, DiagnosticsUserView,
50};
51
52fn path_user_node(room_uuid: &str, user: &DiagnosticsUserView, selected: bool) -> Value {
53    let session_id = user.user_id.path_segment();
54    let title = if selected {
55        format!("{session_id} selected")
56    } else {
57        session_id.to_string()
58    };
59    json!({
60        "id": format!("user:{}:{}", room_uuid, session_id),
61        "title": title,
62        "subtitle": transport_health_label(user.transport.health.as_ref()),
63        "mainStat": format!("{} bps", user.transport.quality_summary.current_incoming_bitrate.total),
64        "secondaryStat": format!("{} pub / {} sub", user.publications.len(), user.subscriptions.len()),
65        "detail__connection": user.transport.connection_id.to_string(),
66        "detail__worker": user.transport.media_worker_id,
67        "color": if selected { "blue" } else { "gray" },
68    })
69}
70
71fn path_worker_node(detail: &DiagnosticsRoomDetail, media_worker_id: usize) -> Value {
72    let users_on_worker = detail
73        .users
74        .iter()
75        .filter(|user| user.transport.media_worker_id == media_worker_id);
76    let mut user_count = 0_usize;
77    let mut publication_count = 0_usize;
78    let mut subscription_count = 0_usize;
79    for user in users_on_worker {
80        user_count = user_count.saturating_add(1);
81        publication_count = publication_count.saturating_add(user.publications.len());
82        subscription_count = subscription_count.saturating_add(user.subscriptions.len());
83    }
84
85    json!({
86        "id": format!("worker:{}", media_worker_id),
87        "title": format!("worker {}", media_worker_id),
88        "subtitle": "media worker",
89        "mainStat": format!("{} users", user_count),
90        "secondaryStat": format!("{} pub / {} sub", publication_count, subscription_count),
91    })
92}
93
94fn path_source_node(room_uuid: &str, source: &DiagnosticsSource) -> Value {
95    let stream_id = stream_id_label(&source.stream_id);
96    json!({
97        "id": format!("source:{}:{}", room_uuid, source.source_id),
98        "title": format!("{} #{}", stream_id, source.source_id),
99        "subtitle": format!("{:?}", source.media_kind).to_lowercase(),
100        "mainStat": format!("{} bps", source.current_incoming_bitrate_bps),
101        "secondaryStat": format!("{} encodings", source.encodings.len()),
102        "detail__owner_user": source.owner_user_id.path_segment(),
103        "detail__stream_id": stream_id,
104        "detail__transport_media_id": source.transport_media_id,
105    })
106}
107
108fn ensure_path_user(
109    nodes: &mut Vec<Value>,
110    seen_nodes: &mut HashSet<String>,
111    detail: &DiagnosticsRoomDetail,
112    user: &DiagnosticsUserView,
113    selected: bool,
114) {
115    let room_uuid = detail.summary.uuid.as_str();
116    let user_id = user.user_id.path_segment();
117    push_unique_node(
118        nodes,
119        seen_nodes,
120        format!("user:{room_uuid}:{user_id}"),
121        path_user_node(room_uuid, user, selected),
122    );
123    push_unique_node(
124        nodes,
125        seen_nodes,
126        format!("worker:{}", user.transport.media_worker_id),
127        path_worker_node(detail, user.transport.media_worker_id),
128    );
129}
130
131fn push_user_transport_edge(
132    edges: &mut Vec<Value>,
133    seen_edges: &mut HashSet<String>,
134    detail: &DiagnosticsRoomDetail,
135    user: &DiagnosticsUserView,
136    direction: &str,
137) {
138    let room_uuid = detail.summary.uuid.as_str();
139    let user_id = user.user_id.path_segment();
140    let edge_id = format!(
141        "transport:{room_uuid}:{user_id}:{}",
142        user.transport.media_worker_id
143    );
144    push_unique_edge(
145        edges,
146        seen_edges,
147        edge_id.clone(),
148        json!({
149            "id": edge_id,
150            "source": format!("user:{}:{}", room_uuid, user_id),
151            "target": format!("worker:{}", user.transport.media_worker_id),
152            "mainStat": "transport",
153            "secondaryStat": transport_health_label(user.transport.health.as_ref()),
154            "detail__direction": direction,
155            "detail__connection": user.transport.connection_id,
156        }),
157    );
158}
159
160fn push_publish_path(
161    nodes: &mut Vec<Value>,
162    edges: &mut Vec<Value>,
163    seen_nodes: &mut HashSet<String>,
164    seen_edges: &mut HashSet<String>,
165    detail: &DiagnosticsRoomDetail,
166    owner: &DiagnosticsUserView,
167    source: &DiagnosticsSource,
168) {
169    let room_uuid = detail.summary.uuid.as_str();
170    push_unique_node(
171        nodes,
172        seen_nodes,
173        format!("source:{room_uuid}:{}", source.source_id),
174        path_source_node(room_uuid, source),
175    );
176    let edge_id = format!("publish:{room_uuid}:{}", source.source_id);
177    push_unique_edge(
178        edges,
179        seen_edges,
180        edge_id.clone(),
181        json!({
182            "id": edge_id,
183            "source": format!("worker:{}", owner.transport.media_worker_id),
184            "target": format!("source:{}:{}", room_uuid, source.source_id),
185            "mainStat": format!("{} upload", stream_id_label(&source.stream_id)),
186            "secondaryStat": format!("{} bps", source.current_incoming_bitrate_bps),
187            "color": stream_id_color(&source.stream_id),
188            "detail__owner_user": owner.user_id.path_segment(),
189            "detail__stream_id": &source.stream_id,
190            "detail__transport_media_id": source.transport_media_id,
191        }),
192    );
193}
194
195fn push_subscription_delivery_path(
196    edges: &mut Vec<Value>,
197    seen_edges: &mut HashSet<String>,
198    detail: &DiagnosticsRoomDetail,
199    receiver: &DiagnosticsUserView,
200    sub: &DiagnosticsSubscription,
201    direction: &str,
202) {
203    let room_uuid = detail.summary.uuid.as_str();
204    let receiver_id = receiver.user_id.path_segment();
205    let deliver_edge_id = format!("deliver:{room_uuid}:{}:{receiver_id}", sub.source_id);
206    push_unique_edge(
207        edges,
208        seen_edges,
209        deliver_edge_id.clone(),
210        json!({
211            "id": deliver_edge_id,
212            "source": format!("source:{}:{}", room_uuid, sub.source_id),
213            "target": format!("worker:{}", receiver.transport.media_worker_id),
214            "mainStat": download_main_stat(sub),
215            "secondaryStat": route_state_label(&sub.state),
216            "color": route_state_color(&sub.state),
217            "detail__direction": direction,
218            "detail__receiver_user": receiver_id,
219            "detail__selector": format!("{:?}", sub.selection.selector),
220            "detail__source_transport_media_id": sub.source_transport_media_id,
221            "detail__consumer_transport_media_id": sub.consumer_transport_media_id,
222        }),
223    );
224
225    let consume_edge_id = format!(
226        "consume:{room_uuid}:{}:{}",
227        sub.source_id,
228        receiver.user_id.path_segment()
229    );
230    push_unique_edge(
231        edges,
232        seen_edges,
233        consume_edge_id.clone(),
234        json!({
235            "id": consume_edge_id,
236            "source": format!("worker:{}", receiver.transport.media_worker_id),
237            "target": format!("user:{}:{}", room_uuid, receiver.user_id.path_segment()),
238            "mainStat": "consume",
239            "secondaryStat": direction,
240            "color": route_state_color(&sub.state),
241        }),
242    );
243}
244
245fn push_inbound_paths(
246    nodes: &mut Vec<Value>,
247    edges: &mut Vec<Value>,
248    seen_nodes: &mut HashSet<String>,
249    seen_edges: &mut HashSet<String>,
250    detail: &DiagnosticsRoomDetail,
251    selected_user: &DiagnosticsUserView,
252) {
253    for sub in &selected_user.subscriptions {
254        let Some(source) = source_by_id(detail, sub.source_id) else {
255            continue;
256        };
257        let Some(owner) = user_by_id(detail, &source.owner_user_id) else {
258            continue;
259        };
260        ensure_path_user(nodes, seen_nodes, detail, owner, false);
261        push_user_transport_edge(edges, seen_edges, detail, owner, "inbound-source");
262        push_publish_path(nodes, edges, seen_nodes, seen_edges, detail, owner, source);
263        push_subscription_delivery_path(edges, seen_edges, detail, selected_user, sub, "inbound");
264    }
265}
266
267fn push_outbound_paths(
268    nodes: &mut Vec<Value>,
269    edges: &mut Vec<Value>,
270    seen_nodes: &mut HashSet<String>,
271    seen_edges: &mut HashSet<String>,
272    detail: &DiagnosticsRoomDetail,
273    selected_user: &DiagnosticsUserView,
274) {
275    let selected_sources = detail
276        .sources
277        .iter()
278        .filter(|source| source.owner_user_id == selected_user.user_id);
279    for source in selected_sources {
280        push_publish_path(
281            nodes,
282            edges,
283            seen_nodes,
284            seen_edges,
285            detail,
286            selected_user,
287            source,
288        );
289        for receiver in &detail.users {
290            for sub in receiver
291                .subscriptions
292                .iter()
293                .filter(|sub| sub.source_id == source.source_id)
294            {
295                ensure_path_user(nodes, seen_nodes, detail, receiver, false);
296                push_user_transport_edge(edges, seen_edges, detail, receiver, "outbound-target");
297                push_subscription_delivery_path(
298                    edges, seen_edges, detail, receiver, sub, "outbound",
299                );
300            }
301        }
302    }
303}
304
305/// Build the Grafana node graph payload for one user's media paths in a room.
306///
307/// The graph shows the selected user's outbound publications through their
308/// source media worker to every receiver, and their inbound subscriptions from
309/// the producer user through the producer and receiver workers. It is a
310/// diagnostics projection of the current routing snapshot, not a transport
311/// control API.
312#[must_use]
313pub fn build_user_graph(detail: &DiagnosticsRoomDetail, requested_user_id: &str) -> Option<Value> {
314    let selected_user = detail
315        .users
316        .iter()
317        .find(|user| user.user_id.path_segment().as_ref() == requested_user_id)?;
318    let mut nodes = Vec::new();
319    let mut edges = Vec::new();
320    let mut seen_nodes = HashSet::new();
321    let mut seen_edges = HashSet::new();
322
323    ensure_path_user(&mut nodes, &mut seen_nodes, detail, selected_user, true);
324    push_user_transport_edge(
325        &mut edges,
326        &mut seen_edges,
327        detail,
328        selected_user,
329        "selected",
330    );
331    push_inbound_paths(
332        &mut nodes,
333        &mut edges,
334        &mut seen_nodes,
335        &mut seen_edges,
336        detail,
337        selected_user,
338    );
339    push_outbound_paths(
340        &mut nodes,
341        &mut edges,
342        &mut seen_nodes,
343        &mut seen_edges,
344        detail,
345        selected_user,
346    );
347
348    Some(json!({
349        "nodes": nodes,
350        "edges": edges,
351    }))
352}