1use std::sync::Arc;
4use std::time::Duration;
5
6use axum::http::header::CONTENT_TYPE;
7use axum::http::{HeaderValue, StatusCode};
8use axum::response::{IntoResponse, Response};
9use prometheus::{
10 Encoder, HistogramOpts, HistogramVec, IntCounter, IntCounterVec, IntGaugeVec, Opts, Registry,
11 TextEncoder,
12};
13use thiserror::Error;
14
15const TEXT_FORMAT: &str = "text/plain; version=0.0.4; charset=utf-8";
16
17#[derive(Debug, Error)]
19pub enum MetricsError {
20 #[error("failed to register prometheus metric: {0}")]
22 Register(#[from] prometheus::Error),
23 #[error("failed to encode prometheus metrics: {0}")]
25 Encode(String),
26}
27
28#[derive(Clone, Debug)]
30pub struct Metrics {
31 inner: Arc<MetricsInner>,
32}
33
34#[derive(Debug)]
35struct MetricsInner {
36 registry: Registry,
37 workflows_started: IntCounterVec,
38 workflows_completed: IntCounterVec,
39 workflows_reopened: IntCounterVec,
40 activities_dispatched: IntCounterVec,
41 activities_completed: IntCounterVec,
42 activity_duration: HistogramVec,
43 store_operation_duration: HistogramVec,
44 connected_workers: IntGaugeVec,
45 inflight_activities: IntGaugeVec,
46 signals_delivered: IntCounterVec,
47 schedules_fired: IntCounterVec,
48 deploy_operations: IntCounterVec,
49 deploy_denied: IntCounterVec,
50 loaded_workflow_versions: IntGaugeVec,
51 activity_lease_record_failures: IntCounter,
54}
55
56impl MetricsInner {
57 fn register_collectors(&self) -> Result<(), prometheus::Error> {
58 self.registry
59 .register(Box::new(self.workflows_started.clone()))?;
60 self.registry
61 .register(Box::new(self.workflows_completed.clone()))?;
62 self.registry
63 .register(Box::new(self.workflows_reopened.clone()))?;
64 self.registry
65 .register(Box::new(self.activities_dispatched.clone()))?;
66 self.registry
67 .register(Box::new(self.activities_completed.clone()))?;
68 self.registry
69 .register(Box::new(self.activity_duration.clone()))?;
70 self.registry
71 .register(Box::new(self.store_operation_duration.clone()))?;
72 self.registry
73 .register(Box::new(self.connected_workers.clone()))?;
74 self.registry
75 .register(Box::new(self.inflight_activities.clone()))?;
76 self.registry
77 .register(Box::new(self.signals_delivered.clone()))?;
78 self.registry
79 .register(Box::new(self.schedules_fired.clone()))?;
80 self.registry
81 .register(Box::new(self.deploy_operations.clone()))?;
82 self.registry
83 .register(Box::new(self.deploy_denied.clone()))?;
84 self.registry
85 .register(Box::new(self.loaded_workflow_versions.clone()))?;
86 self.registry
87 .register(Box::new(self.activity_lease_record_failures.clone()))?;
88 Ok(())
89 }
90}
91
92impl Metrics {
93 pub fn new() -> Result<Self, MetricsError> {
99 let inner = build_metrics_inner()?;
100 inner.register_collectors()?;
101 initialize_default_label_sets(&inner);
102 Ok(Self {
103 inner: Arc::new(inner),
104 })
105 }
106
107 pub fn encode(&self) -> Result<Vec<u8>, MetricsError> {
113 let encoder = TextEncoder::new();
114 let families = self.inner.registry.gather();
115 let mut buffer = Vec::new();
116 encoder
117 .encode(&families, &mut buffer)
118 .map_err(|error| MetricsError::Encode(error.to_string()))?;
119 Ok(buffer)
120 }
121
122 pub fn workflow_started(&self, namespace: &str, workflow_type: &str) {
124 self.inner
125 .workflows_started
126 .with_label_values(&[namespace, workflow_type])
127 .inc();
128 }
129
130 pub fn workflow_completed(&self, namespace: &str, status: &str) {
132 self.inner
133 .workflows_completed
134 .with_label_values(&[namespace, status])
135 .inc();
136 }
137
138 pub fn workflow_reopened(&self, namespace: &str) {
140 self.inner
141 .workflows_reopened
142 .with_label_values(&[namespace])
143 .inc();
144 }
145
146 pub fn activity_dispatched(&self, namespace: &str, activity_type: &str) {
148 self.inner
149 .activities_dispatched
150 .with_label_values(&[namespace, activity_type])
151 .inc();
152 self.inner
153 .inflight_activities
154 .with_label_values(&[namespace])
155 .inc();
156 }
157
158 pub fn activity_completed(
160 &self,
161 namespace: &str,
162 activity_type: &str,
163 outcome: &str,
164 duration: Duration,
165 ) {
166 self.inner
167 .activities_completed
168 .with_label_values(&[namespace, outcome])
169 .inc();
170 self.inner
171 .activity_duration
172 .with_label_values(&[namespace, activity_type])
173 .observe(duration.as_secs_f64());
174 self.inner
175 .inflight_activities
176 .with_label_values(&[namespace])
177 .dec();
178 }
179
180 pub fn activity_abandoned(&self, namespace: &str) {
182 self.inner
183 .inflight_activities
184 .with_label_values(&[namespace])
185 .dec();
186 }
187
188 pub fn store_operation(&self, operation: &str, duration: Duration) {
190 self.inner
191 .store_operation_duration
192 .with_label_values(&[operation])
193 .observe(duration.as_secs_f64());
194 }
195
196 pub fn activity_lease_record_failed(&self) {
198 self.inner.activity_lease_record_failures.inc();
199 }
200
201 pub fn worker_connected(&self, namespace: &str) {
203 self.inner
204 .connected_workers
205 .with_label_values(&[namespace])
206 .inc();
207 }
208
209 pub fn worker_disconnected(&self, namespace: &str) {
211 self.inner
212 .connected_workers
213 .with_label_values(&[namespace])
214 .dec();
215 }
216
217 pub fn signal_delivered(&self, namespace: &str, residency: &str) {
219 self.inner
220 .signals_delivered
221 .with_label_values(&[namespace, residency])
222 .inc();
223 }
224
225 pub fn schedule_fired(&self, namespace: &str) {
227 self.inner
228 .schedules_fired
229 .with_label_values(&[namespace])
230 .inc();
231 }
232
233 pub fn deploy_operation(&self, operation: &str, outcome: &str) {
235 self.inner
236 .deploy_operations
237 .with_label_values(&[operation, outcome])
238 .inc();
239 }
240
241 pub fn deploy_denied(&self, transport: &str) {
243 self.inner
244 .deploy_denied
245 .with_label_values(&[transport])
246 .inc();
247 }
248
249 pub fn set_loaded_workflow_versions(&self, workflow_type: &str, count: i64) {
252 self.inner
253 .loaded_workflow_versions
254 .with_label_values(&[workflow_type])
255 .set(count);
256 }
257
258 #[cfg(test)]
260 pub(crate) fn inflight_activities_value(&self, namespace: &str) -> i64 {
261 self.inner
262 .inflight_activities
263 .with_label_values(&[namespace])
264 .get()
265 }
266
267 #[cfg(test)]
270 pub(crate) fn activities_dispatched_value(&self, namespace: &str, activity_type: &str) -> u64 {
271 self.inner
272 .activities_dispatched
273 .with_label_values(&[namespace, activity_type])
274 .get()
275 }
276
277 #[cfg(test)]
280 pub(crate) fn activities_completed_value(&self, namespace: &str, outcome: &str) -> u64 {
281 self.inner
282 .activities_completed
283 .with_label_values(&[namespace, outcome])
284 .get()
285 }
286}
287
288fn build_workflow_metrics() -> Result<(IntCounterVec, IntCounterVec, IntCounterVec), MetricsError> {
289 let workflows_started = IntCounterVec::new(
290 Opts::new(
291 "aion_workflows_started_total",
292 "Total workflow executions started by namespace and workflow type.",
293 ),
294 &["namespace", "workflow_type"],
295 )?;
296 let workflows_completed = IntCounterVec::new(
297 Opts::new(
298 "aion_workflows_completed_total",
299 "Total workflow executions that reached a terminal status by namespace and status.",
300 ),
301 &["namespace", "status"],
302 )?;
303 let workflows_reopened = IntCounterVec::new(
304 Opts::new(
305 "aion_workflows_reopened_total",
306 "Total failed workflow runs reopened, by namespace.",
307 ),
308 &["namespace"],
309 )?;
310 Ok((workflows_started, workflows_completed, workflows_reopened))
311}
312
313fn build_metrics_inner() -> Result<MetricsInner, MetricsError> {
314 let registry = Registry::new();
315 let (workflows_started, workflows_completed, workflows_reopened) = build_workflow_metrics()?;
316 let activities_dispatched = IntCounterVec::new(
317 Opts::new(
318 "aion_activities_dispatched_total",
319 "Total activities dispatched to workers by namespace and activity type.",
320 ),
321 &["namespace", "activity_type"],
322 )?;
323 let activities_completed = IntCounterVec::new(
324 Opts::new(
325 "aion_activities_completed_total",
326 "Total activity results received by namespace and outcome.",
327 ),
328 &["namespace", "outcome"],
329 )?;
330 let activity_duration = HistogramVec::new(
331 HistogramOpts::new(
332 "aion_activity_duration_seconds",
333 "Wall-clock activity execution latency from dispatch to result by namespace and activity type.",
334 )
335 .buckets(vec![
336 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0, 30.0,
337 ]),
338 &["namespace", "activity_type"],
339 )?;
340 let store_operation_duration = HistogramVec::new(
341 HistogramOpts::new(
342 "aion_store_operation_duration_seconds",
343 "Store operation latency by operation.",
344 )
345 .buckets(vec![
346 0.001, 0.0025, 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0,
347 ]),
348 &["operation"],
349 )?;
350 let connected_workers = IntGaugeVec::new(
351 Opts::new(
352 "aion_connected_workers",
353 "Current connected worker streams by namespace.",
354 ),
355 &["namespace"],
356 )?;
357 let inflight_activities = IntGaugeVec::new(
358 Opts::new(
359 "aion_inflight_activities",
360 "Current dispatched activities awaiting worker completion by namespace.",
361 ),
362 &["namespace"],
363 )?;
364 let signals_delivered = IntCounterVec::new(
365 Opts::new(
366 "aion_signals_delivered_total",
367 "Total signals delivered by namespace and residency classification.",
368 ),
369 &["namespace", "residency"],
370 )?;
371 let schedules_fired = IntCounterVec::new(
372 Opts::new(
373 "aion_schedules_fired_total",
374 "Total schedule timer evaluations that started a workflow by namespace.",
375 ),
376 &["namespace"],
377 )?;
378 let (deploy_operations, deploy_denied, loaded_workflow_versions) = build_deploy_metrics()?;
379 let activity_lease_record_failures = IntCounter::new(
380 "aion_activity_lease_record_failures_total",
381 "Total activity leases this server knew the worker for and failed to record; each is an \
382 attempt the history reads as unattributed.",
383 )?;
384
385 Ok(MetricsInner {
386 registry,
387 workflows_started,
388 workflows_completed,
389 workflows_reopened,
390 activities_dispatched,
391 activities_completed,
392 activity_duration,
393 store_operation_duration,
394 connected_workers,
395 inflight_activities,
396 signals_delivered,
397 schedules_fired,
398 deploy_operations,
399 deploy_denied,
400 loaded_workflow_versions,
401 activity_lease_record_failures,
402 })
403}
404
405fn build_deploy_metrics() -> Result<(IntCounterVec, IntCounterVec, IntGaugeVec), MetricsError> {
408 let deploy_operations = IntCounterVec::new(
409 Opts::new(
410 "aion_deploy_operations_total",
411 "Total deploy API mutations by operation and outcome class.",
412 ),
413 &["operation", "outcome"],
414 )?;
415 let deploy_denied = IntCounterVec::new(
416 Opts::new(
417 "aion_deploy_denied_total",
418 "Total deploy API authorization denials by transport.",
419 ),
420 &["transport"],
421 )?;
422 let loaded_workflow_versions = IntGaugeVec::new(
423 Opts::new(
424 "aion_loaded_workflow_versions",
425 "Currently loaded package versions per workflow type.",
426 ),
427 &["workflow_type"],
428 )?;
429 Ok((deploy_operations, deploy_denied, loaded_workflow_versions))
430}
431
432fn initialize_default_label_sets(inner: &MetricsInner) {
435 for operation in ["append", "read_history", "list_active", "list_workflow_ids"] {
436 inner
437 .store_operation_duration
438 .with_label_values(&[operation]);
439 }
440 inner
441 .activity_duration
442 .with_label_values(&["default", "default"]);
443}
444
445pub async fn metrics_handler(
447 axum::extract::State(metrics): axum::extract::State<Metrics>,
448) -> Response {
449 match metrics.encode() {
450 Ok(body) => {
451 let mut response = body.into_response();
452 response
453 .headers_mut()
454 .insert(CONTENT_TYPE, HeaderValue::from_static(TEXT_FORMAT));
455 response
456 }
457 Err(error) => (StatusCode::INTERNAL_SERVER_ERROR, error.to_string()).into_response(),
458 }
459}