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 #[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#[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}