Skip to main content

af_workflow/
metrics.rs

1//! Dependency-free scheduler metrics and readiness.
2
3use std::sync::atomic::{AtomicI64, AtomicU64, Ordering};
4
5const READY_MISSED_PASSES: i64 = 3;
6
7pub struct RuntimeMetrics {
8    prefix: String,
9    passes_total: AtomicU64,
10    pass_failures_total: AtomicU64,
11    claimed_total: AtomicU64,
12    evaluation_failures_total: AtomicU64,
13    overruns_total: AtomicU64,
14    backlog_events_total: AtomicU64,
15    last_pass_unix_ms: AtomicI64,
16    last_pass_duration_ms: AtomicU64,
17    active_instances: AtomicI64,
18    due_instances: AtomicI64,
19    stalled_instances: AtomicI64,
20    started_unix_ms: AtomicI64,
21}
22
23impl RuntimeMetrics {
24    pub fn new(now_ms: i64, prefix: impl Into<String>) -> Self {
25        Self {
26            prefix: prefix.into(),
27            passes_total: AtomicU64::new(0),
28            pass_failures_total: AtomicU64::new(0),
29            claimed_total: AtomicU64::new(0),
30            evaluation_failures_total: AtomicU64::new(0),
31            overruns_total: AtomicU64::new(0),
32            backlog_events_total: AtomicU64::new(0),
33            last_pass_unix_ms: AtomicI64::new(0),
34            last_pass_duration_ms: AtomicU64::new(0),
35            active_instances: AtomicI64::new(0),
36            due_instances: AtomicI64::new(0),
37            stalled_instances: AtomicI64::new(0),
38            started_unix_ms: AtomicI64::new(now_ms),
39        }
40    }
41
42    pub fn prefix(&self) -> &str {
43        &self.prefix
44    }
45
46    pub fn record_pass(&self, now_ms: i64, duration_ms: u64, claimed: u64, failed: u64) {
47        self.passes_total.fetch_add(1, Ordering::Relaxed);
48        self.claimed_total.fetch_add(claimed, Ordering::Relaxed);
49        self.evaluation_failures_total
50            .fetch_add(failed, Ordering::Relaxed);
51        self.last_pass_unix_ms.store(now_ms, Ordering::Relaxed);
52        self.last_pass_duration_ms
53            .store(duration_ms, Ordering::Relaxed);
54    }
55
56    pub fn record_pass_failure(&self) {
57        self.pass_failures_total.fetch_add(1, Ordering::Relaxed);
58    }
59
60    pub fn record_overrun(&self) {
61        self.overruns_total.fetch_add(1, Ordering::Relaxed);
62    }
63
64    pub fn record_backlog(&self) {
65        self.backlog_events_total.fetch_add(1, Ordering::Relaxed);
66    }
67
68    pub fn set_gauges(&self, active: i64, due: i64, stalled: i64) {
69        self.active_instances.store(active, Ordering::Relaxed);
70        self.due_instances.store(due, Ordering::Relaxed);
71        self.stalled_instances.store(stalled, Ordering::Relaxed);
72    }
73
74    pub fn readiness(&self, now_ms: i64, interval_secs: u64) -> Readiness {
75        let active = self.active_instances.load(Ordering::Relaxed);
76        let last = self.last_pass_unix_ms.load(Ordering::Relaxed);
77        let reference = if last > 0 {
78            last
79        } else {
80            self.started_unix_ms.load(Ordering::Relaxed)
81        };
82        let age_ms = now_ms.saturating_sub(reference).max(0);
83        let budget_ms = (interval_secs as i64)
84            .saturating_mul(1_000)
85            .saturating_mul(READY_MISSED_PASSES);
86        let ready = active == 0 || age_ms <= budget_ms;
87        Readiness {
88            ready,
89            reason: (!ready).then(|| {
90                format!(
91                    "{active} active instance(s) but no completed pass in {}s",
92                    age_ms / 1_000
93                )
94            }),
95            last_pass_age_secs: age_ms / 1_000,
96            active,
97        }
98    }
99
100    pub fn render_prometheus(&self, now_ms: i64) -> String {
101        let last = self.last_pass_unix_ms.load(Ordering::Relaxed);
102        let age = if last > 0 {
103            now_ms.saturating_sub(last).max(0) / 1_000
104        } else {
105            -1
106        };
107        let values = [
108            (
109                "passes_total",
110                "counter",
111                self.passes_total.load(Ordering::Relaxed) as i64,
112            ),
113            (
114                "pass_failures_total",
115                "counter",
116                self.pass_failures_total.load(Ordering::Relaxed) as i64,
117            ),
118            (
119                "claimed_total",
120                "counter",
121                self.claimed_total.load(Ordering::Relaxed) as i64,
122            ),
123            (
124                "evaluation_failures_total",
125                "counter",
126                self.evaluation_failures_total.load(Ordering::Relaxed) as i64,
127            ),
128            (
129                "overruns_total",
130                "counter",
131                self.overruns_total.load(Ordering::Relaxed) as i64,
132            ),
133            (
134                "backlog_events_total",
135                "counter",
136                self.backlog_events_total.load(Ordering::Relaxed) as i64,
137            ),
138            (
139                "last_pass_duration_ms",
140                "gauge",
141                self.last_pass_duration_ms.load(Ordering::Relaxed) as i64,
142            ),
143            ("last_pass_age_seconds", "gauge", age),
144            (
145                "active_instances",
146                "gauge",
147                self.active_instances.load(Ordering::Relaxed),
148            ),
149            (
150                "due_instances",
151                "gauge",
152                self.due_instances.load(Ordering::Relaxed),
153            ),
154            (
155                "stalled_instances",
156                "gauge",
157                self.stalled_instances.load(Ordering::Relaxed),
158            ),
159        ];
160        values
161            .into_iter()
162            .map(|(suffix, kind, value)| {
163                let name = format!("{}_{suffix}", self.prefix);
164                format!("# TYPE {name} {kind}\n{name} {value}\n")
165            })
166            .collect()
167    }
168}
169
170#[derive(Debug, Clone, PartialEq, Eq)]
171pub struct Readiness {
172    pub ready: bool,
173    pub reason: Option<String>,
174    pub last_pass_age_secs: i64,
175    pub active: i64,
176}
177
178#[cfg(test)]
179mod tests {
180    use super::*;
181
182    #[test]
183    fn readiness_detects_an_idle_active_runtime() {
184        let metrics = RuntimeMetrics::new(0, "jobs");
185        metrics.set_gauges(2, 2, 0);
186        assert!(!metrics.readiness(31_000, 10).ready);
187        metrics.record_pass(30_000, 25, 2, 1);
188        assert!(metrics.readiness(31_000, 10).ready);
189    }
190
191    #[test]
192    fn prometheus_output_is_namespaced_and_complete() {
193        let metrics = RuntimeMetrics::new(0, "jobs");
194        metrics.set_gauges(2, 1, 1);
195        metrics.record_pass_failure();
196        metrics.record_overrun();
197        metrics.record_backlog();
198        let output = metrics.render_prometheus(5_000);
199        assert!(output.contains("jobs_active_instances 2"));
200        assert!(output.contains("jobs_pass_failures_total 1"));
201        assert!(output.contains("jobs_overruns_total 1"));
202        assert!(output.contains("jobs_backlog_events_total 1"));
203        assert_eq!(metrics.prefix(), "jobs");
204    }
205}