saluki_app/metrics/
mod.rs

1//! Metrics.
2
3use std::{sync::OnceLock, time::Duration};
4
5use async_trait::async_trait;
6use metrics::{gauge, Gauge, Level};
7use saluki_common::sync::shutdown::ShutdownHandle;
8use saluki_core::{
9    observability::metrics::initialize_shared_metrics_state,
10    runtime::{self, InitializationError, Supervisable, Supervisor, SupervisorFuture},
11};
12use saluki_error::{ErrorContext as _, GenericError};
13use saluki_metrics::static_metrics;
14use tokio::{runtime::Handle, select, time::sleep};
15
16mod api;
17pub use self::api::{MetricsAPIHandler, MetricsOverrideWorker};
18
19/// Initializes the metrics subsystem for `metrics`.
20///
21/// The given prefix is used to namespace all metrics that are emitted by the application, and is prepended to all
22/// metrics, followed by a period (for example, `<prefix>.<metric name>`). The given default level seeds the runtime
23/// filter and is what the filter is restored to when [`MetricsAPIHandler`]'s reset route is invoked.
24///
25/// Returns a [`Supervisor`] holding the subsystem's background workers. The caller must arrange for it to run --
26/// typically by adding it to a parent supervisor -- or internal metrics are never flushed to subscribers.
27///
28/// # Errors
29///
30/// If the metrics subsystem was already initialized, or its supervisor can't be constructed, an error will be
31/// returned.
32pub(crate) async fn initialize_metrics(
33    metrics_prefix: impl Into<String>, default_level: Level,
34) -> Result<Supervisor, GenericError> {
35    // We forward to the implementation in `saluki_core` so that we can have this crate be the collection point of all
36    // helpers/types that are specific to generic application setup/initialization.
37    //
38    // The implementation itself has to live in `saluki_core`, however, to have access to all of the underlying types
39    // that are created and used to install the global recorder, such that they need not be exposed publicly.
40    let (filter_handle, flusher) =
41        saluki_core::observability::metrics::initialize_metrics(metrics_prefix.into(), default_level).await?;
42
43    let override_processor = MetricsOverrideWorker::new(filter_handle);
44
45    // Subscribe to the registry now rather than when the worker starts, so that flushes landing between here and the
46    // supervisor running are still folded into the shared state.
47    let reflector = initialize_shared_metrics_state()?;
48
49    // Capture the current runtime handle eagerly so the runtime metrics worker measures the runtime that owns
50    // bootstrap, regardless of where the worker future eventually executes under the supervisor.
51    let runtime = RuntimeMetricsWorker::new("primary", Handle::current());
52
53    let mut supervisor = Supervisor::new("metrics").error_context("Failed to construct metrics supervisor.")?;
54    supervisor.add_worker(flusher);
55    supervisor.add_worker(runtime);
56    supervisor.add_worker(override_processor);
57
58    // Transient rather than permanent: the reflector worker treats an exhausted source as terminal and returns
59    // normally, so restarting it would only hand it the same dead source again.
60    supervisor.add_worker(runtime::supervisable(reflector).transient().build());
61
62    Ok(supervisor)
63}
64
65/// Emits the startup metrics for the application.
66///
67/// This is generally meant to be called after the application has been initialized, in order to indicate the
68/// application has completed start-up and is now running.
69///
70/// Must be called after the metrics subsystem has been initialized.
71pub fn emit_startup_metrics() {
72    // We hold the handle for the life of the process so it doesn't get idle reaped.
73    static RUNNING: OnceLock<Gauge> = OnceLock::new();
74
75    let app_details = saluki_metadata::get_app_details();
76    let app_version = if app_details.is_dev_build() {
77        format!("{}-dev-{}", app_details.version().raw(), app_details.git_hash(),)
78    } else {
79        app_details.version().raw().to_string()
80    };
81
82    // Emit a "running" metric to indicate that the application is running.
83    let running = RUNNING.get_or_init(|| gauge!("running", "version" => app_version));
84    running.set(1.0);
85}
86
87/// Collects Tokio runtime metrics from the given runtime handle.
88///
89/// All metrics generated will include a `runtime_id` label which maps to the given runtime ID. This allows for
90/// differentiating between multiple runtimes that may be running in the same process.
91pub async fn collect_runtime_metrics(runtime_id: &str, handle: Handle) {
92    // Grab the total number of runtime workers to properly initialize/register our metrics.
93    let runtime_metrics = RuntimeMetrics::with_workers(runtime_id, handle.metrics().num_workers());
94
95    // With our metrics registered, enter the main loop where we periodically scrape the metrics.
96    loop {
97        let latest_runtime_metrics = handle.metrics();
98        runtime_metrics.update(&latest_runtime_metrics);
99
100        sleep(Duration::from_secs(5)).await;
101    }
102}
103
104/// A worker that periodically collects Tokio runtime metrics.
105///
106/// The runtime is captured at construction time so that the metrics measured always describe the runtime that
107/// owned the bootstrap, regardless of where this worker eventually executes under a supervisor.
108pub struct RuntimeMetricsWorker {
109    runtime_id: String,
110    handle: Handle,
111}
112
113impl RuntimeMetricsWorker {
114    /// Creates a new `RuntimeMetricsWorker` that collects metrics from the given runtime.
115    pub fn new<S: Into<String>>(runtime_id: S, handle: Handle) -> Self {
116        Self {
117            runtime_id: runtime_id.into(),
118            handle,
119        }
120    }
121}
122
123#[async_trait]
124impl Supervisable for RuntimeMetricsWorker {
125    fn name(&self) -> &str {
126        "tokio-runtime-metrics-collector"
127    }
128
129    async fn initialize(&self, process_shutdown: ShutdownHandle) -> Result<SupervisorFuture, InitializationError> {
130        let runtime_id = self.runtime_id.clone();
131        let handle = self.handle.clone();
132
133        Ok(Box::pin(async move {
134            select! {
135                _ = process_shutdown => {},
136                _ = collect_runtime_metrics(&runtime_id, handle) => {},
137            }
138
139            Ok(())
140        }))
141    }
142}
143
144#[static_metrics(prefix = runtime_worker, labels(runtime_id, worker_idx))]
145#[derive(Clone)]
146struct WorkerMetrics {
147    #[metric(level = trace)]
148    local_queue_depth: Gauge,
149    #[metric(level = trace)]
150    local_schedule_count: Gauge,
151    #[metric(level = trace)]
152    mean_poll_time: Gauge,
153    #[metric(level = trace)]
154    noop_count: Gauge,
155    #[metric(level = trace)]
156    overflow_count: Gauge,
157    #[metric(level = trace)]
158    park_count: Gauge,
159    #[metric(level = trace)]
160    park_unpark_count: Gauge,
161    #[metric(level = trace)]
162    poll_count: Gauge,
163    #[metric(level = trace)]
164    steal_count: Gauge,
165    #[metric(level = trace)]
166    steal_operations: Gauge,
167    #[metric(level = trace)]
168    total_busy_duration: Gauge,
169}
170
171impl WorkerMetrics {
172    fn with_worker_idx(runtime_id: &str, worker_idx: usize) -> Self {
173        Self::new(runtime_id, worker_idx)
174    }
175
176    fn update(&self, worker_idx: usize, metrics: &tokio::runtime::RuntimeMetrics) {
177        self.local_queue_depth()
178            .set(metrics.worker_local_queue_depth(worker_idx) as f64);
179        self.local_schedule_count()
180            .set(metrics.worker_local_schedule_count(worker_idx) as f64);
181        self.mean_poll_time()
182            .set(metrics.worker_mean_poll_time(worker_idx).as_nanos() as f64);
183        self.noop_count().set(metrics.worker_noop_count(worker_idx) as f64);
184        self.overflow_count()
185            .set(metrics.worker_overflow_count(worker_idx) as f64);
186        self.park_count().set(metrics.worker_park_count(worker_idx) as f64);
187        self.park_unpark_count()
188            .set(metrics.worker_park_unpark_count(worker_idx) as f64);
189        self.poll_count().set(metrics.worker_poll_count(worker_idx) as f64);
190        self.steal_count().set(metrics.worker_steal_count(worker_idx) as f64);
191        self.steal_operations()
192            .set(metrics.worker_steal_operations(worker_idx) as f64);
193        self.total_busy_duration()
194            .set(metrics.worker_total_busy_duration(worker_idx).as_nanos() as f64);
195    }
196}
197
198#[static_metrics(prefix = runtime, labels(runtime_id))]
199#[derive(Clone)]
200struct GlobalRuntimeMetrics {
201    #[metric(level = debug)]
202    num_alive_tasks: Gauge,
203    #[metric(level = debug)]
204    blocking_queue_depth: Gauge,
205    #[metric(level = debug)]
206    budget_forced_yield_count: Gauge,
207    #[metric(level = debug)]
208    global_queue_depth: Gauge,
209    #[metric(level = debug)]
210    io_driver_fd_deregistered_count: Gauge,
211    #[metric(level = debug)]
212    io_driver_fd_registered_count: Gauge,
213    #[metric(level = debug)]
214    io_driver_ready_count: Gauge,
215    #[metric(level = debug)]
216    num_blocking_threads: Gauge,
217    #[metric(level = debug)]
218    num_idle_blocking_threads: Gauge,
219    #[metric(level = debug)]
220    num_workers: Gauge,
221    #[metric(level = debug)]
222    remote_schedule_count: Gauge,
223    #[metric(level = debug)]
224    spawned_tasks_count: Gauge,
225}
226
227struct RuntimeMetrics {
228    global: GlobalRuntimeMetrics,
229    workers: Vec<WorkerMetrics>,
230}
231
232impl RuntimeMetrics {
233    fn with_workers(runtime_id: &str, workers_len: usize) -> Self {
234        let mut workers = Vec::with_capacity(workers_len);
235        for i in 0..workers_len {
236            workers.push(WorkerMetrics::with_worker_idx(runtime_id, i));
237        }
238
239        Self {
240            global: GlobalRuntimeMetrics::new(runtime_id),
241            workers,
242        }
243    }
244
245    fn update(&self, metrics: &tokio::runtime::RuntimeMetrics) {
246        self.global.num_alive_tasks().set(metrics.num_alive_tasks() as f64);
247        self.global
248            .blocking_queue_depth()
249            .set(metrics.blocking_queue_depth() as f64);
250        self.global
251            .budget_forced_yield_count()
252            .set(metrics.budget_forced_yield_count() as f64);
253        self.global
254            .global_queue_depth()
255            .set(metrics.global_queue_depth() as f64);
256        self.global
257            .io_driver_fd_deregistered_count()
258            .set(metrics.io_driver_fd_deregistered_count() as f64);
259        self.global
260            .io_driver_fd_registered_count()
261            .set(metrics.io_driver_fd_registered_count() as f64);
262        self.global
263            .io_driver_ready_count()
264            .set(metrics.io_driver_ready_count() as f64);
265        self.global
266            .num_blocking_threads()
267            .set(metrics.num_blocking_threads() as f64);
268        self.global
269            .num_idle_blocking_threads()
270            .set(metrics.num_idle_blocking_threads() as f64);
271        self.global.num_workers().set(metrics.num_workers() as f64);
272        self.global
273            .remote_schedule_count()
274            .set(metrics.remote_schedule_count() as f64);
275        self.global
276            .spawned_tasks_count()
277            .set(metrics.spawned_tasks_count() as f64);
278
279        for (worker_idx, worker) in self.workers.iter().enumerate() {
280            worker.update(worker_idx, metrics);
281        }
282    }
283}