Skip to main content

nmbrs_metrics/reporters/
per_instance.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Per-instance metric snapshot reporter — appends one JSONL
5//! record per snapshot tick to a separate file for each
6//! `(metric_name, labels)` combination.
7//!
8//! ## Layout
9//!
10//! All files live in a single configured directory. Filenames
11//! encode the metric name and every label pair:
12//!
13//! ```text
14//! <name>__<key1>_<value1>__<key2>_<value2>.jsonl
15//! ```
16//!
17//! A metric instance with no labels writes to `<name>.jsonl`.
18//! Characters that aren't `[A-Za-z0-9_-]` get replaced with
19//! `_` so the path stays portable across filesystems and
20//! shell-safe.
21//!
22//! ## Record format
23//!
24//! Each line is one JSON object with `ts` (ms since epoch),
25//! `name`, `labels`, `type`, and a type-specific value
26//! payload. Distinct record shapes per metric type are
27//! distinguished by the `type` discriminator so a downstream
28//! consumer can parse without per-file metadata.
29
30use std::fs;
31use std::io::Write;
32use std::path::{Path, PathBuf};
33
34use crate::labels::Labels;
35use crate::scheduler::Reporter;
36use crate::snapshot::{BucketBound, MetricSet, MetricValue};
37
38/// Per-instance JSONL reporter.
39///
40/// **No file-handle cache.** The previous implementation kept
41/// every `(metric_name, labels)` combination's `File` alive in a
42/// `HashMap` for the entire run, which monotonically grew the
43/// open-FD count — for workloads with many distinct label
44/// tuples (e.g. `phase × profile × optimize_for × k × r`), the
45/// process would exhaust its `ulimit -n`. Each `report()` call
46/// now opens, writes, and closes the target file per record.
47/// At the 30s cadence the metrics scheduler runs this on, the
48/// open() cost is dwarfed by the rest of the snapshot path.
49pub struct PerInstanceReporter {
50    dir: PathBuf,
51}
52
53impl PerInstanceReporter {
54    /// Construct a reporter rooted at `dir`. The directory
55    /// is created if missing; an existing directory is
56    /// reused so consecutive sessions can co-locate output
57    /// when desired.
58    pub fn new(dir: impl AsRef<Path>) -> Result<Self, String> {
59        let dir = dir.as_ref().to_path_buf();
60        fs::create_dir_all(&dir)
61            .map_err(|e| format!("create per-instance metrics dir {:?}: {e}", dir))?;
62        Ok(Self { dir })
63    }
64}
65
66impl Reporter for PerInstanceReporter {
67    fn report(&mut self, snapshot: &MetricSet) {
68        let now_ms = std::time::SystemTime::now()
69            .duration_since(std::time::UNIX_EPOCH)
70            .map(|d| d.as_millis() as i64)
71            .unwrap_or(0);
72
73        for family in snapshot.families() {
74            let name = family.name();
75            for metric in family.metrics() {
76                let labels = metric.labels();
77                let Some(point) = metric.point() else {
78                    continue;
79                };
80                let line = render_record(now_ms, name, labels, point.value());
81                let key = instance_filename(name, labels);
82                let path = self.dir.join(format!("{key}.jsonl"));
83                // Open-write-close per record. The File is
84                // dropped at the end of the block, closing
85                // the FD immediately — bounded resource use
86                // regardless of how many distinct instances
87                // the workload produces.
88                let result = std::fs::OpenOptions::new()
89                    .create(true)
90                    .append(true)
91                    .open(&path)
92                    .and_then(|mut f| writeln!(f, "{line}"));
93                if let Err(e) = result {
94                    crate::diag::warn(&format!(
95                        "warning: per-instance metrics write failed for {key}: {e}"
96                    ));
97                }
98            }
99        }
100    }
101
102    fn flush(&mut self) {
103        // No cache, nothing to flush. Each record's open-write-
104        // close path already commits before the FD is dropped.
105    }
106}
107
108/// Build the safe filename stem for one metric instance.
109/// `<name>__<k1>_<v1>__<k2>_<v2>` — label pairs are kept in
110/// declaration order so consecutive snapshots write to the
111/// same file. Sanitisation runs per token so a `/` or `,` in
112/// a label value can't escape the metrics directory.
113fn instance_filename(name: &str, labels: &Labels) -> String {
114    let mut out = sanitize(name);
115    for (k, v) in labels.iter() {
116        out.push_str("__");
117        out.push_str(&sanitize(k));
118        out.push('_');
119        out.push_str(&sanitize(v));
120    }
121    out
122}
123
124/// Replace anything outside `[A-Za-z0-9_-]` with `_`. Empty
125/// input becomes `_` so the resulting filename always has a
126/// non-empty token between separators.
127fn sanitize(s: &str) -> String {
128    if s.is_empty() {
129        return "_".to_string();
130    }
131    let mut out = String::with_capacity(s.len());
132    for c in s.chars() {
133        if c.is_ascii_alphanumeric() || c == '_' || c == '-' {
134            out.push(c);
135        } else {
136            out.push('_');
137        }
138    }
139    out
140}
141
142/// Render one JSON record for the metric's current value.
143/// Hand-rolled so this crate stays free of a serde_json
144/// dependency; the surface here is small enough that a tiny
145/// escape helper covers it.
146/// One JSONL record for a metric instance at a tick.
147///
148/// `pub(crate)` so the single-file metrics log renders byte-identical records to
149/// the per-instance files: the two differ only in how records are routed to
150/// files, never in their shape, so a consumer can parse either with one reader.
151pub(crate) fn render_record(
152    now_ms: i64,
153    name: &str,
154    labels: &Labels,
155    value: &MetricValue,
156) -> String {
157    let mut out = String::with_capacity(128);
158    out.push('{');
159    write_kv(&mut out, "ts", &now_ms.to_string(), false);
160    out.push(',');
161    write_kv(&mut out, "name", name, true);
162    out.push(',');
163    out.push_str("\"labels\":{");
164    let mut first = true;
165    for (k, v) in labels.iter() {
166        if !first {
167            out.push(',');
168        }
169        first = false;
170        write_kv(&mut out, k, v, true);
171    }
172    out.push('}');
173    match value {
174        MetricValue::Counter(c) => {
175            out.push(',');
176            write_kv(&mut out, "type", "counter", true);
177            out.push(',');
178            write_kv(&mut out, "count", &c.cumulative.to_string(), false);
179        }
180        MetricValue::Gauge(g) => {
181            out.push(',');
182            write_kv(&mut out, "type", "gauge", true);
183            out.push(',');
184            write_kv(&mut out, "value", &format_f64(g.value), false);
185        }
186        MetricValue::Histogram(h) => {
187            let r = &h.reservoir;
188            out.push(',');
189            write_kv(&mut out, "type", "histogram", true);
190            out.push(',');
191            write_kv(&mut out, "count", &r.len().to_string(), false);
192            out.push(',');
193            write_kv(&mut out, "min", &format_f64(r.min() as f64), false);
194            out.push(',');
195            write_kv(&mut out, "max", &format_f64(r.max() as f64), false);
196            out.push(',');
197            write_kv(&mut out, "mean", &format_f64(r.mean()), false);
198            out.push(',');
199            write_kv(&mut out, "stddev", &format_f64(r.stdev()), false);
200            for (label, q) in &[
201                ("p50", 0.50),
202                ("p75", 0.75),
203                ("p90", 0.90),
204                ("p95", 0.95),
205                ("p98", 0.98),
206                ("p99", 0.99),
207                ("p999", 0.999),
208            ] {
209                out.push(',');
210                write_kv(
211                    &mut out,
212                    label,
213                    &format_f64(r.value_at_quantile(*q) as f64),
214                    false,
215                );
216            }
217        }
218        MetricValue::BucketedHistogram(h) => {
219            out.push(',');
220            write_kv(&mut out, "type", "bucketed_histogram", true);
221            out.push_str(",\"buckets\":[");
222            let mut first = true;
223            for (bound, count) in &h.buckets {
224                if !first {
225                    out.push(',');
226                }
227                first = false;
228                let le = match bound {
229                    BucketBound::Finite(v) => format_f64(*v as f64),
230                    BucketBound::PositiveInfinity => "\"+Inf\"".to_string(),
231                };
232                out.push_str(&format!("{{\"le\":{le},\"count\":{count}}}"));
233            }
234            out.push(']');
235        }
236        MetricValue::Info(_) => {
237            out.push(',');
238            write_kv(&mut out, "type", "info", true);
239            out.push(',');
240            write_kv(&mut out, "value", "1", false);
241        }
242        MetricValue::StateSet(s) => {
243            out.push(',');
244            write_kv(&mut out, "type", "stateset", true);
245            out.push_str(",\"states\":{");
246            let mut first = true;
247            for (state, active) in &s.states {
248                if !first {
249                    out.push(',');
250                }
251                first = false;
252                write_kv(
253                    &mut out,
254                    state,
255                    if *active { "true" } else { "false" },
256                    false,
257                );
258            }
259            out.push('}');
260        }
261    }
262    out.push('}');
263    out
264}
265
266pub(crate) fn write_kv(out: &mut String, key: &str, value: &str, quote_value: bool) {
267    out.push('"');
268    escape_into(out, key);
269    out.push_str("\":");
270    if quote_value {
271        out.push('"');
272        escape_into(out, value);
273        out.push('"');
274    } else {
275        out.push_str(value);
276    }
277}
278
279fn escape_into(out: &mut String, s: &str) {
280    for c in s.chars() {
281        match c {
282            '"' => out.push_str("\\\""),
283            '\\' => out.push_str("\\\\"),
284            '\n' => out.push_str("\\n"),
285            '\r' => out.push_str("\\r"),
286            '\t' => out.push_str("\\t"),
287            c if (c as u32) < 0x20 => out.push_str(&format!("\\u{:04x}", c as u32)),
288            c => out.push(c),
289        }
290    }
291}
292
293/// Plain JSON-friendly float formatting: NaN / ±Inf serialise
294/// as `null` (so a downstream parser doesn't choke), finite
295/// values render with `{:?}` so a sensible precision survives
296/// without sticking trailing zeros on whole numbers.
297fn format_f64(v: f64) -> String {
298    if !v.is_finite() {
299        return "null".to_string();
300    }
301    if v == v.trunc() && v.abs() < 1e16 {
302        format!("{}", v as i64)
303    } else {
304        format!("{v}")
305    }
306}
307
308#[cfg(test)]
309mod tests {
310    use super::*;
311    use crate::labels::Labels;
312    use crate::snapshot::MetricSet;
313    use std::time::{Duration, Instant};
314
315    #[test]
316    fn writes_one_file_per_instance() {
317        let dir = std::env::temp_dir().join("nb_per_instance_test");
318        let _ = fs::remove_dir_all(&dir);
319        let mut reporter = PerInstanceReporter::new(&dir).unwrap();
320
321        let mut snap = MetricSet::new(Duration::from_secs(1));
322        let ann = Labels::of("phase", "ann");
323        let pvs = Labels::of("phase", "pvs");
324        snap.insert_counter("ops_total", ann, 42, Instant::now());
325        snap.insert_counter("ops_total", pvs, 17, Instant::now());
326        reporter.report(&snap);
327        reporter.flush();
328
329        let ann_path = dir.join("ops_total__phase_ann.jsonl");
330        let pvs_path = dir.join("ops_total__phase_pvs.jsonl");
331        assert!(ann_path.exists(), "expected per-instance file for ann");
332        assert!(pvs_path.exists(), "expected per-instance file for pvs");
333
334        let ann_line = fs::read_to_string(&ann_path).unwrap();
335        assert!(ann_line.contains("\"count\":42"));
336        assert!(ann_line.contains("\"phase\":\"ann\""));
337        let pvs_line = fs::read_to_string(&pvs_path).unwrap();
338        assert!(pvs_line.contains("\"count\":17"));
339
340        let _ = fs::remove_dir_all(&dir);
341    }
342
343    #[test]
344    fn label_values_sanitised_for_path() {
345        let labels = Labels::of("name", "path/with..slashes,and,commas");
346        let stem = instance_filename("metric", &labels);
347        assert!(!stem.contains('/'));
348        assert!(!stem.contains(','));
349        assert!(!stem.contains('.'));
350        assert!(stem.starts_with("metric__name_"));
351    }
352
353    #[test]
354    fn appending_preserves_history_across_ticks() {
355        let dir = std::env::temp_dir().join("nb_per_instance_append");
356        let _ = fs::remove_dir_all(&dir);
357        let mut reporter = PerInstanceReporter::new(&dir).unwrap();
358
359        let labels = Labels::of("phase", "ann");
360        for n in [1u64, 2, 3] {
361            let mut snap = MetricSet::new(Duration::from_secs(1));
362            snap.insert_counter("ops_total", labels.clone(), n, Instant::now());
363            reporter.report(&snap);
364            reporter.flush();
365        }
366        let body = fs::read_to_string(dir.join("ops_total__phase_ann.jsonl")).unwrap();
367        let lines: Vec<&str> = body.lines().collect();
368        assert_eq!(lines.len(), 3);
369        assert!(lines[0].contains("\"count\":1"));
370        assert!(lines[2].contains("\"count\":3"));
371        let _ = fs::remove_dir_all(&dir);
372    }
373}