1use std::{future::pending, num::NonZeroU64, sync::Mutex, time::Duration};
2
3use async_trait::async_trait;
4use ddsketch::DDSketch;
5use hashbrown::{hash_map::Entry, HashMap};
6use saluki_common::time::get_unix_timestamp;
7use saluki_core::{
8 accounting::{MemoryBounds, MemoryBoundsBuilder, UsageExpr},
9 components::{transforms::*, BuildContext},
10 data_model::event::{
11 metric::{context::Context, *},
12 Event, EventType,
13 },
14 observability::ComponentMetricsExt as _,
15 topology::{interconnect::BufferedDispatcher, EventsBuffer, OutputDefinition},
16};
17use saluki_error::{generic_error, GenericError};
18use saluki_metrics::MetricsBuilder;
19use smallvec::SmallVec;
20use stringtheory::MetaString;
21use tokio::{
22 select,
23 sync::{mpsc, oneshot},
24 time::interval_at,
25};
26use tracing::{debug, error, info, trace, warn};
27
28mod telemetry;
29use self::telemetry::Telemetry;
30
31mod config;
32pub use self::config::HistogramConfiguration;
33use self::config::HistogramStatistic;
34
35const CONTEXT_SNAPSHOT_REQUEST_CHANNEL_CAPACITY: usize = 1;
36
37#[derive(Clone, Copy, Debug, Eq, PartialEq)]
39pub enum AggregateMetricType {
40 Counter,
42
43 Rate,
45
46 Gauge,
48
49 Set,
51
52 Histogram,
54
55 Distribution,
57}
58
59impl From<&MetricValues> for AggregateMetricType {
60 fn from(values: &MetricValues) -> Self {
61 match values {
62 MetricValues::Counter(_) => Self::Counter,
63 MetricValues::Rate(_, _) => Self::Rate,
64 MetricValues::Gauge(_) => Self::Gauge,
65 MetricValues::Set(_) => Self::Set,
66 MetricValues::Histogram(_) => Self::Histogram,
67 MetricValues::Distribution(_) => Self::Distribution,
68 }
69 }
70}
71
72#[derive(Clone, Debug, Eq, PartialEq)]
76pub struct AggregateContextSnapshotEntry {
77 context: Context,
78 metric_type: AggregateMetricType,
79 unit: MetaString,
80}
81
82impl AggregateContextSnapshotEntry {
83 pub fn context(&self) -> &Context {
85 &self.context
86 }
87
88 pub fn metric_type(&self) -> AggregateMetricType {
90 self.metric_type
91 }
92
93 pub fn unit(&self) -> Option<&str> {
95 if self.unit.is_empty() {
96 None
97 } else {
98 Some(&self.unit)
99 }
100 }
101
102 #[cfg(any(test, feature = "test-util"))]
104 pub fn for_test(context: Context, metric_type: AggregateMetricType, unit: MetaString) -> Self {
105 Self {
106 context,
107 metric_type,
108 unit,
109 }
110 }
111}
112
113type AggregateContextSnapshot = Vec<AggregateContextSnapshotEntry>;
114type AggregateContextSnapshotRequest = oneshot::Sender<AggregateContextSnapshot>;
115type AggregateContextSnapshotRequestReceiver = mpsc::Receiver<AggregateContextSnapshotRequest>;
116
117#[derive(Clone, Debug)]
122pub struct AggregateContextSnapshotHandle {
123 requests: mpsc::Sender<AggregateContextSnapshotRequest>,
124}
125
126impl AggregateContextSnapshotHandle {
127 pub async fn snapshot(&self) -> Result<Vec<AggregateContextSnapshotEntry>, GenericError> {
136 let (response_tx, response_rx) = oneshot::channel();
137 self.requests
138 .send(response_tx)
139 .await
140 .map_err(|_| generic_error!("aggregate context snapshot owner is unavailable"))?;
141
142 response_rx
143 .await
144 .map_err(|_| generic_error!("aggregate context snapshot owner stopped before responding"))
145 }
146}
147
148pub fn aggregate_context_snapshot_channel() -> (AggregateContextSnapshotHandle, AggregateContextSnapshotReceiver) {
154 let (requests, receiver) = mpsc::channel(CONTEXT_SNAPSHOT_REQUEST_CHANNEL_CAPACITY);
155 (
156 AggregateContextSnapshotHandle { requests },
157 AggregateContextSnapshotReceiver {
158 receiver: Mutex::new(Some(receiver)),
159 },
160 )
161}
162
163#[inline]
164fn send_context_snapshot_if_open(
165 response: AggregateContextSnapshotRequest, build_snapshot: impl FnOnce() -> AggregateContextSnapshot,
166) {
167 if response.is_closed() {
168 return;
169 }
170
171 let _ = response.send(build_snapshot());
172}
173
174#[cfg(any(test, feature = "test-util"))]
176pub struct AggregateContextSnapshotPendingResponse {
177 response: AggregateContextSnapshotRequest,
178}
179
180#[cfg(any(test, feature = "test-util"))]
181impl AggregateContextSnapshotPendingResponse {
182 pub fn respond(self, snapshot: Vec<AggregateContextSnapshotEntry>) {
186 let _ = self.response.send(snapshot);
187 }
188}
189
190#[cfg(any(test, feature = "test-util"))]
192pub struct AggregateContextSnapshotResponder {
193 receiver: AggregateContextSnapshotRequestReceiver,
194}
195
196#[cfg(any(test, feature = "test-util"))]
197impl AggregateContextSnapshotResponder {
198 pub async fn respond(&mut self, snapshot: Vec<AggregateContextSnapshotEntry>) -> Result<(), GenericError> {
206 self.receive().await?.respond(snapshot);
207 Ok(())
208 }
209
210 pub async fn receive(&mut self) -> Result<AggregateContextSnapshotPendingResponse, GenericError> {
219 let response = self
220 .receiver
221 .recv()
222 .await
223 .ok_or_else(|| generic_error!("aggregate context snapshot request channel is closed"))?;
224 Ok(AggregateContextSnapshotPendingResponse { response })
225 }
226
227 pub async fn stop_after_receiving(&mut self) -> Result<(), GenericError> {
236 drop(self.receive().await?);
237 Ok(())
238 }
239}
240
241#[cfg(any(test, feature = "test-util"))]
243pub fn aggregate_context_snapshot_channel_for_test(
244) -> (AggregateContextSnapshotHandle, AggregateContextSnapshotResponder) {
245 let (requests, receiver) = mpsc::channel(CONTEXT_SNAPSHOT_REQUEST_CHANNEL_CAPACITY);
246 (
247 AggregateContextSnapshotHandle { requests },
248 AggregateContextSnapshotResponder { receiver },
249 )
250}
251
252pub struct AggregateContextSnapshotReceiver {
257 receiver: Mutex<Option<AggregateContextSnapshotRequestReceiver>>,
258}
259
260impl AggregateContextSnapshotReceiver {
261 fn take_receiver(&self) -> Result<AggregateContextSnapshotRequestReceiver, GenericError> {
262 let mut receiver = self
263 .receiver
264 .lock()
265 .map_err(|_| generic_error!("aggregate context snapshot receiver lock is poisoned"))?;
266 receiver
267 .take()
268 .ok_or_else(|| generic_error!("aggregate context snapshot receiver has already been taken"))
269 }
270}
271
272pub struct AggregateConfiguration {
292 pub window_duration_seconds: NonZeroU64,
298
299 pub primary_flush_interval: Duration,
304
305 pub context_limit: usize,
314
315 pub flush_open_windows: bool,
325
326 pub counter_expiry_seconds: Option<u64>,
339
340 pub hist_config: HistogramConfiguration,
342
343 pub context_snapshot_receiver: AggregateContextSnapshotReceiver,
348}
349
350#[cfg(test)]
351impl AggregateConfiguration {
352 fn for_test() -> Self {
357 let (_handle, context_snapshot_receiver) = aggregate_context_snapshot_channel();
358
359 Self {
360 window_duration_seconds: NonZeroU64::new(10).expect("not zero"),
361 primary_flush_interval: Duration::from_secs(15),
362 context_limit: 1_000_000,
363 flush_open_windows: false,
364 counter_expiry_seconds: Some(300),
365 hist_config: HistogramConfiguration::default(),
366 context_snapshot_receiver,
367 }
368 }
369}
370
371#[async_trait]
372impl TransformBuilder for AggregateConfiguration {
373 async fn build(&self, context: BuildContext) -> Result<Box<dyn Transform + Send>, GenericError> {
374 let context_snapshot_requests = self.context_snapshot_receiver.take_receiver()?;
375 let metrics_builder = MetricsBuilder::from_component_context(context.component_context());
376 let telemetry = Telemetry::new(&metrics_builder);
377
378 let state = AggregationState::new(
379 self.window_duration_seconds,
380 self.context_limit,
381 self.counter_expiry_seconds.filter(|s| *s != 0).map(Duration::from_secs),
382 self.hist_config.clone(),
383 telemetry.clone(),
384 );
385
386 Ok(Box::new(Aggregate {
387 state,
388 telemetry,
389 primary_flush_interval: self.primary_flush_interval,
390 flush_open_windows: self.flush_open_windows,
391 context_snapshot_requests: Some(context_snapshot_requests),
392 }))
393 }
394
395 fn input_event_type(&self) -> EventType {
396 EventType::Metric
397 }
398
399 fn outputs(&self) -> &[OutputDefinition<EventType>] {
400 static OUTPUTS: &[OutputDefinition<EventType>] = &[OutputDefinition::default_output(EventType::Metric)];
401 OUTPUTS
402 }
403}
404
405impl MemoryBounds for AggregateConfiguration {
406 fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
407 builder
416 .minimum()
417 .with_single_value::<Aggregate>("component struct");
419 builder
420 .firm()
421 .with_expr(UsageExpr::product(
423 "aggregation state map",
424 UsageExpr::sum(
425 "context map entry",
426 UsageExpr::struct_size::<Context>("context"),
427 UsageExpr::struct_size::<AggregatedMetric>("aggregated metric"),
428 ),
429 UsageExpr::config("aggregate_context_limit", self.context_limit),
430 ))
431 .with_expr(UsageExpr::product(
433 "retained context snapshot",
434 UsageExpr::struct_size::<AggregateContextSnapshotEntry>("snapshot entry"),
435 UsageExpr::config("aggregate_context_limit", self.context_limit),
436 ));
437 }
438}
439
440pub struct Aggregate {
441 state: AggregationState,
442 telemetry: Telemetry,
443 primary_flush_interval: Duration,
444 flush_open_windows: bool,
445 context_snapshot_requests: Option<AggregateContextSnapshotRequestReceiver>,
446}
447
448#[async_trait]
449impl Transform for Aggregate {
450 async fn run(mut self: Box<Self>, mut context: TransformContext) -> Result<(), GenericError> {
451 let mut health = context.take_health_handle();
452
453 let mut primary_flush = interval_at(
454 tokio::time::Instant::now() + self.primary_flush_interval,
455 self.primary_flush_interval,
456 );
457 let mut final_primary_flush = false;
458
459 health.mark_ready();
460 debug!("Aggregation transform started.");
461
462 loop {
463 select! {
464 _ = health.live() => continue,
465 _ = primary_flush.tick() => {
466 if !self.state.is_empty() {
470 debug!("Flushing aggregated metrics...");
471
472 let should_flush_open_windows = final_primary_flush && self.flush_open_windows;
473
474 let was_breached = self.state.context_limit_breached();
476
477 let mut dispatcher = context.dispatcher().buffered().expect("default output should always exist");
478 let flush = async {
479 if let Err(e) = self.state.flush(get_unix_timestamp(), should_flush_open_windows, &mut dispatcher).await {
480 error!(error = %e, "Failed to flush aggregation state.");
481 }
482
483 self.telemetry.increment_flushes();
484
485 if was_breached && !self.state.context_limit_breached() {
487 info!("Context limit no longer exceeded, metrics are being accepted again.");
488 }
489
490 match dispatcher.flush().await {
491 Ok(aggregated_events) => debug!(aggregated_events, "Dispatched events."),
492 Err(e) => error!(error = %e, "Failed to flush aggregated events."),
493 }
494 };
495 tokio::pin!(flush);
496
497 loop {
501 select! {
502 biased;
503 _ = health.live() => {},
504 _ = &mut flush => break,
505 }
506 }
507 }
508
509 if final_primary_flush {
511 debug!("All aggregation complete.");
512 break
513 }
514 },
515 snapshot_request = receive_context_snapshot_request(&mut self.context_snapshot_requests) => {
516 match snapshot_request {
517 Some(response) => {
518 send_context_snapshot_if_open(response, || self.state.snapshot_contexts());
519 }
520 None => self.context_snapshot_requests = None,
521 }
522 },
523 maybe_events = context.events().next(), if !final_primary_flush => match maybe_events {
524 Some(events) => {
525 trace!(events_len = events.len(), "Received events.");
526
527 let current_time = get_unix_timestamp();
528
529 for event in events {
530 if let Some(metric) = event.try_into_metric() {
531 let was_breached = self.state.context_limit_breached();
532 if !self.state.insert(current_time, metric) {
533 trace!("Dropping metric due to context limit.");
534 if !was_breached {
535 warn!(context_limit = self.state.context_limit, "Context limit reached, \
537 dropping metrics. Consider increasing `aggregate_context_limit`.");
538 }
539 self.telemetry.increment_events_dropped();
540 }
541 }
542 }
543 },
544 None => {
545 final_primary_flush = true;
548 primary_flush.reset_immediately();
549
550 debug!("Aggregation transform stopping...");
551 }
552 },
553 }
554 }
555
556 debug!("Aggregation transform stopped.");
557
558 Ok(())
559 }
560}
561
562async fn receive_context_snapshot_request(
563 receiver: &mut Option<AggregateContextSnapshotRequestReceiver>,
564) -> Option<AggregateContextSnapshotRequest> {
565 match receiver {
566 Some(receiver) => receiver.recv().await,
567 None => pending().await,
568 }
569}
570
571#[derive(Clone)]
572struct AggregatedMetric {
573 values: MetricValues,
574 metadata: MetricMetadata,
575 last_seen: u64,
576}
577
578struct AggregationState {
579 contexts: HashMap<Context, AggregatedMetric, foldhash::quality::RandomState>,
580 contexts_remove_buf: Vec<Context>,
581 context_limit: usize,
582 bucket_width_secs: NonZeroU64,
583 counter_expire_secs: Option<NonZeroU64>,
584 last_flush: u64,
585 hist_config: HistogramConfiguration,
586 telemetry: Telemetry,
587 context_limit_breached: bool,
590}
591
592impl AggregationState {
593 fn new(
594 bucket_width_secs: NonZeroU64, context_limit: usize, counter_expiration: Option<Duration>,
595 hist_config: HistogramConfiguration, telemetry: Telemetry,
596 ) -> Self {
597 let counter_expire_secs = counter_expiration.map(|d| d.as_secs()).and_then(NonZeroU64::new);
598
599 Self {
600 contexts: HashMap::default(),
601 contexts_remove_buf: Vec::new(),
602 context_limit,
603 bucket_width_secs,
604 counter_expire_secs,
605 last_flush: 0,
606 hist_config,
607 telemetry,
608 context_limit_breached: false,
609 }
610 }
611
612 fn is_empty(&self) -> bool {
613 self.contexts.is_empty()
614 }
615
616 fn snapshot_contexts(&self) -> Vec<AggregateContextSnapshotEntry> {
617 let mut snapshot = Vec::with_capacity(self.contexts.len());
618 for (context, aggregated) in &self.contexts {
619 snapshot.push(AggregateContextSnapshotEntry {
620 context: context.clone(),
621 metric_type: AggregateMetricType::from(&aggregated.values),
622 unit: aggregated.metadata.unit.clone(),
623 });
624 }
625 snapshot
626 }
627
628 fn insert(&mut self, timestamp: u64, metric: Metric) -> bool {
629 if !self.contexts.contains_key(metric.context()) && self.contexts.len() >= self.context_limit {
631 self.context_limit_breached = true;
632 return false;
633 }
634
635 let (context, mut values, metadata) = metric.into_parts();
636
637 let bucket_ts = align_to_bucket_start(timestamp, self.bucket_width_secs);
642 values.collapse_non_timestamped(bucket_ts);
643
644 trace!(
645 bucket_ts,
646 kind = values.as_str(),
647 "Inserting metric into aggregation state."
648 );
649
650 match self.contexts.entry(context) {
653 Entry::Occupied(mut entry) => {
654 let aggregated = entry.get_mut();
655
656 aggregated.last_seen = timestamp;
658 aggregated.values.merge(values);
659 }
660 Entry::Vacant(entry) => {
661 self.telemetry.increment_contexts(entry.key(), &values);
662
663 entry.insert(AggregatedMetric {
664 values,
665 metadata,
666 last_seen: timestamp,
667 });
668
669 saluki_antithesis::always_le!(
670 self.contexts.len(),
671 self.context_limit,
672 "aggregate context map within context_limit",
673 { "len": self.contexts.len(), "limit": self.context_limit }
674 );
675 }
676 }
677
678 true
679 }
680
681 async fn flush(
682 &mut self, current_time: u64, flush_open_buckets: bool, dispatcher: &mut BufferedDispatcher<'_, EventsBuffer>,
683 ) -> Result<(), GenericError> {
684 let bucket_width_secs = self.bucket_width_secs;
685 let counter_expire_secs = self.counter_expire_secs.map(|d| d.get()).unwrap_or(0);
686
687 let split_timestamp = align_to_bucket_start(current_time, bucket_width_secs).saturating_sub(1);
690
691 let mut zero_value_buckets = SmallVec::<[(u64, MetricValues); 4]>::new();
696 if self.last_flush != 0 {
697 let start = align_to_bucket_start(self.last_flush, bucket_width_secs);
698
699 saluki_antithesis::always_ge!(
704 current_time,
705 self.last_flush,
706 "aggregate flush wall-clock did not move backward",
707 { "current_time": current_time, "last_flush": self.last_flush }
708 );
709 saluki_antithesis::always_le!(
712 current_time.saturating_sub(self.last_flush) / bucket_width_secs.get(),
713 10_000,
714 "aggregate zero-value bucket span bounded across a flush",
715 {
716 "current_time": current_time,
717 "last_flush": self.last_flush,
718 "bucket_width_secs": bucket_width_secs.get()
719 }
720 );
721
722 for bucket_start in (start..current_time).step_by(bucket_width_secs.get() as usize) {
723 if is_bucket_closed(current_time, bucket_start, bucket_width_secs, flush_open_buckets) {
724 zero_value_buckets.push((bucket_start, MetricValues::counter((bucket_start, 0.0))));
725 }
726 }
727
728 saluki_antithesis::sometimes!(
730 !zero_value_buckets.is_empty(),
731 "aggregate flush generated zero-value counter buckets",
732 { "count": zero_value_buckets.len() }
733 );
734 }
735
736 debug!(timestamp = current_time, "Flushing buckets.");
738
739 for (context, am) in self.contexts.iter_mut() {
740 let should_expire_if_empty = match &am.values {
749 MetricValues::Counter(..) => {
750 saluki_antithesis::always_le!(
751 am.last_seen,
752 u64::MAX - counter_expire_secs,
753 "aggregate counter expiry add does not overflow",
754 { "last_seen": am.last_seen, "counter_expire_secs": counter_expire_secs }
755 );
756 counter_expire_secs != 0 && am.last_seen.saturating_add(counter_expire_secs) < current_time
757 }
758 _ => true,
759 };
760
761 if let MetricValues::Counter(..) = &mut am.values {
767 let expires_at = am.last_seen.saturating_add(counter_expire_secs);
768 for (zv_bucket_start, zero_value) in &zero_value_buckets {
769 if expires_at > *zv_bucket_start {
770 am.values.merge(zero_value.clone());
771 } else {
772 break;
775 }
776 }
777 }
778
779 if let Some(closed_bucket_values) = am.values.split_at_timestamp(split_timestamp) {
788 self.telemetry.increment_flushed(&closed_bucket_values);
789
790 transform_and_push_metric(
792 context.clone(),
793 closed_bucket_values,
794 am.metadata.clone(),
795 bucket_width_secs,
796 &self.hist_config,
797 dispatcher,
798 )
799 .await?;
800 }
801
802 if am.values.is_empty() && should_expire_if_empty {
803 self.telemetry.decrement_contexts(context, &am.values);
804 self.contexts_remove_buf.push(context.clone());
805 }
806 }
807
808 let contexts_len_before = self.contexts.len();
810 for context in self.contexts_remove_buf.drain(..) {
811 self.contexts.remove(&context);
812 }
813 let contexts_len_after = self.contexts.len();
814
815 let contexts_delta = contexts_len_before.saturating_sub(contexts_len_after);
816 let target_contexts_capacity = contexts_len_after.saturating_add(contexts_delta / 2);
817 self.contexts.shrink_to(target_contexts_capacity);
818
819 if self.context_limit_breached && self.contexts.len() < self.context_limit {
820 self.context_limit_breached = false;
821 }
822
823 self.last_flush = current_time;
824
825 Ok(())
826 }
827
828 fn context_limit_breached(&self) -> bool {
829 self.context_limit_breached
830 }
831}
832
833async fn transform_and_push_metric(
834 context: Context, mut values: MetricValues, metadata: MetricMetadata, bucket_width_secs: NonZeroU64,
835 hist_config: &HistogramConfiguration, dispatcher: &mut BufferedDispatcher<'_, EventsBuffer>,
836) -> Result<(), GenericError> {
837 let bucket_width = Duration::from_secs(bucket_width_secs.get());
838
839 match values {
840 MetricValues::Histogram(ref mut points) => {
843 if hist_config.copy_to_distribution() {
845 let sketch_points = points
846 .into_iter()
847 .map(|(ts, hist)| {
848 let mut sketch = DDSketch::default();
849 for sample in hist.samples() {
850 sketch.insert_n(sample.value.into_inner(), sample.weight.0 as u64);
851 }
852 (ts, sketch)
853 })
854 .collect::<SketchPoints>();
855 let distribution_values = MetricValues::distribution(sketch_points);
856 let metric_context = if !hist_config.copy_to_distribution_prefix().is_empty() {
857 context.with_name(format!(
858 "{}{}",
859 hist_config.copy_to_distribution_prefix(),
860 context.name()
861 ))
862 } else {
863 context.clone()
864 };
865 let new_metric = Metric::from_parts(metric_context, distribution_values, metadata.clone());
866 dispatcher.push(Event::Metric(new_metric)).await?;
867 }
868 let mut sorted_points = Vec::new();
873 for (ts, h) in points {
874 sorted_points.push((ts, h.summary_view()));
875 }
876
877 for statistic in hist_config.statistics() {
878 let new_points = sorted_points
879 .iter()
880 .map(|(ts, hs)| (*ts, statistic.value_from_histogram(hs)))
881 .collect::<ScalarPoints>();
882
883 let new_values = if statistic.is_rate_statistic() {
884 MetricValues::rate(new_points, bucket_width)
885 } else {
886 MetricValues::gauge(new_points)
887 };
888
889 let new_metadata = if matches!(statistic, HistogramStatistic::Count) {
891 metadata.clone().with_unit(MetaString::empty())
892 } else {
893 metadata.clone()
894 };
895
896 let new_context = context.with_name(format!("{}.{}", context.name(), statistic.suffix()));
897 let new_metric = Metric::from_parts(new_context, new_values, new_metadata);
898 dispatcher.push(Event::Metric(new_metric)).await?;
899 }
900
901 Ok(())
902 }
903
904 values => {
907 let adjusted_values = counter_values_to_rate(values, bucket_width_secs);
908
909 let metric = Metric::from_parts(context, adjusted_values, metadata);
910 dispatcher.push(Event::Metric(metric)).await
911 }
912 }
913}
914
915fn counter_values_to_rate(values: MetricValues, interval_secs: NonZeroU64) -> MetricValues {
916 match values {
917 MetricValues::Counter(points) => MetricValues::rate(points, Duration::from_secs(interval_secs.get())),
918 values => values,
919 }
920}
921
922const fn align_to_bucket_start(timestamp: u64, bucket_width_secs: NonZeroU64) -> u64 {
923 timestamp - (timestamp % bucket_width_secs.get())
924}
925
926const fn is_bucket_closed(
927 current_time: u64, bucket_start: u64, bucket_width_secs: NonZeroU64, flush_open_buckets: bool,
928) -> bool {
929 (bucket_start + bucket_width_secs.get() - 1) < current_time || flush_open_buckets
948}
949
950#[cfg(test)]
954mod tests {
955 use std::{cell::Cell, mem::size_of};
956
957 use float_cmp::ApproxEqRatio as _;
958 use saluki_core::{
959 accounting::{ComponentRegistry, MemoryLimiter},
960 components::{
961 destinations::{Destination, DestinationBuilder, DestinationContext},
962 sources::{Source, SourceBuilder, SourceContext},
963 ComponentContext,
964 },
965 data_model::tags::{Tag, TagSet},
966 health::HealthRegistry,
967 runtime::{state::ResourceRegistry, Supervisor},
968 support::SubsystemIdentifier,
969 topology::{interconnect::Dispatcher, EventsDispatcher, OutputDefinition, OutputName, TopologyBlueprint},
970 };
971 use saluki_metrics::test::TestRecorder;
972 use stringtheory::MetaString;
973 use tokio::sync::{mpsc, oneshot};
974
975 use super::config::HistogramStatistic;
976 use super::*;
977
978 const BUCKET_WIDTH_SECS: NonZeroU64 = NonZeroU64::new(10).expect("not zero");
979 const BUCKET_WIDTH: Duration = Duration::from_secs(BUCKET_WIDTH_SECS.get());
980 const COUNTER_EXPIRE_SECS: u64 = 20;
981 const COUNTER_EXPIRE: Option<Duration> = Some(Duration::from_secs(COUNTER_EXPIRE_SECS));
982
983 const fn bucket_ts(step: u64) -> u64 {
985 align_to_bucket_start(insert_ts(step), BUCKET_WIDTH_SECS)
986 }
987
988 const fn insert_ts(step: u64) -> u64 {
990 (BUCKET_WIDTH_SECS.get() * (step + 1)) - 2
991 }
992
993 const fn flush_ts(step: u64) -> u64 {
995 BUCKET_WIDTH_SECS.get() * (step + 1)
996 }
997
998 struct ControlledMetricSource {
999 events: mpsc::Receiver<Event>,
1000 }
1001
1002 #[async_trait]
1003 impl Source for ControlledMetricSource {
1004 async fn run(mut self: Box<Self>, mut context: SourceContext) -> Result<(), GenericError> {
1005 let shutdown = context.take_shutdown_handle();
1006 tokio::pin!(shutdown);
1007 let mut events_open = true;
1008
1009 loop {
1010 select! {
1011 _ = &mut shutdown => break,
1012 maybe_event = self.events.recv(), if events_open => match maybe_event {
1013 Some(event) => context.dispatcher().dispatch_one(event).await?,
1014 None => events_open = false,
1015 }
1016 }
1017 }
1018 Ok(())
1019 }
1020 }
1021
1022 struct ControlledMetricSourceBuilder {
1023 events: Mutex<Option<mpsc::Receiver<Event>>>,
1024 outputs: Vec<OutputDefinition<EventType>>,
1025 }
1026
1027 #[async_trait]
1028 impl SourceBuilder for ControlledMetricSourceBuilder {
1029 fn outputs(&self) -> &[OutputDefinition<EventType>] {
1030 &self.outputs
1031 }
1032
1033 async fn build(&self, _context: BuildContext) -> Result<Box<dyn Source + Send>, GenericError> {
1034 let events = self
1035 .events
1036 .lock()
1037 .map_err(|_| generic_error!("controlled metric source receiver lock is poisoned"))?
1038 .take()
1039 .ok_or_else(|| generic_error!("controlled metric source receiver has already been taken"))?;
1040 Ok(Box::new(ControlledMetricSource { events }))
1041 }
1042 }
1043
1044 impl MemoryBounds for ControlledMetricSourceBuilder {
1045 fn specify_bounds(&self, _builder: &mut MemoryBoundsBuilder) {}
1046 }
1047
1048 struct DrainingMetricDestination;
1049
1050 #[async_trait]
1051 impl Destination for DrainingMetricDestination {
1052 async fn run(self: Box<Self>, mut context: DestinationContext) -> Result<(), GenericError> {
1053 while context.events().next().await.is_some() {}
1054 Ok(())
1055 }
1056 }
1057
1058 struct DrainingMetricDestinationBuilder;
1059
1060 #[async_trait]
1061 impl DestinationBuilder for DrainingMetricDestinationBuilder {
1062 fn input_event_type(&self) -> EventType {
1063 EventType::Metric
1064 }
1065
1066 async fn build(&self, _context: BuildContext) -> Result<Box<dyn Destination + Send>, GenericError> {
1067 Ok(Box::new(DrainingMetricDestination))
1068 }
1069 }
1070
1071 impl MemoryBounds for DrainingMetricDestinationBuilder {
1072 fn specify_bounds(&self, _builder: &mut MemoryBoundsBuilder) {}
1073 }
1074
1075 struct DispatcherReceiver {
1076 receiver: mpsc::Receiver<EventsBuffer>,
1077 }
1078
1079 impl DispatcherReceiver {
1080 fn collect_next(&mut self) -> Vec<Metric> {
1081 match self.receiver.try_recv() {
1082 Ok(event_buffer) => {
1083 let mut metrics = event_buffer
1084 .into_iter()
1085 .filter_map(|event| event.try_into_metric())
1086 .collect::<Vec<Metric>>();
1087
1088 metrics.sort_by(|a, b| a.context().name().cmp(b.context().name()));
1089 metrics
1090 }
1091 Err(_) => Vec::new(),
1092 }
1093 }
1094 }
1095
1096 fn build_basic_dispatcher() -> (EventsDispatcher, DispatcherReceiver) {
1098 let context = ComponentContext::test_transform("test");
1099 let mut dispatcher = Dispatcher::new(context);
1100
1101 let (buffer_tx, buffer_rx) = mpsc::channel(1);
1102 dispatcher.add_output(OutputName::Default).unwrap();
1103 dispatcher
1104 .attach_sender_to_output(&OutputName::Default, buffer_tx)
1105 .unwrap();
1106
1107 (dispatcher, DispatcherReceiver { receiver: buffer_rx })
1108 }
1109
1110 async fn get_flushed_metrics(timestamp: u64, state: &mut AggregationState) -> Vec<Metric> {
1111 let (dispatcher, mut dispatcher_receiver) = build_basic_dispatcher();
1112 let mut buffered_dispatcher = dispatcher.buffered().expect("default output should always exist");
1113
1114 state
1116 .flush(timestamp, true, &mut buffered_dispatcher)
1117 .await
1118 .expect("should not fail to flush aggregation state");
1119
1120 buffered_dispatcher
1123 .flush()
1124 .await
1125 .expect("should not fail to flush buffered sender");
1126
1127 dispatcher_receiver.collect_next()
1128 }
1129
1130 macro_rules! compare_points {
1131 (scalar, $expected:expr, $actual:expr, $error_ratio:literal) => {
1132 for (idx, (expected_value, actual_value)) in $expected.into_iter().zip($actual.into_iter()).enumerate() {
1133 let (expected_ts, expected_point) = expected_value;
1134 let (actual_ts, actual_point) = actual_value;
1135
1136 assert_eq!(
1137 expected_ts, actual_ts,
1138 "timestamp for value #{} does not match: {:?} (expected) vs {:?} (actual)",
1139 idx, expected_ts, actual_ts
1140 );
1141 assert!(
1142 expected_point.approx_eq_ratio(&actual_point, $error_ratio),
1143 "point for value #{} does not match: {} (expected) vs {} (actual)",
1144 idx,
1145 expected_point,
1146 actual_point
1147 );
1148 }
1149 };
1150 (distribution, $expected:expr, $actual:expr) => {
1151 for (idx, (expected_value, actual_value)) in $expected.into_iter().zip($actual.into_iter()).enumerate() {
1152 let (expected_ts, expected_sketch) = expected_value;
1153 let (actual_ts, actual_sketch) = actual_value;
1154
1155 assert_eq!(
1156 expected_ts, actual_ts,
1157 "timestamp for value #{} does not match: {:?} (expected) vs {:?} (actual)",
1158 idx, expected_ts, actual_ts
1159 );
1160 assert_eq!(
1161 expected_sketch, actual_sketch,
1162 "sketch for value #{} does not match: {:?} (expected) vs {:?} (actual)",
1163 idx, expected_sketch, actual_sketch
1164 );
1165 }
1166 };
1167 }
1168
1169 macro_rules! assert_flushed_scalar_metric {
1170 ($original:expr, $actual:expr, [$($ts:expr => $value:expr),+]) => {
1171 assert_flushed_scalar_metric!($original, $actual, [$($ts => $value),+], error_ratio => 0.000001);
1172 };
1173 ($original:expr, $actual:expr, [$($ts:expr => $value:expr),+], error_ratio => $error_ratio:literal) => {
1174 let actual_metric = $actual;
1175
1176 assert_eq!($original.context(), actual_metric.context(), "expected context ({}) and actual context ({}) do not match", $original.context(), actual_metric.context());
1177
1178 let expected_points = ScalarPoints::from([$(($ts, $value)),+]);
1179
1180 match actual_metric.values() {
1181 MetricValues::Counter(ref actual_points) | MetricValues::Gauge(ref actual_points) | MetricValues::Rate(ref actual_points, _) => {
1182 assert_eq!(expected_points.len(), actual_points.len(), "expected and actual values have different number of points");
1183 compare_points!(scalar, expected_points, actual_points, $error_ratio);
1184 },
1185 _ => panic!("only counters, rates, and gauges are supported in assert_flushed_scalar_metric"),
1186 }
1187 };
1188 }
1189
1190 macro_rules! assert_flushed_distribution_metric {
1191 ($original:expr, $actual:expr, [$($ts:expr => $value:expr),+]) => {
1192 assert_flushed_distribution_metric!($original, $actual, [$($ts => $value),+], error_ratio => 0.000001);
1193 };
1194 ($original:expr, $actual:expr, [$($ts:expr => $value:expr),+], error_ratio => $error_ratio:literal) => {
1195 let actual_metric = $actual;
1196
1197 assert_eq!($original.context(), actual_metric.context());
1198
1199 match actual_metric.values() {
1200 MetricValues::Distribution(ref actual_points) => {
1201 let expected_points = SketchPoints::from([$(($ts, $value)),+]);
1202 assert_eq!(expected_points.len(), actual_points.len(), "expected and actual values have different number of points");
1203
1204 compare_points!(distribution, &expected_points, actual_points);
1205 },
1206 _ => panic!("only distributions are supported in assert_flushed_distribution_metric"),
1207 }
1208 };
1209 }
1210
1211 #[test]
1212 fn aggregate_metric_type_matches_every_metric_shape() {
1213 let cases = [
1214 (MetricValues::counter(1.0), AggregateMetricType::Counter),
1215 (
1216 MetricValues::rate(1.0, Duration::from_secs(10)),
1217 AggregateMetricType::Rate,
1218 ),
1219 (MetricValues::gauge(1.0), AggregateMetricType::Gauge),
1220 (MetricValues::set("value"), AggregateMetricType::Set),
1221 (MetricValues::histogram([1.0]), AggregateMetricType::Histogram),
1222 (
1223 MetricValues::distribution(&[1.0][..]),
1224 AggregateMetricType::Distribution,
1225 ),
1226 ];
1227
1228 for (values, expected) in cases {
1229 assert_eq!(AggregateMetricType::from(&values), expected);
1230 }
1231 }
1232
1233 #[test]
1234 fn snapshot_contexts_preserves_full_context_shape_and_unit() {
1235 let mut state = AggregationState::new(
1236 BUCKET_WIDTH_SECS,
1237 10,
1238 COUNTER_EXPIRE,
1239 HistogramConfiguration::default(),
1240 Telemetry::noop(),
1241 );
1242
1243 let histogram_context = Context::from_static_parts("request.duration", &["env:prod"])
1244 .with_host(Some(MetaString::from_static("host-a")))
1245 .with_origin_tags(TagSet::from(Tag::from_static("container:one")));
1246 let gauge_context = Context::from_static_parts("request.duration", &["env:prod"])
1247 .with_host(Some(MetaString::from_static("host-b")))
1248 .with_origin_tags(TagSet::from(Tag::from_static("container:two")));
1249 let histogram = Metric::from_parts(
1250 histogram_context.clone(),
1251 MetricValues::histogram([12.0]),
1252 MetricMetadata::default().with_unit(MetaString::from_static("millisecond")),
1253 );
1254 let gauge = Metric::gauge(gauge_context.clone(), 2.0);
1255
1256 assert!(state.insert(insert_ts(1), histogram));
1257 assert!(state.insert(insert_ts(1), gauge));
1258
1259 let mut snapshot = state.snapshot_contexts();
1260 assert_eq!(snapshot.len(), 2);
1261 assert!(snapshot.capacity() >= state.contexts.len());
1262 snapshot.sort_by(|a, b| a.context().host().cmp(&b.context().host()));
1263
1264 assert_eq!(snapshot[0].context(), &histogram_context);
1265 assert_eq!(snapshot[0].context().tags().len(), 1);
1266 assert_eq!(snapshot[0].context().origin_tags().len(), 1);
1267 assert_eq!(snapshot[0].metric_type(), AggregateMetricType::Histogram);
1268 assert_eq!(snapshot[0].unit(), Some("millisecond"));
1269
1270 assert_eq!(snapshot[1].context(), &gauge_context);
1271 assert_eq!(snapshot[1].metric_type(), AggregateMetricType::Gauge);
1272 assert_eq!(snapshot[1].unit(), None);
1273 }
1274
1275 #[tokio::test]
1276 async fn snapshot_contexts_follows_ordinary_context_lifecycle() {
1277 let mut state = AggregationState::new(
1278 BUCKET_WIDTH_SECS,
1279 10,
1280 COUNTER_EXPIRE,
1281 HistogramConfiguration::default(),
1282 Telemetry::noop(),
1283 );
1284 let context = Context::from_static_name("active.gauge");
1285
1286 assert!(state.insert(insert_ts(1), Metric::gauge(context.clone(), 1.0)));
1287 let active_snapshot = state.snapshot_contexts();
1288 assert_eq!(active_snapshot.len(), 1);
1289 assert_eq!(active_snapshot[0].context(), &context);
1290 assert_eq!(active_snapshot[0].metric_type(), AggregateMetricType::Gauge);
1291
1292 let _ = get_flushed_metrics(flush_ts(1), &mut state).await;
1293 assert!(state.snapshot_contexts().is_empty());
1294 }
1295
1296 #[tokio::test]
1297 async fn snapshot_contexts_retains_idle_counter_until_expiry() {
1298 let mut state = AggregationState::new(
1299 BUCKET_WIDTH_SECS,
1300 10,
1301 COUNTER_EXPIRE,
1302 HistogramConfiguration::default(),
1303 Telemetry::noop(),
1304 );
1305 let context = Context::from_static_name("sparse.counter");
1306
1307 assert!(state.insert(insert_ts(1), Metric::counter(context.clone(), 1.0)));
1308 let _ = get_flushed_metrics(flush_ts(1), &mut state).await;
1309 let first_snapshot = state.snapshot_contexts();
1310 assert_eq!(first_snapshot.len(), 1);
1311 assert_eq!(first_snapshot[0].context(), &context);
1312 assert_eq!(first_snapshot[0].metric_type(), AggregateMetricType::Counter);
1313
1314 let _ = get_flushed_metrics(flush_ts(2), &mut state).await;
1315 let second_snapshot = state.snapshot_contexts();
1316 assert_eq!(second_snapshot.len(), 1);
1317 assert_eq!(second_snapshot[0].context(), &context);
1318 assert_eq!(second_snapshot[0].metric_type(), AggregateMetricType::Counter);
1319
1320 let _ = get_flushed_metrics(flush_ts(3), &mut state).await;
1321 assert!(state.snapshot_contexts().is_empty());
1322 }
1323
1324 #[test]
1325 fn canceled_snapshot_response_skips_snapshot_construction() {
1326 let mut state = AggregationState::new(
1327 BUCKET_WIDTH_SECS,
1328 10,
1329 COUNTER_EXPIRE,
1330 HistogramConfiguration::default(),
1331 Telemetry::noop(),
1332 );
1333 assert!(state.insert(
1334 insert_ts(1),
1335 Metric::gauge(Context::from_static_name("canceled.snapshot"), 1.0),
1336 ));
1337 let (response, receiver) = oneshot::channel();
1338 drop(receiver);
1339 let snapshot_calls = Cell::new(0);
1340
1341 send_context_snapshot_if_open(response, || {
1342 snapshot_calls.set(snapshot_calls.get() + 1);
1343 state.snapshot_contexts()
1344 });
1345
1346 assert_eq!(snapshot_calls.get(), 0);
1347 }
1348
1349 #[test]
1350 fn open_snapshot_response_constructs_once_and_returns_exact_entries() {
1351 let mut state = AggregationState::new(
1352 BUCKET_WIDTH_SECS,
1353 10,
1354 COUNTER_EXPIRE,
1355 HistogramConfiguration::default(),
1356 Telemetry::noop(),
1357 );
1358 let context = Context::from_static_name("open.snapshot");
1359 assert!(state.insert(insert_ts(1), Metric::gauge(context.clone(), 1.0)));
1360 let expected = vec![AggregateContextSnapshotEntry {
1361 context,
1362 metric_type: AggregateMetricType::Gauge,
1363 unit: MetaString::empty(),
1364 }];
1365 let (response, mut receiver) = oneshot::channel();
1366 let snapshot_calls = Cell::new(0);
1367
1368 send_context_snapshot_if_open(response, || {
1369 snapshot_calls.set(snapshot_calls.get() + 1);
1370 state.snapshot_contexts()
1371 });
1372
1373 assert_eq!(snapshot_calls.get(), 1);
1374 assert_eq!(
1375 receiver.try_recv().expect("open requester should receive a snapshot"),
1376 expected
1377 );
1378 }
1379
1380 #[test]
1381 fn aggregate_memory_bounds_include_peak_context_snapshot() {
1382 let context_limit = 17;
1383 let config = AggregateConfiguration {
1384 context_limit,
1385 ..AggregateConfiguration::for_test()
1386 };
1387 let registry = ComponentRegistry::default();
1388 config.specify_bounds(&mut registry.bounds_builder(&SubsystemIdentifier::from_dotted("test")));
1389 let bounds = registry.as_bounds();
1390
1391 let expected_minimum = size_of::<Aggregate>();
1392 let aggregation_state_bytes = context_limit * (size_of::<Context>() + size_of::<AggregatedMetric>());
1393 let context_snapshot_bytes = context_limit * size_of::<AggregateContextSnapshotEntry>();
1394
1395 assert_eq!(bounds.total_minimum_required_bytes(), expected_minimum);
1396 assert_eq!(
1397 bounds.total_firm_limit_bytes(),
1398 expected_minimum + aggregation_state_bytes + context_snapshot_bytes
1399 );
1400 }
1401
1402 #[tokio::test]
1403 async fn production_owner_loop_serves_snapshots_and_stops_after_cancellation() {
1404 tokio::time::timeout(Duration::from_secs(5), async {
1405 let (snapshot_handle, context_snapshot_receiver) = aggregate_context_snapshot_channel();
1406 let config = AggregateConfiguration {
1407 primary_flush_interval: Duration::from_secs(60),
1408 context_snapshot_receiver,
1409 ..AggregateConfiguration::for_test()
1410 };
1411
1412 let available_request_capacity = snapshot_handle.requests.capacity();
1413 let canceled_request = tokio::spawn({
1414 let snapshot_handle = snapshot_handle.clone();
1415 async move { snapshot_handle.snapshot().await }
1416 });
1417 while snapshot_handle.requests.capacity() == available_request_capacity {
1418 tokio::task::yield_now().await;
1419 }
1420 canceled_request.abort();
1421 assert!(canceled_request
1422 .await
1423 .expect_err("snapshot requester should be canceled")
1424 .is_cancelled());
1425
1426 let (events_tx, events_rx) = mpsc::channel(1);
1427 let source = ControlledMetricSourceBuilder {
1428 events: Mutex::new(Some(events_rx)),
1429 outputs: vec![OutputDefinition::default_output(EventType::Metric)],
1430 };
1431 let component_registry = ComponentRegistry::default();
1432 let mut blueprint = TopologyBlueprint::new("aggregate_snapshot_owner", &component_registry);
1433 blueprint
1434 .add_source("source", source)
1435 .expect("controlled source should be accepted")
1436 .add_transform("aggregate", config)
1437 .expect("aggregate transform should be accepted")
1438 .add_destination("destination", DrainingMetricDestinationBuilder)
1439 .expect("draining destination should be accepted");
1440 blueprint
1441 .connect_components_in_order(["source", "aggregate", "destination"])
1442 .expect("test topology should connect");
1443 blueprint
1444 .with_health_registry(HealthRegistry::new())
1445 .with_memory_limiter(MemoryLimiter::noop())
1446 .with_resource_registry(ResourceRegistry::new())
1447 .with_ambient_worker_pool();
1448
1449 let mut supervisor =
1450 Supervisor::new("aggregate-snapshot-owner").expect("test supervisor should be created");
1451 supervisor.add_worker(blueprint);
1452 let (shutdown_tx, shutdown_rx) = oneshot::channel::<()>();
1453 let topology_task = tokio::spawn(async move { supervisor.run_with_shutdown(shutdown_rx).await });
1454
1455 let retained_context = Context::from_static_name("owner.loop.mixed.counter");
1456 let mixed_values = ScalarPoints::from_iter([(None, 2.0), (NonZeroU64::new(insert_ts(1)), 3.0)]);
1457 events_tx
1458 .send(Event::Metric(Metric::counter(retained_context.clone(), mixed_values)))
1459 .await
1460 .expect("controlled source should accept a mixed event");
1461
1462 let snapshot = loop {
1463 let snapshot = snapshot_handle
1464 .snapshot()
1465 .await
1466 .expect("running aggregate should fulfill snapshots");
1467 if snapshot.iter().any(|entry| entry.context() == &retained_context) {
1468 break snapshot;
1469 }
1470 tokio::task::yield_now().await;
1471 };
1472 assert_eq!(
1473 snapshot.len(),
1474 1,
1475 "only the mixed metric's retained portion belongs in state"
1476 );
1477 let entry = &snapshot[0];
1478 assert_eq!(entry.context(), &retained_context);
1479 assert_eq!(entry.metric_type(), AggregateMetricType::Counter);
1480 assert_eq!(entry.unit(), None);
1481
1482 drop(events_tx);
1483 drop(snapshot_handle);
1484 shutdown_tx.send(()).expect("test topology should still be running");
1485 let topology_result = topology_task.await.expect("topology task should not panic");
1486 assert!(
1487 topology_result.is_ok(),
1488 "topology should stop cleanly: {topology_result:?}"
1489 );
1490 })
1491 .await
1492 .expect("production aggregate owner loop should complete without spinning or hanging");
1493 }
1494
1495 #[tokio::test]
1496 async fn snapshot_handle_round_trips_through_responder() {
1497 let (handle, mut responder) = aggregate_context_snapshot_channel_for_test();
1498 let expected = vec![AggregateContextSnapshotEntry::for_test(
1499 Context::from_static_name("round.trip"),
1500 AggregateMetricType::Gauge,
1501 MetaString::from_static("widget"),
1502 )];
1503 let snapshot_task = tokio::spawn(async move { handle.snapshot().await });
1504
1505 responder
1506 .respond(expected.clone())
1507 .await
1508 .expect("responder should receive and fulfill a snapshot request");
1509
1510 let actual = snapshot_task
1511 .await
1512 .expect("snapshot task should complete")
1513 .expect("snapshot request should succeed");
1514 assert_eq!(actual, expected);
1515 }
1516
1517 #[tokio::test]
1518 async fn snapshot_handle_reports_dropped_receiver() {
1519 let (handle, responder) = aggregate_context_snapshot_channel_for_test();
1520 drop(responder);
1521
1522 let error = handle
1523 .snapshot()
1524 .await
1525 .expect_err("snapshot should fail after its owner is dropped");
1526 assert!(error.to_string().contains("unavailable"));
1527 }
1528
1529 #[tokio::test]
1530 async fn snapshot_responder_stops_after_accepting_request_and_cancels_response() {
1531 let (handle, mut responder) = aggregate_context_snapshot_channel_for_test();
1532 let snapshot_task = tokio::spawn(async move { handle.snapshot().await });
1533
1534 responder
1535 .stop_after_receiving()
1536 .await
1537 .expect("responder should accept the snapshot request before stopping");
1538
1539 let error = snapshot_task
1540 .await
1541 .expect("snapshot task should complete")
1542 .expect_err("snapshot should fail when the accepted response is dropped");
1543 assert!(error.to_string().contains("stopped before responding"));
1544 }
1545
1546 #[tokio::test]
1547 async fn aggregate_configuration_receiver_can_only_be_taken_once() {
1548 let config = AggregateConfiguration::for_test();
1549 let first = config.build(BuildContext::test_transform("aggregate_one")).await;
1550 assert!(first.is_ok());
1551
1552 let second = config.build(BuildContext::test_transform("aggregate_two")).await;
1553 let error = match second {
1554 Ok(_) => panic!("second build should not take the snapshot receiver again"),
1555 Err(error) => error,
1556 };
1557 assert!(error.to_string().contains("already been taken"));
1558 }
1559
1560 #[tokio::test]
1561 async fn snapshot_responder_reports_closed_request_channel() {
1562 let (handle, mut responder) = aggregate_context_snapshot_channel_for_test();
1563 drop(handle);
1564
1565 let error = responder
1566 .respond(Vec::new())
1567 .await
1568 .expect_err("responder should fail when every request handle is dropped");
1569 assert!(error.to_string().contains("closed"));
1570 }
1571
1572 #[tokio::test]
1573 async fn snapshot_responder_ignores_canceled_requester() {
1574 let (handle, mut responder) = aggregate_context_snapshot_channel_for_test();
1575 let snapshot_task = tokio::spawn(async move { handle.snapshot().await });
1576 let pending_response = responder
1577 .receive()
1578 .await
1579 .expect("responder should accept the snapshot request");
1580
1581 snapshot_task.abort();
1582 let _ = snapshot_task.await;
1583
1584 pending_response.respond(Vec::new());
1585 }
1586
1587 #[test]
1588 fn bucket_is_closed() {
1589 let cases = [
1592 (1000, 995, BUCKET_WIDTH_SECS, false, false),
1594 (1000, 995, BUCKET_WIDTH_SECS, true, true),
1595 (1000, 1000, BUCKET_WIDTH_SECS, false, false),
1597 (1000, 1000, BUCKET_WIDTH_SECS, true, true),
1598 (1010, 1000, BUCKET_WIDTH_SECS, false, true),
1600 (1010, 1000, BUCKET_WIDTH_SECS, true, true),
1601 ];
1602
1603 for (current_time, bucket_start, bucket_width_secs, flush_open_buckets, expected) in cases {
1604 let expected_reason = if expected {
1605 "closed, was open"
1606 } else {
1607 "open, was closed"
1608 };
1609
1610 assert_eq!(
1611 is_bucket_closed(current_time, bucket_start, bucket_width_secs, flush_open_buckets),
1612 expected,
1613 "expected bucket to be {} (current_time={}, bucket_start={}, bucket_width={}, flush_open_buckets={})",
1614 expected_reason,
1615 current_time,
1616 bucket_start,
1617 bucket_width_secs,
1618 flush_open_buckets
1619 );
1620 }
1621 }
1622
1623 #[tokio::test]
1624 async fn context_limit() {
1625 let mut state = AggregationState::new(
1627 BUCKET_WIDTH_SECS,
1628 2,
1629 COUNTER_EXPIRE,
1630 HistogramConfiguration::default(),
1631 Telemetry::noop(),
1632 );
1633
1634 let input_metrics = [
1637 Metric::gauge("metric1", 1.0),
1638 Metric::gauge("metric2", 2.0),
1639 Metric::gauge("metric3", 3.0),
1640 Metric::gauge("metric4", 4.0),
1641 ];
1642
1643 assert!(!state.context_limit_breached());
1644
1645 assert!(state.insert(insert_ts(1), input_metrics[0].clone()));
1646 assert!(state.insert(insert_ts(1), input_metrics[1].clone()));
1647 assert!(!state.context_limit_breached());
1648
1649 assert!(!state.insert(insert_ts(1), input_metrics[2].clone()));
1650 assert!(state.context_limit_breached());
1651 assert!(!state.insert(insert_ts(1), input_metrics[3].clone()));
1652 assert!(state.context_limit_breached());
1653
1654 let flushed_metrics = get_flushed_metrics(flush_ts(1), &mut state).await;
1657 assert_eq!(flushed_metrics.len(), 2);
1658 assert_eq!(input_metrics[0].context(), flushed_metrics[0].context());
1659 assert_eq!(input_metrics[1].context(), flushed_metrics[1].context());
1660 assert!(!state.context_limit_breached());
1661
1662 assert!(state.insert(insert_ts(2), input_metrics[2].clone()));
1665 assert!(state.insert(insert_ts(2), input_metrics[3].clone()));
1666
1667 let flushed_metrics = get_flushed_metrics(flush_ts(2), &mut state).await;
1668 assert_eq!(flushed_metrics.len(), 2);
1669 assert_eq!(input_metrics[2].context(), flushed_metrics[0].context());
1670 assert_eq!(input_metrics[3].context(), flushed_metrics[1].context());
1671 }
1672
1673 #[tokio::test]
1674 async fn context_limit_with_zero_value_counters() {
1675 let mut state = AggregationState::new(
1677 BUCKET_WIDTH_SECS,
1678 2,
1679 COUNTER_EXPIRE,
1680 HistogramConfiguration::default(),
1681 Telemetry::noop(),
1682 );
1683
1684 let input_metrics = [
1686 Metric::counter("metric1", 1.0),
1687 Metric::counter("metric2", 2.0),
1688 Metric::counter("metric3", 3.0),
1689 ];
1690
1691 assert!(state.insert(insert_ts(1), input_metrics[0].clone()));
1692 assert!(state.insert(insert_ts(1), input_metrics[1].clone()));
1693
1694 let flushed_metrics = get_flushed_metrics(flush_ts(1), &mut state).await;
1696 assert_eq!(flushed_metrics.len(), 2);
1697 assert_flushed_scalar_metric!(&input_metrics[0], &flushed_metrics[0], [bucket_ts(1) => 1.0]);
1698 assert_flushed_scalar_metric!(&input_metrics[1], &flushed_metrics[1], [bucket_ts(1) => 2.0]);
1699
1700 let flushed_metrics = get_flushed_metrics(flush_ts(2), &mut state).await;
1702 assert_eq!(flushed_metrics.len(), 2);
1703 assert_flushed_scalar_metric!(&input_metrics[0], &flushed_metrics[0], [bucket_ts(2) => 0.0]);
1704 assert_flushed_scalar_metric!(&input_metrics[1], &flushed_metrics[1], [bucket_ts(2) => 0.0]);
1705
1706 assert!(!state.insert(insert_ts(3), input_metrics[2].clone()));
1708
1709 let flushed_metrics = get_flushed_metrics(flush_ts(3), &mut state).await;
1711 assert_eq!(flushed_metrics.len(), 2);
1712 assert_flushed_scalar_metric!(&input_metrics[0], &flushed_metrics[0], [bucket_ts(3) => 0.0]);
1713 assert_flushed_scalar_metric!(&input_metrics[1], &flushed_metrics[1], [bucket_ts(3) => 0.0]);
1714
1715 let flushed_metrics = get_flushed_metrics(flush_ts(4), &mut state).await;
1718 assert_eq!(flushed_metrics.len(), 0);
1719
1720 assert!(state.insert(insert_ts(5), input_metrics[2].clone()));
1722
1723 let flushed_metrics = get_flushed_metrics(flush_ts(5), &mut state).await;
1724 assert_eq!(flushed_metrics.len(), 1);
1725 assert_flushed_scalar_metric!(&input_metrics[2], &flushed_metrics[0], [bucket_ts(5) => 3.0]);
1726 }
1727
1728 #[tokio::test]
1729 async fn zero_value_counters() {
1730 let mut state = AggregationState::new(
1732 BUCKET_WIDTH_SECS,
1733 10,
1734 COUNTER_EXPIRE,
1735 HistogramConfiguration::default(),
1736 Telemetry::noop(),
1737 );
1738
1739 let input_metrics = [Metric::counter("metric1", 1.0), Metric::counter("metric2", 2.0)];
1741
1742 assert!(state.insert(insert_ts(1), input_metrics[0].clone()));
1743 assert!(state.insert(insert_ts(1), input_metrics[1].clone()));
1744
1745 let flushed_metrics = get_flushed_metrics(flush_ts(1), &mut state).await;
1747 assert_eq!(flushed_metrics.len(), 2);
1748 assert_flushed_scalar_metric!(&input_metrics[0], &flushed_metrics[0], [bucket_ts(1) => 1.0]);
1749 assert_flushed_scalar_metric!(&input_metrics[1], &flushed_metrics[1], [bucket_ts(1) => 2.0]);
1750
1751 let flushed_metrics = get_flushed_metrics(flush_ts(2), &mut state).await;
1753 assert_eq!(flushed_metrics.len(), 2);
1754 assert_flushed_scalar_metric!(&input_metrics[0], &flushed_metrics[0], [bucket_ts(2) => 0.0]);
1755 assert_flushed_scalar_metric!(&input_metrics[1], &flushed_metrics[1], [bucket_ts(2) => 0.0]);
1756
1757 assert!(state.insert(insert_ts(4), input_metrics[0].clone()));
1759 assert!(state.insert(insert_ts(4), input_metrics[1].clone()));
1760
1761 let flushed_metrics = get_flushed_metrics(flush_ts(4), &mut state).await;
1764 assert_eq!(flushed_metrics.len(), 2);
1765 assert_flushed_scalar_metric!(&input_metrics[0], &flushed_metrics[0], [bucket_ts(3) => 0.0, bucket_ts(4) => 1.0]);
1766 assert_flushed_scalar_metric!(&input_metrics[1], &flushed_metrics[1], [bucket_ts(3) => 0.0, bucket_ts(4) => 2.0]);
1767
1768 let flushed_metrics = get_flushed_metrics(flush_ts(7), &mut state).await;
1772 assert_eq!(flushed_metrics.len(), 2);
1773 assert_flushed_scalar_metric!(&input_metrics[0], &flushed_metrics[0], [bucket_ts(5) => 0.0, bucket_ts(6) => 0.0]);
1774 assert_flushed_scalar_metric!(&input_metrics[1], &flushed_metrics[1], [bucket_ts(5) => 0.0, bucket_ts(6) => 0.0]);
1775 }
1776
1777 #[tokio::test]
1778 async fn merge_identical_timestamped_values_on_flush() {
1779 let mut state = AggregationState::new(
1781 BUCKET_WIDTH_SECS,
1782 10,
1783 COUNTER_EXPIRE,
1784 HistogramConfiguration::default(),
1785 Telemetry::noop(),
1786 );
1787
1788 let input_metric = Metric::counter("metric1", [1.0, 2.0, 3.0, 4.0, 5.0]);
1790
1791 assert!(state.insert(insert_ts(1), input_metric.clone()));
1792
1793 let flushed_metrics = get_flushed_metrics(flush_ts(1), &mut state).await;
1796 assert_eq!(flushed_metrics.len(), 1);
1797 assert_flushed_scalar_metric!(&input_metric, &flushed_metrics[0], [bucket_ts(1) => 15.0]);
1798 }
1799
1800 #[tokio::test]
1801 async fn histogram_statistics() {
1802 let hist_config = HistogramConfiguration::from_statistics(
1804 &[
1805 HistogramStatistic::Count,
1806 HistogramStatistic::Sum,
1807 HistogramStatistic::Percentile {
1808 q: 0.5,
1809 suffix: "p50".into(),
1810 },
1811 ],
1812 false,
1813 "".into(),
1814 );
1815 let mut state = AggregationState::new(BUCKET_WIDTH_SECS, 10, COUNTER_EXPIRE, hist_config, Telemetry::noop());
1816
1817 let input_metric = Metric::histogram("metric1", [1.0, 2.0, 3.0, 4.0, 5.0]);
1819 assert!(state.insert(insert_ts(1), input_metric.clone()));
1820
1821 let flushed_metrics = get_flushed_metrics(flush_ts(1), &mut state).await;
1824 assert_eq!(flushed_metrics.len(), 3);
1825
1826 let count_metric = Metric::rate("metric1.count", 0.0, Duration::from_secs(BUCKET_WIDTH_SECS.get()));
1829 let sum_metric = Metric::gauge("metric1.sum", 0.0);
1830 let p50_metric = Metric::gauge("metric1.p50", 0.0);
1831
1832 assert_flushed_scalar_metric!(count_metric, &flushed_metrics[0], [bucket_ts(1) => 5.0]);
1835 assert_flushed_scalar_metric!(p50_metric, &flushed_metrics[1], [bucket_ts(1) => 3.0], error_ratio => 0.0025);
1836 assert_flushed_scalar_metric!(sum_metric, &flushed_metrics[2], [bucket_ts(1) => 15.0]);
1837 }
1838
1839 #[tokio::test]
1840 async fn histogram_statistics_unit_propagation() {
1841 let hist_config = HistogramConfiguration::from_statistics(
1843 &[
1844 HistogramStatistic::Count,
1845 HistogramStatistic::Sum,
1846 HistogramStatistic::Percentile {
1847 q: 0.5,
1848 suffix: "p50".into(),
1849 },
1850 ],
1851 false,
1852 "".into(),
1853 );
1854 let mut state = AggregationState::new(BUCKET_WIDTH_SECS, 10, COUNTER_EXPIRE, hist_config, Telemetry::noop());
1855
1856 let context = Context::from_static_parts("metric1", &[]);
1858 let metadata = MetricMetadata::default().with_unit(MetaString::from_static("millisecond"));
1859 let input_metric = Metric::from_parts(
1860 context,
1861 MetricValues::histogram([1.0_f64, 2.0, 3.0, 4.0, 5.0]),
1862 metadata,
1863 );
1864 assert!(state.insert(insert_ts(1), input_metric));
1865
1866 let flushed_metrics = get_flushed_metrics(flush_ts(1), &mut state).await;
1867 assert_eq!(flushed_metrics.len(), 3);
1868
1869 for metric in &flushed_metrics {
1872 let name = metric.context().name();
1873 if name.ends_with(".count") {
1874 assert_eq!(
1875 metric.metadata().unit(),
1876 None,
1877 "flushed metric '{}' should be dimensionless",
1878 name
1879 );
1880 } else {
1881 assert_eq!(
1882 metric.metadata().unit(),
1883 Some("millisecond"),
1884 "flushed metric '{}' should carry unit='millisecond'",
1885 name
1886 );
1887 }
1888 }
1889 }
1890
1891 #[tokio::test]
1892 async fn distributions() {
1893 let mut state = AggregationState::new(
1895 BUCKET_WIDTH_SECS,
1896 10,
1897 COUNTER_EXPIRE,
1898 HistogramConfiguration::default(),
1899 Telemetry::noop(),
1900 );
1901
1902 let values = [1.0, 2.0, 3.0, 4.0, 5.0];
1904 let input_metric = Metric::distribution("metric1", &values[..]);
1905
1906 assert!(state.insert(insert_ts(1), input_metric.clone()));
1907
1908 let flushed_metrics = get_flushed_metrics(flush_ts(1), &mut state).await;
1910 assert_eq!(flushed_metrics.len(), 1);
1911
1912 assert_flushed_distribution_metric!(&input_metric, &flushed_metrics[0], [bucket_ts(1) => &values[..]]);
1913 }
1914
1915 #[tokio::test]
1916 async fn histogram_copy_to_distribution() {
1917 let hist_config = HistogramConfiguration::from_statistics(
1918 &[
1919 HistogramStatistic::Count,
1920 HistogramStatistic::Sum,
1921 HistogramStatistic::Percentile {
1922 q: 0.5,
1923 suffix: "p50".into(),
1924 },
1925 ],
1926 true,
1927 "dist_prefix.".into(),
1928 );
1929 let mut state = AggregationState::new(BUCKET_WIDTH_SECS, 10, COUNTER_EXPIRE, hist_config, Telemetry::noop());
1930
1931 let values = [1.0, 2.0, 3.0, 4.0, 5.0];
1933 let input_metric = Metric::histogram("metric1", values);
1934 assert!(state.insert(insert_ts(1), input_metric.clone()));
1935
1936 let flushed_metrics = get_flushed_metrics(flush_ts(1), &mut state).await;
1939 assert_eq!(flushed_metrics.len(), 4);
1940
1941 let count_metric = Metric::rate("metric1.count", 0.0, BUCKET_WIDTH);
1944 let sum_metric = Metric::gauge("metric1.sum", 0.0);
1945 let p50_metric = Metric::gauge("metric1.p50", 0.0);
1946 let expected_distribution = Metric::distribution("dist_prefix.metric1", &values[..]);
1947
1948 assert_flushed_distribution_metric!(expected_distribution, &flushed_metrics[0], [bucket_ts(1) => &values[..]]);
1951 assert_flushed_scalar_metric!(count_metric, &flushed_metrics[1], [bucket_ts(1) => 5.0]);
1952 assert_flushed_scalar_metric!(p50_metric, &flushed_metrics[2], [bucket_ts(1) => 3.0], error_ratio => 0.0025);
1953 assert_flushed_scalar_metric!(sum_metric, &flushed_metrics[3], [bucket_ts(1) => 15.0]);
1954 }
1955
1956 #[tokio::test]
1957 async fn nonaggregated_counters_to_rate() {
1958 let counter_value = 42.0;
1959
1960 let mut state = AggregationState::new(
1962 BUCKET_WIDTH_SECS,
1963 10,
1964 COUNTER_EXPIRE,
1965 HistogramConfiguration::default(),
1966 Telemetry::noop(),
1967 );
1968
1969 let input_metric = Metric::counter("metric1", counter_value);
1971 assert!(state.insert(insert_ts(1), input_metric.clone()));
1972
1973 let flushed_metrics = get_flushed_metrics(flush_ts(1), &mut state).await;
1976 assert_eq!(flushed_metrics.len(), 1);
1977 let flushed_metric = &flushed_metrics[0];
1978
1979 assert_flushed_scalar_metric!(&input_metric, flushed_metric, [bucket_ts(1) => counter_value]);
1980 assert_eq!(flushed_metric.values().as_str(), "rate");
1981 }
1982
1983 #[tokio::test]
1984 async fn telemetry() {
1985 let recorder = TestRecorder::default();
1992 let _local = metrics::set_default_local_recorder(&recorder);
1993
1994 let builder = MetricsBuilder::default();
1995 let telemetry = Telemetry::new(&builder);
1996
1997 let mut state = AggregationState::new(
1998 BUCKET_WIDTH_SECS,
1999 2,
2000 COUNTER_EXPIRE,
2001 HistogramConfiguration::default(),
2002 telemetry,
2003 );
2004
2005 assert_eq!(recorder.gauge("aggregate_active_contexts"), Some(0.0));
2007 assert_eq!(
2008 recorder.counter(("component_events_dropped_total", &[("intentional", "true")])),
2009 Some(0)
2010 );
2011 for metric_type in &["counter", "gauge", "rate", "set", "histogram", "distribution"] {
2012 assert_eq!(
2013 recorder.gauge(("aggregate_active_contexts_by_type", &[("metric_type", *metric_type)])),
2014 Some(0.0)
2015 );
2016 }
2017
2018 assert!(state.insert(insert_ts(1), Metric::counter("metric1", 42.0)));
2020 assert_eq!(recorder.gauge("aggregate_active_contexts"), Some(1.0));
2021 assert_eq!(
2022 recorder.gauge(("aggregate_active_contexts_by_type", &[("metric_type", "counter")])),
2023 Some(1.0)
2024 );
2025
2026 assert!(state.insert(insert_ts(1), Metric::gauge("metric2", (insert_ts(1), 42.0))));
2028 assert_eq!(recorder.gauge("aggregate_active_contexts"), Some(2.0));
2029 assert_eq!(
2030 recorder.gauge(("aggregate_active_contexts_by_type", &[("metric_type", "gauge")])),
2031 Some(1.0)
2032 );
2033
2034 assert!(!state.insert(insert_ts(1), Metric::counter("metric3", 42.0)));
2036 assert_eq!(recorder.gauge("aggregate_active_contexts"), Some(2.0));
2037 assert_eq!(
2038 recorder.gauge(("aggregate_active_contexts_by_type", &[("metric_type", "counter")])),
2039 Some(1.0)
2040 );
2041
2042 let _ = get_flushed_metrics(flush_ts(1), &mut state).await;
2045 assert_eq!(recorder.gauge("aggregate_active_contexts"), Some(1.0));
2046 assert_eq!(
2047 recorder.gauge(("aggregate_active_contexts_by_type", &[("metric_type", "counter")])),
2048 Some(1.0)
2049 );
2050 assert_eq!(
2051 recorder.gauge(("aggregate_active_contexts_by_type", &[("metric_type", "gauge")])),
2052 Some(0.0)
2053 );
2054 }
2055}