1use core::{fmt, ops::ControlFlow};
6
7use emit_core::{
8 and::And,
9 event::{Event, ToEvent},
10 extent::{Extent, ToExtent},
11 or::Or,
12 path::Path,
13 props::{ErasedProps, Props},
14 str::{Str, ToStr},
15 template::{self, Template},
16 timestamp::Timestamp,
17 value::{ToValue, Value},
18 well_known::{
19 KEY_DIST_COUNT, KEY_DIST_EXP_SCALE, KEY_DIST_MAX, KEY_DIST_MIN, KEY_DIST_SUM, KEY_EVT_KIND,
20 KEY_METRIC_AGG, KEY_METRIC_DESCRIPTION, KEY_METRIC_NAME, KEY_METRIC_UNIT, KEY_METRIC_VALUE,
21 },
22};
23
24#[cfg(feature = "alloc")]
25use emit_core::well_known::KEY_DIST_EXP_BUCKETS;
26
27use crate::kind::Kind;
28
29pub use self::{sampler::Sampler, source::Source};
30
31pub struct Metric<'a, P> {
39 mdl: Path<'a>,
40 extent: Option<Extent>,
41 tpl: Option<Template<'a>>,
42 props: P,
43}
44
45impl<'a, P> Metric<'a, P> {
46 pub fn new(mdl: impl Into<Path<'a>>, extent: impl ToExtent, props: P) -> Self {
56 Metric {
57 mdl: mdl.into(),
58 extent: extent.to_extent(),
59 tpl: None,
60 props,
61 }
62 }
63
64 pub fn mdl(&self) -> &Path<'a> {
68 &self.mdl
69 }
70
71 pub fn with_mdl(mut self, mdl: impl Into<Path<'a>>) -> Self {
75 self.mdl = mdl.into();
76 self
77 }
78
79 pub fn extent(&self) -> Option<&Extent> {
83 self.extent.as_ref()
84 }
85
86 pub fn with_extent(mut self, extent: impl ToExtent) -> Self {
90 self.extent = extent.to_extent();
91 self
92 }
93
94 pub fn ts(&self) -> Option<&Timestamp> {
100 self.extent.as_ref().map(|extent| extent.as_point())
101 }
102
103 pub fn ts_start(&self) -> Option<&Timestamp> {
109 self.extent
110 .as_ref()
111 .and_then(|extent| extent.as_range())
112 .map(|span| &span.start)
113 }
114
115 pub fn tpl(&self) -> &Template<'a> {
119 self.tpl.as_ref().unwrap_or(&TEMPLATE)
120 }
121
122 pub fn with_tpl(mut self, tpl: impl Into<Template<'a>>) -> Self {
126 self.tpl = Some(tpl.into());
127 self
128 }
129
130 pub fn props(&self) -> &P {
134 &self.props
135 }
136
137 pub fn props_mut(&mut self) -> &mut P {
141 &mut self.props
142 }
143
144 pub fn with_props<U>(self, props: U) -> Metric<'a, U> {
148 Metric {
149 mdl: self.mdl,
150 extent: self.extent,
151 tpl: self.tpl,
152 props,
153 }
154 }
155
156 pub fn map_props<U>(self, map: impl FnOnce(P) -> U) -> Metric<'a, U> {
160 Metric {
161 mdl: self.mdl,
162 extent: self.extent,
163 tpl: self.tpl,
164 props: map(self.props),
165 }
166 }
167}
168
169impl<'a, P: Props> Metric<'a, P> {
170 pub fn name(&self) -> Option<Str<'_>> {
174 self.props.pull(KEY_METRIC_NAME)
175 }
176
177 pub fn with_name(
181 self,
182 name: impl Into<Str<'a>>,
183 ) -> Metric<'a, And<(&'static str, Str<'a>), P>> {
184 self.map_props(|props| (KEY_METRIC_NAME, name.into()).and_props(props))
185 }
186
187 pub fn description(&self) -> Option<Str<'_>> {
191 self.props.pull(KEY_METRIC_DESCRIPTION)
192 }
193
194 pub fn with_description(
198 self,
199 description: impl Into<Str<'a>>,
200 ) -> Metric<'a, And<(&'static str, Str<'a>), P>> {
201 self.map_props(|props| (KEY_METRIC_DESCRIPTION, description.into()).and_props(props))
202 }
203
204 pub fn agg(&self) -> Option<Str<'_>> {
210 self.props.pull(KEY_METRIC_AGG)
211 }
212
213 pub fn with_agg(self, agg: impl Into<Str<'a>>) -> Metric<'a, And<(&'static str, Str<'a>), P>> {
219 self.map_props(|props| (KEY_METRIC_AGG, agg.into()).and_props(props))
220 }
221
222 pub fn unit(&self) -> Option<Str<'_>> {
226 self.props.pull(KEY_METRIC_UNIT)
227 }
228
229 pub fn with_unit(
233 self,
234 unit: impl Into<Str<'a>>,
235 ) -> Metric<'a, And<(&'static str, Str<'a>), P>> {
236 self.map_props(|props| (KEY_METRIC_UNIT, unit.into()).and_props(props))
237 }
238
239 pub fn value(&self) -> Option<Value<'_>> {
243 self.props.get(KEY_METRIC_VALUE)
244 }
245
246 pub fn with_value(
250 self,
251 value: impl Into<Value<'a>>,
252 ) -> Metric<'a, And<(&'static str, Value<'a>), P>> {
253 self.map_props(|props| (KEY_METRIC_VALUE, value.into()).and_props(props))
254 }
255
256 pub fn dist_min(&self) -> Option<Value<'_>> {
260 self.props.get(KEY_DIST_MIN)
261 }
262
263 pub fn with_dist_min(
267 self,
268 dist_min: impl Into<Value<'a>>,
269 ) -> Metric<'a, And<(&'static str, Value<'a>), P>> {
270 self.map_props(|props| (KEY_DIST_MIN, dist_min.into()).and_props(props))
271 }
272
273 pub fn dist_max(&self) -> Option<Value<'_>> {
277 self.props.get(KEY_DIST_MAX)
278 }
279
280 pub fn with_dist_max(
284 self,
285 dist_max: impl Into<Value<'a>>,
286 ) -> Metric<'a, And<(&'static str, Value<'a>), P>> {
287 self.map_props(|props| (KEY_DIST_MAX, dist_max.into()).and_props(props))
288 }
289
290 pub fn dist_count(&self) -> Option<Value<'_>> {
294 self.props.get(KEY_DIST_COUNT)
295 }
296
297 pub fn with_dist_count(
301 self,
302 dist_count: impl Into<Value<'a>>,
303 ) -> Metric<'a, And<(&'static str, Value<'a>), P>> {
304 self.map_props(|props| (KEY_DIST_COUNT, dist_count.into()).and_props(props))
305 }
306
307 pub fn dist_sum(&self) -> Option<Value<'_>> {
311 self.props.get(KEY_DIST_SUM)
312 }
313
314 pub fn with_dist_sum(
318 self,
319 dist_sum: impl Into<Value<'a>>,
320 ) -> Metric<'a, And<(&'static str, Value<'a>), P>> {
321 self.map_props(|props| (KEY_DIST_SUM, dist_sum.into()).and_props(props))
322 }
323
324 pub fn dist_exp_scale(&self) -> Option<i32> {
328 self.props.pull(KEY_DIST_EXP_SCALE)
329 }
330
331 pub fn with_dist_exp_scale(
335 self,
336 dist_exp_scale: impl Into<i32>,
337 ) -> Metric<'a, And<(&'static str, i32), P>> {
338 self.map_props(|props| (KEY_DIST_EXP_SCALE, dist_exp_scale.into()).and_props(props))
339 }
340
341 #[cfg(feature = "alloc")]
345 pub fn dist_exp_buckets(&self) -> Option<exp::BucketSet> {
346 self.props.pull(KEY_DIST_EXP_BUCKETS)
347 }
348
349 #[cfg(feature = "alloc")]
353 pub fn with_dist_exp_buckets(
354 self,
355 dist_exp_buckets: impl Into<exp::BucketSet>,
356 ) -> Metric<'a, And<(&'static str, exp::BucketSet), P>> {
357 self.map_props(|props| (KEY_DIST_EXP_BUCKETS, dist_exp_buckets.into()).and_props(props))
358 }
359}
360
361impl<'a, P: Props> fmt::Debug for Metric<'a, P> {
362 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
363 fmt::Debug::fmt(&self.to_event(), f)
364 }
365}
366
367impl<'a, P: Props> ToEvent for Metric<'a, P> {
368 type Props<'b>
369 = &'b Self
370 where
371 Self: 'b;
372
373 fn to_event<'b>(&'b self) -> Event<'b, Self::Props<'b>> {
374 Event::new(
375 self.mdl.by_ref(),
376 self.tpl().by_ref(),
377 self.extent.clone(),
378 self,
379 )
380 }
381}
382
383impl<'a, P: Props> Metric<'a, P> {
384 pub fn by_ref<'b>(&'b self) -> Metric<'b, &'b P> {
388 Metric {
389 mdl: self.mdl.by_ref(),
390 extent: self.extent.clone(),
391 tpl: self.tpl.as_ref().map(|tpl| tpl.by_ref()),
392 props: &self.props,
393 }
394 }
395
396 pub fn erase<'b>(&'b self) -> Metric<'b, &'b dyn ErasedProps> {
400 Metric {
401 mdl: self.mdl.by_ref(),
402 extent: self.extent.clone(),
403 tpl: self.tpl.as_ref().map(|tpl| tpl.by_ref()),
404 props: &self.props,
405 }
406 }
407}
408
409impl<'a, P> ToExtent for Metric<'a, P> {
410 fn to_extent(&self) -> Option<Extent> {
411 self.extent.clone()
412 }
413}
414
415impl<'a, P: Props> Props for Metric<'a, P> {
416 fn for_each<'kv, F: FnMut(Str<'kv>, Value<'kv>) -> ControlFlow<()>>(
417 &'kv self,
418 mut for_each: F,
419 ) -> ControlFlow<()> {
420 for_each(KEY_EVT_KIND.to_str(), Kind::Metric.to_value())?;
421
422 self.props.for_each(for_each)
423 }
424}
425
426const TEMPLATE_PARTS: &'static [template::Part<'static>] = &[
428 template::Part::hole("metric_agg"),
429 template::Part::text(" of "),
430 template::Part::hole("metric_name"),
431 template::Part::text(" is "),
432 template::Part::hole("metric_value"),
433];
434
435static TEMPLATE: Template<'static> = Template::new(TEMPLATE_PARTS);
436
437pub mod source {
438 use self::sampler::ErasedSampler;
445
446 use super::*;
447
448 pub trait Source {
452 fn sample_metrics<S: sampler::Sampler>(&self, sampler: S);
456
457 fn and_sample<U>(self, other: U) -> And<Self, U>
461 where
462 Self: Sized,
463 {
464 And::new(self, other)
465 }
466 }
467
468 impl<'a, T: Source + ?Sized> Source for &'a T {
469 fn sample_metrics<S: sampler::Sampler>(&self, sampler: S) {
470 (**self).sample_metrics(sampler)
471 }
472 }
473
474 impl<T: Source> Source for Option<T> {
475 fn sample_metrics<S: sampler::Sampler>(&self, sampler: S) {
476 if let Some(source) = self {
477 source.sample_metrics(sampler);
478 }
479 }
480 }
481
482 #[cfg(feature = "alloc")]
483 impl<'a, T: Source + ?Sized + 'a> Source for alloc::boxed::Box<T> {
484 fn sample_metrics<S: sampler::Sampler>(&self, sampler: S) {
485 (**self).sample_metrics(sampler)
486 }
487 }
488
489 #[cfg(feature = "alloc")]
490 impl<'a, T: Source + ?Sized + 'a> Source for alloc::sync::Arc<T> {
491 fn sample_metrics<S: sampler::Sampler>(&self, sampler: S) {
492 (**self).sample_metrics(sampler)
493 }
494 }
495
496 impl<T: Source, U: Source> Source for And<T, U> {
497 fn sample_metrics<S: sampler::Sampler>(&self, sampler: S) {
498 self.left().sample_metrics(&sampler);
499 self.right().sample_metrics(&sampler);
500 }
501 }
502
503 impl<T: Source, U: Source> Source for Or<T, U> {
504 fn sample_metrics<S: sampler::Sampler>(&self, sampler: S) {
505 self.left().sample_metrics(&sampler);
506 self.right().sample_metrics(&sampler);
507 }
508 }
509
510 impl<'a, P: Props> Source for Metric<'a, P> {
511 fn sample_metrics<S: Sampler>(&self, sampler: S) {
512 sampler.metric(self.by_ref());
513 }
514 }
515
516 pub struct FromFn<F = fn(&mut dyn ErasedSampler)>(F);
522
523 pub const fn from_fn<F: Fn(&mut dyn ErasedSampler)>(source: F) -> FromFn<F> {
527 FromFn::new(source)
528 }
529
530 impl<F> FromFn<F> {
531 pub const fn new(source: F) -> Self {
535 FromFn(source)
536 }
537 }
538
539 impl<F: Fn(&mut dyn ErasedSampler)> Source for FromFn<F> {
540 fn sample_metrics<S: sampler::Sampler>(&self, mut sampler: S) {
541 (self.0)(&mut sampler)
542 }
543 }
544
545 mod internal {
546 use super::*;
547
548 pub trait DispatchSource {
549 fn dispatch_sample_metrics(&self, sampler: &dyn sampler::ErasedSampler);
550 }
551
552 pub trait SealedSource {
553 fn erase_source(&self) -> crate::internal::Erased<&dyn DispatchSource>;
554 }
555 }
556
557 pub trait ErasedSource: internal::SealedSource {}
563
564 impl<T: Source> ErasedSource for T {}
565
566 impl<T: Source> internal::SealedSource for T {
567 fn erase_source(&self) -> crate::internal::Erased<&dyn internal::DispatchSource> {
568 crate::internal::Erased(self)
569 }
570 }
571
572 impl<T: Source> internal::DispatchSource for T {
573 fn dispatch_sample_metrics(&self, sampler: &dyn sampler::ErasedSampler) {
574 self.sample_metrics(sampler)
575 }
576 }
577
578 impl<'a> Source for dyn ErasedSource + 'a {
579 fn sample_metrics<S: sampler::Sampler>(&self, sampler: S) {
580 self.erase_source().0.dispatch_sample_metrics(&sampler)
581 }
582 }
583
584 impl<'a> Source for dyn ErasedSource + Send + Sync + 'a {
585 fn sample_metrics<S: sampler::Sampler>(&self, sampler: S) {
586 (self as &(dyn ErasedSource + 'a)).sample_metrics(sampler)
587 }
588 }
589
590 #[cfg(test)]
591 mod tests {
592 use super::*;
593 use std::cell::Cell;
594
595 #[test]
596 fn source_sample_emit() {
597 struct MySource;
598
599 impl Source for MySource {
600 fn sample_metrics<S: Sampler>(&self, sampler: S) {
601 sampler.metric(Metric::new(
602 Path::new_raw("test"),
603 crate::Empty,
604 crate::Empty,
605 ));
606
607 sampler.metric(Metric::new(
608 Path::new_raw("test"),
609 crate::Empty,
610 crate::Empty,
611 ));
612 }
613 }
614
615 let calls = Cell::new(0);
616
617 MySource.sample_metrics(sampler::from_fn(|_| {
618 calls.set(calls.get() + 1);
619 }));
620
621 assert_eq!(2, calls.get());
622 }
623
624 #[test]
625 fn and_sample() {
626 let calls = Cell::new(0);
627
628 from_fn(|sampler| {
629 sampler.metric(Metric::new(
630 Path::new_raw("test"),
631 crate::Empty,
632 crate::Empty,
633 ));
634 })
635 .and_sample(from_fn(|sampler| {
636 sampler.metric(Metric::new(
637 Path::new_raw("test"),
638 crate::Empty,
639 crate::Empty,
640 ));
641 }))
642 .sample_metrics(sampler::from_fn(|_| {
643 calls.set(calls.get() + 1);
644 }));
645
646 assert_eq!(2, calls.get());
647 }
648
649 #[test]
650 fn from_fn_source() {
651 let calls = Cell::new(0);
652
653 from_fn(|sampler| {
654 sampler.metric(Metric::new(
655 Path::new_raw("test"),
656 crate::Empty,
657 crate::Empty,
658 ));
659
660 sampler.metric(Metric::new(
661 Path::new_raw("test"),
662 crate::Empty,
663 crate::Empty,
664 ));
665 })
666 .sample_metrics(sampler::from_fn(|_| {
667 calls.set(calls.get() + 1);
668 }));
669
670 assert_eq!(2, calls.get());
671 }
672
673 #[test]
674 fn erased_source() {
675 let source = from_fn(|sampler| {
676 sampler.metric(Metric::new(
677 Path::new_raw("test"),
678 crate::Empty,
679 crate::Empty,
680 ));
681
682 sampler.metric(Metric::new(
683 Path::new_raw("test"),
684 crate::Empty,
685 crate::Empty,
686 ));
687 });
688
689 let source = &source as &dyn ErasedSource;
690
691 let calls = Cell::new(0);
692
693 source.sample_metrics(sampler::from_fn(|_| {
694 calls.set(calls.get() + 1);
695 }));
696
697 assert_eq!(2, calls.get());
698 }
699
700 #[test]
701 fn metric_as_source() {
702 let sampler = sampler::from_fn(|metric| {
703 assert_eq!("metric", metric.name().unwrap().to_string());
704 assert_eq!("count", metric.agg().unwrap().to_string());
705 });
706
707 let metric = Metric::new(
708 Path::new_raw("test"),
709 crate::Empty,
710 [("metric_name", "metric"), ("metric_agg", "count")],
711 );
712
713 metric.sample_metrics(sampler);
714 }
715 }
716}
717
718#[cfg(feature = "alloc")]
719mod alloc_support {
720 use super::*;
721
722 use alloc::{boxed::Box, vec::Vec};
723 use core::ops::Range;
724
725 use crate::{
726 clock::{Clock, ErasedClock},
727 emitter::Emitter,
728 metric::source::{ErasedSource, Source},
729 };
730
731 pub struct Reporter {
749 sources: Vec<Box<dyn ErasedSource + Send + Sync>>,
750 clock: ReporterClock,
751 }
752
753 impl Reporter {
754 pub const fn new() -> Self {
761 Reporter {
762 sources: Vec::new(),
763 clock: ReporterClock::Default,
764 }
765 }
766
767 pub fn normalize_with_clock(
771 &mut self,
772 clock: impl Clock + Send + Sync + 'static,
773 ) -> &mut Self {
774 self.clock = ReporterClock::Other(Some(Box::new(clock)));
775
776 self
777 }
778
779 pub fn without_normalization(&mut self) -> &mut Self {
783 self.clock = ReporterClock::Other(None);
784
785 self
786 }
787
788 pub fn add_source(&mut self, source: impl Source + Send + Sync + 'static) -> &mut Self {
792 self.sources.push(Box::new(source));
793
794 self
795 }
796
797 pub fn sample_metrics<S: sampler::Sampler>(&self, sampler: S) {
801 let sampler = TimeNormalizer::new(self.clock.now(), sampler);
802
803 for source in &self.sources {
804 source.sample_metrics(&sampler);
805 }
806 }
807
808 pub fn emit_metrics<E: Emitter>(&self, emitter: E) {
812 self.sample_metrics(sampler::from_emitter(emitter).with_sampled_at(self.clock.now()))
813 }
814 }
815
816 impl Source for Reporter {
817 fn sample_metrics<S: sampler::Sampler>(&self, sampler: S) {
818 self.sample_metrics(sampler)
819 }
820 }
821
822 struct TimeNormalizer<S> {
823 now: Option<Timestamp>,
824 inner: S,
825 }
826
827 impl<S> TimeNormalizer<S> {
828 fn new(now: Option<Timestamp>, sampler: S) -> TimeNormalizer<S> {
829 TimeNormalizer {
830 now,
831 inner: sampler,
832 }
833 }
834 }
835
836 impl<S: Sampler> Sampler for TimeNormalizer<S> {
837 fn metric<P: Props>(&self, metric: Metric<P>) {
838 if let Some(now) = self.now {
839 let extent = metric.extent();
840
841 let extent = if let Some(range) = extent.and_then(|extent| extent.as_range()) {
842 normalize_range(now, range)
844 .map(Extent::range)
845 .unwrap_or_else(|| Extent::range(range.clone()))
847 } else {
848 Extent::point(now)
850 };
851
852 self.inner.metric(metric.with_extent(extent))
853 } else {
854 self.inner.metric(metric)
855 }
856 }
857
858 fn sampled_at(&self) -> Option<Timestamp> {
859 self.now
860 }
861 }
862
863 fn normalize_range(now: Timestamp, range: &Range<Timestamp>) -> Option<Range<Timestamp>> {
864 let len = range.end.duration_since(range.start)?;
867 let start = now.checked_sub(len)?;
868
869 Some(start..now)
870 }
871
872 enum ReporterClock {
873 Default,
874 Other(Option<Box<dyn ErasedClock + Send + Sync>>),
875 }
876
877 impl Clock for ReporterClock {
878 fn now(&self) -> Option<Timestamp> {
879 match self {
880 ReporterClock::Default => crate::platform::DefaultClock::new().now(),
881 ReporterClock::Other(clock) => clock.now(),
882 }
883 }
884 }
885
886 #[cfg(test)]
887 mod tests {
888 use super::*;
889 use std::time::Duration;
890
891 #[cfg(all(
892 target_arch = "wasm32",
893 target_vendor = "unknown",
894 target_os = "unknown"
895 ))]
896 use wasm_bindgen_test::*;
897
898 #[test]
899 fn reporter_is_send_sync() {
900 fn check<T: Send + Sync>() {}
901
902 check::<Reporter>();
903 }
904
905 #[test]
906 #[cfg(not(miri))]
907 #[cfg_attr(
908 all(
909 target_arch = "wasm32",
910 target_vendor = "unknown",
911 target_os = "unknown"
912 ),
913 wasm_bindgen_test
914 )]
915 fn reporter_sample() {
916 use std::cell::Cell;
917
918 let mut reporter = Reporter::new();
919
920 reporter
921 .add_source(source::from_fn(|sampler| {
922 sampler.metric(Metric::new(
923 Path::new_raw("test"),
924 crate::Empty,
925 crate::Empty,
926 ));
927 }))
928 .add_source(source::from_fn(|sampler| {
929 sampler.metric(Metric::new(
930 Path::new_raw("test"),
931 crate::Empty,
932 crate::Empty,
933 ));
934 }));
935
936 let calls = Cell::new(0);
937
938 reporter.sample_metrics(sampler::from_fn(|_| {
939 calls.set(calls.get() + 1);
940 }));
941
942 assert_eq!(2, calls.get());
943 }
944
945 struct TestClock(Option<Timestamp>);
946
947 impl Clock for TestClock {
948 fn now(&self) -> Option<Timestamp> {
949 self.0
950 }
951 }
952
953 #[test]
954 #[cfg(all(feature = "std", not(miri)))]
955 #[cfg_attr(
956 all(
957 target_arch = "wasm32",
958 target_vendor = "unknown",
959 target_os = "unknown"
960 ),
961 wasm_bindgen_test
962 )]
963 fn reporter_normalize_std() {
964 let mut reporter = Reporter::new();
965
966 reporter.add_source(source::from_fn(|sampler| {
967 sampler.metric(Metric::new(
968 Path::new_raw("test"),
969 crate::Empty,
970 crate::Empty,
971 ));
972 }));
973
974 reporter.sample_metrics(sampler::from_fn(|metric| {
975 assert!(metric.extent().is_some());
976 }));
977 }
978
979 #[wasm_bindgen_test]
980 #[cfg(all(feature = "web", not(miri)))]
981 #[cfg(all(
982 target_arch = "wasm32",
983 target_vendor = "unknown",
984 target_os = "unknown"
985 ))]
986 fn reporter_normalize_web() {
987 let mut reporter = Reporter::new();
988
989 reporter.add_source(source::from_fn(|sampler| {
990 sampler.metric(Metric::new(
991 Path::new_raw("test"),
992 crate::Empty,
993 crate::Empty,
994 ));
995 }));
996
997 reporter.sample_metrics(sampler::from_fn(|metric| {
998 assert!(metric.extent().is_some());
999 }));
1000 }
1001
1002 #[test]
1003 #[cfg_attr(
1004 all(
1005 target_arch = "wasm32",
1006 target_vendor = "unknown",
1007 target_os = "unknown"
1008 ),
1009 wasm_bindgen_test
1010 )]
1011 fn reporter_normalize_empty_extent() {
1012 let mut reporter = Reporter::new();
1013
1014 reporter.normalize_with_clock(TestClock(Some(Timestamp::MIN)));
1015
1016 reporter.add_source(source::from_fn(|sampler| {
1017 sampler.metric(Metric::new(
1018 Path::new_raw("test"),
1019 crate::Empty,
1020 crate::Empty,
1021 ));
1022 }));
1023
1024 reporter.sample_metrics(sampler::from_fn(|metric| {
1025 assert_eq!(Timestamp::MIN, metric.extent().unwrap().as_point());
1026 }));
1027 }
1028
1029 #[test]
1030 #[cfg_attr(
1031 all(
1032 target_arch = "wasm32",
1033 target_vendor = "unknown",
1034 target_os = "unknown"
1035 ),
1036 wasm_bindgen_test
1037 )]
1038 fn reporter_normalize_point_extent() {
1039 let mut reporter = Reporter::new();
1040
1041 reporter.normalize_with_clock(TestClock(Some(
1042 Timestamp::from_unix(Duration::from_secs(37)).unwrap(),
1043 )));
1044
1045 reporter.add_source(source::from_fn(|sampler| {
1046 sampler.metric(Metric::new(
1047 Path::new_raw("test"),
1048 crate::Empty,
1049 crate::Empty,
1050 ));
1051 }));
1052
1053 reporter.sample_metrics(sampler::from_fn(|metric| {
1054 assert_eq!(
1055 Timestamp::from_unix(Duration::from_secs(37)).unwrap(),
1056 metric.extent().unwrap().as_point()
1057 );
1058 }));
1059 }
1060
1061 #[test]
1062 #[cfg_attr(
1063 all(
1064 target_arch = "wasm32",
1065 target_vendor = "unknown",
1066 target_os = "unknown"
1067 ),
1068 wasm_bindgen_test
1069 )]
1070 fn reporter_normalize_range_extent() {
1071 let mut reporter = Reporter::new();
1072
1073 reporter.normalize_with_clock(TestClock(Some(
1074 Timestamp::from_unix(Duration::from_secs(350)).unwrap(),
1075 )));
1076
1077 reporter.add_source(source::from_fn(|sampler| {
1078 sampler.metric(Metric::new(
1079 Path::new_raw("test"),
1080 Timestamp::from_unix(Duration::from_secs(100)).unwrap()
1081 ..Timestamp::from_unix(Duration::from_secs(200)).unwrap(),
1082 crate::Empty,
1083 ));
1084 }));
1085
1086 reporter.sample_metrics(sampler::from_fn(|metric| {
1087 assert_eq!(
1088 Timestamp::from_unix(Duration::from_secs(250)).unwrap()
1089 ..Timestamp::from_unix(Duration::from_secs(350)).unwrap(),
1090 metric.extent().unwrap().as_range().unwrap().clone()
1091 );
1092 }));
1093 }
1094 }
1095}
1096
1097#[cfg(feature = "alloc")]
1098pub use self::alloc_support::*;
1099
1100pub mod sampler {
1101 use emit_core::{
1108 clock::Clock, ctxt::Ctxt, emitter::Emitter, empty::Empty, filter::Filter, rng::Rng,
1109 runtime::Runtime,
1110 };
1111
1112 use super::*;
1113
1114 pub trait Sampler {
1118 fn metric<P: Props>(&self, metric: Metric<P>);
1122
1123 fn sampled_at(&self) -> Option<Timestamp> {
1129 None
1130 }
1131
1132 fn with_sampled_at(self, now: Option<Timestamp>) -> WithSampledAt<Self>
1136 where
1137 Self: Sized,
1138 {
1139 WithSampledAt::new(self, now)
1140 }
1141 }
1142
1143 impl<'a, T: Sampler + ?Sized> Sampler for &'a T {
1144 fn metric<P: Props>(&self, metric: Metric<P>) {
1145 (**self).metric(metric)
1146 }
1147
1148 fn sampled_at(&self) -> Option<Timestamp> {
1149 (**self).sampled_at()
1150 }
1151 }
1152
1153 impl Sampler for Empty {
1154 fn metric<P: Props>(&self, _: Metric<P>) {}
1155 }
1156
1157 pub struct WithSampledAt<S> {
1161 sampler: S,
1162 now: Option<Timestamp>,
1163 }
1164
1165 impl<S> WithSampledAt<S> {
1166 pub const fn new(sampler: S, now: Option<Timestamp>) -> Self {
1170 WithSampledAt { sampler, now }
1171 }
1172 }
1173
1174 impl<S: Sampler> Sampler for WithSampledAt<S> {
1175 fn metric<P: Props>(&self, metric: Metric<P>) {
1176 self.sampler.metric(metric)
1177 }
1178
1179 fn sampled_at(&self) -> Option<Timestamp> {
1180 self.now
1181 }
1182 }
1183
1184 pub struct FromEmitter<E>(E);
1192
1193 impl<E: Emitter> Sampler for FromEmitter<E> {
1194 fn metric<P: Props>(&self, metric: Metric<P>) {
1195 self.0.emit(metric)
1196 }
1197 }
1198
1199 impl<E> FromEmitter<E> {
1200 pub const fn new(emitter: E) -> Self {
1204 FromEmitter(emitter)
1205 }
1206 }
1207
1208 pub const fn from_emitter<E: Emitter>(emitter: E) -> FromEmitter<E> {
1214 FromEmitter(emitter)
1215 }
1216
1217 pub struct FromRuntime<'a, E, F, C, T, R>(&'a Runtime<E, F, C, T, R>);
1225
1226 impl<'a, E: Emitter, F: Filter, C: Ctxt, T: Clock, R: Rng> Sampler
1227 for FromRuntime<'a, E, F, C, T, R>
1228 {
1229 fn metric<P: Props>(&self, metric: Metric<P>) {
1230 self.0.emit(metric)
1231 }
1232 }
1233
1234 impl<'a, E, F, C, T, R> FromRuntime<'a, E, F, C, T, R> {
1235 pub const fn new(rt: &'a Runtime<E, F, C, T, R>) -> Self {
1239 FromRuntime(rt)
1240 }
1241 }
1242
1243 pub const fn from_runtime<'a, E: Emitter, F: Filter, C: Ctxt, T: Clock, R: Rng>(
1249 rt: &'a Runtime<E, F, C, T, R>,
1250 ) -> FromRuntime<'a, E, F, C, T, R> {
1251 FromRuntime(rt)
1252 }
1253
1254 pub struct FromFn<F = fn(Metric<&dyn ErasedProps>)>(F);
1260
1261 pub const fn from_fn<F: Fn(Metric<&dyn ErasedProps>)>(f: F) -> FromFn<F> {
1265 FromFn(f)
1266 }
1267
1268 impl<F> FromFn<F> {
1269 pub const fn new(sampler: F) -> FromFn<F> {
1273 FromFn(sampler)
1274 }
1275 }
1276
1277 impl<F: Fn(Metric<&dyn ErasedProps>)> Sampler for FromFn<F> {
1278 fn metric<P: Props>(&self, metric: Metric<P>) {
1279 (self.0)(metric.erase())
1280 }
1281 }
1282
1283 mod internal {
1284 use super::*;
1285
1286 pub trait DispatchSampler {
1287 fn dispatch_metric(&self, metric: Metric<&dyn ErasedProps>);
1288
1289 fn dispatch_sampled_at(&self) -> Option<Timestamp>;
1290 }
1291
1292 pub trait SealedSampler {
1293 fn erase_sampler(&self) -> crate::internal::Erased<&dyn DispatchSampler>;
1294 }
1295 }
1296
1297 pub trait ErasedSampler: internal::SealedSampler {}
1303
1304 impl<T: Sampler> ErasedSampler for T {}
1305
1306 impl<T: Sampler> internal::SealedSampler for T {
1307 fn erase_sampler(&self) -> crate::internal::Erased<&dyn internal::DispatchSampler> {
1308 crate::internal::Erased(self)
1309 }
1310 }
1311
1312 impl<T: Sampler> internal::DispatchSampler for T {
1313 fn dispatch_metric(&self, metric: Metric<&dyn ErasedProps>) {
1314 self.metric(metric)
1315 }
1316
1317 fn dispatch_sampled_at(&self) -> Option<Timestamp> {
1318 self.sampled_at()
1319 }
1320 }
1321
1322 impl<'a> Sampler for dyn ErasedSampler + 'a {
1323 fn metric<P: Props>(&self, metric: Metric<P>) {
1324 self.erase_sampler().0.dispatch_metric(metric.erase())
1325 }
1326
1327 fn sampled_at(&self) -> Option<Timestamp> {
1328 self.erase_sampler().0.dispatch_sampled_at()
1329 }
1330 }
1331
1332 impl<'a> Sampler for dyn ErasedSampler + Send + Sync + 'a {
1333 fn metric<P: Props>(&self, metric: Metric<P>) {
1334 (self as &(dyn ErasedSampler + 'a)).metric(metric)
1335 }
1336
1337 fn sampled_at(&self) -> Option<Timestamp> {
1338 (self as &(dyn ErasedSampler + 'a)).sampled_at()
1339 }
1340 }
1341
1342 #[cfg(test)]
1343 mod tests {
1344 use super::*;
1345
1346 use emit_core::{emitter, runtime::Runtime};
1347
1348 use std::cell::Cell;
1349
1350 #[test]
1351 fn from_fn_sampler() {
1352 let called = Cell::new(false);
1353
1354 let sampler = from_fn(|metric| {
1355 assert_eq!("test", metric.name().unwrap());
1356
1357 called.set(true);
1358 });
1359
1360 sampler.metric(Metric::new(
1361 Path::new_raw("test"),
1362 Empty,
1363 ("metric_name", "test"),
1364 ));
1365
1366 assert!(called.get());
1367 }
1368
1369 #[test]
1370 fn erased_sampler() {
1371 let called = Cell::new(false);
1372
1373 let sampler = from_fn(|metric| {
1374 assert_eq!("test", metric.name().unwrap());
1375
1376 called.set(true);
1377 });
1378
1379 let sampler = &sampler as &dyn ErasedSampler;
1380
1381 sampler.metric(Metric::new(
1382 Path::new_raw("test"),
1383 crate::Empty,
1384 ("metric_name", "test"),
1385 ));
1386
1387 assert!(called.get());
1388 }
1389
1390 #[test]
1391 fn from_runtime_sampler() {
1392 let called = Cell::new(false);
1393
1394 let rt = Runtime::default().with_emitter(emitter::from_fn(|_| {
1395 called.set(true);
1396 }));
1397
1398 let sampler = from_runtime(&rt);
1399
1400 sampler.metric(Metric::new(
1401 Path::new_raw("test"),
1402 crate::Empty,
1403 crate::Empty,
1404 ));
1405
1406 assert!(called.get());
1407 }
1408 }
1409}
1410
1411pub mod exp {
1412 use crate::{
1417 platform::libm,
1418 value::{FromValue, ToValue, Value},
1419 };
1420
1421 use core::{cmp, fmt, hash, str::FromStr};
1422
1423 #[derive(Debug)]
1427 pub struct ParsePointError {}
1428
1429 impl fmt::Display for ParsePointError {
1430 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1431 write!(f, "the input was not a valid point")
1432 }
1433 }
1434
1435 #[cfg(feature = "std")]
1436 impl std::error::Error for ParsePointError {}
1437
1438 #[derive(Clone, Copy)]
1446 #[repr(transparent)]
1447 pub struct Point(f64);
1448
1449 impl Point {
1450 pub const fn new(value: f64) -> Self {
1454 Point(value)
1455 }
1456
1457 pub fn try_from_str(s: &str) -> Result<Self, ParsePointError> {
1461 Ok(Point::new(s.parse().map_err(|_| ParsePointError {})?))
1462 }
1463
1464 pub const fn get(&self) -> f64 {
1468 self.0
1469 }
1470
1471 pub const fn is_sign_positive(&self) -> bool {
1475 self.get().is_sign_positive()
1476 }
1477
1478 pub const fn is_sign_negative(&self) -> bool {
1482 self.get().is_sign_negative()
1483 }
1484
1485 pub const fn is_zero_bucket(&self) -> bool {
1491 self.get() == 0.0
1492 }
1493
1494 pub const fn is_positive_bucket(&self) -> bool {
1500 self.is_indexable() && self.is_sign_positive()
1501 }
1502
1503 pub const fn is_negative_bucket(&self) -> bool {
1509 self.is_indexable() && self.is_sign_negative()
1510 }
1511
1512 pub const fn is_indexable(&self) -> bool {
1521 let value = self.get();
1522
1523 value != 0.0 && value.is_finite()
1524 }
1525 }
1526
1527 impl From<f64> for Point {
1528 fn from(value: f64) -> Self {
1529 Point::new(value)
1530 }
1531 }
1532
1533 impl From<Point> for f64 {
1534 fn from(value: Point) -> Self {
1535 value.get()
1536 }
1537 }
1538
1539 impl PartialEq for Point {
1540 fn eq(&self, other: &Self) -> bool {
1541 self.cmp(other) == cmp::Ordering::Equal
1542 }
1543 }
1544
1545 impl Eq for Point {}
1546
1547 impl PartialOrd for Point {
1548 fn partial_cmp(&self, other: &Self) -> Option<cmp::Ordering> {
1549 Some(self.cmp(other))
1550 }
1551 }
1552
1553 impl Ord for Point {
1554 fn cmp(&self, other: &Self) -> cmp::Ordering {
1555 libm::cmp(self.get()).cmp(&libm::cmp(other.get()))
1556 }
1557 }
1558
1559 impl hash::Hash for Point {
1560 fn hash<H: hash::Hasher>(&self, state: &mut H) {
1561 libm::cmp(self.get()).hash(state)
1562 }
1563 }
1564
1565 impl fmt::Debug for Point {
1566 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1567 fmt::Debug::fmt(&self.get(), f)
1568 }
1569 }
1570
1571 impl fmt::Display for Point {
1572 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1573 fmt::Display::fmt(&self.get(), f)
1574 }
1575 }
1576
1577 impl FromStr for Point {
1578 type Err = ParsePointError;
1579
1580 fn from_str(s: &str) -> Result<Self, Self::Err> {
1581 Self::try_from_str(s)
1582 }
1583 }
1584
1585 #[cfg(feature = "sval")]
1586 impl sval::Value for Point {
1587 fn stream<'sval, S: sval::Stream<'sval> + ?Sized>(
1588 &'sval self,
1589 stream: &mut S,
1590 ) -> sval::Result {
1591 stream.f64(self.get())
1592 }
1593 }
1594
1595 #[cfg(feature = "serde")]
1596 impl serde::Serialize for Point {
1597 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1598 where
1599 S: serde::Serializer,
1600 {
1601 serializer.serialize_f64(self.get())
1602 }
1603 }
1604
1605 impl ToValue for Point {
1606 fn to_value(&self) -> Value<'_> {
1607 Value::capture_display(self)
1608 }
1609 }
1610
1611 impl<'v> FromValue<'v> for Point {
1612 fn from_value(value: Value<'v>) -> Option<Self>
1613 where
1614 Self: Sized,
1615 {
1616 value
1617 .downcast_ref::<Point>()
1618 .copied()
1619 .or_else(|| f64::from_value(value).map(Point::new))
1620 }
1621 }
1622
1623 pub const fn gamma(scale: i32) -> f64 {
1635 libm::pow(2.0, libm::pow(2.0, -(scale as f64)))
1636 }
1637
1638 pub const fn midpoint(value: f64, scale: i32) -> Point {
1658 let sign = value.signum();
1659 let value = value.abs();
1660
1661 if value == 0.0 {
1662 return Point::new(value);
1663 }
1664
1665 let gamma = gamma(scale);
1666
1667 let index = libm::ceil(libm::log(value, gamma));
1668
1669 let lower = libm::pow(gamma, index - 1.0);
1670 let upper = lower * gamma;
1671
1672 Point::new(sign * lower.midpoint(upper))
1673 }
1674
1675 #[cfg(feature = "alloc")]
1676 mod alloc_support {
1677 use super::*;
1678
1679 use emit_core::{
1680 props::Props,
1681 str::{Str, ToStr},
1682 well_known::{
1683 KEY_DIST_COUNT, KEY_DIST_EXP_BUCKETS, KEY_DIST_EXP_SCALE, KEY_DIST_MAX,
1684 KEY_DIST_MIN, KEY_DIST_SUM,
1685 },
1686 };
1687
1688 use crate::core::{cmp, ops::ControlFlow};
1689
1690 pub mod bucket_set {
1691 use emit_core::value::{FromValue, ToValue, Value};
1696
1697 use crate::{
1698 alloc::collections::{BTreeMap, btree_map},
1699 buf::{find, trim, trim_start},
1700 core::{
1701 fmt::{self, Write as _},
1702 str::FromStr,
1703 },
1704 metric::exp::Point,
1705 };
1706
1707 #[derive(Debug)]
1711 pub struct ParseBucketSetError {}
1712
1713 impl fmt::Display for ParseBucketSetError {
1714 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1715 write!(f, "the input was not a valid bucket set")
1716 }
1717 }
1718
1719 #[cfg(feature = "std")]
1720 impl std::error::Error for ParseBucketSetError {}
1721
1722 #[derive(Clone, PartialEq, Eq)]
1728 pub struct BucketSet {
1729 total: u64,
1730 buckets: BTreeMap<Point, u64>,
1731 }
1732
1733 impl fmt::Debug for BucketSet {
1734 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
1735 fmt::Display::fmt(self, f)
1736 }
1737 }
1738
1739 impl fmt::Display for BucketSet {
1740 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
1741 f.write_char('[')?;
1742
1743 let mut first = true;
1744 for (k, v) in &self.buckets {
1745 if !first {
1746 f.write_char(',')?;
1747 }
1748 first = false;
1749
1750 f.write_char('[')?;
1751 fmt::Display::fmt(k, f)?;
1752 f.write_char(',')?;
1753 fmt::Display::fmt(v, f)?;
1754 f.write_char(']')?;
1755 }
1756
1757 f.write_char(']')
1758 }
1759 }
1760
1761 impl BucketSet {
1762 pub fn new() -> Self {
1768 BucketSet {
1769 buckets: BTreeMap::new(),
1770 total: 0,
1771 }
1772 }
1773
1774 pub fn try_from_str(s: &str) -> Result<Self, ParseBucketSetError> {
1778 Self::try_from_slice(s.as_bytes())
1779 }
1780
1781 fn try_from_slice(mut s: &[u8]) -> Result<Self, ParseBucketSetError> {
1782 let mut set = BucketSet::new();
1783
1784 if s.len() < 2 {
1785 return Err(ParseBucketSetError {});
1787 }
1788
1789 let container_end = match (s.first(), s.last()) {
1791 (Some(&b'['), Some(&b']')) => b']',
1792 (Some(&b'('), Some(&b')')) => b')',
1793 (Some(&b'{'), Some(&b'}')) => b'}',
1794 _ => return Err(ParseBucketSetError {}),
1795 };
1796 s = &s[1..];
1797 s = trim_start(s);
1798
1799 let mut first = true;
1800 while s.len() > 1 {
1801 if !first {
1803 if s.first() != Some(&b',') {
1804 return Err(ParseBucketSetError {});
1806 }
1807 s = &s[1..];
1808 s = trim_start(s);
1809 }
1810 first = false;
1811
1812 let (key_start_skip, key_end, value_end): (
1814 usize,
1815 &[(u8, u8)],
1816 &[(u8, u8)],
1817 ) = match s.first() {
1818 Some(&b'[') => {
1820 (1, &[(b',', 1u8), (b':', 1u8), (b'=', 1u8)], &[(b']', 1u8)])
1821 }
1822 Some(&b'(') => {
1824 (1, &[(b',', 1u8), (b':', 1u8), (b'=', 1u8)], &[(b')', 1u8)])
1825 }
1826 Some(&b'{') => {
1828 (1, &[(b',', 1u8), (b':', 1u8), (b'=', 1u8)], &[(b'}', 1u8)])
1829 }
1830 _ => (
1832 0,
1833 &[(b':', 1u8), (b'=', 1u8)],
1834 &[(b',', 0u8), (container_end, 0u8)],
1835 ),
1836 };
1837 s = &s[key_start_skip..];
1838
1839 let Some((key_end, key_end_skip)) = find(s, key_end) else {
1841 return Err(ParseBucketSetError {});
1843 };
1844
1845 let key = str::from_utf8(trim(&s[..key_end]))
1846 .map_err(|_| ParseBucketSetError {})?;
1847 s = &s[key_end + key_end_skip..];
1848
1849 let Some((value_end, value_end_skip)) = find(s, value_end) else {
1851 return Err(ParseBucketSetError {});
1853 };
1854
1855 let value = str::from_utf8(trim(&s[..value_end]))
1856 .map_err(|_| ParseBucketSetError {})?;
1857 s = &s[value_end + value_end_skip..];
1858
1859 let key = key.parse().map_err(|_| ParseBucketSetError {})?;
1861 let value = value.parse().map_err(|_| ParseBucketSetError {})?;
1862
1863 set.total = set
1864 .total
1865 .checked_add(value)
1866 .ok_or_else(|| ParseBucketSetError {})?;
1867 if set.buckets.insert(key, value).is_some() {
1868 return Err(ParseBucketSetError {});
1870 }
1871
1872 s = trim_start(s);
1873 }
1874
1875 if s.len() != 1 {
1876 return Err(ParseBucketSetError {});
1878 }
1879
1880 Ok(set)
1881 }
1882
1883 pub fn observe(&mut self, value: Point) {
1895 self.observe_all(value, 1)
1896 }
1897
1898 pub fn observe_all(&mut self, value: Point, count: u64) {
1910 let entry = self.buckets.entry(value).or_default();
1911
1912 *entry = entry.checked_add(count).unwrap_or_else(|| {
1913 panic!("adding {count} observations would overflow bucket")
1914 });
1915 self.total = self.total.checked_add(count).unwrap_or_else(|| {
1916 panic!("adding {count} observations would overflow total")
1917 });
1918 }
1919
1920 pub fn remap(&mut self, mut map: impl FnMut(Point) -> Point) {
1930 let mut remapped = BTreeMap::<Point, u64>::new();
1931
1932 for (value, count) in &self.buckets {
1933 let entry = remapped.entry(map(*value)).or_default();
1934
1935 *entry = entry.checked_add(*count).unwrap_or_else(|| {
1936 panic!("adding {count} observations would overflow bucket")
1937 });
1938 }
1939
1940 self.buckets = remapped;
1941 }
1942
1943 pub fn len(&self) -> usize {
1949 self.buckets.len()
1950 }
1951
1952 pub fn total(&self) -> u64 {
1956 self.total
1957 }
1958
1959 pub fn clear(&mut self) {
1963 self.buckets.clear();
1964 self.total = 0;
1965 }
1966
1967 pub fn get(&self, value: Point) -> Option<u64> {
1971 self.buckets.get(&value).copied()
1972 }
1973
1974 pub fn first(&self) -> Option<(Point, u64)> {
1978 self.buckets.first_key_value().map(|(k, v)| (*k, *v))
1979 }
1980
1981 pub fn last(&self) -> Option<(Point, u64)> {
1985 self.buckets.last_key_value().map(|(k, v)| (*k, *v))
1986 }
1987
1988 pub fn iter(&self) -> Iter<'_> {
1992 Iter(self.buckets.iter())
1993 }
1994 }
1995
1996 impl<'a> IntoIterator for &'a BucketSet {
1997 type IntoIter = Iter<'a>;
1998 type Item = (Point, u64);
1999
2000 fn into_iter(self) -> Self::IntoIter {
2001 self.iter()
2002 }
2003 }
2004
2005 impl<'a> FromIterator<(Point, u64)> for BucketSet {
2006 fn from_iter<I: IntoIterator<Item = (Point, u64)>>(iter: I) -> Self {
2007 let mut set = BucketSet::new();
2008 set.extend(iter);
2009
2010 set
2011 }
2012 }
2013
2014 impl<'a> Extend<(Point, u64)> for BucketSet {
2015 fn extend<I: IntoIterator<Item = (Point, u64)>>(&mut self, iter: I) {
2016 for (value, count) in iter {
2017 self.observe_all(value, count);
2018 }
2019 }
2020 }
2021
2022 pub struct Iter<'a>(btree_map::Iter<'a, Point, u64>);
2028
2029 impl<'a> Iterator for Iter<'a> {
2030 type Item = (Point, u64);
2031
2032 fn next(&mut self) -> Option<Self::Item> {
2033 self.0.next().map(|(k, v)| (*k, *v))
2034 }
2035 }
2036
2037 #[cfg(feature = "sval")]
2038 impl sval::Value for BucketSet {
2039 fn stream<'sval, S: sval::Stream<'sval> + ?Sized>(
2040 &'sval self,
2041 stream: &mut S,
2042 ) -> sval::Result {
2043 stream.seq_begin(Some(self.buckets.len()))?;
2044
2045 for bucket in &self.buckets {
2046 stream.value_computed(&bucket)?;
2047 }
2048
2049 stream.seq_end()
2050 }
2051 }
2052
2053 #[cfg(feature = "serde")]
2054 impl serde::Serialize for BucketSet {
2055 fn serialize<S: serde::Serializer>(
2056 &self,
2057 serializer: S,
2058 ) -> Result<S::Ok, S::Error> {
2059 use serde::ser::SerializeSeq as _;
2060
2061 let mut seq = serializer.serialize_seq(Some(self.buckets.len()))?;
2062
2063 for bucket in &self.buckets {
2064 seq.serialize_element(&bucket)?;
2065 }
2066
2067 seq.end()
2068 }
2069 }
2070
2071 impl FromStr for BucketSet {
2072 type Err = ParseBucketSetError;
2073
2074 fn from_str(s: &str) -> Result<Self, Self::Err> {
2075 Self::try_from_str(s)
2076 }
2077 }
2078
2079 impl ToValue for BucketSet {
2080 fn to_value(&self) -> Value<'_> {
2081 #[cfg(feature = "sval")]
2082 {
2083 Value::capture_sval(self)
2084 }
2085 #[cfg(all(feature = "serde", not(feature = "sval")))]
2086 {
2087 Value::capture_serde(self)
2088 }
2089 #[cfg(all(not(feature = "serde"), not(feature = "sval")))]
2090 {
2091 Value::capture_display(self)
2092 }
2093 }
2094 }
2095
2096 impl<'a> FromValue<'a> for BucketSet {
2097 fn from_value(v: Value<'a>) -> Option<Self> {
2098 if let Some(buckets) = v.downcast_ref::<Self>() {
2099 return Some(buckets.clone());
2100 }
2101
2102 #[cfg(feature = "sval")]
2103 {
2104 if let Some(buckets) = from_sval(v.by_ref()) {
2105 return Some(buckets);
2106 }
2107 }
2108
2109 #[cfg(all(not(feature = "sval"), feature = "serde"))]
2110 {
2111 if let Some(buckets) = from_serde(v.by_ref()) {
2112 return Some(buckets);
2113 }
2114 }
2115
2116 v.parse()
2117 }
2118 }
2119
2120 #[cfg(any(feature = "sval", feature = "serde"))]
2121 #[derive(Default)]
2122 struct Extract {
2123 depth: usize,
2124 buckets: BTreeMap<Point, u64>,
2125 count: u64,
2126 next_midpoint: Option<f64>,
2127 next_count: Option<u64>,
2128 }
2129
2130 #[derive(Debug)]
2131 #[cfg(any(feature = "sval", feature = "serde"))]
2132 struct Incompatible;
2133
2134 #[cfg(any(feature = "sval", feature = "serde"))]
2135 impl Extract {
2136 fn push(
2137 &mut self,
2138 midpoint: impl FnOnce() -> Option<f64>,
2139 count: impl FnOnce() -> Option<u64>,
2140 ) -> Result<(), Incompatible> {
2141 if self.depth == 2 {
2142 if self.next_midpoint.is_none() {
2143 self.next_midpoint = midpoint();
2144
2145 return Ok(());
2146 }
2147
2148 if self.next_count.is_none() {
2149 self.next_count = count();
2150
2151 return Ok(());
2152 }
2153 }
2154
2155 Err(Incompatible)
2156 }
2157
2158 fn apply(&mut self) -> Result<(), Incompatible> {
2159 if self.depth == 2 {
2160 let midpoint = self.next_midpoint.take().ok_or(Incompatible)?;
2161 let count = self.next_count.take().ok_or(Incompatible)?;
2162
2163 let entry = self.buckets.entry(Point::new(midpoint)).or_default();
2164 *entry = entry.checked_add(count).ok_or_else(|| Incompatible)?;
2165
2166 self.count = self.count.checked_add(count).ok_or_else(|| Incompatible)?;
2167
2168 Ok(())
2169 } else {
2170 Ok(())
2171 }
2172 }
2173
2174 fn down(&mut self) -> Result<(), Incompatible> {
2175 self.depth += 1;
2176
2177 if self.depth > 2 {
2178 Err(Incompatible)
2179 } else {
2180 Ok(())
2181 }
2182 }
2183
2184 fn up(&mut self) -> Result<(), Incompatible> {
2185 self.apply()?;
2186 self.depth -= 1;
2187
2188 Ok(())
2189 }
2190
2191 fn end(self) -> BucketSet {
2192 BucketSet {
2193 buckets: self.buckets,
2194 total: self.count,
2195 }
2196 }
2197 }
2198
2199 #[cfg(feature = "sval")]
2200 fn from_sval(value: Value) -> Option<BucketSet> {
2201 #[allow(non_local_definitions)]
2202 impl From<Incompatible> for sval::Error {
2203 fn from(_: Incompatible) -> sval::Error {
2204 sval::Error::new()
2205 }
2206 }
2207
2208 #[allow(non_local_definitions)]
2209 impl<'sval> sval::Stream<'sval> for Extract {
2210 fn null(&mut self) -> sval::Result {
2211 sval::error()
2212 }
2213
2214 fn bool(&mut self, _: bool) -> sval::Result {
2215 sval::error()
2216 }
2217
2218 fn text_begin(&mut self, _: Option<usize>) -> sval::Result {
2219 sval::error()
2220 }
2221
2222 fn text_fragment_computed(&mut self, _: &str) -> sval::Result {
2223 sval::error()
2224 }
2225
2226 fn text_end(&mut self) -> sval::Result {
2227 sval::error()
2228 }
2229
2230 fn i64(&mut self, value: i64) -> sval::Result {
2231 Ok(self.push(|| Some(value as f64), || value.try_into().ok())?)
2232 }
2233
2234 fn u64(&mut self, value: u64) -> sval::Result {
2235 Ok(self.push(|| Some(value as f64), || Some(value))?)
2236 }
2237
2238 fn i128(&mut self, value: i128) -> sval::Result {
2239 Ok(self.push(|| Some(value as f64), || value.try_into().ok())?)
2240 }
2241
2242 fn u128(&mut self, value: u128) -> sval::Result {
2243 Ok(self.push(|| Some(value as f64), || value.try_into().ok())?)
2244 }
2245
2246 fn f64(&mut self, value: f64) -> sval::Result {
2247 Ok(self.push(|| Some(value), || Some(value as u64))?)
2248 }
2249
2250 fn seq_begin(&mut self, _: Option<usize>) -> sval::Result {
2251 Ok(self.down()?)
2252 }
2253
2254 fn seq_value_begin(&mut self) -> sval::Result {
2255 Ok(())
2256 }
2257
2258 fn seq_value_end(&mut self) -> sval::Result {
2259 Ok(())
2260 }
2261
2262 fn seq_end(&mut self) -> sval::Result {
2263 Ok(self.up()?)
2264 }
2265 }
2266
2267 let mut extract = Extract::default();
2268 sval::stream(&mut extract, &value).ok()?;
2269
2270 Some(extract.end())
2271 }
2272
2273 #[cfg(all(not(feature = "sval"), feature = "serde"))]
2274 fn from_serde(value: Value) -> Option<BucketSet> {
2275 use serde::Serialize as _;
2276
2277 #[allow(non_local_definitions)]
2278 impl fmt::Display for Incompatible {
2279 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
2280 f.write_str("incompatible")
2281 }
2282 }
2283
2284 #[allow(non_local_definitions)]
2285 impl serde::ser::StdError for Incompatible {}
2286
2287 #[allow(non_local_definitions)]
2288 impl serde::ser::Error for Incompatible {
2289 fn custom<T>(_: T) -> Self
2290 where
2291 T: fmt::Display,
2292 {
2293 Incompatible
2294 }
2295 }
2296
2297 #[allow(non_local_definitions)]
2298 impl<'a> serde::Serializer for &'a mut Extract {
2299 type Ok = ();
2300 type Error = Incompatible;
2301 type SerializeSeq = Self;
2302 type SerializeTuple = Self;
2303 type SerializeTupleStruct = Self;
2304 type SerializeTupleVariant = Self;
2305 type SerializeMap = Self;
2306 type SerializeStruct = Self;
2307 type SerializeStructVariant = Self;
2308
2309 fn serialize_bool(self, _: bool) -> Result<Self::Ok, Self::Error> {
2310 Err(Incompatible)
2311 }
2312
2313 fn serialize_i8(self, value: i8) -> Result<Self::Ok, Self::Error> {
2314 self.push(|| Some(value as f64), || value.try_into().ok())
2315 }
2316
2317 fn serialize_i16(self, value: i16) -> Result<Self::Ok, Self::Error> {
2318 self.push(|| Some(value as f64), || value.try_into().ok())
2319 }
2320
2321 fn serialize_i32(self, value: i32) -> Result<Self::Ok, Self::Error> {
2322 self.push(|| Some(value as f64), || value.try_into().ok())
2323 }
2324
2325 fn serialize_i64(self, value: i64) -> Result<Self::Ok, Self::Error> {
2326 self.push(|| Some(value as f64), || value.try_into().ok())
2327 }
2328
2329 fn serialize_u8(self, value: u8) -> Result<Self::Ok, Self::Error> {
2330 self.push(|| Some(value as f64), || value.try_into().ok())
2331 }
2332
2333 fn serialize_u16(self, value: u16) -> Result<Self::Ok, Self::Error> {
2334 self.push(|| Some(value as f64), || value.try_into().ok())
2335 }
2336
2337 fn serialize_u32(self, value: u32) -> Result<Self::Ok, Self::Error> {
2338 self.push(|| Some(value as f64), || value.try_into().ok())
2339 }
2340
2341 fn serialize_u64(self, value: u64) -> Result<Self::Ok, Self::Error> {
2342 self.push(|| Some(value as f64), || value.try_into().ok())
2343 }
2344
2345 fn serialize_u128(self, value: u128) -> Result<Self::Ok, Self::Error> {
2346 self.push(|| Some(value as f64), || value.try_into().ok())
2347 }
2348
2349 fn serialize_i128(self, value: i128) -> Result<Self::Ok, Self::Error> {
2350 self.push(|| Some(value as f64), || value.try_into().ok())
2351 }
2352
2353 fn serialize_f32(self, value: f32) -> Result<Self::Ok, Self::Error> {
2354 self.push(|| Some(value as f64), || Some(value as u64))
2355 }
2356
2357 fn serialize_f64(self, value: f64) -> Result<Self::Ok, Self::Error> {
2358 self.push(|| Some(value as f64), || Some(value as u64))
2359 }
2360
2361 fn serialize_char(self, _: char) -> Result<Self::Ok, Self::Error> {
2362 Err(Incompatible)
2363 }
2364
2365 fn serialize_str(self, _: &str) -> Result<Self::Ok, Self::Error> {
2366 Err(Incompatible)
2367 }
2368
2369 fn serialize_bytes(self, _: &[u8]) -> Result<Self::Ok, Self::Error> {
2370 Err(Incompatible)
2371 }
2372
2373 fn serialize_none(self) -> Result<Self::Ok, Self::Error> {
2374 Err(Incompatible)
2375 }
2376
2377 fn serialize_some<T>(self, value: &T) -> Result<Self::Ok, Self::Error>
2378 where
2379 T: ?Sized + serde::Serialize,
2380 {
2381 value.serialize(self)
2382 }
2383
2384 fn serialize_unit(self) -> Result<Self::Ok, Self::Error> {
2385 Err(Incompatible)
2386 }
2387
2388 fn serialize_unit_struct(
2389 self,
2390 name: &'static str,
2391 ) -> Result<Self::Ok, Self::Error> {
2392 name.serialize(self)
2393 }
2394
2395 fn serialize_unit_variant(
2396 self,
2397 _: &'static str,
2398 _: u32,
2399 variant: &'static str,
2400 ) -> Result<Self::Ok, Self::Error> {
2401 variant.serialize(self)
2402 }
2403
2404 fn serialize_newtype_struct<T>(
2405 self,
2406 _: &'static str,
2407 value: &T,
2408 ) -> Result<Self::Ok, Self::Error>
2409 where
2410 T: ?Sized + serde::Serialize,
2411 {
2412 value.serialize(self)
2413 }
2414
2415 fn serialize_newtype_variant<T>(
2416 self,
2417 _: &'static str,
2418 _: u32,
2419 _: &'static str,
2420 value: &T,
2421 ) -> Result<Self::Ok, Self::Error>
2422 where
2423 T: ?Sized + serde::Serialize,
2424 {
2425 value.serialize(self)
2426 }
2427
2428 fn serialize_seq(
2429 self,
2430 _: Option<usize>,
2431 ) -> Result<Self::SerializeSeq, Self::Error> {
2432 self.down()?;
2433
2434 Ok(self)
2435 }
2436
2437 fn serialize_tuple(
2438 self,
2439 _: usize,
2440 ) -> Result<Self::SerializeTuple, Self::Error> {
2441 self.down()?;
2442
2443 Ok(self)
2444 }
2445
2446 fn serialize_tuple_struct(
2447 self,
2448 _: &'static str,
2449 _: usize,
2450 ) -> Result<Self::SerializeTupleStruct, Self::Error> {
2451 self.down()?;
2452
2453 Ok(self)
2454 }
2455
2456 fn serialize_tuple_variant(
2457 self,
2458 _: &'static str,
2459 _: u32,
2460 _: &'static str,
2461 _: usize,
2462 ) -> Result<Self::SerializeTupleVariant, Self::Error> {
2463 self.down()?;
2464
2465 Ok(self)
2466 }
2467
2468 fn serialize_map(
2469 self,
2470 _: Option<usize>,
2471 ) -> Result<Self::SerializeMap, Self::Error> {
2472 self.down()?;
2473
2474 Ok(self)
2475 }
2476
2477 fn serialize_struct(
2478 self,
2479 _: &'static str,
2480 _: usize,
2481 ) -> Result<Self::SerializeStruct, Self::Error> {
2482 self.down()?;
2483
2484 Ok(self)
2485 }
2486
2487 fn serialize_struct_variant(
2488 self,
2489 _: &'static str,
2490 _: u32,
2491 _: &'static str,
2492 _: usize,
2493 ) -> Result<Self::SerializeStructVariant, Self::Error> {
2494 self.down()?;
2495
2496 Ok(self)
2497 }
2498 }
2499
2500 #[allow(non_local_definitions)]
2501 impl<'a> serde::ser::SerializeSeq for &'a mut Extract {
2502 type Ok = ();
2503 type Error = Incompatible;
2504
2505 fn serialize_element<T>(&mut self, value: &T) -> Result<(), Self::Error>
2506 where
2507 T: ?Sized + serde::Serialize,
2508 {
2509 value.serialize(&mut **self)
2510 }
2511
2512 fn end(self) -> Result<Self::Ok, Self::Error> {
2513 self.up()
2514 }
2515 }
2516
2517 #[allow(non_local_definitions)]
2518 impl<'a> serde::ser::SerializeTuple for &'a mut Extract {
2519 type Ok = ();
2520 type Error = Incompatible;
2521
2522 fn serialize_element<T>(&mut self, value: &T) -> Result<(), Self::Error>
2523 where
2524 T: ?Sized + serde::Serialize,
2525 {
2526 value.serialize(&mut **self)
2527 }
2528
2529 fn end(self) -> Result<Self::Ok, Self::Error> {
2530 self.up()
2531 }
2532 }
2533
2534 #[allow(non_local_definitions)]
2535 impl<'a> serde::ser::SerializeTupleStruct for &'a mut Extract {
2536 type Ok = ();
2537 type Error = Incompatible;
2538
2539 fn serialize_field<T>(&mut self, value: &T) -> Result<(), Self::Error>
2540 where
2541 T: ?Sized + serde::Serialize,
2542 {
2543 value.serialize(&mut **self)
2544 }
2545
2546 fn end(self) -> Result<Self::Ok, Self::Error> {
2547 self.up()
2548 }
2549 }
2550
2551 #[allow(non_local_definitions)]
2552 impl<'a> serde::ser::SerializeTupleVariant for &'a mut Extract {
2553 type Ok = ();
2554 type Error = Incompatible;
2555
2556 fn serialize_field<T>(&mut self, value: &T) -> Result<(), Self::Error>
2557 where
2558 T: ?Sized + serde::Serialize,
2559 {
2560 value.serialize(&mut **self)
2561 }
2562
2563 fn end(self) -> Result<Self::Ok, Self::Error> {
2564 self.up()
2565 }
2566 }
2567
2568 #[allow(non_local_definitions)]
2569 impl<'a> serde::ser::SerializeMap for &'a mut Extract {
2570 type Ok = ();
2571 type Error = Incompatible;
2572
2573 fn serialize_key<T>(&mut self, key: &T) -> Result<(), Self::Error>
2574 where
2575 T: ?Sized + serde::Serialize,
2576 {
2577 self.down()?;
2578 key.serialize(&mut **self)
2579 }
2580
2581 fn serialize_value<T>(&mut self, value: &T) -> Result<(), Self::Error>
2582 where
2583 T: ?Sized + serde::Serialize,
2584 {
2585 value.serialize(&mut **self)?;
2586 self.up()
2587 }
2588
2589 fn end(self) -> Result<Self::Ok, Self::Error> {
2590 self.up()
2591 }
2592 }
2593
2594 #[allow(non_local_definitions)]
2595 impl<'a> serde::ser::SerializeStruct for &'a mut Extract {
2596 type Ok = ();
2597 type Error = Incompatible;
2598
2599 fn serialize_field<T>(
2600 &mut self,
2601 key: &'static str,
2602 value: &T,
2603 ) -> Result<(), Self::Error>
2604 where
2605 T: ?Sized + serde::Serialize,
2606 {
2607 self.down()?;
2608 key.serialize(&mut **self)?;
2609 value.serialize(&mut **self)?;
2610 self.up()
2611 }
2612
2613 fn end(self) -> Result<Self::Ok, Self::Error> {
2614 self.up()
2615 }
2616 }
2617
2618 #[allow(non_local_definitions)]
2619 impl<'a> serde::ser::SerializeStructVariant for &'a mut Extract {
2620 type Ok = ();
2621 type Error = Incompatible;
2622
2623 fn serialize_field<T>(
2624 &mut self,
2625 key: &'static str,
2626 value: &T,
2627 ) -> Result<(), Self::Error>
2628 where
2629 T: ?Sized + serde::Serialize,
2630 {
2631 self.down()?;
2632 key.serialize(&mut **self)?;
2633 value.serialize(&mut **self)?;
2634 self.up()
2635 }
2636
2637 fn end(self) -> Result<Self::Ok, Self::Error> {
2638 self.up()
2639 }
2640 }
2641
2642 let mut extract = Extract::default();
2643 value.serialize(&mut extract).ok()?;
2644
2645 Some(extract.end())
2646 }
2647
2648 #[cfg(test)]
2649 mod tests {
2650 use super::*;
2651
2652 use std::collections::{BTreeMap, BTreeSet};
2653
2654 #[test]
2655 fn bucket_set_observe() {
2656 let mut set = BucketSet::new();
2657
2658 assert_eq!(0, set.len());
2659
2660 set.observe(Point::new(0.0));
2661 set.observe_all(Point::new(0.0), 2);
2662 set.observe_all(Point::new(1.0), 2);
2663
2664 assert_eq!(2, set.len());
2665 assert_eq!(3, set.get(Point::new(0.0)).unwrap());
2666 assert_eq!(2, set.get(Point::new(1.0)).unwrap());
2667
2668 assert_eq!((Point::new(0.0), 3), set.first().unwrap());
2669 assert_eq!((Point::new(1.0), 2), set.last().unwrap());
2670 }
2671
2672 #[test]
2673 fn bucket_set_remap() {
2674 let mut set = BucketSet::new();
2675
2676 set.observe_all(Point::new(0.0), 3);
2677 set.observe_all(Point::new(1.0), 2);
2678
2679 assert_eq!(2, set.len());
2680
2681 set.remap(|_| Point::new(2.0));
2682
2683 assert_eq!(1, set.len());
2684 assert_eq!(5, set.get(Point::new(2.0)).unwrap());
2685 }
2686
2687 #[test]
2688 fn bucket_set_roundtrip() {
2689 for case in [
2690 BucketSet::new(),
2691 {
2692 let mut set = BucketSet::new();
2693 set.observe(Point::new(0.0));
2694 set
2695 },
2696 {
2697 let mut set = BucketSet::new();
2698 set.observe(Point::new(0.0));
2699 set.observe(Point::new(1.0));
2700 set
2701 },
2702 ] {
2703 let fmt = case.to_string();
2704 assert_eq!(Some(case), BucketSet::try_from_str(&fmt).ok(), "{fmt}");
2705 }
2706 }
2707
2708 #[test]
2709 fn bucket_set_from_iter() {
2710 let mut set = BucketSet::from_iter([
2711 (Point::new(0.0), 3),
2712 (Point::new(0.0), 2),
2713 (Point::new(1.0), 2),
2714 ]);
2715
2716 assert_eq!(5, set.get(Point::new(0.0)).unwrap());
2717 assert_eq!(2, set.get(Point::new(1.0)).unwrap());
2718
2719 set.extend([(Point::new(1.0), 3), (Point::new(2.0), 2)]);
2720
2721 assert_eq!(5, set.get(Point::new(1.0)).unwrap());
2722 assert_eq!(2, set.get(Point::new(2.0)).unwrap());
2723 }
2724
2725 #[test]
2726 fn bucket_set_parse() {
2727 for (case, expected) in [
2728 (format!("{:?}", ([[1, 1], [2, 2]])), {
2729 let mut set = BucketSet::new();
2730 set.observe_all(Point::new(1.0), 1);
2731 set.observe_all(Point::new(2.0), 2);
2732 set
2733 }),
2734 (format!("{:?}", ([(1.0, 1), (2.0, 2)])), {
2735 let mut set = BucketSet::new();
2736 set.observe_all(Point::new(1.0), 1);
2737 set.observe_all(Point::new(2.0), 2);
2738 set
2739 }),
2740 (
2741 format!("{:?}", {
2742 let mut set = BTreeSet::new();
2743 set.insert((1, 1));
2744 set.insert((2, 2));
2745 set
2746 }),
2747 {
2748 let mut set = BucketSet::new();
2749 set.observe_all(Point::new(1.0), 1);
2750 set.observe_all(Point::new(2.0), 2);
2751 set
2752 },
2753 ),
2754 (
2755 format!("{:?}", {
2756 let mut set = BTreeMap::new();
2757 set.insert(1, 1);
2758 set.insert(2, 2);
2759 set
2760 }),
2761 {
2762 let mut set = BucketSet::new();
2763 set.observe_all(Point::new(1.0), 1);
2764 set.observe_all(Point::new(2.0), 2);
2765 set
2766 },
2767 ),
2768 ("[ [ 1 , 1 ] , [ 2 , 2 ] ]".to_string(), {
2769 let mut set = BucketSet::new();
2770 set.observe_all(Point::new(1.0), 1);
2771 set.observe_all(Point::new(2.0), 2);
2772 set
2773 }),
2774 ("[ 1 : 1 , 2 : 2 ]".to_string(), {
2775 let mut set = BucketSet::new();
2776 set.observe_all(Point::new(1.0), 1);
2777 set.observe_all(Point::new(2.0), 2);
2778 set
2779 }),
2780 ] {
2781 assert_eq!(
2782 Some(expected),
2783 BucketSet::try_from_str(&case).ok(),
2784 "{case}"
2785 );
2786 }
2787 }
2788
2789 #[test]
2790 fn bucket_set_parse_exotic() {
2791 for case in ["[[inf,1]]", "[[nan,1]]", "[(1, 1), [2, 1], {3: 1}]"] {
2792 assert!(BucketSet::try_from_str(case).is_ok());
2793 }
2794 }
2795
2796 #[test]
2797 fn bucket_set_to_from_value() {
2798 for case in [{
2799 let mut set = BucketSet::new();
2800 set.observe_all(Point::new(1.0), 1);
2801 set.observe_all(Point::new(2.0), 2);
2802 set
2803 }] {
2804 assert_eq!(case, BucketSet::from_value(case.to_value()).unwrap());
2805 }
2806 }
2807
2808 #[test]
2809 fn bucket_set_from_value_string() {
2810 for (case, expected) in [("[[1.0,1],[2.0,2]]", {
2811 let mut set = BucketSet::new();
2812 set.observe_all(Point::new(1.0), 1);
2813 set.observe_all(Point::new(2.0), 2);
2814 set
2815 })] {
2816 assert_eq!(expected, Value::from(case).cast().unwrap());
2817 }
2818 }
2819
2820 #[test]
2821 fn bucket_set_from_value_structured() {
2822 #[cfg(feature = "sval")]
2823 trait CaseSval: sval::Value {}
2824 #[cfg(feature = "sval")]
2825 impl<T: sval::Value> CaseSval for T {}
2826 #[cfg(not(feature = "sval"))]
2827 trait CaseSval {}
2828 #[cfg(not(feature = "sval"))]
2829 impl<T> CaseSval for T {}
2830
2831 #[cfg(feature = "serde")]
2832 trait CaseSerde: serde::Serialize {}
2833 #[cfg(feature = "serde")]
2834 impl<T: serde::Serialize> CaseSerde for T {}
2835 #[cfg(not(feature = "serde"))]
2836 trait CaseSerde {}
2837 #[cfg(not(feature = "serde"))]
2838 impl<T> CaseSerde for T {}
2839
2840 trait Case: CaseSval + CaseSerde + fmt::Debug {}
2841 impl<T: fmt::Debug + CaseSval + CaseSerde> Case for T {}
2842
2843 fn case(case: &impl Case, expected: &BucketSet) {
2844 assert_eq!(expected, &Value::from_debug(case).cast().unwrap());
2845
2846 #[cfg(feature = "sval")]
2847 {
2848 assert_eq!(expected, &Value::from_sval(case).cast().unwrap());
2849 }
2850
2851 #[cfg(feature = "serde")]
2852 {
2853 assert_eq!(expected, &Value::from_serde(case).cast().unwrap());
2854 }
2855 }
2856
2857 let mut set = BucketSet::new();
2858 set.observe_all(Point::new(1.0), 1);
2859 set.observe_all(Point::new(2.0), 2);
2860
2861 case(&[[1, 1], [2, 2]], &set);
2862 case(&[(1.0, 1), (2.0, 2)], &set);
2863 case(
2864 &{
2865 let mut set = BTreeSet::new();
2866 set.insert((1, 1));
2867 set.insert((2, 2));
2868 set
2869 },
2870 &set,
2871 );
2872 case(
2873 &{
2874 let mut set = BTreeMap::new();
2875 set.insert(1, 1);
2876 set.insert(2, 2);
2877 set
2878 },
2879 &set,
2880 );
2881 }
2882
2883 #[test]
2884 fn err_bucket_set_invalid() {
2885 for case in [
2886 "",
2887 "<>",
2888 "[1, 1]",
2889 "1, 1",
2890 "[[1, 1]], [[2, 1]]",
2891 "[[]]",
2892 "[}",
2893 "[[1,1}]",
2894 "[[1 1]]",
2895 "[[1, 1] [2, 1]]",
2896 "[,]",
2897 "[[,]]",
2898 "[[1,]]",
2899 "[[,1]]",
2900 "[[:]]",
2901 "[[1:]]",
2902 "[[:1]]",
2903 "[[1, -1]]",
2904 "[[1, 1.0]]",
2905 "[[1, 0xff]]",
2906 "[[1, ff]]",
2907 "{1.2789: 11111111111111111111, 2789: 11111111111111111111, 2 \0: \0: 2}",
2908 ] {
2909 assert!(BucketSet::try_from_str(case).is_err(), "{case}");
2910 }
2911 }
2912 }
2913 }
2914
2915 pub use self::bucket_set::BucketSet;
2916
2917 pub struct Distribution {
2933 max_buckets: usize,
2934 max_scale: i32,
2935 scale: i32,
2936 sum: Option<f64>,
2937 min: Option<f64>,
2938 max: Option<f64>,
2939 buckets: BucketSet,
2940 }
2941
2942 impl Default for Distribution {
2943 fn default() -> Self {
2944 Self::new(Self::DEFAULT_MAX_SCALE, Self::DEFAULT_MAX_BUCKETS)
2945 }
2946 }
2947
2948 impl Distribution {
2949 pub const DEFAULT_MAX_SCALE: i32 = 20;
2953
2954 pub const DEFAULT_MAX_BUCKETS: usize = 160;
2958
2959 pub fn new(max_scale: i32, max_buckets: usize) -> Self {
2965 Distribution {
2966 max_buckets,
2967 max_scale,
2968 scale: max_scale,
2969 min: None,
2970 max: None,
2971 sum: None,
2972 buckets: BucketSet::new(),
2973 }
2974 }
2975
2976 pub fn observe(&mut self, raw_value: f64) {
2986 self.observe_all(raw_value, 1)
2987 }
2988
2989 pub fn observe_all(&mut self, raw_value: f64, count: u64) {
2999 self.buckets
3000 .observe_all(midpoint(raw_value, self.scale), count);
3001
3002 self.min = self
3004 .min
3005 .map(|min| cmp::min_by(min, raw_value, |a, b| a.total_cmp(b)))
3006 .or(Some(raw_value));
3007 self.max = self
3008 .max
3009 .map(|max| cmp::max_by(max, raw_value, |a, b| a.total_cmp(b)))
3010 .or(Some(raw_value));
3011 self.sum = self.sum.map(|sum| sum + raw_value).or(Some(raw_value));
3012
3013 if self.buckets.len() > self.max_buckets {
3016 self.scale -= 1;
3017 self.buckets
3018 .remap(|value| midpoint(value.get(), self.scale));
3019 }
3020 }
3021
3022 pub fn reset(&mut self) {
3028 let Distribution {
3029 max_scale,
3030 max_buckets: _,
3031 scale,
3032 min,
3033 max,
3034 sum,
3035 buckets,
3036 } = self;
3037
3038 buckets.clear();
3039 *min = None;
3040 *max = None;
3041 *sum = None;
3042 *scale = *max_scale;
3043 }
3044
3045 pub fn count(&self) -> u64 {
3051 self.buckets.total()
3052 }
3053
3054 pub fn min(&self) -> Option<f64> {
3060 self.min
3061 }
3062
3063 pub fn max(&self) -> Option<f64> {
3069 self.max
3070 }
3071
3072 pub fn sum(&self) -> Option<f64> {
3078 self.sum
3079 }
3080
3081 pub fn scale(&self) -> i32 {
3085 self.scale
3086 }
3087
3088 pub fn buckets(&self) -> &BucketSet {
3092 &self.buckets
3093 }
3094
3095 pub fn max_buckets(&self) -> usize {
3099 self.max_buckets
3100 }
3101
3102 pub fn max_scale(&self) -> i32 {
3106 self.max_scale
3107 }
3108 }
3109
3110 impl Props for Distribution {
3111 fn for_each<'kv, F: FnMut(Str<'kv>, Value<'kv>) -> ControlFlow<()>>(
3112 &'kv self,
3113 mut for_each: F,
3114 ) -> ControlFlow<()> {
3115 for_each(KEY_DIST_EXP_SCALE.to_str(), self.scale().into())?;
3116 for_each(KEY_DIST_EXP_BUCKETS.to_str(), self.buckets().to_value())?;
3117
3118 for_each(KEY_DIST_COUNT.to_str(), self.count().into())?;
3119
3120 if let Some(sum) = self.sum() {
3121 for_each(KEY_DIST_SUM.to_str(), sum.into())?;
3122 }
3123 if let Some(min) = self.min() {
3124 for_each(KEY_DIST_MIN.to_str(), min.into())?;
3125 }
3126 if let Some(max) = self.max() {
3127 for_each(KEY_DIST_MAX.to_str(), max.into())?;
3128 }
3129
3130 ControlFlow::Continue(())
3131 }
3132 }
3133
3134 #[cfg(test)]
3135 mod tests {
3136 use super::*;
3137
3138 #[test]
3139 fn distribution_observe() {
3140 let mut distribution = Distribution::new(10, 10);
3141
3142 assert_eq!(distribution.max_scale(), distribution.scale());
3143 assert_eq!(0, distribution.buckets().len());
3144 assert_eq!(None, distribution.min());
3145 assert_eq!(None, distribution.max());
3146 assert_eq!(None, distribution.sum());
3147 assert_eq!(0, distribution.count());
3148
3149 distribution.observe(1.0);
3150 distribution.observe(1.0);
3151
3152 assert_eq!(
3153 2,
3154 distribution
3155 .buckets()
3156 .get(midpoint(1.0, distribution.max_scale()))
3157 .unwrap()
3158 );
3159 assert_eq!(1, distribution.buckets().len());
3160 assert_eq!(Some(1.0), distribution.min());
3161 assert_eq!(Some(1.0), distribution.max());
3162 assert_eq!(Some(2.0), distribution.sum());
3163 assert_eq!(2, distribution.count());
3164
3165 distribution.reset();
3166
3167 assert_eq!(distribution.max_scale(), distribution.scale());
3168 assert_eq!(0, distribution.buckets().len());
3169 assert_eq!(None, distribution.min());
3170 assert_eq!(None, distribution.max());
3171 assert_eq!(None, distribution.sum());
3172 assert_eq!(0, distribution.count());
3173 }
3174
3175 #[test]
3176 fn distribution_rescale() {
3177 let mut distribution = Distribution::new(10, 10);
3178
3179 for i in 0..100 {
3180 distribution.observe(i as f64);
3181 }
3182
3183 assert!(distribution.buckets().len() <= distribution.max_buckets());
3184 assert!(distribution.scale() < distribution.max_scale());
3185
3186 distribution.reset();
3187
3188 assert_eq!(distribution.max_scale(), distribution.scale());
3189 }
3190 }
3191 }
3192
3193 #[cfg(feature = "alloc")]
3194 pub use self::alloc_support::*;
3195
3196 #[cfg(test)]
3197 mod tests {
3198 use core::f64::consts::PI;
3199
3200 use super::*;
3201
3202 #[test]
3203 fn point_cmp() {
3204 let mut values = vec![
3205 Point::new(1.0),
3206 Point::new(f64::NAN),
3207 Point::new(0.0),
3208 Point::new(f64::NEG_INFINITY),
3209 Point::new(-1.0),
3210 Point::new(-0.0),
3211 Point::new(f64::INFINITY),
3212 ];
3213
3214 values.sort();
3215
3216 assert_eq!(
3217 vec![
3218 Point::new(f64::NEG_INFINITY),
3219 Point::new(-1.0),
3220 Point::new(-0.0),
3221 Point::new(0.0),
3222 Point::new(1.0),
3223 Point::new(f64::INFINITY),
3224 Point::new(f64::NAN)
3225 ],
3226 &*values
3227 );
3228 }
3229
3230 #[test]
3231 fn point_roundtrip() {
3232 let point = Point::new(1.0);
3233
3234 assert_eq!(point, Point::try_from_str(&point.to_string()).unwrap());
3235 }
3236
3237 #[test]
3238 fn point_is_indexable() {
3239 for (case, indexable) in [
3240 (Point::new(0.0), false),
3241 (Point::new(-0.0), false),
3242 (Point::new(f64::INFINITY), false),
3243 (Point::new(f64::NEG_INFINITY), false),
3244 (Point::new(f64::NAN), false),
3245 (Point::new(f64::EPSILON), true),
3246 (Point::new(-f64::EPSILON), true),
3247 (Point::new(f64::MIN), true),
3248 (Point::new(f64::MAX), true),
3249 ] {
3250 assert_eq!(indexable, case.is_indexable());
3251 }
3252 }
3253
3254 #[test]
3255 fn point_is_bucket() {
3256 for (case, zero, neg, pos) in [
3257 (Point::new(0.0), true, false, false),
3258 (Point::new(-0.0), true, false, false),
3259 (Point::new(f64::INFINITY), false, false, false),
3260 (Point::new(f64::NEG_INFINITY), false, false, false),
3261 (Point::new(f64::NAN), false, false, false),
3262 (Point::new(f64::EPSILON), false, false, true),
3263 (Point::new(-f64::EPSILON), false, true, false),
3264 (Point::new(f64::MIN), false, true, false),
3265 (Point::new(f64::MAX), false, false, true),
3266 ] {
3267 assert_eq!(zero, case.is_zero_bucket());
3268 assert_eq!(neg, case.is_negative_bucket());
3269 assert_eq!(pos, case.is_positive_bucket());
3270 }
3271 }
3272
3273 #[cfg(feature = "sval")]
3274 #[test]
3275 fn point_stream() {
3276 sval_test::assert_tokens(&Point::new(3.1), &[sval_test::Token::F64(3.1)]);
3277 }
3278
3279 #[cfg(feature = "serde")]
3280 #[test]
3281 fn point_serialize() {
3282 serde_test::assert_ser_tokens(&Point::new(3.1), &[serde_test::Token::F64(3.1)]);
3283 }
3284
3285 #[test]
3286 fn point_to_from_value() {
3287 let point = Point::new(3.1);
3288
3289 assert_eq!(point, Point::from_value(point.to_value()).unwrap());
3290 }
3291
3292 #[test]
3293 fn compute_midpoints() {
3294 let cases = [
3295 0.0f64,
3296 PI,
3297 PI * 100.0f64,
3298 PI * 1000.0f64,
3299 -0.0f64,
3300 -PI,
3301 -(PI * 100.0f64),
3302 -(PI * 1000.0f64),
3303 f64::INFINITY,
3304 f64::NEG_INFINITY,
3305 f64::NAN,
3306 ];
3307 for (scale, expected) in [
3308 (
3309 0i32,
3310 [
3311 0.0f64,
3312 3.0f64,
3313 384.0f64,
3314 3072.0f64,
3315 0.0f64,
3316 -3.0f64,
3317 -384.0f64,
3318 -3072.0f64,
3319 f64::INFINITY,
3320 f64::NEG_INFINITY,
3321 f64::NAN,
3322 ],
3323 ),
3324 (
3325 2i32,
3326 [
3327 0.0f64,
3328 3.0960063928805233f64,
3329 333.2378467041041f64,
3330 3170.3105463096517f64,
3331 0.0f64,
3332 -3.0960063928805233f64,
3333 -333.2378467041041f64,
3334 -3170.3105463096517f64,
3335 f64::INFINITY,
3336 f64::NEG_INFINITY,
3337 f64::NAN,
3338 ],
3339 ),
3340 (
3341 4i32,
3342 [
3343 0.0f64,
3344 3.152701157357188f64,
3345 311.17631066575086f64,
3346 3091.493858659732f64,
3347 0.0f64,
3348 -3.152701157357188f64,
3349 -311.17631066575086f64,
3350 -3091.493858659732f64,
3351 f64::INFINITY,
3352 f64::NEG_INFINITY,
3353 f64::NAN,
3354 ],
3355 ),
3356 (
3357 8i32,
3358 [
3359 0.0f64,
3360 3.1391891212579424f64,
3361 314.0658342072582f64,
3362 3145.6489181930947f64,
3363 0.0f64,
3364 -3.1391891212579424f64,
3365 -314.0658342072582f64,
3366 -3145.6489181930947f64,
3367 f64::INFINITY,
3368 f64::NEG_INFINITY,
3369 f64::NAN,
3370 ],
3371 ),
3372 (
3373 16i32,
3374 [
3375 0.0f64,
3376 3.141594303685526f64,
3377 314.1602303152259f64,
3378 3141.606302893263f64,
3379 0.0f64,
3380 -3.141594303685526f64,
3381 -314.1602303152259f64,
3382 -3141.606302893263f64,
3383 f64::INFINITY,
3384 f64::NEG_INFINITY,
3385 f64::NAN,
3386 ],
3387 ),
3388 ] {
3389 for (case, expected) in cases.iter().copied().zip(expected.iter().copied()) {
3390 let actual = midpoint(case, scale);
3391 let roundtrip = midpoint(actual.get(), scale);
3392
3393 if expected.is_nan() && actual.get().is_nan() && roundtrip.get().is_nan() {
3394 continue;
3395 }
3396
3397 assert_eq!(
3398 expected.to_bits(),
3399 actual.get().to_bits(),
3400 "expected midpoint({case}, {scale}) to be {expected}, but got {actual}"
3401 );
3402
3403 assert_eq!(
3404 actual.get().to_bits(),
3405 roundtrip.get().to_bits(),
3406 "expected midpoint(midpoint({case}, {scale}), {scale}) to roundtrip to {actual}, but got {roundtrip}"
3407 );
3408 }
3409 }
3410 }
3411 }
3412}
3413
3414mod delta {
3415 use super::*;
3416
3417 use core::mem;
3418
3419 use crate::Timestamp;
3420
3421 pub struct Delta<T> {
3436 start: Option<Timestamp>,
3437 value: T,
3438 }
3439
3440 impl<T> Delta<T> {
3441 pub fn new(start: Option<Timestamp>, initial: T) -> Self {
3445 Delta {
3446 start,
3447 value: initial,
3448 }
3449 }
3450
3451 pub fn new_default(start: Option<Timestamp>) -> Self
3455 where
3456 T: Default,
3457 {
3458 Self::new(start, Default::default())
3459 }
3460
3461 pub fn current_start(&self) -> Option<&Timestamp> {
3465 self.start.as_ref()
3466 }
3467
3468 pub fn current_value_mut(&mut self) -> &mut T {
3472 &mut self.value
3473 }
3474
3475 pub fn current_value(&self) -> &T {
3479 &self.value
3480 }
3481
3482 pub fn advance(&mut self, end: Option<Timestamp>) -> (Option<Extent>, &mut T) {
3491 let start = mem::replace(&mut self.start, end);
3492
3493 let extent = (start..end).to_extent();
3494
3495 (extent, &mut self.value)
3496 }
3497
3498 pub fn advance_default(&mut self, end: Option<Timestamp>) -> (Option<Extent>, T)
3507 where
3508 T: Default,
3509 {
3510 let (extent, value) = self.advance(end);
3511
3512 (extent, mem::take(value))
3513 }
3514 }
3515
3516 #[cfg(test)]
3517 mod tests {
3518 use super::*;
3519
3520 use core::time::Duration;
3521
3522 #[test]
3523 fn delta_advance() {
3524 let mut delta = Delta::new(Some(Timestamp::MIN), 0);
3525
3526 *delta.current_value_mut() += 1;
3527
3528 let (extent, value) = delta.advance(Some(Timestamp::MIN + Duration::from_secs(1)));
3529 let extent = extent.unwrap();
3530 let range = extent.as_range().unwrap();
3531
3532 assert_eq!(
3533 Timestamp::MIN..Timestamp::MIN + Duration::from_secs(1),
3534 *range
3535 );
3536 assert_eq!(1, *value);
3537 }
3538 }
3539}
3540
3541pub use self::delta::*;
3542
3543#[cfg(test)]
3544mod tests {
3545 use super::*;
3546 use std::time::Duration;
3547
3548 use crate::{Timestamp, Value};
3549
3550 #[test]
3551 fn metric_new() {
3552 let metric = Metric::new(
3553 Path::new_raw("test"),
3554 Timestamp::from_unix(Duration::from_secs(1)),
3555 [
3556 ("metric_prop", Value::from(true)),
3557 ("metric_name", Value::from("my metric")),
3558 ("metric_value", Value::from(42)),
3559 ("metric_description", Value::from("my description")),
3560 ("metric_agg", Value::from("count")),
3561 ],
3562 );
3563
3564 assert_eq!("test", metric.mdl());
3565 assert_eq!(
3566 Timestamp::from_unix(Duration::from_secs(1)).unwrap(),
3567 metric.extent().unwrap().as_point()
3568 );
3569 assert_eq!("my metric", metric.name().unwrap());
3570 assert_eq!("my description", metric.description().unwrap());
3571 assert_eq!("count", metric.agg().unwrap());
3572 assert_eq!(42, metric.value().to_value().cast::<i32>().unwrap());
3573 assert_eq!(true, metric.props().pull::<bool, _>("metric_prop").unwrap());
3574
3575 let metric = metric
3576 .with_name("my metric 2")
3577 .with_description("my description 2")
3578 .with_agg("last")
3579 .with_value(17)
3580 .with_unit("ms")
3581 .with_dist_min(0)
3582 .with_dist_max(100)
3583 .with_dist_sum(1000)
3584 .with_dist_count(42);
3585
3586 assert_eq!("my metric 2", metric.name().unwrap());
3587 assert_eq!("my description 2", metric.description().unwrap());
3588 assert_eq!("last", metric.agg().unwrap());
3589 assert_eq!(17, metric.value().to_value().cast::<i32>().unwrap());
3590 assert_eq!("ms", metric.unit().unwrap());
3591
3592 assert_eq!(0, metric.dist_min().to_value().cast::<i32>().unwrap());
3593 assert_eq!(100, metric.dist_max().to_value().cast::<i32>().unwrap());
3594 assert_eq!(1000, metric.dist_sum().to_value().cast::<i32>().unwrap());
3595 assert_eq!(42, metric.dist_count().to_value().cast::<i32>().unwrap());
3596
3597 #[cfg(feature = "alloc")]
3598 {
3599 let set = exp::BucketSet::from_iter([
3600 (exp::Point::new(0.0), 3),
3601 (exp::Point::new(0.0), 2),
3602 (exp::Point::new(1.0), 2),
3603 ]);
3604
3605 let metric = metric
3606 .with_dist_exp_scale(-1)
3607 .with_dist_exp_buckets(set.clone());
3608
3609 assert_eq!(-1, metric.dist_exp_scale().unwrap());
3610 assert_eq!(set, metric.dist_exp_buckets().unwrap());
3611 }
3612 }
3613
3614 #[test]
3615 fn metric_to_event() {
3616 let metric = Metric::new(
3617 Path::new_raw("test"),
3618 Timestamp::from_unix(Duration::from_secs(1)),
3619 [
3620 ("metric_prop", Value::from(true)),
3621 ("metric_name", Value::from("my metric")),
3622 ("metric_agg", Value::from("count")),
3623 ("metric_value", Value::from(42)),
3624 ],
3625 );
3626
3627 let evt = metric.to_event();
3628
3629 assert_eq!("test", evt.mdl());
3630 assert_eq!(
3631 Timestamp::from_unix(Duration::from_secs(1)).unwrap(),
3632 evt.extent().unwrap().as_point()
3633 );
3634 assert_eq!("count of my metric is 42", evt.msg().to_string());
3635 assert_eq!("count", evt.props().pull::<Str, _>(KEY_METRIC_AGG).unwrap());
3636 assert_eq!(42, evt.props().pull::<i32, _>(KEY_METRIC_VALUE).unwrap());
3637 assert_eq!(
3638 "my metric",
3639 evt.props().pull::<Str, _>(KEY_METRIC_NAME).unwrap()
3640 );
3641 assert_eq!(true, evt.props().pull::<bool, _>("metric_prop").unwrap());
3642 assert_eq!(
3643 Kind::Metric,
3644 evt.props().pull::<Kind, _>(KEY_EVT_KIND).unwrap()
3645 );
3646 }
3647
3648 #[test]
3649 fn metric_to_event_uses_tpl() {
3650 assert_eq!(
3651 "test",
3652 Metric::new(
3653 Path::new_raw("test"),
3654 Timestamp::from_unix(Duration::from_secs(1)),
3655 ("metric_prop", true),
3656 )
3657 .with_tpl(Template::literal("test"))
3658 .to_event()
3659 .msg()
3660 .to_string(),
3661 );
3662 }
3663
3664 #[test]
3665 fn metric_to_extent() {
3666 for (case, expected) in [
3667 (
3668 Some(Timestamp::from_unix(Duration::from_secs(1)).unwrap()),
3669 Some(Extent::point(
3670 Timestamp::from_unix(Duration::from_secs(1)).unwrap(),
3671 )),
3672 ),
3673 (None, None),
3674 ] {
3675 let metric = Metric::new(Path::new_raw("test"), case, ("metric_prop", true));
3676
3677 let extent = metric.to_extent();
3678
3679 assert_eq!(
3680 expected.map(|extent| extent.as_range().cloned()),
3681 extent.map(|extent| extent.as_range().cloned())
3682 );
3683 }
3684 }
3685}