saluki_components/transforms/aggregate/
mod.rs

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/// The shape of metric values retained by the aggregate transform.
38#[derive(Clone, Copy, Debug, Eq, PartialEq)]
39pub enum AggregateMetricType {
40    /// Counter values.
41    Counter,
42
43    /// Rate values.
44    Rate,
45
46    /// Gauge values.
47    Gauge,
48
49    /// Set values.
50    Set,
51
52    /// Histogram values.
53    Histogram,
54
55    /// Distribution values.
56    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/// A retained metric context and its small aggregation metadata.
73///
74/// Cloning an entry shares the underlying context name and tags rather than copying their contents.
75#[derive(Clone, Debug, Eq, PartialEq)]
76pub struct AggregateContextSnapshotEntry {
77    context: Context,
78    metric_type: AggregateMetricType,
79    unit: MetaString,
80}
81
82impl AggregateContextSnapshotEntry {
83    /// Returns the retained metric context.
84    pub fn context(&self) -> &Context {
85        &self.context
86    }
87
88    /// Returns the shape of the retained metric values.
89    pub fn metric_type(&self) -> AggregateMetricType {
90        self.metric_type
91    }
92
93    /// Returns the unit attached to the retained metric values, if one is set.
94    pub fn unit(&self) -> Option<&str> {
95        if self.unit.is_empty() {
96            None
97        } else {
98            Some(&self.unit)
99        }
100    }
101
102    /// Creates a snapshot entry for downstream test and benchmark fixtures.
103    #[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/// A handle for requesting retained-context snapshots from an aggregate transform.
118///
119/// Snapshot construction runs on the aggregate owner task. The returned entries contain shared context handles and
120/// small metric metadata, allowing callers to perform heavier processing after the owner resumes ingestion.
121#[derive(Clone, Debug)]
122pub struct AggregateContextSnapshotHandle {
123    requests: mpsc::Sender<AggregateContextSnapshotRequest>,
124}
125
126impl AggregateContextSnapshotHandle {
127    /// Requests the aggregate transform's current retained contexts.
128    ///
129    /// The returned snapshot uses O(context count) memory. After delivery, the caller owns this memory, so callers that
130    /// retain snapshots must include their retained size in their own memory accounting.
131    ///
132    /// # Errors
133    ///
134    /// Returns an error if the aggregate owner is unavailable or stops before responding.
135    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
148/// Creates a retained-context snapshot handle and its owner-side receiver.
149///
150/// The receiver belongs to the aggregate transform, and is supplied to it through
151/// [`AggregateConfiguration::context_snapshot_receiver`]. The handle belongs to whoever requests snapshots, and can be
152/// cloned to let several callers request them.
153pub 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/// An accepted retained-context snapshot request for test fixtures.
175#[cfg(any(test, feature = "test-util"))]
176pub struct AggregateContextSnapshotPendingResponse {
177    response: AggregateContextSnapshotRequest,
178}
179
180#[cfg(any(test, feature = "test-util"))]
181impl AggregateContextSnapshotPendingResponse {
182    /// Responds to the accepted snapshot request with the supplied entries.
183    ///
184    /// If the requester was canceled after the request was accepted, the response is discarded.
185    pub fn respond(self, snapshot: Vec<AggregateContextSnapshotEntry>) {
186        let _ = self.response.send(snapshot);
187    }
188}
189
190/// A responder for retained-context snapshot test fixtures.
191#[cfg(any(test, feature = "test-util"))]
192pub struct AggregateContextSnapshotResponder {
193    receiver: AggregateContextSnapshotRequestReceiver,
194}
195
196#[cfg(any(test, feature = "test-util"))]
197impl AggregateContextSnapshotResponder {
198    /// Waits for one snapshot request and responds with the supplied entries.
199    ///
200    /// A canceled requester is treated as a successful no-op delivery.
201    ///
202    /// # Errors
203    ///
204    /// Returns an error if the request channel closes before a request arrives.
205    pub async fn respond(&mut self, snapshot: Vec<AggregateContextSnapshotEntry>) -> Result<(), GenericError> {
206        self.receive().await?.respond(snapshot);
207        Ok(())
208    }
209
210    /// Waits for one snapshot request and returns its pending response.
211    ///
212    /// The returned response lets tests deterministically control whether the owner responds, stops, or outlives a
213    /// canceled requester after accepting the request.
214    ///
215    /// # Errors
216    ///
217    /// Returns an error if the request channel closes before a request arrives.
218    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    /// Waits for one snapshot request and stops without responding.
228    ///
229    /// This accepts the request from the owner channel before dropping its one-shot response sender, allowing tests to
230    /// distinguish an owner that stops mid-request from an owner whose request channel is unavailable.
231    ///
232    /// # Errors
233    ///
234    /// Returns an error if the request channel closes before a request arrives.
235    pub async fn stop_after_receiving(&mut self) -> Result<(), GenericError> {
236        drop(self.receive().await?);
237        Ok(())
238    }
239}
240
241/// Creates a retained-context snapshot handle and owner-side responder for tests.
242#[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
252/// The owner side of a retained-context snapshot channel.
253///
254/// The aggregate transform takes the receiver out of this holder when it is built. A receiver can only be taken once,
255/// so an [`AggregateConfiguration`] holding one can only build a single transform.
256pub 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
272/// Aggregate transform.
273///
274/// Aggregates metrics into fixed-size windows, flushing them at a regular interval.
275///
276/// ## Zero-value counters
277///
278/// When metrics are aggregated and then flushed, they're typically removed entirely from the aggregation state. Unless
279/// they're updated again, they won't be emitted again. However, for counters, a slightly different approach is
280/// taken by tracking "zero-value" counters.
281///
282/// Counters are aggregated and flushed normally. However, when flushed, counters are added to a list of "zero-value"
283/// counters, and if those counters aren't updated again, the transform emits a copy of the counter with a value of
284/// zero. It does this until the counter is updated again, or the zero-value counter expires (no updates), whichever
285/// comes first.
286///
287/// This provides a continuity in the output of a counter, from the perspective of a downstream system, when counters
288/// are otherwise sparse. The expiration period is configurable, and allows a trade-off in how sparse/infrequent the
289/// updates to counters can be versus how long it takes for counters that don't exist anymore to actually cease to be
290/// emitted.
291pub struct AggregateConfiguration {
292    /// Size of the aggregation window, in seconds.
293    ///
294    /// Metrics are aggregated into fixed-size windows, such that all updates to the same metric within a window are
295    /// aggregated into a single metric. The window size controls how efficiently metrics are aggregated, and in turn,
296    /// how many data points are emitted downstream.
297    pub window_duration_seconds: NonZeroU64,
298
299    /// How often to flush buckets.
300    ///
301    /// This represents a trade-off between the savings in network bandwidth (sending fewer requests to downstream
302    /// systems, etc) and the frequency of updates (how often updates to a metric are emitted).
303    pub primary_flush_interval: Duration,
304
305    /// Maximum number of contexts to aggregate per window.
306    ///
307    /// A context is the unique combination of a metric name and its set of tags. For example,
308    /// `metric.name.here{tag1=A,tag2=B}` represents a single context, and would be different than
309    /// `metric.name.here{tag1=A,tag2=C}`.
310    ///
311    /// When the maximum number of contexts is reached in the current aggregation window, additional metrics are dropped
312    /// until the next window starts.
313    pub context_limit: usize,
314
315    /// Whether to flush open buckets when stopping the transform.
316    ///
317    /// Normally, open buckets (a bucket whose end hasn't yet occurred) aren't flushed when the transform is stopped.
318    /// This is done to avoid the chance of flushing a partial window, restarting the process, and then flushing the
319    /// same window again. Downstream systems sometimes can't cope with this gracefully, as there is no way to
320    /// determine that it's an incremental update, and so they treat it as an absolute update, overwriting the
321    /// previously flushed value.
322    ///
323    /// In cases where flushing all outstanding data is paramount, this can be enabled.
324    pub flush_open_windows: bool,
325
326    /// How long to keep idle counters alive after they've been flushed, in seconds.
327    ///
328    /// When metrics are flushed, they're removed from the aggregation state. However, if a counter expiration is set,
329    /// counters will be kept alive in an "idle" state. For as long as a counter is idle, but not yet expired, a zero
330    /// value will be emitted for it during each flush. This allows more gracefully handling sparse counters, where
331    /// updates are infrequent but leaving gaps in the time series would be undesirable from a user experience
332    /// perspective.
333    ///
334    /// After a counter has been idle (no updates) for longer than the expiry period, it will be completely removed and
335    /// no further zero values will be emitted.
336    ///
337    /// A value of `0`, or `None`, disables idle counter keep-alive.
338    pub counter_expiry_seconds: Option<u64>,
339
340    /// Statistics to calculate over histograms, and how to copy them to distributions.
341    pub hist_config: HistogramConfiguration,
342
343    /// Owner side of the channel used to serve retained-context snapshot requests.
344    ///
345    /// This is runtime wiring rather than configuration: it carries no settings, and is created by the caller with
346    /// [`aggregate_context_snapshot_channel`] so that the caller keeps the matching handle.
347    pub context_snapshot_receiver: AggregateContextSnapshotReceiver,
348}
349
350#[cfg(test)]
351impl AggregateConfiguration {
352    /// Creates a fixture configuration, for tests that exercise aggregation behavior rather than configuration.
353    ///
354    /// The snapshot handle paired with the configuration's receiver is dropped. Tests that request snapshots build the
355    /// channel with [`aggregate_context_snapshot_channel`] and override `context_snapshot_receiver`.
356    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        // TODO: While we account for the aggregation state map accurately, what we don't currently account for is the
408        // fact that a metric could have multiple distinct values. For the common pipeline of metrics in via DogStatsD,
409        // this generally shouldn't be a problem because the values don't have a timestamp, so they get aggregated into
410        // the same bucket, leading to two values per `MetricValues` at most, which is already baked into the size of
411        // `MetricValues` due to using `SmallVec`.
412        //
413        // However, there could be many more values in a single metric, and we don't account for that.
414
415        builder
416            .minimum()
417            // Capture the size of the heap allocation when the component is built.
418            .with_single_value::<Aggregate>("component struct");
419        builder
420            .firm()
421            // Account for the aggregation state map, where we map contexts to the merged metric.
422            .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            // A snapshot is constructed while the aggregation state remains live, so its peak allocation is additive.
432            .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                    // We've reached the end of the current window. Flush our aggregation state and forward the metrics
467                    // onwards. Regardless of whether any metrics were aggregated, we always update the aggregation
468                    // state to track the start time of the current aggregation window.
469                    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                        // Remember if the context limit had been surpassed before this flush.
475                        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 flush recovered us from a breach, log the recovery.
486                            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                        // A flush can wait a long time for space in a full downstream channel. Keep answering liveness
498                        // probes while it waits. Probes go first so that a waiting probe is seen on every wakeup, even
499                        // when the flush would use up the task's cooperative budget.
500                        loop {
501                            select! {
502                                biased;
503                                _ = health.live() => {},
504                                _ = &mut flush => break,
505                            }
506                        }
507                    }
508
509                    // If this is the final flush, we break out of the loop.
510                    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                                        // First drop since the last recovery — emit a single warning.
536                                        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                        // We've reached the end of our input stream, so mark ourselves for a final flush and reset the
546                        // interval so it ticks immediately on the next loop iteration.
547                        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    /// Tracks whether the context limit has been breached. Starts out as `false`. Set to `true` on the first dropped
588    /// metric. Reset to `false` when the context count drops below the limit during flush.
589    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 we haven't seen this context yet, and it would put us over the limit to insert it, then return early.
630        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        // Collapse all non-timestamped values into a single timestamped value.
638        //
639        // We do this pre-aggregation step because unless we're merging into an existing context, we'll end up with
640        // however many values were in the original metric instead of full aggregated values.
641        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        // If we're already tracking this context, update the last seen time and merge the new values into the existing
651        // values. Otherwise, create a new entry.
652        match self.contexts.entry(context) {
653            Entry::Occupied(mut entry) => {
654                let aggregated = entry.get_mut();
655
656                // We ignore metadata changes within a flush interval to keep things simple.
657                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        // We want our split timestamp to be before the start of the current bucket, which ensures any timestamp that is
688        // less than or equal to `split_timestamp` resides in a closed bucket.
689        let split_timestamp = align_to_bucket_start(current_time, bucket_width_secs).saturating_sub(1);
690
691        // Calculate the buckets we need to potentially generate zero-value counters for.
692        //
693        // We only need to do this if we've flushed before, since we won't have any knowledge of which counters are idle
694        // or not until that happens.
695        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            // Clock-skew guards. Bucketing reads the wall clock while the flush cadence is monotonic, so a wall-clock
700            // jump is not bounded by the flush interval. A backward jump empties the zero-value range (a silent counter
701            // gap); a forward jump makes the loop below run once per bucket across the whole jumped span — O(jump) work
702            // and allocation. Assert before the loop so a flood fails fast rather than after the damage is done.
703            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            // The 10_000 bound is generous. A default 15s flush over a 10s bucket yields one or two buckets. The bound
710            // trips only on a multi-hour wall-clock jump, never on a slow-but-sane flush.
711            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            // Anti-vacuity anchor: prove the idle-counter zero-value path actually runs in some timeline.
729            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        // Iterate over each context we're tracking, and flush any values that are in buckets which are now closed.
737        debug!(timestamp = current_time, "Flushing buckets.");
738
739        for (context, am) in self.contexts.iter_mut() {
740            // Figure out if we should remove this metric or not if it has no values in open buckets.
741            //
742            // We have a special carve-out for counters here, which we have the ability to keep alive after they are
743            // flushed, based on a configured expiration period. This allows us to continue emitting a zero value for
744            // counters when they're idle, which can make them appear "live" in downstream systems, even when they're
745            // not.
746            //
747            // This is useful for sparsely-updated counters.
748            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 we're dealing with a counter, we'll merge in our calculated set of zero values. We only merge in the
762            // values that represent now-closed buckets.
763            //
764            // This is also safe to do even when there are real values in those buckets since adding zero to anything is
765            // a no-op from the perspective of what we end up flushing, and it doesn't mess with the "last seen" time.
766            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                        // Since zero-value buckets are in order, we can break early if this bucket is past the
773                        // expiration cutoff of the counter. No other bucket will be within the expiration range.
774                        break;
775                    }
776                }
777            }
778
779            // Finally, figure out if the current metric can be removed.
780            //
781            // For any metric with values that are in open buckets, we split off the values that are in closed buckets
782            // and keep the metric alive. When all the values are in closed buckets, or there are no values, we'll
783            // remove the metric if `should_remove_if_empty` is `true`.
784            //
785            // This means we'll always remove all-closed/empty non-counter metrics, and we _may_ remove all-closed/empty
786            // counters.
787            if let Some(closed_bucket_values) = am.values.split_at_timestamp(split_timestamp) {
788                self.telemetry.increment_flushed(&closed_bucket_values);
789
790                // We got some closed bucket values, so flush those out.
791                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        // Remove any contexts that were marked as needing to be removed.
809        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        // If we're dealing with a histogram, we calculate a configured set of aggregates/percentiles from it, and emit
841        // them as individual metrics.
842        MetricValues::Histogram(ref mut points) => {
843            // Convert histogram to distribution
844            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            // We collect our histogram points in their "summary" view, which sorts the underlying samples allowing
869            // proper quantile queries to be answered, hence our "sorted" points. We do it this way because rather than
870            // sort every time we insert, or cloning the points, we only sort when a summary view is constructed, which
871            // requires mutable access to sort the samples in-place.
872            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                // Counts are dimensionless, so clear any unit inherited from the input histogram.
890                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        // If we're not dealing with a histogram, then all we need to worry about is converting counters to rates before
905        // forwarding our single, aggregated metric.
906        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    // A bucket is considered "closed" if the current time is greater than the end of the bucket, or if
930    // `flush_open_buckets` is `true`.
931    //
932    // Buckets represent a half-open interval, where the start is inclusive and the end is exclusive. This means that
933    // for a bucket start of 10, and a width of 10, the bucket is 10 seconds "wide", and its start and end are 10 and
934    // 20, with the 20 excluded, or [10, 20) in interval notation. Simply put, if we have a timestamp of 10, or anything
935    // smaller than 20, we would consider it to fall within the bucket... but 20 or more would be outside of the bucket.
936    //
937    // We can also represent this visually:
938    //
939    // <--------- bucket 1 ----------> <--------- bucket 2 ----------> <--------- bucket 3 ---------->
940    // [10 11 12 13 14 15 16 17 18 19] [20 21 22 23 24 25 26 27 28 29] [30 31 32 33 34 35 36 37 38 39]
941    //
942    // We can see that each bucket is 10 seconds wide (10 elements, one for each second), and that their ends are
943    // effectively `start + width - 1`. This means that for any of these buckets to be considered "closed", the current
944    // time has to be _greater_ than `start + width - 1`. For example, if the current time is 19, then no buckets are
945    // closed, and if the current time is 29, then bucket 1 is closed but buckets 2 and 3 are still open, and if the
946    // current time is 30, then both buckets 1 and 2 are closed, but bucket 3 is still open.
947    (bucket_start + bucket_width_secs.get() - 1) < current_time || flush_open_buckets
948}
949
950// TODO: One thing we ought to consider is a property test, specifically a state machine property test, where we
951// generate a randomized offset to start time from, a bucket width, flush interval, and operations, and so on... and
952// then we run it to make sure that we are always generating sequential timestamps for data points, etc.
953#[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    /// Gets the bucket start timestamp for the given step.
984    const fn bucket_ts(step: u64) -> u64 {
985        align_to_bucket_start(insert_ts(step), BUCKET_WIDTH_SECS)
986    }
987
988    /// Gets the insert timestamp for the given step.
989    const fn insert_ts(step: u64) -> u64 {
990        (BUCKET_WIDTH_SECS.get() * (step + 1)) - 2
991    }
992
993    /// Gets the flush timestamp for the given step.
994    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    /// Constructs a basic `Dispatcher` with a fixed-size event buffer.
1097    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        // Flush the metrics to an event buffer.
1115        state
1116            .flush(timestamp, true, &mut buffered_dispatcher)
1117            .await
1118            .expect("should not fail to flush aggregation state");
1119
1120        // Flush our buffered dispatcher, which should ensure that the event buffer is sent out, and then read it from the
1121        // receiver:
1122        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        // Cases are defined as:
1590        // (current time, bucket start, bucket width, flush open buckets, expected result)
1591        let cases = [
1592            // Bucket goes from [995, 1005), current time of 1000, so bucket is open.
1593            (1000, 995, BUCKET_WIDTH_SECS, false, false),
1594            (1000, 995, BUCKET_WIDTH_SECS, true, true),
1595            // Bucket goes from [1000, 1010), current time of 1000, so bucket is open.
1596            (1000, 1000, BUCKET_WIDTH_SECS, false, false),
1597            (1000, 1000, BUCKET_WIDTH_SECS, true, true),
1598            // Bucket goes from [1000, 1010), current time of 1010, so bucket is closed.
1599            (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        // Create our aggregation state with a context limit of 2.
1626        let mut state = AggregationState::new(
1627            BUCKET_WIDTH_SECS,
1628            2,
1629            COUNTER_EXPIRE,
1630            HistogramConfiguration::default(),
1631            Telemetry::noop(),
1632        );
1633
1634        // Create four unique gauges, and insert all of them. The third and fourth should fail because we've reached
1635        // the context limit.
1636        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        // We should only see the first two gauges after flushing. The flush should also clear the breached flag since
1655        // contexts drop below the limit.
1656        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        // We should be able to insert the third and fourth gauges now as the first two have been flushed, and along
1663        // with them, their contexts should no longer be tracked in the aggregation state:
1664        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        // We test here to ensure that zero-value counters contribute to the context limit.
1676        let mut state = AggregationState::new(
1677            BUCKET_WIDTH_SECS,
1678            2,
1679            COUNTER_EXPIRE,
1680            HistogramConfiguration::default(),
1681            Telemetry::noop(),
1682        );
1683
1684        // Create our input metrics.
1685        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        // Flush the aggregation state, and observe they're both present.
1695        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        // Flush _again_ to ensure that we then emit zero-value variants for both counters.
1701        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        // Now try to insert a third counter, which should fail because we've reached the context limit.
1707        assert!(!state.insert(insert_ts(3), input_metrics[2].clone()));
1708
1709        // Flush the aggregation state, and observe that we only see the two original counters.
1710        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        // With a fourth flush interval, the two counters should now have expired, and thus be dropped and no longer
1716        // contributing to the context limit.
1717        let flushed_metrics = get_flushed_metrics(flush_ts(4), &mut state).await;
1718        assert_eq!(flushed_metrics.len(), 0);
1719
1720        // Now we should be able to insert the third counter, and it should be the only one present after flushing.
1721        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        // We're testing that we properly emit and expire zero-value counters in all relevant scenarios.
1731        let mut state = AggregationState::new(
1732            BUCKET_WIDTH_SECS,
1733            10,
1734            COUNTER_EXPIRE,
1735            HistogramConfiguration::default(),
1736            Telemetry::noop(),
1737        );
1738
1739        // Create two unique counters, and insert both of them.
1740        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        // Flush the aggregation state, and observe they're both present.
1746        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        // Perform our second flush, which should have them as zero-value counters.
1752        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        // Now, we'll pretend to skip a flush period and add updates to them again after that.
1758        assert!(state.insert(insert_ts(4), input_metrics[0].clone()));
1759        assert!(state.insert(insert_ts(4), input_metrics[1].clone()));
1760
1761        // Flush the aggregation state, and observe that we have two zero-value counters for the flush period we
1762        // skipped, but that we see them appear again in the fourth flush period.
1763        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        // Now we'll skip multiple flush periods and ensure that we emit zero-value counters up until the point they
1769        // expire. As our zero-value counter expiration is 20 seconds, this is two flush periods, so we skip by three
1770        // flush periods, and we should only see the counters emitted for the first two.
1771        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        // We're testing that we properly emit and expire zero-value counters in all relevant scenarios.
1780        let mut state = AggregationState::new(
1781            BUCKET_WIDTH_SECS,
1782            10,
1783            COUNTER_EXPIRE,
1784            HistogramConfiguration::default(),
1785            Telemetry::noop(),
1786        );
1787
1788        // Create one multi-value counter, and insert it.
1789        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        // Flush the aggregation state, and observe the metric is present _and_ that we've properly merged all of the
1794        // values within the same timestamp.
1795        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        // We're testing that we properly emit individual metrics (min, max, sum, etc) for a histogram.
1803        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        // Create one multi-value histogram and insert it.
1818        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        // Flush the aggregation state, and observe that we've emitted all of the configured distribution statistics in
1822        // the form of three metrics: count, sum, and p50.
1823        let flushed_metrics = get_flushed_metrics(flush_ts(1), &mut state).await;
1824        assert_eq!(flushed_metrics.len(), 3);
1825
1826        // Create versions of the metric for each of the statistics we're expecting to emit. The values themselves don't
1827        // matter here, but we do need a `Metric` for it to compare the context to.
1828        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        // We use a less strict error ratio (how much the expected vs actual) for the percentile check, as we generally
1833        // expect the value to be somewhat off the exact value due to the lossy nature of `DDSketch`.
1834        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        // We're testing that the unit from the input histogram metadata propagates to all flushed output metrics.
1842        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        // Build a histogram with unit = "millisecond", simulating what arrives from a DogStatsD `ms` metric.
1857        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        // Counts are dimensionless: the `.count` series must drop the unit while all other
1870        // aggregate series carry the unit from the input histogram.
1871        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        // We're testing that we pass through distributions untouched.
1894        let mut state = AggregationState::new(
1895            BUCKET_WIDTH_SECS,
1896            10,
1897            COUNTER_EXPIRE,
1898            HistogramConfiguration::default(),
1899            Telemetry::noop(),
1900        );
1901
1902        // Create one multi-value distribution, with server-side aggregation, and insert it.
1903        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        // Flush the aggregation state, and observe that we've emitted the original distribution.
1909        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        // Create one multi-value histogram and insert it.
1932        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        // Flush the aggregation state, and observe that we've emitted all of the configured distribution statistics in
1937        // the form of three metrics: count, sum, and p50 as well as the additional metric from copying the histogram.
1938        let flushed_metrics = get_flushed_metrics(flush_ts(1), &mut state).await;
1939        assert_eq!(flushed_metrics.len(), 4);
1940
1941        // Create versions of the metric for each of the statistics we're expecting to emit. The values themselves don't
1942        // matter here, but we do need a `Metric` for it to compare the context to.
1943        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        // We use a less strict error ratio (how much the expected vs actual) for the percentile check, as we generally
1949        // expect the value to be somewhat off the exact value due to the lossy nature of `DDSketch`.
1950        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        // Create a basic aggregation state.
1961        let mut state = AggregationState::new(
1962            BUCKET_WIDTH_SECS,
1963            10,
1964            COUNTER_EXPIRE,
1965            HistogramConfiguration::default(),
1966            Telemetry::noop(),
1967        );
1968
1969        // Create a simple non-aggregated counter, and insert it.
1970        let input_metric = Metric::counter("metric1", counter_value);
1971        assert!(state.insert(insert_ts(1), input_metric.clone()));
1972
1973        // Flush the aggregation state, and observe that we've emitted the expected counter and that it has the right
1974        // value, but specifically that it's a rate with an interval that matches our configured bucket width:
1975        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        // TODO: We don't check `component_events_dropped_total` here as it's set directly in the aggregate
1986        // component future rather than `AggregationState`, which is harder to drive overall and would have
1987        // required even more boilerplate.
1988        //
1989        // Leaving that as a future improvement.
1990
1991        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        // Make sure our telemetry is registered at default values.
2006        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        // Insert a counter with a non-timestamped value.
2019        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        // Insert a gauge with a timestamped value.
2027        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        // We've reached our context limit at this point, so the next metric should not be inserted.
2035        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        // Now let's flush the state which should flush the gauge entirely, reducing the context count, but not flush
2043        // the counter, since it'll be in zero-value mode.
2044        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}