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