commonware_runtime/tokio/
telemetry.rs1use 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
21pub struct Logs {
23 pub level: Level,
25
26 pub json: bool,
32}
33
34pub fn init(context: Context, logs: Logs, metrics: Option<SocketAddr>, traces: Option<Config>) {
41 let log_layer = tracing_subscriber::fmt::layer()
43 .with_line_number(true)
44 .with_thread_ids(true)
45 .with_file(true);
46
47 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 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 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 tracing::subscriber::set_global_default(registry).expect("Failed to set subscriber");
87
88 if let Some(cfg) = metrics {
90 context.child("metrics").spawn(move |context| async move {
91 let listener = TcpListener::bind(cfg)
97 .await
98 .expect("Failed to bind metrics server");
99
100 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(listener, app.into_make_service())
120 .await
121 .expect("Could not serve metrics");
122 });
123 }
124}