1use 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}