1use parking_lot::RwLock;
18use std::collections::HashMap;
19use std::sync::Arc;
20
21pub struct Summary {
26 name: String,
27 help: String,
28 quantiles: Vec<f64>,
30 samples: Arc<RwLock<Vec<f64>>>,
32 sum: Arc<RwLock<f64>>,
34 count: Arc<RwLock<u64>>,
36}
37
38impl Summary {
39 pub fn new(name: impl Into<String>, help: impl Into<String>, quantiles: Vec<f64>) -> Self {
46 Self {
47 name: name.into(),
48 help: help.into(),
49 quantiles,
50 samples: Arc::new(RwLock::new(Vec::new())),
51 sum: Arc::new(RwLock::new(0.0)),
52 count: Arc::new(RwLock::new(0)),
53 }
54 }
55
56 pub fn observe(&self, value: f64) {
58 let mut samples = self.samples.write();
59 let pos = samples.partition_point(|&v| v < value);
60 samples.insert(pos, value);
61
62 let mut sum = self.sum.write();
63 *sum += value;
64
65 let mut count = self.count.write();
66 *count += 1;
67 }
68
69 pub fn quantile(&self, q: f64) -> Option<f64> {
74 if !(0.0..=1.0).contains(&q) {
75 return None;
76 }
77 let samples = self.samples.read();
78 if samples.is_empty() {
79 return None;
80 }
81 let n = samples.len();
82 let rank = ((q * n as f64).ceil() as usize).max(1).min(n);
83 Some(samples[rank - 1])
84 }
85
86 pub fn quantiles(&self) -> Vec<(f64, Option<f64>)> {
88 self.quantiles
89 .iter()
90 .map(|&q| (q, self.quantile(q)))
91 .collect()
92 }
93
94 pub fn count(&self) -> u64 {
96 *self.count.read()
97 }
98
99 pub fn sum(&self) -> f64 {
101 *self.sum.read()
102 }
103
104 pub fn name(&self) -> &str {
106 &self.name
107 }
108
109 pub fn render(&self) -> String {
111 let samples = self.samples.read();
112 let sum = *self.sum.read();
113 let count = *self.count.read();
114
115 let mut output = String::new();
116 output.push_str(&format!("# HELP {} {}\n", self.name, self.help));
117 output.push_str(&format!("# TYPE {} summary\n", self.name));
118
119 for &q in &self.quantiles {
120 let value = if samples.is_empty() {
121 0.0
122 } else {
123 let n = samples.len();
124 let rank = ((q * n as f64).ceil() as usize).max(1).min(n);
125 samples[rank - 1]
126 };
127 output.push_str(&format!("{}{{quantile=\"{}\"}} {}\n", self.name, q, value));
128 }
129
130 output.push_str(&format!("{}_sum {}\n", self.name, sum));
131 output.push_str(&format!("{}_count {}\n", self.name, count));
132 output
133 }
134
135 pub fn reset(&self) {
137 let mut samples = self.samples.write();
138 samples.clear();
139 *self.sum.write() = 0.0;
140 *self.count.write() = 0;
141 }
142}
143
144pub struct LabeledHistogram {
149 name: String,
151 help: String,
153 buckets: Vec<f64>,
155 series: RwLock<HashMap<String, LabeledSeries>>,
157}
158
159#[derive(Debug, Clone)]
161struct LabeledSeries {
162 labels: Vec<(String, String)>,
164 counts: Vec<u64>,
166 sum: f64,
168 count: u64,
170}
171
172impl LabeledSeries {
173 fn new(labels: Vec<(String, String)>, bucket_count: usize) -> Self {
174 Self {
175 labels,
176 counts: vec![0; bucket_count],
177 sum: 0.0,
178 count: 0,
179 }
180 }
181
182 fn label_key(labels: &[(String, String)]) -> String {
183 labels
184 .iter()
185 .map(|(k, v)| format!("{}=\"{}\"", k, v.replace('"', "\\\"")))
186 .collect::<Vec<_>>()
187 .join(",")
188 }
189}
190
191impl LabeledHistogram {
192 pub fn new(name: impl Into<String>, help: impl Into<String>, mut buckets: Vec<f64>) -> Self {
199 buckets.sort_by(|a, b| a.partial_cmp(b).unwrap_or(std::cmp::Ordering::Equal));
200 if !buckets.contains(&f64::INFINITY) {
201 buckets.push(f64::INFINITY);
202 }
203 Self {
204 name: name.into(),
205 help: help.into(),
206 buckets,
207 series: RwLock::new(HashMap::new()),
208 }
209 }
210
211 pub fn observe(&self, labels: &HashMap<String, String>, value: f64) {
217 let mut sorted_labels: Vec<(String, String)> =
218 labels.iter().map(|(k, v)| (k.clone(), v.clone())).collect();
219 sorted_labels.sort_by(|a, b| a.0.cmp(&b.0));
220
221 let key = LabeledSeries::label_key(&sorted_labels);
222 let mut series = self.series.write();
223 let entry = series
224 .entry(key)
225 .or_insert_with(|| LabeledSeries::new(sorted_labels.clone(), self.buckets.len()));
226
227 for (i, bucket) in self.buckets.iter().enumerate() {
228 if value <= *bucket {
229 entry.counts[i] += 1;
230 }
231 }
232 entry.sum += value;
233 entry.count += 1;
234 }
235
236 pub fn count(&self, labels: &HashMap<String, String>) -> u64 {
238 let mut sorted_labels: Vec<(String, String)> =
239 labels.iter().map(|(k, v)| (k.clone(), v.clone())).collect();
240 sorted_labels.sort_by(|a, b| a.0.cmp(&b.0));
241 let key = LabeledSeries::label_key(&sorted_labels);
242 self.series.read().get(&key).map(|s| s.count).unwrap_or(0)
243 }
244
245 pub fn label_combination_count(&self) -> usize {
247 self.series.read().len()
248 }
249
250 pub fn render(&self) -> String {
252 let series = self.series.read();
253 let mut output = String::new();
254 output.push_str(&format!("# HELP {} {}\n", self.name, self.help));
255 output.push_str(&format!("# TYPE {} histogram\n", self.name));
256
257 for s in series.values() {
258 let label_str = LabeledSeries::label_key(&s.labels);
259 for (i, bucket) in self.buckets.iter().enumerate() {
260 if *bucket == f64::INFINITY {
261 output.push_str(&format!(
262 "{}_bucket{{{},le=\"+Inf\"}} {}\n",
263 self.name, label_str, s.counts[i]
264 ));
265 } else {
266 output.push_str(&format!(
267 "{}_bucket{{{},le=\"{}\"}} {}\n",
268 self.name, label_str, bucket, s.counts[i]
269 ));
270 }
271 }
272 output.push_str(&format!("{}_sum{{{}}} {}\n", self.name, label_str, s.sum));
273 output.push_str(&format!(
274 "{}_count{{{}}} {}\n",
275 self.name, label_str, s.count
276 ));
277 }
278
279 output
280 }
281}
282
283#[derive(Debug, Clone)]
285pub struct PushgatewayConfig {
286 pub endpoint: String,
288 pub job: String,
290 pub instance: Option<String>,
292}
293
294impl Default for PushgatewayConfig {
295 fn default() -> Self {
296 Self {
297 endpoint: "http://localhost:9091".to_string(),
298 job: "sz-orm".to_string(),
299 instance: None,
300 }
301 }
302}
303
304pub struct PushgatewayExporter {
321 config: PushgatewayConfig,
322 pushed: RwLock<Vec<PushSnapshot>>,
324}
325
326#[derive(Debug, Clone)]
328pub struct PushSnapshot {
329 pub timestamp_ms: i64,
331 pub metrics_text: String,
333 pub job: String,
335 pub instance: Option<String>,
337}
338
339impl PushgatewayExporter {
340 pub fn new(config: PushgatewayConfig) -> Self {
342 Self {
343 config,
344 pushed: RwLock::new(Vec::new()),
345 }
346 }
347
348 pub fn push(&self, metrics_text: impl Into<String>) -> Result<(), String> {
362 let text = metrics_text.into();
363 let snapshot = PushSnapshot {
364 timestamp_ms: current_timestamp_ms(),
365 metrics_text: text.clone(),
366 job: self.config.job.clone(),
367 instance: self.config.instance.clone(),
368 };
369
370 self.pushed.write().push(snapshot);
372
373 #[cfg(feature = "push-gateway")]
375 {
376 return self.push_http(&text);
377 }
378
379 #[cfg(not(feature = "push-gateway"))]
381 Ok(())
382 }
383
384 #[cfg(feature = "push-gateway")]
386 fn push_http(&self, text: &str) -> Result<(), String> {
387 let mut url = format!(
389 "{}/metrics/job/{}",
390 self.config.endpoint.trim_end_matches('/'),
391 url_encode(&self.config.job)
392 );
393 if let Some(ref instance) = self.config.instance {
394 url.push_str(&format!("/instance/{}", url_encode(instance)));
395 }
396
397 let client = reqwest::blocking::Client::builder()
400 .timeout(std::time::Duration::from_secs(10))
401 .build()
402 .map_err(|e| format!("reqwest client build failed: {}", e))?;
403
404 let resp = client
405 .put(&url)
406 .header("Content-Type", "text/plain; version=0.0.4; charset=utf-8")
407 .body(text.to_string())
408 .send()
409 .map_err(|e| format!("push to {} failed: {}", url, e))?;
410
411 let status = resp.status();
412 if status.is_success() {
413 Ok(())
414 } else {
415 let body = resp.text().unwrap_or_default();
416 Err(format!(
417 "push to {} returned non-2xx status {}: {}",
418 url, status, body
419 ))
420 }
421 }
422
423 pub fn push_from_registry(&self, registry: &crate::MetricsRegistry) -> Result<(), String> {
425 let text = registry.render();
426 self.push(text)
427 }
428
429 pub fn snapshots(&self) -> Vec<PushSnapshot> {
431 self.pushed.read().clone()
432 }
433
434 pub fn push_count(&self) -> usize {
436 self.pushed.read().len()
437 }
438
439 pub fn clear(&self) {
441 self.pushed.write().clear();
442 }
443
444 pub fn config(&self) -> &PushgatewayConfig {
446 &self.config
447 }
448}
449
450fn current_timestamp_ms() -> i64 {
451 use std::time::{SystemTime, UNIX_EPOCH};
452 SystemTime::now()
453 .duration_since(UNIX_EPOCH)
454 .unwrap_or_default()
455 .as_millis() as i64
456}
457
458#[allow(dead_code)]
462fn url_encode(s: &str) -> String {
463 let mut out = String::with_capacity(s.len());
464 for byte in s.bytes() {
465 match byte {
466 b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'_' | b'.' | b'~' => {
467 out.push(byte as char);
468 }
469 _ => {
470 out.push_str(&format!("%{:02X}", byte));
471 }
472 }
473 }
474 out
475}
476
477#[cfg(test)]
478mod tests {
479 use super::*;
480
481 #[test]
484 fn test_url_encode_alphanumeric() {
485 assert_eq!(url_encode("job1"), "job1");
486 assert_eq!(url_encode("my-job"), "my-job");
487 assert_eq!(url_encode("job.test"), "job.test");
488 assert_eq!(url_encode("job_test"), "job_test");
489 assert_eq!(url_encode("job~test"), "job~test");
490 }
491
492 #[test]
493 fn test_url_encode_special_chars() {
494 assert_eq!(url_encode("my job"), "my%20job");
496 assert_eq!(url_encode("a/b"), "a%2Fb");
498 assert_eq!(url_encode("任务"), "%E4%BB%BB%E5%8A%A1");
500 }
501
502 #[test]
503 fn test_url_encode_empty() {
504 assert_eq!(url_encode(""), "");
505 }
506
507 #[test]
510 fn test_pushgateway_memory_mode_records_snapshot() {
511 let exporter = PushgatewayExporter::new(PushgatewayConfig::default());
513 let result = exporter.push("# HELP test_metric\n");
514 assert!(result.is_ok());
515 assert_eq!(exporter.push_count(), 1);
516 let snap = &exporter.snapshots()[0];
517 assert_eq!(snap.metrics_text, "# HELP test_metric\n");
518 assert_eq!(snap.job, "sz-orm");
519 }
520
521 #[test]
524 fn test_summary_new_empty() {
525 let s = Summary::new("latency", "latency summary", vec![0.5, 0.9, 0.99]);
526 assert_eq!(s.count(), 0);
527 assert_eq!(s.sum(), 0.0);
528 assert!(s.quantile(0.5).is_none());
529 }
530
531 #[test]
532 fn test_summary_observe_single() {
533 let s = Summary::new("latency", "help", vec![0.5]);
534 s.observe(1.5);
535 assert_eq!(s.count(), 1);
536 assert!((s.sum() - 1.5).abs() < 1e-9);
537 assert!((s.quantile(0.5).unwrap() - 1.5).abs() < 1e-9);
538 }
539
540 #[test]
541 fn test_summary_observe_multiple_p50() {
542 let s = Summary::new("latency", "help", vec![0.5]);
543 for v in [1.0, 2.0, 3.0, 4.0, 5.0] {
544 s.observe(v);
545 }
546 assert!((s.quantile(0.5).unwrap() - 3.0).abs() < 1e-9);
548 }
549
550 #[test]
551 fn test_summary_observe_multiple_p99() {
552 let s = Summary::new("latency", "help", vec![0.99]);
553 for v in [1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0, 100.0] {
554 s.observe(v);
555 }
556 assert!((s.quantile(0.99).unwrap() - 100.0).abs() < 1e-9);
558 }
559
560 #[test]
561 fn test_summary_quantile_out_of_range() {
562 let s = Summary::new("latency", "help", vec![0.5]);
563 s.observe(1.0);
564 assert!(s.quantile(-0.1).is_none());
565 assert!(s.quantile(1.1).is_none());
566 }
567
568 #[test]
569 fn test_summary_quantile_empty() {
570 let s = Summary::new("latency", "help", vec![0.5]);
571 assert!(s.quantile(0.5).is_none());
572 }
573
574 #[test]
575 fn test_summary_quantile_p0_and_p1() {
576 let s = Summary::new("latency", "help", vec![]);
577 for v in [10.0, 20.0, 30.0] {
578 s.observe(v);
579 }
580 assert!((s.quantile(0.0).unwrap() - 10.0).abs() < 1e-9);
582 assert!((s.quantile(1.0).unwrap() - 30.0).abs() < 1e-9);
584 }
585
586 #[test]
587 fn test_summary_quantiles_all() {
588 let s = Summary::new("latency", "help", vec![0.5, 0.9, 0.99]);
589 for v in [1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0, 10.0] {
590 s.observe(v);
591 }
592 let qs = s.quantiles();
593 assert_eq!(qs.len(), 3);
594 assert!(qs.iter().all(|(_, v)| v.is_some()));
595
596 let qmap: Vec<(f64, f64)> = qs
598 .iter()
599 .filter_map(|(q, v)| v.map(|val| (*q, val)))
600 .collect();
601 let lookup = |target: f64| -> f64 {
602 qmap.iter()
603 .find(|(q, _)| (*q - target).abs() < 1e-9)
604 .map(|(_, v)| *v)
605 .unwrap_or(f64::NAN)
606 };
607
608 let p50 = lookup(0.5);
609 assert!(
610 (5.0..=6.0).contains(&p50),
611 "p50 应在 5-6 之间,实际: {}",
612 p50
613 );
614
615 let p90 = lookup(0.9);
616 assert!(
617 (9.0..=10.0).contains(&p90),
618 "p90 应在 9-10 之间,实际: {}",
619 p90
620 );
621
622 let p99 = lookup(0.99);
623 assert!(
624 (9.0..=10.0).contains(&p99),
625 "p99 应在 9-10 之间,实际: {}",
626 p99
627 );
628 }
629
630 #[test]
631 fn test_summary_unsorted_input_stays_sorted() {
632 let s = Summary::new("latency", "help", vec![0.5]);
633 s.observe(50.0);
634 s.observe(10.0);
635 s.observe(30.0);
636 assert!((s.quantile(0.5).unwrap() - 30.0).abs() < 1e-9);
638 }
639
640 #[test]
641 fn test_summary_render_contains_type() {
642 let s = Summary::new("latency", "latency help", vec![0.5, 0.99]);
643 s.observe(1.0);
644 let output = s.render();
645 assert!(output.contains("# HELP latency latency help"));
646 assert!(output.contains("# TYPE latency summary"));
647 assert!(output.contains("latency{quantile=\"0.5\"}"));
648 assert!(output.contains("latency{quantile=\"0.99\"}"));
649 assert!(output.contains("latency_sum"));
650 assert!(output.contains("latency_count"));
651 }
652
653 #[test]
654 fn test_summary_render_empty_shows_zero() {
655 let s = Summary::new("latency", "help", vec![0.5]);
656 let output = s.render();
657 assert!(output.contains("latency{quantile=\"0.5\"} 0"));
659 assert!(output.contains("latency_count 0"));
660 }
661
662 #[test]
663 fn test_summary_reset() {
664 let s = Summary::new("latency", "help", vec![0.5]);
665 s.observe(1.0);
666 s.observe(2.0);
667 assert_eq!(s.count(), 2);
668
669 s.reset();
670 assert_eq!(s.count(), 0);
671 assert!((s.sum() - 0.0).abs() < 1e-9);
672 assert!(s.quantile(0.5).is_none());
673 }
674
675 #[test]
676 fn test_summary_name() {
677 let s = Summary::new("my_metric", "help", vec![0.5]);
678 assert_eq!(s.name(), "my_metric");
679 }
680
681 #[test]
684 fn test_labeled_histogram_new() {
685 let h = LabeledHistogram::new("requests", "help", vec![0.1, 0.5, 1.0]);
686 assert_eq!(h.label_combination_count(), 0);
687 }
688
689 #[test]
690 fn test_labeled_histogram_observe_single_label() {
691 let h = LabeledHistogram::new("requests", "help", vec![0.1, 0.5, 1.0]);
692 let mut labels = HashMap::new();
693 labels.insert("method".to_string(), "GET".to_string());
694
695 h.observe(&labels, 0.3);
696 assert_eq!(h.count(&labels), 1);
697 assert_eq!(h.label_combination_count(), 1);
698 }
699
700 #[test]
701 fn test_labeled_histogram_observe_multiple_labels() {
702 let h = LabeledHistogram::new("requests", "help", vec![0.1, 0.5, 1.0]);
703
704 let mut get_labels = HashMap::new();
705 get_labels.insert("method".to_string(), "GET".to_string());
706
707 let mut post_labels = HashMap::new();
708 post_labels.insert("method".to_string(), "POST".to_string());
709
710 h.observe(&get_labels, 0.1);
711 h.observe(&get_labels, 0.2);
712 h.observe(&post_labels, 0.5);
713
714 assert_eq!(h.count(&get_labels), 2);
715 assert_eq!(h.count(&post_labels), 1);
716 assert_eq!(h.label_combination_count(), 2);
717 }
718
719 #[test]
720 fn test_labeled_histogram_label_order_independent() {
721 let h = LabeledHistogram::new("requests", "help", vec![0.1, 1.0]);
722
723 let mut labels1 = HashMap::new();
724 labels1.insert("a".to_string(), "1".to_string());
725 labels1.insert("b".to_string(), "2".to_string());
726
727 let mut labels2 = HashMap::new();
728 labels2.insert("b".to_string(), "2".to_string());
729 labels2.insert("a".to_string(), "1".to_string());
730
731 h.observe(&labels1, 0.5);
732 assert_eq!(h.count(&labels2), 1);
734 assert_eq!(h.label_combination_count(), 1);
735 }
736
737 #[test]
738 fn test_labeled_histogram_count_missing_labels() {
739 let h = LabeledHistogram::new("requests", "help", vec![0.1, 1.0]);
740 let labels = HashMap::new();
741 assert_eq!(h.count(&labels), 0);
742 }
743
744 #[test]
745 fn test_labeled_histogram_render_contains_labels() {
746 let h = LabeledHistogram::new("requests", "request help", vec![0.1, 1.0]);
747 let mut labels = HashMap::new();
748 labels.insert("method".to_string(), "GET".to_string());
749 h.observe(&labels, 0.05);
750
751 let output = h.render();
752 assert!(output.contains("# HELP requests request help"));
753 assert!(output.contains("# TYPE requests histogram"));
754 assert!(output.contains("method=\"GET\""));
755 assert!(output.contains("requests_count"));
756 assert!(output.contains("requests_sum"));
757 }
758
759 #[test]
760 fn test_labeled_histogram_render_inf_bucket() {
761 let h = LabeledHistogram::new("req", "help", vec![0.1]);
762 let labels = HashMap::new();
763 h.observe(&labels, 0.05);
764 h.observe(&labels, 5.0);
765 let output = h.render();
766 assert!(output.contains("le=\"+Inf\""));
767 }
768
769 #[test]
772 fn test_pushgateway_config_default() {
773 let config = PushgatewayConfig::default();
774 assert_eq!(config.endpoint, "http://localhost:9091");
775 assert_eq!(config.job, "sz-orm");
776 assert!(config.instance.is_none());
777 }
778
779 #[test]
780 fn test_pushgateway_exporter_new() {
781 let exporter = PushgatewayExporter::new(PushgatewayConfig::default());
782 assert_eq!(exporter.push_count(), 0);
783 assert!(exporter.snapshots().is_empty());
784 }
785
786 #[test]
787 fn test_pushgateway_push_text() {
788 let exporter = PushgatewayExporter::new(PushgatewayConfig::default());
789 exporter.push("metric1 1\n").unwrap();
790 exporter.push("metric2 2\n").unwrap();
791
792 assert_eq!(exporter.push_count(), 2);
793 let snaps = exporter.snapshots();
794 assert_eq!(snaps.len(), 2);
795 assert_eq!(snaps[0].metrics_text, "metric1 1\n");
796 assert_eq!(snaps[1].metrics_text, "metric2 2\n");
797 }
798
799 #[test]
800 fn test_pushgateway_push_from_registry() {
801 let registry = crate::MetricsRegistry::new();
802 let counter = registry.register_counter("test_total", "test");
803 counter.inc();
804
805 let exporter = PushgatewayExporter::new(PushgatewayConfig::default());
806 exporter.push_from_registry(®istry).unwrap();
807
808 assert_eq!(exporter.push_count(), 1);
809 let snap = &exporter.snapshots()[0];
810 assert!(snap.metrics_text.contains("test_total"));
811 }
812
813 #[test]
814 fn test_pushgateway_snapshot_has_metadata() {
815 let config = PushgatewayConfig {
816 endpoint: "http://push:9091".to_string(),
817 job: "myjob".to_string(),
818 instance: Some("inst1".to_string()),
819 };
820 let exporter = PushgatewayExporter::new(config);
821 exporter.push("m 1\n").unwrap();
822
823 let snap = &exporter.snapshots()[0];
824 assert_eq!(snap.job, "myjob");
825 assert_eq!(snap.instance, Some("inst1".to_string()));
826 assert!(snap.timestamp_ms > 0);
827 }
828
829 #[test]
830 fn test_pushgateway_clear() {
831 let exporter = PushgatewayExporter::new(PushgatewayConfig::default());
832 exporter.push("m 1\n").unwrap();
833 assert_eq!(exporter.push_count(), 1);
834
835 exporter.clear();
836 assert_eq!(exporter.push_count(), 0);
837 }
838
839 #[test]
840 fn test_pushgateway_config_access() {
841 let config = PushgatewayConfig {
842 job: "custom".to_string(),
843 ..Default::default()
844 };
845 let exporter = PushgatewayExporter::new(config);
846 assert_eq!(exporter.config().job, "custom");
847 }
848}