Skip to main content

o_sfu_telemetry/metrics/
descriptor.rs

1use std::fmt::Write as _;
2
3use o_sfu_model::WebSocketCloseCode;
4
5#[cfg(any(test, feature = "test-support"))]
6use super::snapshot::{RuntimeMetricsSnapshot, SnapshotWriter};
7use super::{
8    catalog::RuntimeMetrics,
9    counter::{
10        CounterFamily, ExportedMetricLabel, Histogram, HistogramBucketLabel, HistogramFamily,
11        MetricBucketLabel, UpDownCounterFamily,
12    },
13    labels::{
14        ControlPlaneDurationBucket, ExportedMetricLabelPair, HttpRoute, RtcRelayEnqueueResult,
15        RtcTransportIoFailure,
16    },
17    rtc::RtcMetricsSnapshot,
18    rtp::{RtpMetricsSnapshot, RtpTrafficSnapshot},
19};
20
21#[cfg(test)]
22#[path = "TESTS/descriptor.rs"]
23mod tests;
24
25macro_rules! metric_catalog {
26    (@kind Counter) => { "counter" };
27    (@kind Gauge) => { "gauge" };
28    (@kind Histogram) => { "histogram" };
29    (@metadata $name:literal, $help:literal, $kind:ident) => {
30        concat!("# HELP ", $name, " ", $help, "\n# TYPE ", $name, " ", metric_catalog!(@kind $kind), "\n")
31    };
32    ($($id:ident {
33        name: $name:literal,
34        help: $help:literal,
35        kind: $kind:ident,
36        samples: |$metrics:ident, $capture:ident, $output:ident| $samples:expr
37    }),+ $(,)?) => {
38        /// Names every exported Prometheus metric family.
39        #[derive(Debug, Clone, Copy, PartialEq, Eq)]
40        pub enum MetricName {
41            $(
42                #[doc = concat!("`", $name, "`\n\n", $help)]
43                $id,
44            )+
45        }
46
47        #[cfg(test)]
48        pub(crate) const METRIC_FAMILY_COUNT: usize = [$(MetricName::$id),+].len();
49
50        const PROMETHEUS_METADATA_CAPACITY: usize = 0 $(
51            + metric_catalog!(@metadata $name, $help, $kind).len()
52        )+;
53
54        fn export(
55            metrics: &RuntimeMetrics,
56            room_gauges: RoomGaugeValues,
57            destination: &mut MetricDestination<'_>,
58        ) {
59            let capture = MetricCapture {
60                room_gauges,
61                rtp: metrics.rtp_metrics.snapshot(),
62                rtc: metrics.rtc_metrics.snapshot(),
63            };
64            $(
65                {
66                    let $metrics = metrics;
67                    let $capture = &capture;
68                    let $output = &mut MetricOutput::new(
69                        MetricDescriptor {
70                            #[cfg(any(test, feature = "test-support"))]
71                            id: MetricName::$id,
72                            name: $name,
73                        },
74                        metric_catalog!(@metadata $name, $help, $kind),
75                        destination,
76                    );
77                    let _ = ($metrics, $capture);
78                    $samples
79                }
80            )+
81        }
82    };
83}
84
85#[derive(Clone, Copy)]
86struct MetricDescriptor {
87    #[cfg(any(test, feature = "test-support"))]
88    id: MetricName,
89    name: &'static str,
90}
91
92struct MetricCapture {
93    room_gauges: RoomGaugeValues,
94    rtp: RtpMetricsSnapshot,
95    rtc: RtcMetricsSnapshot,
96}
97
98/// Room counts supplied to one export and saturated during gauge encoding.
99#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
100pub struct RoomGaugeValues {
101    pub rooms: usize,
102    pub users: usize,
103    pub publications: usize,
104    pub subscriptions: usize,
105    pub recording_rooms: usize,
106}
107
108#[derive(Debug, Clone, Copy, PartialEq, Eq)]
109pub(super) enum MetricLabelValue {
110    Text(&'static str),
111    Number(usize),
112}
113
114#[derive(Debug, Clone, Copy, PartialEq, Eq)]
115pub(super) struct MetricLabel {
116    pub(super) name: &'static str,
117    pub(super) value: MetricLabelValue,
118}
119
120impl MetricLabel {
121    const fn text(name: &'static str, value: &'static str) -> Self {
122        Self {
123            name,
124            value: MetricLabelValue::Text(value),
125        }
126    }
127
128    const fn number(name: &'static str, value: usize) -> Self {
129        Self {
130            name,
131            value: MetricLabelValue::Number(value),
132        }
133    }
134}
135
136enum MetricDestination<'a> {
137    Prometheus(&'a mut String),
138    #[cfg(any(test, feature = "test-support"))]
139    Snapshot(&'a mut SnapshotWriter),
140}
141
142struct MetricOutput<'a, 'b> {
143    family: MetricDescriptor,
144    destination: &'a mut MetricDestination<'b>,
145}
146
147impl<'a, 'b> MetricOutput<'a, 'b> {
148    fn new(
149        family: MetricDescriptor,
150        metadata: &str,
151        destination: &'a mut MetricDestination<'b>,
152    ) -> Self {
153        match destination {
154            MetricDestination::Prometheus(output) => output.push_str(metadata),
155            #[cfg(any(test, feature = "test-support"))]
156            MetricDestination::Snapshot(_) => {}
157        }
158        Self {
159            family,
160            destination,
161        }
162    }
163
164    fn counter(&mut self, labels: &[MetricLabel], value: u64) {
165        match self.destination {
166            MetricDestination::Prometheus(output) => {
167                append_sample_name(output, self.family.name, labels);
168                let _ = writeln!(output, " {value}");
169            }
170            #[cfg(any(test, feature = "test-support"))]
171            MetricDestination::Snapshot(snapshot) => {
172                snapshot.counter(self.family.id, Box::from(labels), value);
173            }
174        }
175    }
176
177    fn gauge(&mut self, labels: &[MetricLabel], value: i64) {
178        match self.destination {
179            MetricDestination::Prometheus(output) => {
180                append_sample_name(output, self.family.name, labels);
181                let _ = writeln!(output, " {value}");
182            }
183            #[cfg(any(test, feature = "test-support"))]
184            MetricDestination::Snapshot(snapshot) => {
185                snapshot.gauge(self.family.id, Box::from(labels), value);
186            }
187        }
188    }
189
190    fn histogram<B>(
191        &mut self,
192        labels: &[MetricLabel],
193        load_bucket: impl Fn(B) -> u64,
194        load_count: impl Fn() -> u64,
195        load_sum_micros: impl Fn() -> u64,
196    ) where
197        B: MetricBucketLabel,
198    {
199        let descriptor = self.family;
200        match self.destination {
201            MetricDestination::Prometheus(output) => {
202                let mut floor = 0;
203                for bucket in B::VARIANTS {
204                    let value = floor.max(load_bucket(*bucket));
205                    floor = value;
206                    output.push_str(descriptor.name);
207                    output.push_str("_bucket");
208                    append_labels(
209                        output,
210                        labels,
211                        Some(MetricLabel::text("le", bucket.upper_bound())),
212                    );
213                    let _ = writeln!(output, " {value}");
214                }
215                let count = floor.max(load_count());
216                let sum_micros = load_sum_micros();
217                output.push_str(descriptor.name);
218                output.push_str("_bucket");
219                append_labels(output, labels, Some(MetricLabel::text("le", "+Inf")));
220                let _ = writeln!(output, " {count}");
221                output.push_str(descriptor.name);
222                output.push_str("_sum");
223                append_labels(output, labels, None);
224                output.push(' ');
225                append_seconds_from_micros(output, sum_micros);
226                output.push('\n');
227                output.push_str(descriptor.name);
228                output.push_str("_count");
229                append_labels(output, labels, None);
230                let _ = writeln!(output, " {count}");
231            }
232            #[cfg(any(test, feature = "test-support"))]
233            MetricDestination::Snapshot(snapshot) => {
234                let mut floor = 0;
235                let buckets = B::VARIANTS
236                    .iter()
237                    .map(|bucket| {
238                        let value = floor.max(load_bucket(*bucket));
239                        floor = value;
240                        (bucket.upper_bound(), value)
241                    })
242                    .collect();
243                snapshot.histogram(
244                    descriptor.id,
245                    Box::from(labels),
246                    buckets,
247                    floor.max(load_count()),
248                    load_sum_micros(),
249                );
250            }
251        }
252    }
253}
254
255pub(crate) fn render_prometheus_text(metrics: &RuntimeMetrics, gauges: RoomGaugeValues) -> String {
256    let mut text = String::with_capacity(PROMETHEUS_METADATA_CAPACITY);
257    export(
258        metrics,
259        gauges,
260        &mut MetricDestination::Prometheus(&mut text),
261    );
262    text
263}
264
265#[cfg(any(test, feature = "test-support"))]
266pub(super) fn build_snapshot(metrics: &RuntimeMetrics) -> RuntimeMetricsSnapshot {
267    let mut snapshot = SnapshotWriter::default();
268    export(
269        metrics,
270        RoomGaugeValues::default(),
271        &mut MetricDestination::Snapshot(&mut snapshot),
272    );
273    snapshot.finish()
274}
275
276fn gauge_count(value: usize) -> i64 {
277    i64::try_from(value).unwrap_or(i64::MAX)
278}
279
280metric_catalog! {
281    HttpConnectionRejectionsTotal {
282        name: "osfu_http_connection_rejections_total",
283        help: "Total TCP connections rejected at the HTTP listener connection cap.",
284        kind: Counter,
285        samples: |metrics, capture, output| output.counter(&[], metrics.http_connection_rejections.load())
286    },
287    HttpNoopRequestsTotal {
288        name: "osfu_http_noop_requests_total",
289        help: "Total HTTP requests served by /v1/noop.",
290        kind: Counter,
291        samples: |metrics, capture, output| output.counter(&[], metrics.http_requests.load(HttpRoute::Noop))
292    },
293    HttpStatsRequestsTotal {
294        name: "osfu_http_stats_requests_total",
295        help: "Total HTTP requests served by /v1/stats.",
296        kind: Counter,
297        samples: |metrics, capture, output| output.counter(&[], metrics.http_requests.load(HttpRoute::Stats))
298    },
299    HttpRoomRequestsTotal {
300        name: "osfu_http_room_requests_total",
301        help: "Total HTTP requests received by /v1/channel.",
302        kind: Counter,
303        samples: |metrics, capture, output| output.counter(&[], metrics.http_requests.load(HttpRoute::Room))
304    },
305    HttpRoomResponsesTotal {
306        name: "osfu_http_room_responses_total",
307        help: "Total HTTP /v1/channel responses by status.",
308        kind: Counter,
309        samples: |metrics, capture, output| write_counter_family(output, &metrics.http_room_responses, "status")
310    },
311    HttpDisconnectRequestsTotal {
312        name: "osfu_http_disconnect_requests_total",
313        help: "Total HTTP requests received by /v1/disconnect.",
314        kind: Counter,
315        samples: |metrics, capture, output| output.counter(&[], metrics.http_requests.load(HttpRoute::Disconnect))
316    },
317    HttpDisconnectResponsesTotal {
318        name: "osfu_http_disconnect_responses_total",
319        help: "Total HTTP /v1/disconnect responses by status.",
320        kind: Counter,
321        samples: |metrics, capture, output| write_counter_family(output, &metrics.http_disconnect_responses, "status")
322    },
323    HttpMetricsRequestsTotal {
324        name: "osfu_http_metrics_requests_total",
325        help: "Total HTTP requests served by /metrics.",
326        kind: Counter,
327        samples: |metrics, capture, output| output.counter(&[], metrics.http_requests.load(HttpRoute::Metrics))
328    },
329    HttpInflightRequests {
330        name: "osfu_http_inflight_requests",
331        help: "Current in-flight HTTP requests by route.",
332        kind: Gauge,
333        samples: |metrics, capture, output| write_up_down_counter_family(output, &metrics.http_inflight_requests, "route")
334    },
335    HttpRequestDurationSeconds {
336        name: "osfu_http_request_duration_seconds",
337        help: "HTTP request duration by route.",
338        kind: Histogram,
339        samples: |metrics, capture, output| write_histogram_family(output, &metrics.http_request_duration, "route")
340    },
341    WsConnectionsTotal {
342        name: "osfu_ws_connections_total",
343        help: "Total websocket connections observed at each handshake stage.",
344        kind: Counter,
345        samples: |metrics, capture, output| write_counter_family(output, &metrics.ws_connections, "stage")
346    },
347    WsPreAuthRejectionsTotal {
348        name: "osfu_ws_pre_auth_rejections_total",
349        help: "Total websocket upgrades rejected by pre-auth admission limit.",
350        kind: Counter,
351        samples: |metrics, capture, output| write_counter_family(output, &metrics.ws_pre_auth_rejections, "limit")
352    },
353    WsHandshakeRejectionsTotal {
354        name: "osfu_ws_handshake_rejections_total",
355        help: "Total websocket handshake rejections by close code bucket.",
356        kind: Counter,
357        samples: |metrics, capture, output| {
358            counter(output,
359                [("close_code", WebSocketCloseCode::AuthTimeout.label_value())],
360                metrics.ws_handshake_rejections.load(WebSocketCloseCode::AuthTimeout),
361            );
362            counter(output,
363                [("close_code", WebSocketCloseCode::AuthFailed.label_value())],
364                metrics.ws_handshake_rejections.load(WebSocketCloseCode::AuthFailed),
365            );
366            counter(output,
367                [("close_code", WebSocketCloseCode::ProtocolError.label_value())],
368                metrics.ws_handshake_rejections.load(WebSocketCloseCode::ProtocolError),
369            );
370            counter(output,
371                [("close_code", WebSocketCloseCode::RoomFull.label_value())],
372                metrics.ws_handshake_rejections.load(WebSocketCloseCode::RoomFull),
373            );
374            counter(output, [("close_code", "error")], metrics.ws_handshake_rejections_other.load());
375        }
376    },
377    WsStartupFailuresTotal {
378        name: "osfu_ws_startup_failures_total",
379        help: "Total websocket startup failures before the steady-state user loop.",
380        kind: Counter,
381        samples: |metrics, capture, output| write_counter_family(output, &metrics.ws_startup_failures, "kind")
382    },
383    WsHandshakeDurationSeconds {
384        name: "osfu_ws_handshake_duration_seconds",
385        help: "Websocket handshake duration from upgrade to user readiness or rejection.",
386        kind: Histogram,
387        samples: |metrics, capture, output| control_plane_histogram(output, &metrics.ws_handshake_duration)
388    },
389    WsAuthDurationSeconds {
390        name: "osfu_ws_auth_duration_seconds",
391        help: "Websocket authentication duration from first auth wait through token validation.",
392        kind: Histogram,
393        samples: |metrics, capture, output| control_plane_histogram(output, &metrics.ws_auth_duration)
394    },
395    WsUserInitializeDurationSeconds {
396        name: "osfu_ws_user_initialize_duration_seconds",
397        help: "Websocket user initialization duration after room admission.",
398        kind: Histogram,
399        samples: |metrics, capture, output| control_plane_histogram(output, &metrics.ws_user_initialize_duration)
400    },
401    WsUserLoopsStartedTotal {
402        name: "osfu_ws_user_loops_started_total",
403        help: "Total websocket user loops started after a successful join.",
404        kind: Counter,
405        samples: |metrics, capture, output| output.counter(&[], metrics.ws_user_loops_started.load())
406    },
407    WsUserLoopExitsTotal {
408        name: "osfu_ws_user_loop_exits_total",
409        help: "Total websocket user loop exits by reason.",
410        kind: Counter,
411        samples: |metrics, capture, output| write_counter_family(output, &metrics.ws_user_loop_exits, "reason")
412    },
413    WsBusBatchesTotal {
414        name: "osfu_ws_bus_batches_total",
415        help: "Total websocket signaling batches processed by direction.",
416        kind: Counter,
417        samples: |metrics, capture, output| write_counter_family(output, &metrics.ws_bus_batches, "direction")
418    },
419    WsBusEnvelopesTotal {
420        name: "osfu_ws_bus_envelopes_total",
421        help: "Total websocket signaling envelopes processed by direction.",
422        kind: Counter,
423        samples: |metrics, capture, output| write_counter_family(output, &metrics.ws_bus_envelopes, "direction")
424    },
425    WsBusParseFailuresTotal {
426        name: "osfu_ws_bus_parse_failures_total",
427        help: "Total websocket signaling parse failures.",
428        kind: Counter,
429        samples: |metrics, capture, output| output.counter(&[], metrics.ws_bus_parse_failures.load())
430    },
431    WsBusFailuresTotal {
432        name: "osfu_ws_bus_failures_total",
433        help: "Total websocket signaling failures by kind.",
434        kind: Counter,
435        samples: |metrics, capture, output| write_counter_family(output, &metrics.ws_bus_failures, "kind")
436    },
437    WsBusClientFramesTotal {
438        name: "osfu_ws_bus_client_frames_total",
439        help: "Total client websocket signaling frames by kind.",
440        kind: Counter,
441        samples: |metrics, capture, output| write_counter_family(output, &metrics.ws_bus_client_frames, "kind")
442    },
443    WsOutboundQueuedMessages {
444        name: "osfu_ws_outbound_queued_messages",
445        help: "Current websocket outbound room messages waiting in per-user queues.",
446        kind: Gauge,
447        samples: |metrics, capture, output| output.gauge(&[], metrics.ws_outbound_queued_messages.load())
448    },
449    WsOutboundQueueOverflowsTotal {
450        name: "osfu_ws_outbound_queue_overflows_total",
451        help: "Total websocket users marked for slow-consumer shutdown after outbound queue overflow.",
452        kind: Counter,
453        samples: |metrics, capture, output| output.counter(&[], metrics.ws_outbound_queue_overflows.load())
454    },
455    RoomsActive {
456        name: "osfu_rooms_active",
457        help: "Current number of active rooms owned by this runtime.",
458        kind: Gauge,
459        samples: |metrics, capture, output| output.gauge(&[], gauge_count(capture.room_gauges.rooms))
460    },
461    UsersActive {
462        name: "osfu_users_active",
463        help: "Current number of active room users owned by this runtime.",
464        kind: Gauge,
465        samples: |metrics, capture, output| output.gauge(&[], gauge_count(capture.room_gauges.users))
466    },
467    PublicationsActive {
468        name: "osfu_publications_active",
469        help: "Current number of committed or pending published media entries owned by this runtime.",
470        kind: Gauge,
471        samples: |metrics, capture, output| output.gauge(&[], gauge_count(capture.room_gauges.publications))
472    },
473    SubscriptionIntentEvictionsTotal {
474        name: "osfu_subscription_intent_evictions_total",
475        help: "Total absent publisher targets evicted from bounded receiver subscription intent.",
476        kind: Counter,
477        samples: |metrics, capture, output| output.counter(&[], metrics.subscription_intent_evictions.load())
478    },
479    SubscriptionsActive {
480        name: "osfu_subscriptions_active",
481        help: "Current number of committed or pending consumer subscriptions owned by this runtime.",
482        kind: Gauge,
483        samples: |metrics, capture, output| output.gauge(&[], gauge_count(capture.room_gauges.subscriptions))
484    },
485    TransportUsersActive {
486        name: "osfu_transport_users_active",
487        help: "Current number of active RTC transport users on this runtime.",
488        kind: Gauge,
489        samples: |metrics, capture, output| output.gauge(&[], metrics.active_transport_users.load())
490    },
491    RecordingActionsTotal {
492        name: "osfu_recording_actions_total",
493        help: "Total recording control actions by action and outcome.",
494        kind: Counter,
495        samples: |metrics, capture, output| write_label_pair_counter_family(output, &metrics.recording_actions)
496    },
497    RecordingRoomsActive {
498        name: "osfu_recording_rooms_active",
499        help: "Current number of rooms with an active recording user.",
500        kind: Gauge,
501        samples: |metrics, capture, output| output.gauge(&[], gauge_count(capture.room_gauges.recording_rooms))
502    },
503    RecordingCapturedPacketsTotal {
504        name: "osfu_recording_captured_packets_total",
505        help: "Total packets accepted by the recording capture path.",
506        kind: Counter,
507        samples: |metrics, capture, output| output.counter(&[], metrics.recording_captured_packets.load())
508    },
509    RecordingCapturedStreamsTotal {
510        name: "osfu_recording_captured_streams_total",
511        help: "Total unique media streams first seen by the recording capture path.",
512        kind: Counter,
513        samples: |metrics, capture, output| output.counter(&[], metrics.recording_captured_streams.load())
514    },
515    TransportHealthUsers {
516        name: "osfu_transport_health_users",
517        help: "Current number of transport users by observed health state.",
518        kind: Gauge,
519        samples: |metrics, capture, output| write_up_down_counter_family(output, &metrics.transport_health_users, "state")
520    },
521    TransportHealthTransitionsTotal {
522        name: "osfu_transport_health_transitions_total",
523        help: "Total transport health-state transitions observed from the transport adapter.",
524        kind: Counter,
525        samples: |metrics, capture, output| write_label_pair_counter_family(output, &metrics.transport_health_transitions)
526    },
527    RtpPacketsTotal {
528        name: "osfu_rtp_packets_total",
529        help: "Total RTP packets processed by flow direction.",
530        kind: Counter,
531        samples: |metrics, capture, output| write_snapshot_counters(output,
532            &capture.rtp.traffic,
533            "direction",
534            RtpTrafficSnapshot::packets
535        )
536    },
537    RtpPayloadBytesTotal {
538        name: "osfu_rtp_payload_bytes_total",
539        help: "Total RTP payload bytes processed by flow direction.",
540        kind: Counter,
541        samples: |metrics, capture, output| write_snapshot_counters(output,
542            &capture.rtp.traffic,
543            "direction",
544            RtpTrafficSnapshot::payload_bytes
545        )
546    },
547    RtpForwardedPacketsTotal {
548        name: "osfu_rtp_forwarded_packets_total",
549        help: "Total RTP packet fan-out operations by forwarding destination.",
550        kind: Counter,
551        samples: |metrics, capture, output| write_snapshot_counters(output,
552            &capture.rtp.traffic,
553            "destination",
554            RtpTrafficSnapshot::forwarded_packets
555        )
556    },
557    RtpForwardedPayloadBytesTotal {
558        name: "osfu_rtp_forwarded_payload_bytes_total",
559        help: "Total RTP payload bytes fanned out by forwarding destination.",
560        kind: Counter,
561        samples: |metrics, capture, output| write_snapshot_counters(output,
562            &capture.rtp.traffic,
563            "destination",
564            RtpTrafficSnapshot::forwarded_payload_bytes
565        )
566    },
567    WorkerRtpPacketsTotal {
568        name: "osfu_worker_rtp_packets_total",
569        help: "Total RTP packets processed by media worker and flow direction.",
570        kind: Counter,
571        samples: |metrics, capture, output| write_rtp_worker_counters(output,
572            &capture.rtp,
573            "direction",
574            RtpTrafficSnapshot::packets
575        )
576    },
577    WorkerRtpPayloadBytesTotal {
578        name: "osfu_worker_rtp_payload_bytes_total",
579        help: "Total RTP payload bytes processed by media worker and flow direction.",
580        kind: Counter,
581        samples: |metrics, capture, output| write_rtp_worker_counters(output,
582            &capture.rtp,
583            "direction",
584            RtpTrafficSnapshot::payload_bytes
585        )
586    },
587    WorkerRtpForwardedPacketsTotal {
588        name: "osfu_worker_rtp_forwarded_packets_total",
589        help: "Total RTP packet fan-out operations by media worker and forwarding destination.",
590        kind: Counter,
591        samples: |metrics, capture, output| write_rtp_worker_counters(output,
592            &capture.rtp,
593            "destination",
594            RtpTrafficSnapshot::forwarded_packets
595        )
596    },
597    WorkerRtpForwardedPayloadBytesTotal {
598        name: "osfu_worker_rtp_forwarded_payload_bytes_total",
599        help: "Total RTP payload bytes fanned out by media worker and forwarding destination.",
600        kind: Counter,
601        samples: |metrics, capture, output| write_rtp_worker_counters(output,
602            &capture.rtp,
603            "destination",
604            RtpTrafficSnapshot::forwarded_payload_bytes
605        )
606    },
607    RtpRelayOverloadDropsTotal {
608        name: "osfu_rtp_relay_overload_drops_total",
609        help: "Total RTP relay packets dropped because the bounded relay mailbox was full.",
610        kind: Counter,
611        samples: |metrics, capture, output| write_counter_family(output, &metrics.rtp_relay_overload_drops, "destination")
612    },
613    RtpDecoderRefreshesTotal {
614        name: "osfu_rtp_decoder_refreshes_total",
615        help: "Total decoder-refresh RTP packets observed by source scope.",
616        kind: Counter,
617        samples: |metrics, capture, output| write_snapshot_counters(output,
618            &capture.rtp,
619            "scope",
620            RtpMetricsSnapshot::decoder_refreshes
621        )
622    },
623    TransportIceStateChangesTotal {
624        name: "osfu_transport_ice_state_changes_total",
625        help: "Total RTC ICE state-change events observed from the transport adapter.",
626        kind: Counter,
627        samples: |metrics, capture, output| write_counter_family(output, &metrics.transport_ice_state_changes, "state")
628    },
629    TransportDtlsConnectedTotal {
630        name: "osfu_transport_dtls_connected_total",
631        help: "Total RTC DTLS-connected events observed from the transport adapter.",
632        kind: Counter,
633        samples: |metrics, capture, output| output.counter(&[], metrics.transport_dtls_connected.load())
634    },
635    TransportUserLifetimeSeconds {
636        name: "osfu_transport_user_lifetime_seconds",
637        help: "Lifetime of closed RTC transport users observed at cold-path teardown.",
638        kind: Histogram,
639        samples: |metrics, capture, output| output.histogram(
640            &[],
641            |bucket| metrics.transport_user_lifetime_buckets.load(bucket),
642            || metrics.transport_user_lifetime_count.load(),
643            || metrics.transport_user_lifetime_sum_micros.load(),
644        )
645    },
646    MediaQualitySamplesTotal {
647        name: "osfu_media_quality_samples_total",
648        help: "Total sampled transport-quality events by str0m stats source.",
649        kind: Counter,
650        samples: |metrics, capture, output| write_counter_family(output, &metrics.media_quality_samples, "sample")
651    },
652    MediaQualityRttSeconds {
653        name: "osfu_media_quality_rtt_seconds",
654        help: "Sampled RTC round-trip time by str0m stats source.",
655        kind: Histogram,
656        samples: |metrics, capture, output| write_histogram_family(output, &metrics.media_quality_rtt, "sample")
657    },
658    MediaQualityLossPpmObservedTotal {
659        name: "osfu_media_quality_loss_ppm_observed_total",
660        help: "Sum of sampled packet loss observations in parts per million by direction.",
661        kind: Counter,
662        samples: |metrics, capture, output| write_counter_family(output, &metrics.media_quality_loss_ppm_observed, "direction")
663    },
664    MediaQualityLossObservationsTotal {
665        name: "osfu_media_quality_loss_observations_total",
666        help: "Total packet loss observations by direction.",
667        kind: Counter,
668        samples: |metrics, capture, output| write_counter_family(output, &metrics.media_quality_loss_observations, "direction")
669    },
670    MediaQualityBweBpsObservedTotal {
671        name: "osfu_media_quality_bwe_bps_observed_total",
672        help: "Sum of sampled peer egress bandwidth estimates in bits per second.",
673        kind: Counter,
674        samples: |metrics, capture, output| output.counter(&[], metrics.media_quality_bwe_bps_observed.load())
675    },
676    MediaQualityBweObservationsTotal {
677        name: "osfu_media_quality_bwe_observations_total",
678        help: "Total peer egress bandwidth estimate observations.",
679        kind: Counter,
680        samples: |metrics, capture, output| output.counter(&[], metrics.media_quality_bwe_observations.load())
681    },
682    MediaQualityJitterRtpTimestampUnitsObservedTotal {
683        name: "osfu_media_quality_jitter_rtp_timestamp_units_observed_total",
684        help: "Sum of sampled remote egress jitter observations in RTP timestamp units.",
685        kind: Counter,
686        samples: |metrics, capture, output| output.counter(
687            &[],
688            metrics.media_quality_jitter_rtp_timestamp_units_observed.load(),
689        )
690    },
691    MediaQualityJitterObservationsTotal {
692        name: "osfu_media_quality_jitter_observations_total",
693        help: "Total remote egress jitter observations.",
694        kind: Counter,
695        samples: |metrics, capture, output| output.counter(&[], metrics.media_quality_jitter_observations.load())
696    },
697    TransportCleanupFailuresTotal {
698        name: "osfu_transport_cleanup_failures_total",
699        help: "Total terminal transport cleanup failures.",
700        kind: Counter,
701        samples: |metrics, capture, output| counter(output,
702            [("kind", "terminal")],
703            metrics.transport_cleanup_failures.load()
704        )
705    },
706    RtcDatagramRoutesTotal {
707        name: "osfu_rtc_datagram_routes_total",
708        help: "Total RTC UDP datagrams accepted by routing path.",
709        kind: Counter,
710        samples: |metrics, capture, output| write_snapshot_counters(output,
711            &capture.rtc,
712            "path",
713            RtcMetricsSnapshot::datagram_routes
714        )
715    },
716    RtcDatagramDropsTotal {
717        name: "osfu_rtc_datagram_drops_total",
718        help: "Total RTC UDP datagrams dropped by ingress routing before session delivery.",
719        kind: Counter,
720        samples: |metrics, capture, output| write_snapshot_counters(output,
721            &capture.rtc,
722            "reason",
723            RtcMetricsSnapshot::datagram_drops
724        )
725    },
726    RtcDatagramFallbackScansTotal {
727        name: "osfu_rtc_datagram_fallback_scans_total",
728        help: "Total fallback scans across RTC users for UDP datagram routing.",
729        kind: Counter,
730        samples: |metrics, capture, output| output.counter(&[], capture.rtc.datagram_fallback_scans())
731    },
732    RtcDatagramScanUsersTotal {
733        name: "osfu_rtc_datagram_scan_users_total",
734        help: "Total RTC users examined by UDP fallback scans.",
735        kind: Counter,
736        samples: |metrics, capture, output| output.counter(&[], capture.rtc.datagram_scan_users())
737    },
738    RtcNacksTotal {
739        name: "osfu_rtc_nacks_total",
740        help: "Total Generic NACK feedback events by WebRTC direction.",
741        kind: Counter,
742        samples: |metrics, capture, output| write_snapshot_counters(output,
743            &capture.rtc,
744            "direction",
745            RtcMetricsSnapshot::nacks
746        )
747    },
748    RtcRtxPacketsTotal {
749        name: "osfu_rtc_rtx_packets_total",
750        help: "Total authenticated RTX packets accepted from publishers.",
751        kind: Counter,
752        samples: |metrics, capture, output| counter(output,
753            [("direction", "received_from_publisher")],
754            capture.rtc.rtx_packets_received_from_publisher()
755        )
756    },
757    RtcRtxPayloadBytesTotal {
758        name: "osfu_rtc_rtx_payload_bytes_total",
759        help: "Total de-RTX media payload bytes accepted from publishers.",
760        kind: Counter,
761        samples: |metrics, capture, output| counter(output,
762            [("direction", "received_from_publisher")],
763            capture.rtc.rtx_payload_bytes_received_from_publisher()
764        )
765    },
766    RtcRtcpIngressBudgetDropsTotal {
767        name: "osfu_rtc_rtcp_ingress_budget_drops_total",
768        help: "Total candidate RTCP datagrams dropped by the per-session ingress byte budget.",
769        kind: Counter,
770        samples: |metrics, capture, output| output.counter(&[], capture.rtc.rtcp_ingress_budget_drops())
771    },
772    RtcTransportIoFailuresTotal {
773        name: "osfu_rtc_transport_io_failures_total",
774        help: "Total RTC UDP socket failures by direction and bounded category.",
775        kind: Counter,
776        samples: |metrics, capture, output| write_label_pair_counters::<RtcTransportIoFailure>(
777            output, |failure| capture.rtc.transport_io_failures(failure)
778        )
779    },
780    RtcInputFailuresTotal {
781        name: "osfu_rtc_input_failures_total",
782        help: "Total RTC datagram input failures by bounded category.",
783        kind: Counter,
784        samples: |metrics, capture, output| write_snapshot_counters(output,
785            &capture.rtc, "category", RtcMetricsSnapshot::input_failures
786        )
787    },
788    RtcDrainFailuresTotal {
789        name: "osfu_rtc_drain_failures_total",
790        help: "Total terminal RTC session drain failures by stage.",
791        kind: Counter,
792        samples: |metrics, capture, output| write_snapshot_counters(output,
793            &capture.rtc, "stage", RtcMetricsSnapshot::drain_failures
794        )
795    },
796    RtcWorkerTerminalFailuresTotal {
797        name: "osfu_rtc_worker_terminal_failures_total",
798        help: "Total media workers retired after unexpected terminal failure.",
799        kind: Counter,
800        samples: |metrics, capture, output| output.counter(&[], capture.rtc.worker_terminal_failures())
801    },
802    RtcWorkerObservationTimeoutsTotal {
803        name: "osfu_rtc_worker_observation_timeouts_total",
804        help: "Total bounded media worker observations that timed out by command kind.",
805        kind: Counter,
806        samples: |metrics, capture, output| write_snapshot_counters(output,
807            &capture.rtc, "kind", RtcMetricsSnapshot::worker_observation_timeouts
808        )
809    },
810    RtcOutputBudgetExhaustionsTotal {
811        name: "osfu_rtc_output_budget_exhaustions_total",
812        help: "Total RTC session drains that exhausted the output budget by limit.",
813        kind: Counter,
814        samples: |metrics, capture, output| write_snapshot_counters(output,
815            &capture.rtc,
816            "limit",
817            RtcMetricsSnapshot::output_budget_exhaustions
818        )
819    },
820    RtcOutputBudgetSessionClosesTotal {
821        name: "osfu_rtc_output_budget_session_closes_total",
822        help: "Total RTC sessions closed after output-budget exhaustion.",
823        kind: Counter,
824        samples: |metrics, capture, output| output.counter(&[], capture.rtc.output_budget_session_closes())
825    },
826    RtcProducerSsrcBindingsTotal {
827        name: "osfu_rtc_producer_ssrc_bindings_total",
828        help: "Total packet-triggered producer SSRC binding changes and rejections by outcome.",
829        kind: Counter,
830        samples: |metrics, capture, output| write_snapshot_counters(output,
831            &capture.rtc,
832            "outcome",
833            RtcMetricsSnapshot::producer_ssrc_bindings
834        )
835    },
836    RtcRouteControlTotal {
837        name: "osfu_rtc_route_control_total",
838        help: "Total RTC route-control decisions observed at the transport boundary.",
839        kind: Counter,
840        samples: |metrics, capture, output| write_snapshot_counters(output,
841            &capture.rtc,
842            "outcome",
843            RtcMetricsSnapshot::route_control
844        )
845    },
846    RtcKeyframeRequestsTotal {
847        name: "osfu_rtc_keyframe_requests_total",
848        help: "Total RTC keyframe request tracker outcomes.",
849        kind: Counter,
850        samples: |metrics, capture, output| write_snapshot_counters(output,
851            &capture.rtc,
852            "outcome",
853            RtcMetricsSnapshot::keyframe_requests
854        )
855    },
856    RtcRelayEnqueuesTotal {
857        name: "osfu_rtc_relay_enqueues_total",
858        help: "Total relay enqueue attempts by target kind and outcome.",
859        kind: Counter,
860        samples: |metrics, capture, output| write_label_pair_counters(output, |result: RtcRelayEnqueueResult| {
861            capture.rtc.relay_enqueues(result)
862        })
863    },
864    RtcRelayMailboxDepthSamplesTotal {
865        name: "osfu_rtc_relay_mailbox_depth_samples_total",
866        help: "Total sampled intra-node relay mailbox depth observations.",
867        kind: Counter,
868        samples: |metrics, capture, output| output.counter(&[], capture.rtc.relay_mailbox_depth_samples())
869    },
870    RtcRelayMailboxDepthObservedTotal {
871        name: "osfu_rtc_relay_mailbox_depth_observed_total",
872        help: "Sum of sampled intra-node relay mailbox depths.",
873        kind: Counter,
874        samples: |metrics, capture, output| output.counter(&[], capture.rtc.relay_mailbox_depth_total())
875    },
876    RtcRelayDrainBatchesTotal {
877        name: "osfu_rtc_relay_drain_batches_total",
878        help: "Total non-empty packet-loop relay drain batches.",
879        kind: Counter,
880        samples: |metrics, capture, output| output.counter(&[], capture.rtc.relay_drain_batches())
881    },
882    RtcRelayDrainedPacketsTotal {
883        name: "osfu_rtc_relay_drained_packets_total",
884        help: "Total relay packets drained into packet-loop batches.",
885        kind: Counter,
886        samples: |metrics, capture, output| output.counter(&[], capture.rtc.relay_drained_packets())
887    },
888    RtcRelayDrainCapHitsTotal {
889        name: "osfu_rtc_relay_drain_cap_hits_total",
890        help: "Total relay drain batches that left queued relay packets behind after hitting the per-turn cap.",
891        kind: Counter,
892        samples: |metrics, capture, output| output.counter(&[], capture.rtc.relay_drain_cap_hits())
893    },
894    RtcRemoteControlDropsTotal {
895        name: "osfu_rtc_remote_control_drops_total",
896        help: "Total remote-source control commands dropped before enqueue by command kind.",
897        kind: Counter,
898        samples: |metrics, capture, output| write_snapshot_counters(output,
899            &capture.rtc,
900            "kind",
901            RtcMetricsSnapshot::remote_control_drops
902        )
903    },
904    RtcRemotePacketGateConvergenceTotal {
905        name: "osfu_rtc_remote_packet_gate_convergence_total",
906        help: "Total remote packet-gate convergence retry attempts and successful pending flushes.",
907        kind: Counter,
908        samples: |metrics, capture, output| write_snapshot_counters(output,
909            &capture.rtc,
910            "outcome",
911            RtcMetricsSnapshot::remote_packet_gate_convergence
912        )
913    },
914    SourceSelectionUpdatesTotal {
915        name: "osfu_source_selection_updates_total",
916        help: "Total room-scoped source selector updates accepted by source policy.",
917        kind: Counter,
918        samples: |metrics, capture, output| write_counter_family(output, &metrics.source_selection_updates, "selector")
919    },
920    BudgetSolverOutcomesTotal {
921        name: "osfu_budget_solver_outcomes_total",
922        help: "Total committed receiver video route degradation, pause and resume transitions.",
923        kind: Counter,
924        samples: |metrics, capture, output| write_counter_family(output, &metrics.budget_solver_outcomes, "outcome")
925    },
926}
927
928fn write_counter_family<L>(
929    output: &mut MetricOutput<'_, '_>,
930    family: &CounterFamily<L>,
931    label_name: &'static str,
932) where
933    L: ExportedMetricLabel,
934{
935    for label in L::VARIANTS {
936        counter(
937            output,
938            [(label_name, label.label_value())],
939            family.load(*label),
940        );
941    }
942}
943
944fn write_label_pair_counter_family<L>(output: &mut MetricOutput<'_, '_>, family: &CounterFamily<L>)
945where
946    L: ExportedMetricLabelPair,
947{
948    write_label_pair_counters(output, |label| family.load(label));
949}
950
951fn write_label_pair_counters<L>(output: &mut MetricOutput<'_, '_>, load: impl Fn(L) -> u64)
952where
953    L: ExportedMetricLabelPair,
954{
955    for label in L::VARIANTS {
956        counter(output, label.label_pair(), load(*label));
957    }
958}
959
960fn write_snapshot_counters<S, L>(
961    output: &mut MetricOutput<'_, '_>,
962    snapshot: &S,
963    label_name: &'static str,
964    read: fn(&S, L) -> u64,
965) where
966    L: ExportedMetricLabel,
967{
968    for label in L::VARIANTS {
969        counter(
970            output,
971            [(label_name, label.label_value())],
972            read(snapshot, *label),
973        );
974    }
975}
976
977fn write_rtp_worker_counters<L>(
978    output: &mut MetricOutput<'_, '_>,
979    snapshot: &RtpMetricsSnapshot,
980    label_name: &'static str,
981    read: fn(&RtpTrafficSnapshot, L) -> u64,
982) where
983    L: ExportedMetricLabel,
984{
985    for worker in snapshot.worker_snapshots() {
986        for label in L::VARIANTS {
987            output.counter(
988                &[
989                    MetricLabel::number("media_worker_id", worker.media_worker_id()),
990                    MetricLabel::text(label_name, label.label_value()),
991                ],
992                read(&worker.traffic, *label),
993            );
994        }
995    }
996}
997
998fn write_up_down_counter_family<L>(
999    output: &mut MetricOutput<'_, '_>,
1000    family: &UpDownCounterFamily<L>,
1001    label_name: &'static str,
1002) where
1003    L: ExportedMetricLabel,
1004{
1005    for label in L::VARIANTS {
1006        output.gauge(
1007            &[MetricLabel::text(label_name, label.label_value())],
1008            family.load(*label),
1009        );
1010    }
1011}
1012
1013fn write_histogram_family<L, B>(
1014    output: &mut MetricOutput<'_, '_>,
1015    family: &HistogramFamily<L, B>,
1016    label_name: &'static str,
1017) where
1018    L: ExportedMetricLabel,
1019    B: HistogramBucketLabel,
1020{
1021    for label in L::VARIANTS {
1022        output.histogram(
1023            &[MetricLabel::text(label_name, label.label_value())],
1024            |bucket| family.load_bucket(*label, bucket),
1025            || family.load_count(*label),
1026            || family.load_sum_micros(*label),
1027        );
1028    }
1029}
1030
1031fn counter<const N: usize>(
1032    output: &mut MetricOutput<'_, '_>,
1033    labels: [(&'static str, &'static str); N],
1034    value: u64,
1035) {
1036    output.counter(
1037        &labels.map(|(name, value)| MetricLabel::text(name, value)),
1038        value,
1039    );
1040}
1041
1042fn control_plane_histogram(
1043    output: &mut MetricOutput<'_, '_>,
1044    histogram: &Histogram<ControlPlaneDurationBucket>,
1045) {
1046    output.histogram(
1047        &[],
1048        |bucket| histogram.load_bucket(bucket),
1049        || histogram.load_count(),
1050        || histogram.load_sum_micros(),
1051    );
1052}
1053
1054fn append_sample_name(output: &mut String, name: &str, labels: &[MetricLabel]) {
1055    output.push_str(name);
1056    append_labels(output, labels, None);
1057}
1058
1059fn append_labels(output: &mut String, labels: &[MetricLabel], extra_label: Option<MetricLabel>) {
1060    if labels.is_empty() && extra_label.is_none() {
1061        return;
1062    }
1063    output.push('{');
1064    for (index, label) in labels.iter().chain(extra_label.iter()).enumerate() {
1065        if index != 0 {
1066            output.push(',');
1067        }
1068        output.push_str(label.name);
1069        output.push_str("=\"");
1070        match label.value {
1071            MetricLabelValue::Text(value) => output.push_str(value),
1072            MetricLabelValue::Number(value) => {
1073                let _ = write!(output, "{value}");
1074            }
1075        }
1076        output.push('"');
1077    }
1078    output.push('}');
1079}
1080
1081fn append_seconds_from_micros(output: &mut String, micros: u64) {
1082    let whole_seconds = micros / 1_000_000;
1083    let fractional_micros = micros % 1_000_000;
1084    let _ = write!(output, "{whole_seconds}");
1085    if fractional_micros == 0 {
1086        output.push_str(".0");
1087        return;
1088    }
1089    let _ = write!(output, ".{fractional_micros:06}");
1090    while output.ends_with('0') {
1091        output.pop();
1092    }
1093}