Skip to main content

aion_server/observability/
metrics.rs

1//! Prometheus metrics registry and recording helpers.
2
3use 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/// Prometheus registry construction or exposition error.
18#[derive(Debug, Error)]
19pub enum MetricsError {
20    /// A metric failed to register in the prometheus registry.
21    #[error("failed to register prometheus metric: {0}")]
22    Register(#[from] prometheus::Error),
23    /// The prometheus text encoder failed to encode gathered metric families.
24    #[error("failed to encode prometheus metrics: {0}")]
25    Encode(String),
26}
27
28/// Cloneable server metrics handle backed by a prometheus registry.
29#[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    /// Leases this install knew and failed to record (WA-010 R3). One truth
52    /// with the dispatcher's ledger: both move by one on every loss.
53    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    /// Construct the server metrics registry and register all exported metrics.
94    ///
95    /// # Errors
96    ///
97    /// Returns [`MetricsError::Register`] if prometheus rejects a metric descriptor.
98    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    /// Encode all currently gathered metrics in prometheus text exposition format.
108    ///
109    /// # Errors
110    ///
111    /// Returns [`MetricsError::Encode`] if prometheus cannot encode gathered metrics.
112    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    /// Increment the workflow-start counter.
123    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    /// Increment the workflow-terminal counter.
131    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    /// Increment the workflow-reopen counter when a failed run is reopened.
139    pub fn workflow_reopened(&self, namespace: &str) {
140        self.inner
141            .workflows_reopened
142            .with_label_values(&[namespace])
143            .inc();
144    }
145
146    /// Increment the activity-dispatch counter and in-flight gauge.
147    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    /// Increment the activity-completion counter, observe duration, and decrement in-flight gauge.
159    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    /// Decrement in-flight activity gauge when dispatch fails before a result can arrive.
181    pub fn activity_abandoned(&self, namespace: &str) {
182        self.inner
183            .inflight_activities
184            .with_label_values(&[namespace])
185            .dec();
186    }
187
188    /// Observe a store operation duration.
189    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    /// An activity lease this install knew and failed to record.
197    pub fn activity_lease_record_failed(&self) {
198        self.inner.activity_lease_record_failures.inc();
199    }
200
201    /// Increment connected worker gauge for a namespace.
202    pub fn worker_connected(&self, namespace: &str) {
203        self.inner
204            .connected_workers
205            .with_label_values(&[namespace])
206            .inc();
207    }
208
209    /// Decrement connected worker gauge for a namespace.
210    pub fn worker_disconnected(&self, namespace: &str) {
211        self.inner
212            .connected_workers
213            .with_label_values(&[namespace])
214            .dec();
215    }
216
217    /// Increment signal delivery counter.
218    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    /// Increment schedule-fired counter.
226    pub fn schedule_fired(&self, namespace: &str) {
227        self.inner
228            .schedules_fired
229            .with_label_values(&[namespace])
230            .inc();
231    }
232
233    /// Increment the deploy-operation counter for one mutation outcome.
234    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    /// Increment the deploy-denied counter for a transport.
242    pub fn deploy_denied(&self, transport: &str) {
243        self.inner
244            .deploy_denied
245            .with_label_values(&[transport])
246            .inc();
247    }
248
249    /// Set the loaded-version gauge for one workflow type from the
250    /// post-operation listing.
251    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    /// Current value of the `aion_inflight_activities` gauge for a namespace.
259    #[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    /// Current value of the `aion_activities_dispatched_total` counter for a
268    /// namespace and activity type.
269    #[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    /// Current value of the `aion_activities_completed_total` counter for a
278    /// namespace and outcome label.
279    #[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
405/// Deploy API collectors: mutation counter, denial counter, and the
406/// loaded-version gauge fed from the post-operation listing.
407fn 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
432/// Pre-initialize known label sets so all metric families appear in the
433/// prometheus text output before any workflow or activity traffic occurs.
434fn 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
445/// Axum handler for `/metrics`.
446pub 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}