Skip to main content

o_sfu_core/engine/room/
read_model.rs

1//! Two-phase room diagnostics.
2//!
3//! Capture methods copy room facts and transport keys under the room read guard.
4//! Callers collect transport snapshots after releasing that guard, then
5//! projection methods combine both views. The result is not an atomic room and
6//! transport snapshot.
7
8use 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/// Passive facts required by summary and room-list diagnostics.
61#[derive(Debug)]
62pub struct RoomOverviewCapture {
63    /// Current publication and subscription counts.
64    pub media_counts: RoomMediaCounts,
65    /// Primary media worker assigned to the room.
66    pub primary_media_worker_id: Option<MediaWorkerId>,
67    /// Current room recording facts.
68    pub recording_state: RecordingState,
69    /// Active transport sessions in the room.
70    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/// Passive room user inventory used by room-user and worker diagnostics.
91#[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/// Passive room state required by room detail diagnostics.
163#[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/// Passive room and user facts for one room-scoped user lookup.
225#[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;