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,
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/// Passive facts required by summary and room-list diagnostics.
60#[derive(Debug)]
61pub struct RoomOverviewCapture {
62    /// Current publication and subscription counts.
63    pub media_counts: RoomMediaCounts,
64    /// Primary media worker assigned to the room.
65    pub primary_media_worker_id: Option<MediaWorkerId>,
66    /// Current room recording facts.
67    pub recording_state: RecordingState,
68    /// Active transport sessions in the room.
69    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/// Passive room user inventory used by room-user and worker diagnostics.
90#[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/// Passive room state required by room detail and graph diagnostics.
162#[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/// Passive room and user facts for one room-scoped user lookup.
224#[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;