saluki_app/metrics/
mod.rs1use 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
19pub(crate) async fn initialize_metrics(
33 metrics_prefix: impl Into<String>, default_level: Level,
34) -> Result<Supervisor, GenericError> {
35 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 let reflector = initialize_shared_metrics_state()?;
48
49 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 supervisor.add_worker(runtime::supervisable(reflector).transient().build());
61
62 Ok(supervisor)
63}
64
65pub fn emit_startup_metrics() {
72 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 let running = RUNNING.get_or_init(|| gauge!("running", "version" => app_version));
84 running.set(1.0);
85}
86
87pub async fn collect_runtime_metrics(runtime_id: &str, handle: Handle) {
92 let runtime_metrics = RuntimeMetrics::with_workers(runtime_id, handle.metrics().num_workers());
94
95 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
104pub struct RuntimeMetricsWorker {
109 runtime_id: String,
110 handle: Handle,
111}
112
113impl RuntimeMetricsWorker {
114 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}