1use std::{
9 collections::{BTreeMap, BTreeSet},
10 time::{Duration, Instant},
11};
12
13use o_sfu_rfc::rtp::Ssrc;
14use o_sfu_router::{MediaKind, rtp::MediaFormat};
15use o_sfu_telemetry::diagnostics::{
16 DiagnosticsActiveSpeaker, DiagnosticsActiveSpeakerReason, DiagnosticsActiveSpeakerState,
17 DiagnosticsIncomingBitrate, DiagnosticsMediaKind, DiagnosticsPendingUpgrade,
18 DiagnosticsPolicyPauseReason, DiagnosticsPublication, DiagnosticsQualitySummary,
19 DiagnosticsRouteState, DiagnosticsSource, DiagnosticsSourceEncoding,
20 DiagnosticsSourceSelection, DiagnosticsSourceSelectionReason, DiagnosticsSourceSelector,
21 DiagnosticsSubscription, DiagnosticsTransportCounts, DiagnosticsUserSummary,
22 DiagnosticsUserTransport, DiagnosticsUserView, DiagnosticsVideoLayoutRole,
23 DiagnosticsVideoRoutePriority, DiagnosticsWorkerSummary,
24};
25
26use super::{Room, RoomMediaCounts, media_graph::PendingUpgrade, state::RoomState};
27use crate::{
28 Bitrate,
29 engine::{
30 ConnectionId, MediaWorkerId, RecordingState, UserId, UserInfo,
31 media_transport::{
32 ActiveSpeakerActivityReason, ActiveSpeakerActivityState, ActiveSpeakerSourceDiagnostic,
33 MediaTransport, TransportBitrateSnapshot, TransportHealthSnapshot, TransportMediaId,
34 TransportQualitySample, TransportQualitySnapshot, TransportRidActivity,
35 TransportSessionHealth, TransportSessionKey, TransportSourceActivity,
36 TransportSourceDiagnosticsSnapshot, TransportSourceKey,
37 },
38 observability::diagnostics_transport_health,
39 source_model::{
40 ConsumerSourceSelection, PolicyPauseReason, PublishedSourceDescriptor,
41 SourceEncodingDescriptor, SourceEncodingId, SourceRoomPolicySelector,
42 SourceRoutePriority, SourceSelector, UserStreamId,
43 },
44 },
45};
46
47#[derive(Debug, Clone, Default, PartialEq, Eq)]
48pub struct IncomingBitrateSnapshot {
49 pub total: u64,
50 pub by_stream: BTreeMap<UserStreamId, u64>,
51}
52
53#[derive(Debug, Clone, PartialEq, Eq)]
54pub struct RoomUserStatsSnapshot {
55 pub incoming_bitrate: IncomingBitrateSnapshot,
56 pub count: u64,
57 pub active_stream_counts: BTreeMap<UserStreamId, u64>,
58}
59
60#[derive(Debug)]
62pub struct RoomOverviewCapture {
63 pub media_counts: RoomMediaCounts,
65 pub primary_media_worker_id: Option<MediaWorkerId>,
67 pub recording_state: RecordingState,
69 pub session_keys: Vec<TransportSessionKey>,
71}
72
73impl RoomOverviewCapture {
74 #[must_use]
75 pub fn transport_counts(&self, health: &TransportHealthSnapshot) -> DiagnosticsTransportCounts {
76 let mut counts = DiagnosticsTransportCounts::default();
77 for session_key in &self.session_keys {
78 let count = match health.get(session_key) {
79 Some(TransportSessionHealth::Connected) => &mut counts.connected,
80 Some(TransportSessionHealth::Disconnected) => &mut counts.disconnected,
81 None => &mut counts.unknown,
82 };
83 *count = count.saturating_add(1);
84 }
85 counts.total = self.session_keys.len();
86 counts
87 }
88}
89
90#[derive(Debug)]
92pub struct RoomUsersCapture {
93 primary_media_worker_id: Option<MediaWorkerId>,
94 users: Vec<CapturedUserSummary>,
95}
96
97impl RoomUsersCapture {
98 pub fn session_keys(&self) -> impl Iterator<Item = &TransportSessionKey> {
99 self.users.iter().map(|user| &user.session_key)
100 }
101
102 #[must_use]
103 pub fn into_user_summaries(
104 self,
105 room_id: &str,
106 bitrate: &TransportBitrateSnapshot,
107 health: &TransportHealthSnapshot,
108 stream_ids: [&str; 3],
109 ) -> Vec<DiagnosticsUserSummary> {
110 let bitrate_by_media = bitrate_by_media(bitrate);
111 self.users
112 .into_iter()
113 .map(|user| user_summary(room_id, user, &bitrate_by_media, health, stream_ids))
114 .collect()
115 }
116
117 pub fn add_to_worker_summaries(
118 self,
119 health: &TransportHealthSnapshot,
120 workers: &mut BTreeMap<usize, DiagnosticsWorkerSummary>,
121 ) {
122 if self.users.is_empty() {
123 let worker_id = self
124 .primary_media_worker_id
125 .map_or(0, MediaWorkerId::as_usize);
126 let worker = workers
127 .entry(worker_id)
128 .or_insert_with(|| worker_summary(worker_id));
129 worker.room_count = worker.room_count.saturating_add(1);
130 return;
131 }
132 let mut room_workers = BTreeSet::new();
133 for user in self.users {
134 let worker_id = user.session_key.media_worker_id().as_usize();
135 let worker = workers
136 .entry(worker_id)
137 .or_insert_with(|| worker_summary(worker_id));
138 if room_workers.insert(worker_id) {
139 worker.room_count = worker.room_count.saturating_add(1);
140 }
141 worker.user_count = worker.user_count.saturating_add(1);
142 worker.publication_count = worker
143 .publication_count
144 .saturating_add(user.publications.len());
145 worker.subscription_count = worker
146 .subscription_count
147 .saturating_add(user.subscription_count);
148 match health.get(&user.session_key) {
149 Some(TransportSessionHealth::Connected) => {
150 worker.connected_user_count = worker.connected_user_count.saturating_add(1);
151 }
152 Some(TransportSessionHealth::Disconnected) => {
153 worker.disconnected_user_count =
154 worker.disconnected_user_count.saturating_add(1);
155 }
156 None => worker.unknown_user_count = worker.unknown_user_count.saturating_add(1),
157 }
158 }
159 }
160}
161
162#[derive(Debug)]
164pub struct RoomDetailCapture {
165 overview: RoomOverviewCapture,
166 sources: Vec<CapturedSource>,
167 users: Vec<CapturedUser>,
168}
169
170impl RoomDetailCapture {
171 #[must_use]
172 pub fn session_keys(&self) -> &[TransportSessionKey] {
173 &self.overview.session_keys
174 }
175
176 pub fn source_keys(&self) -> impl Iterator<Item = &TransportSourceKey> {
177 self.sources.iter().map(|source| &source.source_key)
178 }
179
180 #[must_use]
181 pub fn into_views(
182 self,
183 bitrate: &TransportBitrateSnapshot,
184 quality: &TransportQualitySnapshot,
185 health: &TransportHealthSnapshot,
186 source_diagnostics: &TransportSourceDiagnosticsSnapshot,
187 ) -> (
188 RoomOverviewCapture,
189 Vec<DiagnosticsUserView>,
190 Vec<DiagnosticsSource>,
191 ) {
192 let bitrate_by_media = bitrate_by_media(bitrate);
193 let users = self
194 .users
195 .into_iter()
196 .map(|user| user_view(user, &bitrate_by_media, quality, health))
197 .collect();
198 let activity_by_media = source_diagnostics
199 .activity
200 .iter()
201 .map(|activity| (activity.transport_media_id(), activity))
202 .collect::<BTreeMap<_, _>>();
203 let speaker_diagnostics_by_media = source_diagnostics
204 .active_speaker_diagnostics
205 .iter()
206 .map(|speaker| (speaker.transport_media_id(), *speaker))
207 .collect::<BTreeMap<_, _>>();
208 let sources = self
209 .sources
210 .iter()
211 .map(|source| {
212 source_view(
213 source,
214 &bitrate_by_media,
215 &activity_by_media,
216 &speaker_diagnostics_by_media,
217 )
218 })
219 .collect();
220 (self.overview, users, sources)
221 }
222}
223
224#[derive(Debug)]
226pub struct RoomUserCapture {
227 recording_state: RecordingState,
228 user: CapturedUser,
229}
230
231impl RoomUserCapture {
232 #[must_use]
233 pub fn session_key(&self) -> &TransportSessionKey {
234 &self.user.session_key
235 }
236
237 #[must_use]
238 pub fn into_view(
239 self,
240 bitrate: &TransportBitrateSnapshot,
241 quality: &TransportQualitySnapshot,
242 health: &TransportHealthSnapshot,
243 ) -> (RecordingState, DiagnosticsUserView) {
244 let bitrate_by_media = bitrate_by_media(bitrate);
245 (
246 self.recording_state,
247 user_view(self.user, &bitrate_by_media, quality, health),
248 )
249 }
250}
251
252#[derive(Debug)]
253struct CapturedUser {
254 publications: Vec<DiagnosticsPublication>,
255 session_key: TransportSessionKey,
256 subscriptions: Vec<DiagnosticsSubscription>,
257 user_id: UserId,
258 user_info: UserInfo,
259 video_soft_pause_remaining_ms: Option<u64>,
260}
261
262#[derive(Debug)]
263struct CapturedUserSummary {
264 publications: Vec<(UserStreamId, TransportMediaId)>,
265 session_key: TransportSessionKey,
266 subscription_count: usize,
267 user_id: UserId,
268}
269
270#[derive(Debug)]
271struct CapturedSource {
272 active: bool,
273 descriptor: PublishedSourceDescriptor,
274 source_key: TransportSourceKey,
275}
276
277impl Room {
278 pub(crate) async fn session_stats_snapshot(
279 &self,
280 transport: &MediaTransport,
281 ) -> RoomUserStatsSnapshot {
282 let state = self.state.read().await;
283 let session_keys = transport_session_keys(&state);
284 let transport_snapshot = transport.transport_bitrate_snapshot(&session_keys);
285 let mut incoming_bitrate = IncomingBitrateSnapshot {
286 total: transport_snapshot.total.as_bps(),
287 ..Default::default()
288 };
289 for (transport_media_id, bits) in transport_snapshot.per_media {
290 let Some(stream_id) =
291 state.producer_stream_id_for_transport_media_id(transport_media_id)
292 else {
293 continue;
294 };
295 let entry = incoming_bitrate.by_stream.entry(stream_id).or_default();
296 *entry = entry.saturating_add(bits.as_bps());
297 }
298 let (count, active_stream_counts) = state.user_stats_counts();
299 drop(state);
300 RoomUserStatsSnapshot {
301 incoming_bitrate,
302 count,
303 active_stream_counts,
304 }
305 }
306
307 pub async fn diagnostics_overview_capture(&self) -> RoomOverviewCapture {
308 let state = self.state.read().await;
309 overview_capture(&state, transport_session_keys(&state))
310 }
311
312 pub async fn diagnostics_users_capture(&self) -> RoomUsersCapture {
313 let state = self.state.read().await;
314 RoomUsersCapture {
315 primary_media_worker_id: state.assigned_primary_media_worker_id(),
316 users: captured_user_summaries(&state),
317 }
318 }
319
320 pub async fn diagnostics_detail_capture(&self) -> RoomDetailCapture {
321 let state = self.state.read().await;
322 let now = Instant::now();
323 let users = state
324 .transport_user_entries()
325 .filter_map(|(user_id, connection_id)| {
326 captured_user(&state, user_id.clone(), connection_id, now)
327 })
328 .collect::<Vec<_>>();
329 let session_keys = users.iter().map(|user| user.session_key.clone()).collect();
330 let sources = state
331 .topology
332 .published_sources()
333 .map(|source| CapturedSource {
334 active: source.active,
335 descriptor: source.descriptor.clone(),
336 source_key: source.transport.clone(),
337 })
338 .collect();
339 let overview = overview_capture(&state, session_keys);
340 drop(state);
341 RoomDetailCapture {
342 overview,
343 sources,
344 users,
345 }
346 }
347
348 pub async fn diagnostics_user_capture(&self, user_key: &str) -> Option<RoomUserCapture> {
349 let state = self.state.read().await;
350 let now = Instant::now();
351 let (user_id, connection_id) = state
352 .transport_user_entries()
353 .find(|(user_id, _)| user_id.path_segment().as_ref() == user_key)?;
354 Some(RoomUserCapture {
355 recording_state: state.recording_state(),
356 user: captured_user(&state, user_id.clone(), connection_id, now)?,
357 })
358 }
359}
360
361fn overview_capture(
362 state: &RoomState,
363 session_keys: Vec<TransportSessionKey>,
364) -> RoomOverviewCapture {
365 RoomOverviewCapture {
366 media_counts: state.media_counts(),
367 primary_media_worker_id: state.assigned_primary_media_worker_id(),
368 recording_state: state.recording_state(),
369 session_keys,
370 }
371}
372
373fn captured_user_summaries(state: &RoomState) -> Vec<CapturedUserSummary> {
374 state
375 .transport_user_entries()
376 .map(|(user_id, connection_id)| CapturedUserSummary {
377 publications: state
378 .topology
379 .published_sources()
380 .filter(|publication| {
381 publication.descriptor.owner().user_id() == user_id
382 && publication.transport.session_key().connection_id() == connection_id
383 })
384 .map(|publication| {
385 (
386 publication.descriptor.stream_id().clone(),
387 publication.transport.transport_media_id(),
388 )
389 })
390 .collect(),
391 session_key: state.transport_user_key(user_id, connection_id),
392 subscription_count: state
393 .topology
394 .committed_consumer_routes_for_user(user_id)
395 .filter(|route| route.route.consumer_session_key().connection_id() == connection_id)
396 .count()
397 .saturating_add(
398 state
399 .topology
400 .pending_consumer_routes_for_user(user_id)
401 .count(),
402 ),
403 user_id: user_id.clone(),
404 })
405 .collect()
406}
407
408fn captured_user(
409 state: &RoomState,
410 user_id: UserId,
411 connection_id: ConnectionId,
412 now: Instant,
413) -> Option<CapturedUser> {
414 Some(CapturedUser {
415 publications: diagnostics_publications(state, &user_id, connection_id),
416 session_key: state.transport_user_key(&user_id, connection_id),
417 subscriptions: diagnostics_subscriptions(state, &user_id, connection_id, now),
418 video_soft_pause_remaining_ms: state
419 .user_for_connection(&user_id, connection_id)?
420 .video_soft_pause_deadline
421 .map(|deadline| duration_millis(deadline.saturating_duration_since(now))),
422 user_info: state.user_info_snapshot(&user_id)?.1,
423 user_id,
424 })
425}
426
427fn diagnostics_publications(
428 state: &RoomState,
429 user_id: &UserId,
430 connection_id: ConnectionId,
431) -> Vec<DiagnosticsPublication> {
432 state
433 .topology
434 .published_sources()
435 .filter(|publication| {
436 publication.descriptor.owner().user_id() == user_id
437 && publication.transport.session_key().connection_id() == connection_id
438 })
439 .map(|publication| {
440 let source = &publication.descriptor;
441 DiagnosticsPublication {
442 active: publication.active,
443 encoding_ids: source
444 .encodings()
445 .map(|encoding| encoding.encoding_id().as_u64())
446 .collect(),
447 media_kind: media_kind(source.media_kind()),
448 source_id: source.source_id().as_u64(),
449 stream_id: source.stream_id().to_string(),
450 transport_media_id: Some(publication.transport.transport_media_id().as_u64()),
451 }
452 })
453 .collect()
454}
455
456fn diagnostics_subscriptions(
457 state: &RoomState,
458 user_id: &UserId,
459 connection_id: ConnectionId,
460 now: Instant,
461) -> Vec<DiagnosticsSubscription> {
462 let project = |source: &PublishedSourceDescriptor,
463 route_selection: ConsumerSourceSelection,
464 pending_upgrade: Option<PendingUpgrade>,
465 consumer_media: Option<TransportMediaId>,
466 source_media: TransportMediaId,
467 route_state| {
468 let layout = state.diagnostics_video_layout_role(user_id, source);
469 DiagnosticsSubscription {
470 consumer_transport_media_id: consumer_media.map(TransportMediaId::as_u64),
471 layout_priority: layout.map(|role| role.priority().into()),
472 layout_role: layout.map(Into::into),
473 producer_user_id: source.owner().user_id().clone(),
474 selection: selection(source, route_selection, pending_upgrade, now),
475 source_id: source.source_id().as_u64(),
476 source_transport_media_id: Some(source_media.as_u64()),
477 state: route_state,
478 stream_id: source.stream_id().to_string(),
479 }
480 };
481 let mut subscriptions = state
482 .topology
483 .committed_consumer_routes_for_user(user_id)
484 .filter(|route| route.route.consumer_session_key().connection_id() == connection_id)
485 .map(|route| {
486 let source = &route.source.descriptor;
487 let route_state = if route.source.active && route.selection.delivery_active() {
488 DiagnosticsRouteState::Active
489 } else {
490 DiagnosticsRouteState::Inactive
491 };
492 project(
493 source,
494 route.selection,
495 route.pending_upgrade.copied(),
496 Some(route.route.consumer_transport_media_id()),
497 route.route.source_transport_media_id(),
498 route_state,
499 )
500 })
501 .collect::<Vec<_>>();
502 subscriptions.extend(
503 state
504 .topology
505 .pending_consumer_routes_for_user(user_id)
506 .map(|route| {
507 let source = &route.source.descriptor;
508 project(
509 source,
510 route.selection,
511 None,
512 None,
513 route.source.transport.transport_media_id(),
514 DiagnosticsRouteState::Pending,
515 )
516 }),
517 );
518 subscriptions
519}
520
521fn user_summary(
522 room_id: &str,
523 user: CapturedUserSummary,
524 bitrate_by_media: &BTreeMap<u64, u64>,
525 health: &TransportHealthSnapshot,
526 stream_ids: [&str; 3],
527) -> DiagnosticsUserSummary {
528 let [audio_stream_id, camera_stream_id, screen_stream_id] = stream_ids;
529 let mut audio = 0_u64;
530 let mut camera = 0_u64;
531 let mut screen = 0_u64;
532 let mut total = 0_u64;
533 for (stream_id, media_id) in &user.publications {
534 let bitrate = bitrate_by_media
535 .get(&media_id.as_u64())
536 .copied()
537 .unwrap_or_default();
538 total = total.saturating_add(bitrate);
539 let stream = match stream_id.as_str() {
540 value if value == audio_stream_id => &mut audio,
541 value if value == camera_stream_id => &mut camera,
542 value if value == screen_stream_id => &mut screen,
543 _ => continue,
544 };
545 *stream = stream.saturating_add(bitrate);
546 }
547 DiagnosticsUserSummary {
548 audio_incoming_bitrate_bps: audio,
549 camera_incoming_bitrate_bps: camera,
550 connection_id: user.session_key.connection_id().as_u64(),
551 health: health
552 .get(&user.session_key)
553 .copied()
554 .map(diagnostics_transport_health),
555 incoming_bitrate_bps: total,
556 media_worker_id: user.session_key.media_worker_id().as_usize(),
557 publication_count: user.publications.len(),
558 room_id: room_id.to_owned(),
559 screen_incoming_bitrate_bps: screen,
560 subscription_count: user.subscription_count,
561 user_key: user.user_id.path_segment().into_owned(),
562 user_id: user.user_id,
563 }
564}
565
566fn worker_summary(media_worker_id: usize) -> DiagnosticsWorkerSummary {
567 DiagnosticsWorkerSummary {
568 media_worker_id,
569 ..Default::default()
570 }
571}
572
573fn user_view(
574 user: CapturedUser,
575 bitrate_by_media: &BTreeMap<u64, u64>,
576 quality: &TransportQualitySnapshot,
577 health: &TransportHealthSnapshot,
578) -> DiagnosticsUserView {
579 let transport = DiagnosticsUserTransport {
580 video_soft_pause_remaining_ms: user.video_soft_pause_remaining_ms,
581 connection_id: user.session_key.connection_id().as_u64(),
582 health: health
583 .get(&user.session_key)
584 .copied()
585 .map(diagnostics_transport_health),
586 media_worker_id: user.session_key.media_worker_id().as_usize(),
587 quality_summary: quality_summary(
588 incoming_bitrate(&user.publications, bitrate_by_media),
589 quality.get(&user.session_key).copied(),
590 ),
591 };
592 DiagnosticsUserView {
593 publications: user.publications,
594 subscriptions: user.subscriptions,
595 transport,
596 user_id: user.user_id,
597 user_info: user.user_info,
598 }
599}
600
601fn bitrate_by_media(snapshot: &TransportBitrateSnapshot) -> BTreeMap<u64, u64> {
602 snapshot
603 .per_media
604 .iter()
605 .map(|(media, bitrate)| (media.as_u64(), bitrate.as_bps()))
606 .collect()
607}
608
609fn incoming_bitrate(
610 publications: &[DiagnosticsPublication],
611 bitrate_by_media: &BTreeMap<u64, u64>,
612) -> DiagnosticsIncomingBitrate {
613 let mut incoming = DiagnosticsIncomingBitrate::default();
614 for publication in publications {
615 let bitrate = publication
616 .transport_media_id
617 .and_then(|media| bitrate_by_media.get(&media))
618 .copied()
619 .unwrap_or_default();
620 incoming.total = incoming.total.saturating_add(bitrate);
621 let stream = incoming
622 .by_stream_bps
623 .entry(publication.stream_id.clone())
624 .or_default();
625 *stream = stream.saturating_add(bitrate);
626 }
627 incoming
628}
629
630fn quality_summary(
631 current_incoming_bitrate: DiagnosticsIncomingBitrate,
632 quality: Option<TransportQualitySample>,
633) -> DiagnosticsQualitySummary {
634 let quality = quality.unwrap_or_default();
635 DiagnosticsQualitySummary {
636 current_incoming_bitrate,
637 sampled_metrics_available: quality.sample_count > 0,
638 latest_bwe_bps: quality.latest_bwe_bps,
639 rtt_ms: quality.rtt_ms,
640 ingress_loss_ppm: quality.ingress_loss_ppm,
641 egress_loss_ppm: quality.egress_loss_ppm,
642 egress_jitter_rtp_timestamp_units: quality.egress_jitter_rtp_timestamp_units,
643 sample_count: quality.sample_count,
644 }
645}
646
647fn source_view(
648 source: &CapturedSource,
649 bitrate_by_media: &BTreeMap<u64, u64>,
650 activity_by_media: &BTreeMap<TransportMediaId, &TransportSourceActivity>,
651 speaker_diagnostics_by_media: &BTreeMap<TransportMediaId, ActiveSpeakerSourceDiagnostic>,
652) -> DiagnosticsSource {
653 let descriptor = &source.descriptor;
654 let media_id = source.source_key.transport_media_id();
655 let activity = activity_by_media.get(&media_id).copied();
656 DiagnosticsSource {
657 active: source.active,
658 active_speaker: active_speaker(descriptor, media_id, speaker_diagnostics_by_media),
659 current_incoming_bitrate_bps: bitrate_by_media
660 .get(&media_id.as_u64())
661 .copied()
662 .unwrap_or_default(),
663 encodings: descriptor
664 .encodings()
665 .map(|encoding| source_encoding(encoding, activity))
666 .collect(),
667 last_packet_age_ms: activity.map(|value| duration_millis(value.last_packet_age())),
668 last_keyframe_age_ms: activity
669 .and_then(TransportSourceActivity::last_keyframe_age)
670 .map(duration_millis),
671 media_kind: media_kind(descriptor.media_kind()),
672 mid: descriptor.mid().map(|mid| mid.as_str().to_owned()),
673 owner_user_id: descriptor.owner().user_id().clone(),
674 source_id: descriptor.source_id().as_u64(),
675 stream_id: descriptor.stream_id().to_string(),
676 transport_media_id: Some(media_id.as_u64()),
677 video_bitrate_cap_bps: descriptor.policy().video_bitrate_cap().map(Bitrate::as_bps),
678 }
679}
680
681fn source_encoding(
682 encoding: &SourceEncodingDescriptor,
683 activity: Option<&TransportSourceActivity>,
684) -> DiagnosticsSourceEncoding {
685 let format = encoding.negotiated_format();
686 let rid_activity = encoding
687 .rid()
688 .and_then(|rid| rid_activity(activity, rid.as_str()));
689 DiagnosticsSourceEncoding {
690 codec: format.map(|value| value.codec_name().to_owned()),
691 encoding_id: encoding.encoding_id().as_u64(),
692 max_bitrate_bps: encoding.max_bitrate().map(Bitrate::as_bps),
693 resolution_scale: encoding.resolution_scale(),
694 max_framerate: encoding.max_framerate(),
695 policy_role: encoding
696 .policy_role()
697 .map(|role| role.as_wire_value().to_owned()),
698 payload_type: format.map(MediaFormat::payload_type),
699 primary_ssrc: encoding.primary_ssrc().map(Ssrc::value),
700 repair_ssrc: encoding.repair_ssrc().map(Ssrc::value),
701 rid: encoding.rid().map(|rid| rid.as_str().to_owned()),
702 last_packet_age_ms: rid_activity.map(|value| duration_millis(value.last_packet_age())),
703 last_keyframe_age_ms: rid_activity
704 .and_then(TransportRidActivity::last_keyframe_age)
705 .map(duration_millis),
706 }
707}
708
709fn rid_activity<'a>(
710 activity: Option<&'a TransportSourceActivity>,
711 rid: &str,
712) -> Option<&'a TransportRidActivity> {
713 activity?.rids().iter().find(|value| value.rid() == rid)
714}
715
716fn selection(
717 source: &PublishedSourceDescriptor,
718 selection: ConsumerSourceSelection,
719 pending_upgrade: Option<PendingUpgrade>,
720 now: Instant,
721) -> DiagnosticsSourceSelection {
722 let selected_encoding_id = selection.selector().selected_encoding();
723 let budget = selection.budget();
724 let selected_encoding = selected_encoding_id.and_then(|id| source.encoding(id));
725 let (selector, selection_reason) = match selection.selector() {
726 SourceSelector::Open => (
727 DiagnosticsSourceSelector::Open,
728 DiagnosticsSourceSelectionReason::Open,
729 ),
730 SourceSelector::Encoding(_) => (
731 DiagnosticsSourceSelector::Encoding,
732 DiagnosticsSourceSelectionReason::ReceiverAdaptation,
733 ),
734 };
735 DiagnosticsSourceSelection {
736 active: selection.active(),
737 active_video_route_count: budget.active_video_route_count(),
738 latest_receiver_bandwidth_estimate_bps: budget
739 .latest_receiver_bandwidth()
740 .map(Bitrate::as_bps),
741 policy_allows_delivery: selection.policy_allows_delivery(),
742 policy_pause_reason: selection.policy_pause_reason().map(Into::into),
743 pending_upgrade: pending_upgrade.map(|pending| DiagnosticsPendingUpgrade {
744 selector: match pending.selector {
745 SourceSelector::Open => DiagnosticsSourceSelector::Open,
746 SourceSelector::Encoding(_) => DiagnosticsSourceSelector::Encoding,
747 },
748 encoding_id: pending
749 .selector
750 .selected_encoding()
751 .map(SourceEncodingId::as_u64),
752 remaining_ms: duration_millis(pending.deadline.saturating_duration_since(now)),
753 }),
754 selection_reason,
755 selector,
756 selected_estimated_bitrate_bps: selected_encoding
757 .and_then(SourceEncodingDescriptor::max_bitrate)
758 .map(Bitrate::as_bps),
759 selected_video_bitrate_bps: budget.selected_video_bitrate().as_bps(),
760 selected_video_budget_bps: budget.selected_video_budget().map(Bitrate::as_bps),
761 selected_encoding_id: selected_encoding_id.map(SourceEncodingId::as_u64),
762 selected_rid: selected_encoding
763 .and_then(SourceEncodingDescriptor::rid)
764 .map(|rid| rid.as_str().to_owned()),
765 }
766}
767
768fn active_speaker(
769 source: &PublishedSourceDescriptor,
770 media_id: TransportMediaId,
771 diagnostics: &BTreeMap<TransportMediaId, ActiveSpeakerSourceDiagnostic>,
772) -> Option<DiagnosticsActiveSpeaker> {
773 (source.media_kind() == MediaKind::Audio).then(|| {
774 diagnostics
775 .get(&media_id)
776 .copied()
777 .map_or_else(DiagnosticsActiveSpeaker::idle, active_speaker_snapshot)
778 })
779}
780
781fn media_kind(value: MediaKind) -> DiagnosticsMediaKind {
782 match value {
783 MediaKind::Audio => DiagnosticsMediaKind::Audio,
784 MediaKind::Video => DiagnosticsMediaKind::Video,
785 }
786}
787
788fn active_speaker_snapshot(diagnostic: ActiveSpeakerSourceDiagnostic) -> DiagnosticsActiveSpeaker {
789 DiagnosticsActiveSpeaker {
790 state: diagnostic.state().into(),
791 reason: diagnostic.reason().into(),
792 last_audio_level_dbov: diagnostic.last_audio_level_dbov(),
793 confidence_observations: diagnostic.confidence_observations(),
794 hold_remaining_ms: diagnostic.hold_remaining().map(duration_millis),
795 }
796}
797
798fn duration_millis(duration: Duration) -> u64 {
799 u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
800}
801
802fn transport_session_keys(state: &RoomState) -> Vec<TransportSessionKey> {
803 state
804 .transport_user_entries()
805 .map(|(user_id, connection_id)| state.transport_user_key(user_id, connection_id))
806 .collect()
807}
808
809impl From<PolicyPauseReason> for DiagnosticsPolicyPauseReason {
810 fn from(value: PolicyPauseReason) -> Self {
811 match value {
812 PolicyPauseReason::BudgetPressure => Self::BudgetPressure,
813 PolicyPauseReason::HiddenTile => Self::HiddenTile,
814 PolicyPauseReason::OverflowTile => Self::OverflowTile,
815 PolicyPauseReason::MissingUsableLayer => Self::MissingUsableLayer,
816 PolicyPauseReason::AudioSpeakerLimit => Self::AudioSpeakerLimit,
817 PolicyPauseReason::ReceiverDeafened => Self::ReceiverDeafened,
818 PolicyPauseReason::VideoDownloadLimit => Self::VideoDownloadLimit,
819 PolicyPauseReason::SourceBitrateLimit => Self::SourceBitrateLimit,
820 }
821 }
822}
823
824impl From<SourceRoomPolicySelector> for DiagnosticsVideoLayoutRole {
825 fn from(value: SourceRoomPolicySelector) -> Self {
826 match value {
827 SourceRoomPolicySelector::Pinned => Self::Pinned,
828 SourceRoomPolicySelector::Featured => Self::Featured,
829 SourceRoomPolicySelector::ReadableDetail => Self::ReadableDetail,
830 SourceRoomPolicySelector::ActiveSpeaker => Self::ActiveSpeaker,
831 SourceRoomPolicySelector::VisibleThumbnail => Self::VisibleThumbnail,
832 SourceRoomPolicySelector::Hidden => Self::Hidden,
833 SourceRoomPolicySelector::Overflow => Self::Overflow,
834 }
835 }
836}
837
838impl From<SourceRoutePriority> for DiagnosticsVideoRoutePriority {
839 fn from(value: SourceRoutePriority) -> Self {
840 match value {
841 SourceRoutePriority::PinnedOrFeatured => Self::PinnedOrFeatured,
842 SourceRoutePriority::ReadableDetail => Self::ReadableDetail,
843 SourceRoutePriority::ActiveSpeaker => Self::ActiveSpeaker,
844 SourceRoutePriority::VisibleThumbnail => Self::VisibleThumbnail,
845 SourceRoutePriority::HiddenOrOverflow => Self::HiddenOrOverflow,
846 }
847 }
848}
849
850impl From<ActiveSpeakerActivityState> for DiagnosticsActiveSpeakerState {
851 fn from(value: ActiveSpeakerActivityState) -> Self {
852 match value {
853 ActiveSpeakerActivityState::Active => Self::Active,
854 ActiveSpeakerActivityState::Idle => Self::Idle,
855 ActiveSpeakerActivityState::Blocked => Self::Blocked,
856 ActiveSpeakerActivityState::RecentlyExpired => Self::RecentlyExpired,
857 }
858 }
859}
860
861impl From<ActiveSpeakerActivityReason> for DiagnosticsActiveSpeakerReason {
862 fn from(value: ActiveSpeakerActivityReason) -> Self {
863 match value {
864 ActiveSpeakerActivityReason::Vad => Self::Vad,
865 ActiveSpeakerActivityReason::AudioLevel => Self::AudioLevel,
866 ActiveSpeakerActivityReason::AudioLevelWarmup => Self::AudioLevelWarmup,
867 ActiveSpeakerActivityReason::VadFalse => Self::VadFalse,
868 ActiveSpeakerActivityReason::LowNoise => Self::LowNoise,
869 ActiveSpeakerActivityReason::BelowSpeechThreshold => Self::BelowSpeechThreshold,
870 ActiveSpeakerActivityReason::MissingAudioMetadata => Self::MissingAudioMetadata,
871 ActiveSpeakerActivityReason::Expired => Self::Expired,
872 ActiveSpeakerActivityReason::NoMetadata => Self::NoMetadata,
873 }
874 }
875}
876
877#[cfg(test)]
878#[path = "TESTS/read_model.rs"]
879mod tests;