Skip to main content

commonware_runtime/tokio/
telemetry.rs

1//! Utilities for collecting and reporting telemetry data.
2
3use super::{
4    Context,
5    tracing::{Config, export},
6};
7use crate::{Metrics as _, Spawner, Supervisor as _};
8use axum::{
9    Extension, Router,
10    body::Body,
11    http::{Response, StatusCode, header},
12    routing::get,
13    serve,
14};
15use cfg_if::cfg_if;
16use std::{net::SocketAddr, sync::Arc};
17use tokio::net::TcpListener;
18use tracing::Level;
19use tracing_subscriber::{Layer, Registry, filter::filter_fn, layer::SubscriberExt};
20
21/// Log configuration.
22pub struct Logs {
23    /// The level of logging to use.
24    pub level: Level,
25
26    /// Whether to log in JSON format.
27    ///
28    /// This is useful for structured logging in server-based environments.
29    /// If you are running things locally, it is recommended to use
30    /// `json = false` to get a human-readable format.
31    pub json: bool,
32}
33
34/// Initialize telemetry with the given configuration.
35///
36/// If `metrics` is provided, starts serving metrics at the given address at `/metrics`.
37/// If `traces` is provided, enables OpenTelemetry trace export. When `None`, span
38/// callsites in the runtime-provided subscriber are disabled and only log events
39/// are emitted (per [`Logs::level`]).
40pub fn init(context: Context, logs: Logs, metrics: Option<SocketAddr>, traces: Option<Config>) {
41    // Create fmt layer for logging
42    let log_layer = tracing_subscriber::fmt::layer()
43        .with_line_number(true)
44        .with_thread_ids(true)
45        .with_file(true);
46
47    // Set the format to JSON (if specified)
48    let log_layer = if logs.json {
49        log_layer.json().boxed()
50    } else {
51        log_layer.compact().boxed()
52    };
53    let log_layer = match &traces {
54        None => log_layer
55            .with_filter(filter_fn(move |metadata| {
56                metadata.is_event() && *metadata.level() <= logs.level
57            }))
58            .boxed(),
59        Some(_) => log_layer
60            .with_filter(tracing_subscriber::EnvFilter::new(logs.level.to_string()))
61            .boxed(),
62    };
63
64    // Create OpenTelemetry layer for tracing
65    let trace_layer = traces.map(|cfg| {
66        let tracer = export(cfg).expect("Failed to initialize tracer");
67        tracing_opentelemetry::layer()
68            .with_tracer(tracer)
69            .with_filter(tracing_subscriber::EnvFilter::new(logs.level.to_string()))
70    });
71
72    // Create tracing registry.
73    cfg_if! {
74        if #[cfg(feature = "tokio-console")] {
75            let console_layer = console_subscriber::spawn();
76            let registry = Registry::default()
77                .with(log_layer)
78                .with(trace_layer)
79                .with(console_layer);
80        } else {
81            let registry = Registry::default().with(log_layer).with(trace_layer);
82        }
83    }
84
85    // Set the global subscriber
86    tracing::subscriber::set_global_default(registry).expect("Failed to set subscriber");
87
88    // Expose metrics over HTTP
89    if let Some(cfg) = metrics {
90        context.child("metrics").spawn(move |context| async move {
91            // Create a tokio listener for the metrics server.
92            //
93            // We explicitly avoid using a runtime `Listener` because
94            // it will track bandwidth used for metrics and apply a policy
95            // for read/write timeouts fit for a p2p network.
96            let listener = TcpListener::bind(cfg)
97                .await
98                .expect("Failed to bind metrics server");
99
100            // Create a router for the metrics server
101            let shared = Arc::new(context);
102            let app = Router::new()
103                .route(
104                    "/metrics",
105                    get(|Extension(ctx): Extension<Arc<Context>>| async move {
106                        Response::builder()
107                            .status(StatusCode::OK)
108                            .header(header::CONTENT_TYPE, "text/plain; version=0.0.4")
109                            .body(Body::from(ctx.encode()))
110                            .expect("Failed to create response")
111                    }),
112                )
113                .layer(Extension(shared));
114
115            // Serve the metrics over HTTP.
116            //
117            // `serve` will spawn its own tasks using `tokio::spawn` (and there is no way to specify
118            // it to do otherwise). These tasks will not be tracked like metrics spawned using `Spawner`.
119            serve(listener, app.into_make_service())
120                .await
121                .expect("Could not serve metrics");
122        });
123    }
124}