1use std::sync::atomic::{AtomicU64, Ordering};
7use std::sync::Arc;
8
9#[derive(Debug, Default)]
11pub struct TataraMetrics {
12 pub jobs_total: AtomicU64,
14 pub jobs_pending: AtomicU64,
15 pub jobs_running: AtomicU64,
16 pub allocations_total: AtomicU64,
17 pub allocations_running: AtomicU64,
18 pub allocations_failed: AtomicU64,
19 pub nodes_total: AtomicU64,
20 pub nodes_ready: AtomicU64,
21 pub services_registered: AtomicU64,
22
23 pub reconcile_total: AtomicU64,
25 pub reconcile_errors: AtomicU64,
26 pub scheduler_evals: AtomicU64,
27 pub health_probes_executed: AtomicU64,
28 pub health_probes_failed: AtomicU64,
29 pub secrets_fetched: AtomicU64,
30 pub ports_allocated: AtomicU64,
31
32 pub reconcile_duration_ms: AtomicU64,
34 pub scheduler_eval_duration_ms: AtomicU64,
35 pub nix_eval_duration_ms: AtomicU64,
36}
37
38impl TataraMetrics {
39 pub fn new() -> Arc<Self> {
40 Arc::new(Self::default())
41 }
42
43 pub fn render_prometheus(&self) -> String {
45 let mut out = String::with_capacity(2048);
46
47 prom_gauge(
49 &mut out,
50 "tatara_jobs_total",
51 "Total number of jobs",
52 self.jobs_total.load(Ordering::Relaxed),
53 );
54 prom_gauge(
55 &mut out,
56 "tatara_jobs_pending",
57 "Jobs in pending state",
58 self.jobs_pending.load(Ordering::Relaxed),
59 );
60 prom_gauge(
61 &mut out,
62 "tatara_jobs_running",
63 "Jobs in running state",
64 self.jobs_running.load(Ordering::Relaxed),
65 );
66 prom_gauge(
67 &mut out,
68 "tatara_allocations_total",
69 "Total allocations",
70 self.allocations_total.load(Ordering::Relaxed),
71 );
72 prom_gauge(
73 &mut out,
74 "tatara_allocations_running",
75 "Running allocations",
76 self.allocations_running.load(Ordering::Relaxed),
77 );
78 prom_gauge(
79 &mut out,
80 "tatara_allocations_failed",
81 "Failed allocations",
82 self.allocations_failed.load(Ordering::Relaxed),
83 );
84 prom_gauge(
85 &mut out,
86 "tatara_nodes_total",
87 "Total cluster nodes",
88 self.nodes_total.load(Ordering::Relaxed),
89 );
90 prom_gauge(
91 &mut out,
92 "tatara_nodes_ready",
93 "Ready nodes",
94 self.nodes_ready.load(Ordering::Relaxed),
95 );
96 prom_gauge(
97 &mut out,
98 "tatara_services_registered",
99 "Registered service instances",
100 self.services_registered.load(Ordering::Relaxed),
101 );
102
103 prom_counter(
105 &mut out,
106 "tatara_reconcile_total",
107 "Total reconciliation ticks",
108 self.reconcile_total.load(Ordering::Relaxed),
109 );
110 prom_counter(
111 &mut out,
112 "tatara_reconcile_errors_total",
113 "Reconciliation errors",
114 self.reconcile_errors.load(Ordering::Relaxed),
115 );
116 prom_counter(
117 &mut out,
118 "tatara_scheduler_evals_total",
119 "Scheduler evaluation cycles",
120 self.scheduler_evals.load(Ordering::Relaxed),
121 );
122 prom_counter(
123 &mut out,
124 "tatara_health_probes_total",
125 "Health probes executed",
126 self.health_probes_executed.load(Ordering::Relaxed),
127 );
128 prom_counter(
129 &mut out,
130 "tatara_health_probes_failed_total",
131 "Health probes failed",
132 self.health_probes_failed.load(Ordering::Relaxed),
133 );
134 prom_counter(
135 &mut out,
136 "tatara_secrets_fetched_total",
137 "Secrets fetched",
138 self.secrets_fetched.load(Ordering::Relaxed),
139 );
140 prom_counter(
141 &mut out,
142 "tatara_ports_allocated_total",
143 "Ports allocated",
144 self.ports_allocated.load(Ordering::Relaxed),
145 );
146
147 prom_gauge(
149 &mut out,
150 "tatara_reconcile_duration_ms",
151 "Last reconcile duration in ms",
152 self.reconcile_duration_ms.load(Ordering::Relaxed),
153 );
154 prom_gauge(
155 &mut out,
156 "tatara_scheduler_eval_duration_ms",
157 "Last scheduler eval duration in ms",
158 self.scheduler_eval_duration_ms.load(Ordering::Relaxed),
159 );
160 prom_gauge(
161 &mut out,
162 "tatara_nix_eval_duration_ms",
163 "Last nix eval duration in ms",
164 self.nix_eval_duration_ms.load(Ordering::Relaxed),
165 );
166
167 out
168 }
169
170 pub fn inc(&self, field: &AtomicU64) {
171 field.fetch_add(1, Ordering::Relaxed);
172 }
173
174 pub fn set(&self, field: &AtomicU64, value: u64) {
175 field.store(value, Ordering::Relaxed);
176 }
177}
178
179fn prom_gauge(out: &mut String, name: &str, help: &str, value: u64) {
180 out.push_str(&format!(
181 "# HELP {name} {help}\n# TYPE {name} gauge\n{name} {value}\n"
182 ));
183}
184
185fn prom_counter(out: &mut String, name: &str, help: &str, value: u64) {
186 out.push_str(&format!(
187 "# HELP {name} {help}\n# TYPE {name} counter\n{name} {value}\n"
188 ));
189}
190
191#[cfg(test)]
192mod tests {
193 use super::*;
194
195 #[test]
196 fn test_render_prometheus() {
197 let metrics = TataraMetrics::default();
198 metrics.jobs_total.store(5, Ordering::Relaxed);
199 metrics.jobs_running.store(3, Ordering::Relaxed);
200 metrics.reconcile_total.store(100, Ordering::Relaxed);
201
202 let output = metrics.render_prometheus();
203 assert!(output.contains("tatara_jobs_total 5"));
204 assert!(output.contains("tatara_jobs_running 3"));
205 assert!(output.contains("tatara_reconcile_total 100"));
206 assert!(output.contains("# TYPE tatara_jobs_total gauge"));
207 assert!(output.contains("# TYPE tatara_reconcile_total counter"));
208 }
209
210 #[test]
211 fn test_prometheus_format_validity() {
212 let metrics = TataraMetrics::default();
213 let output = metrics.render_prometheus();
214
215 for line in output.lines() {
217 if line.starts_with("# HELP") {
218 assert!(line.len() > 7, "HELP line too short: {}", line);
219 } else if line.starts_with("# TYPE") {
220 assert!(
221 line.contains("gauge") || line.contains("counter"),
222 "TYPE must be gauge or counter: {}",
223 line
224 );
225 } else if !line.is_empty() {
226 let parts: Vec<&str> = line.split_whitespace().collect();
228 assert_eq!(parts.len(), 2, "metric line must be 'name value': {}", line);
229 assert!(
230 parts[1].parse::<u64>().is_ok(),
231 "metric value must be numeric: {}",
232 line
233 );
234 }
235 }
236 }
237
238 #[test]
239 fn test_inc_and_set() {
240 let metrics = TataraMetrics::default();
241 metrics.inc(&metrics.reconcile_total);
242 metrics.inc(&metrics.reconcile_total);
243 metrics.set(&metrics.nix_eval_duration_ms, 42);
244
245 assert_eq!(metrics.reconcile_total.load(Ordering::Relaxed), 2);
246 assert_eq!(metrics.nix_eval_duration_ms.load(Ordering::Relaxed), 42);
247 }
248}