1use 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
38pub struct PerInstanceReporter {
50 dir: PathBuf,
51}
52
53impl PerInstanceReporter {
54 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 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 }
106}
107
108fn 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
124fn 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
142pub(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
293fn 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}