Skip to main content

tatara_engine/metrics/
mod.rs

1//! Prometheus metrics for tatara — jobs, allocations, nodes, reconciler,
2//! scheduler, gossip, and driver health.
3//!
4//! Exposes a `/metrics` endpoint via the REST API router.
5
6use std::sync::atomic::{AtomicU64, Ordering};
7use std::sync::Arc;
8
9/// Central metrics registry for tatara.
10#[derive(Debug, Default)]
11pub struct TataraMetrics {
12    // Gauges
13    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    // Counters
24    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    // Timing (last recorded values in milliseconds)
33    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    /// Render metrics in Prometheus text exposition format.
44    pub fn render_prometheus(&self) -> String {
45        let mut out = String::with_capacity(2048);
46
47        // Gauges
48        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        // Counters
104        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        // Timing gauges (last observed value)
148        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        // Every metric must have HELP and TYPE lines
216        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                // Metric line: name value
227                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}