Skip to main content

teksilo_data/
chart_aggregate.rs

1// SPDX-License-Identifier: MPL-2.0
2// SPDX-FileCopyrightText: 2026 FernTech
3
4//! `ChartAggregate<T>` — a bucket/rollup projection over a [`crate::ChartModel`].
5//!
6//! Wraps a [`ChartModel<T>`](crate::ChartModel) and exposes each series
7//! reduced into fixed-size buckets of `bucket_size` source points, each
8//! bucket collapsed to one [`crate::ChartDatum`] via a [`ChartAggregateFn`]
9//! (`Mean` / `Sum` / `Min` / `Max` / `First` / `Last` / `Custom`) — the
10//! "downsample a long series for display" pattern (a year of daily
11//! sensor readings shown as weekly means, a tick feed shown as 1-minute
12//! bars). Bucket `b` covers source indices `[b*bucket_size,
13//! min((b+1)*bucket_size, n))`; a trailing partial bucket is included. A
14//! bucket's category is its first member's category.
15//!
16//! Unlike [`crate::ChartWindow`] (which reads straight through to the
17//! source), `ChartAggregate` **materializes** its buckets — a bucket's
18//! category is a *clone* of a source point's category, so constructing or
19//! rebuilding a `ChartAggregate<T>` requires `T: Clone`. Once built,
20//! read-only queries (`point_count`, `with_point`, …) need only `T:
21//! 'static`.
22//!
23//! ## Reactivity
24//!
25//! A tail append that doesn't change the bucket count updates the
26//! now-not-yet-full last bucket in place (`PointUpdated`); a tail append
27//! that starts a new bucket finalizes the previous last bucket
28//! (`PointUpdated`) and appends the new one(s) (`PointsInserted`).
29//! Symmetrically, a **tail removal** that doesn't eliminate the last bucket
30//! recomputes it in place (`PointUpdated`, since it lost some of its
31//! points); one that eliminates one or more trailing buckets recomputes the
32//! new last bucket the same way and then drops the buckets beyond it
33//! (`PointsRemoved`). A mid-series insert or removal (front or interior)
34//! falls back to a full per-series rebuild reported as
35//! `SeriesDataReplaced`. A `PointUpdated` recomputes just its own bucket.
36//!
37//! ```ignore
38//! use teksilo_data::{ChartModel, ChartAggregate, ChartAggregateFn};
39//! let model: ChartModel<i32> = ChartModel::new();
40//! let s = model.add_series("daily");
41//! for i in 0..70 {
42//!     model.push_point(s, i, i as f32);
43//! }
44//! let weekly = ChartAggregate::new(model, 7, ChartAggregateFn::Mean);
45//! assert_eq!(weekly.point_count(s), 10); // 70 / 7
46//! ```
47
48use std::cell::{Cell, RefCell};
49use std::collections::HashMap;
50use std::rc::Rc;
51
52use teksilo_core::ObserverHandle;
53use teksilo_core::color_prop::ColorProp;
54
55use crate::chart_change::{ChartChange, SeriesId};
56use crate::chart_model::{ChartDatum, ChartModel};
57
58type SeriesIdsFn = Rc<dyn Fn() -> Vec<SeriesId>>;
59type PointCountFn = Rc<dyn Fn(SeriesId) -> usize>;
60type WithPointFn<T> = Rc<dyn Fn(SeriesId, usize, &dyn Fn(&ChartDatum<T>))>;
61type WithSeriesFn = Rc<dyn Fn(SeriesId, &dyn Fn(&str, Option<&ColorProp>, bool))>;
62type ObserveChartFn = Rc<dyn Fn(Box<dyn Fn(&ChartChange)>) -> ObserverHandle>;
63
64/// A reduction applied to the numeric values within one bucket.
65pub enum ChartAggregateFn {
66    /// Arithmetic mean of the bucket's values (`0.0` for an empty bucket).
67    Mean,
68    /// Sum of the bucket's values.
69    Sum,
70    /// Smallest value in the bucket.
71    Min,
72    /// Largest value in the bucket.
73    Max,
74    /// The bucket's first value.
75    First,
76    /// The bucket's last value.
77    Last,
78    /// A caller-supplied reduction.
79    Custom(Rc<dyn Fn(&[f32]) -> f32>),
80}
81
82impl ChartAggregateFn {
83    /// Apply the reduction to a bucket's values.
84    ///
85    /// On an empty slice, every built-in variant returns `0.0`
86    /// (`Mean`/`Min`/`Max`/`First`/`Last`) or the empty sum (`Sum`, also
87    /// `0.0`) — a uniform, unsurprising convention rather than `Min`/`Max`
88    /// leaking their fold seed (`±INFINITY`) into a chart value. `Custom`
89    /// returns whatever the supplied closure computes for `&[]`. No internal
90    /// caller actually passes an empty slice — `compute_bucket_datum` bails
91    /// out before calling `apply` for an empty bucket — so this only bites a
92    /// direct caller.
93    pub fn apply(&self, values: &[f32]) -> f32 {
94        match self {
95            ChartAggregateFn::Mean => {
96                if values.is_empty() {
97                    0.0
98                } else {
99                    values.iter().sum::<f32>() / values.len() as f32
100                }
101            }
102            ChartAggregateFn::Sum => values.iter().sum(),
103            ChartAggregateFn::Min => {
104                if values.is_empty() {
105                    0.0
106                } else {
107                    values.iter().copied().fold(f32::INFINITY, f32::min)
108                }
109            }
110            ChartAggregateFn::Max => {
111                if values.is_empty() {
112                    0.0
113                } else {
114                    values.iter().copied().fold(f32::NEG_INFINITY, f32::max)
115                }
116            }
117            ChartAggregateFn::First => values.first().copied().unwrap_or(0.0),
118            ChartAggregateFn::Last => values.last().copied().unwrap_or(0.0),
119            ChartAggregateFn::Custom(f) => f(values),
120        }
121    }
122}
123
124impl Clone for ChartAggregateFn {
125    fn clone(&self) -> Self {
126        match self {
127            ChartAggregateFn::Mean => ChartAggregateFn::Mean,
128            ChartAggregateFn::Sum => ChartAggregateFn::Sum,
129            ChartAggregateFn::Min => ChartAggregateFn::Min,
130            ChartAggregateFn::Max => ChartAggregateFn::Max,
131            ChartAggregateFn::First => ChartAggregateFn::First,
132            ChartAggregateFn::Last => ChartAggregateFn::Last,
133            ChartAggregateFn::Custom(f) => ChartAggregateFn::Custom(f.clone()),
134        }
135    }
136}
137
138impl std::fmt::Debug for ChartAggregateFn {
139    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
140        match self {
141            ChartAggregateFn::Mean => f.write_str("Mean"),
142            ChartAggregateFn::Sum => f.write_str("Sum"),
143            ChartAggregateFn::Min => f.write_str("Min"),
144            ChartAggregateFn::Max => f.write_str("Max"),
145            ChartAggregateFn::First => f.write_str("First"),
146            ChartAggregateFn::Last => f.write_str("Last"),
147            ChartAggregateFn::Custom(_) => f.write_str("Custom(..)"),
148        }
149    }
150}
151
152struct ObserverEntry {
153    id: u64,
154    callback: Rc<dyn Fn(&ChartChange)>,
155}
156
157struct ChartAggregateInner<T: 'static> {
158    series_ids_fn: SeriesIdsFn,
159    point_count_fn: PointCountFn,
160    with_point_fn: WithPointFn<T>,
161    with_series_fn: WithSeriesFn,
162    bucket_size: usize,
163    aggregate_fn: ChartAggregateFn,
164    /// Materialized bucket data, per series — NOT a view into the source.
165    buckets: HashMap<SeriesId, Vec<ChartDatum<T>>>,
166    /// First bucket index that may have changed, per series.
167    divergence: HashMap<SeriesId, usize>,
168    observers: Vec<ObserverEntry>,
169    next_observer_id: u64,
170    _upstream_handle: Option<ObserverHandle>,
171}
172
173impl<T: 'static> ChartAggregateInner<T> {
174    fn snapshot_callbacks(&self) -> Vec<Rc<dyn Fn(&ChartChange)>> {
175        self.observers.iter().map(|e| e.callback.clone()).collect()
176    }
177}
178
179/// A bucket/rollup projection over a [`ChartModel<T>`].
180///
181/// See the module documentation for semantics.
182pub struct ChartAggregate<T: 'static> {
183    inner: Rc<RefCell<ChartAggregateInner<T>>>,
184}
185
186/// Compute one bucket's reduced datum from source indices `[start, end)`.
187/// Returns `None` for an empty range (no bucket to emit).
188fn compute_bucket_datum<T: Clone>(
189    with_point_fn: &WithPointFn<T>,
190    series: SeriesId,
191    aggregate_fn: &ChartAggregateFn,
192    start: usize,
193    end: usize,
194) -> Option<ChartDatum<T>> {
195    if start >= end {
196        return None;
197    }
198    let values: RefCell<Vec<f32>> = RefCell::new(Vec::with_capacity(end - start));
199    let category: RefCell<Option<T>> = RefCell::new(None);
200    for i in start..end {
201        (with_point_fn)(series, i, &|d: &ChartDatum<T>| {
202            if category.borrow().is_none() {
203                *category.borrow_mut() = Some(d.category.clone());
204            }
205            values.borrow_mut().push(d.value);
206        });
207    }
208    let category = category.into_inner()?;
209    let values = values.into_inner();
210    Some(ChartDatum::new(category, aggregate_fn.apply(&values)))
211}
212
213fn rebuild_series<T: Clone + 'static>(inner: &mut ChartAggregateInner<T>, series: SeriesId) {
214    let bs = inner.bucket_size;
215    let n = (inner.point_count_fn)(series);
216    let bucket_count = n.div_ceil(bs);
217    let mut list = Vec::with_capacity(bucket_count);
218    for b in 0..bucket_count {
219        let start = b * bs;
220        let end = (start + bs).min(n);
221        if let Some(d) = compute_bucket_datum(
222            &inner.with_point_fn,
223            series,
224            &inner.aggregate_fn,
225            start,
226            end,
227        ) {
228            list.push(d);
229        }
230    }
231    inner.buckets.insert(series, list);
232    inner.divergence.insert(series, 0);
233}
234
235fn rebuild_all<T: Clone + 'static>(inner: &mut ChartAggregateInner<T>) {
236    inner.buckets.clear();
237    inner.divergence.clear();
238    for series in (inner.series_ids_fn)() {
239        rebuild_series(inner, series);
240    }
241}
242
243/// Translate one upstream `ChartChange` into zero or more local changes.
244/// See module docs.
245fn translate<T: Clone + 'static>(
246    inner: &mut ChartAggregateInner<T>,
247    change: &ChartChange,
248) -> Vec<ChartChange> {
249    match change {
250        ChartChange::SeriesInserted { index, series } => {
251            let (index, series) = (*index, *series);
252            inner.buckets.insert(series, Vec::new());
253            inner.divergence.insert(series, 0);
254            vec![ChartChange::SeriesInserted { index, series }]
255        }
256        ChartChange::SeriesRemoved { series } => {
257            let series = *series;
258            inner.buckets.remove(&series);
259            inner.divergence.remove(&series);
260            vec![ChartChange::SeriesRemoved { series }]
261        }
262        ChartChange::SeriesMoved { series, from, to } => vec![ChartChange::SeriesMoved {
263            series: *series,
264            from: *from,
265            to: *to,
266        }],
267        ChartChange::SeriesRenamed { series } => {
268            vec![ChartChange::SeriesRenamed { series: *series }]
269        }
270        ChartChange::SeriesColorChanged { series } => {
271            vec![ChartChange::SeriesColorChanged { series: *series }]
272        }
273        ChartChange::SeriesPatternChanged { series } => {
274            vec![ChartChange::SeriesPatternChanged { series: *series }]
275        }
276        ChartChange::SeriesVisibilityChanged { series } => {
277            vec![ChartChange::SeriesVisibilityChanged { series: *series }]
278        }
279        ChartChange::PointsInserted { series, range } => {
280            let series = *series;
281            let range = range.clone();
282            let bs = inner.bucket_size;
283            let new_total = (inner.point_count_fn)(series);
284            let inserted = range.end - range.start;
285            let old_total = new_total - inserted;
286
287            if range.start != old_total {
288                // Not a tail append — a mid-series insert. Rebuild.
289                rebuild_series(inner, series);
290                return vec![ChartChange::SeriesDataReplaced { series }];
291            }
292
293            let old_bc = old_total.div_ceil(bs);
294            let new_bc = new_total.div_ceil(bs);
295            let mut out = Vec::new();
296
297            if old_bc >= 1 {
298                let start = (old_bc - 1) * bs;
299                let end = (start + bs).min(new_total);
300                let mut wrote = false;
301                if let Some(d) = compute_bucket_datum(
302                    &inner.with_point_fn,
303                    series,
304                    &inner.aggregate_fn,
305                    start,
306                    end,
307                ) && let Some(list) = inner.buckets.get_mut(&series)
308                    && old_bc - 1 < list.len()
309                {
310                    list[old_bc - 1] = d;
311                    wrote = true;
312                }
313                // Only notify/diverge if the write actually happened — a
314                // failed guard (stale `buckets` entry, an empty computed
315                // range) must be a silent no-op, not a claim that a bucket
316                // changed when it didn't.
317                if wrote {
318                    inner.divergence.insert(series, old_bc - 1);
319                    out.push(ChartChange::PointUpdated {
320                        series,
321                        index: old_bc - 1,
322                    });
323                } else {
324                    debug_assert!(
325                        false,
326                        "chart_aggregate: last-bucket update guard failed for an existing bucket"
327                    );
328                }
329            }
330
331            if new_bc > old_bc {
332                for b in old_bc..new_bc {
333                    let start = b * bs;
334                    let end = (start + bs).min(new_total);
335                    if let Some(d) = compute_bucket_datum(
336                        &inner.with_point_fn,
337                        series,
338                        &inner.aggregate_fn,
339                        start,
340                        end,
341                    ) {
342                        inner.buckets.entry(series).or_default().push(d);
343                    }
344                }
345                inner.divergence.insert(series, old_bc.saturating_sub(1));
346                out.push(ChartChange::PointsInserted {
347                    series,
348                    range: old_bc..new_bc,
349                });
350            }
351            out
352        }
353        ChartChange::PointsRemoved { series, range } => {
354            let series = *series;
355            let range = range.clone();
356            let bs = inner.bucket_size;
357            let new_total = (inner.point_count_fn)(series);
358            let removed = range.end - range.start;
359            let old_total = new_total + removed;
360
361            if range.end != old_total {
362                // Not a tail removal — a front/mid-series removal shifts
363                // every point after it into a different bucket. Rebuild.
364                rebuild_series(inner, series);
365                return vec![ChartChange::SeriesDataReplaced { series }];
366            }
367
368            // Tail-removal mirror of the `PointsInserted` tail-append
369            // branch above: buckets entirely below the new total are
370            // untouched (their point range didn't change), the new last
371            // bucket (if any) may have lost some of its points and needs
372            // recomputing, and any buckets entirely beyond the new total
373            // no longer exist.
374            let old_bc = old_total.div_ceil(bs);
375            let new_bc = new_total.div_ceil(bs);
376            let mut out = Vec::new();
377            let mut floor = usize::MAX;
378
379            if new_bc >= 1 {
380                let start = (new_bc - 1) * bs;
381                let end = (start + bs).min(new_total);
382                let mut wrote = false;
383                if let Some(d) = compute_bucket_datum(
384                    &inner.with_point_fn,
385                    series,
386                    &inner.aggregate_fn,
387                    start,
388                    end,
389                ) && let Some(list) = inner.buckets.get_mut(&series)
390                    && new_bc - 1 < list.len()
391                {
392                    list[new_bc - 1] = d;
393                    wrote = true;
394                }
395                // Only notify/diverge if the write actually happened — see
396                // the matching guard in the `PointsInserted` arm above.
397                if wrote {
398                    floor = floor.min(new_bc - 1);
399                    out.push(ChartChange::PointUpdated {
400                        series,
401                        index: new_bc - 1,
402                    });
403                } else {
404                    debug_assert!(
405                        false,
406                        "chart_aggregate: last-bucket update guard failed for an existing bucket after removal"
407                    );
408                }
409            }
410
411            if new_bc < old_bc {
412                if let Some(list) = inner.buckets.get_mut(&series) {
413                    list.truncate(new_bc);
414                }
415                floor = floor.min(new_bc);
416                out.push(ChartChange::PointsRemoved {
417                    series,
418                    range: new_bc..old_bc,
419                });
420            }
421
422            if floor != usize::MAX {
423                inner.divergence.insert(series, floor);
424            }
425            out
426        }
427        ChartChange::PointUpdated { series, index } => {
428            let (series, index) = (*series, *index);
429            let bs = inner.bucket_size;
430            let bucket_index = index / bs;
431            let n = (inner.point_count_fn)(series);
432            let start = bucket_index * bs;
433            let end = (start + bs).min(n);
434            let mut wrote = false;
435            if let Some(d) = compute_bucket_datum(
436                &inner.with_point_fn,
437                series,
438                &inner.aggregate_fn,
439                start,
440                end,
441            ) && let Some(list) = inner.buckets.get_mut(&series)
442                && bucket_index < list.len()
443            {
444                list[bucket_index] = d;
445                wrote = true;
446            }
447            // Only notify/diverge if the write actually happened — see the
448            // matching guard in the `PointsInserted` arm above.
449            if wrote {
450                inner.divergence.insert(series, bucket_index);
451                vec![ChartChange::PointUpdated {
452                    series,
453                    index: bucket_index,
454                }]
455            } else {
456                debug_assert!(
457                    false,
458                    "chart_aggregate: point-update guard failed for an existing bucket"
459                );
460                vec![]
461            }
462        }
463        ChartChange::SeriesDataReplaced { series } => {
464            let series = *series;
465            rebuild_series(inner, series);
466            vec![ChartChange::SeriesDataReplaced { series }]
467        }
468        ChartChange::Reset => {
469            rebuild_all(inner);
470            vec![ChartChange::Reset]
471        }
472    }
473}
474
475fn translate_and_notify<T: Clone + 'static>(
476    inner: &Rc<RefCell<ChartAggregateInner<T>>>,
477    change: &ChartChange,
478) {
479    let (changes, callbacks) = {
480        let mut guard = inner.borrow_mut();
481        let changes = translate(&mut guard, change);
482        (changes, guard.snapshot_callbacks())
483    };
484    for c in &changes {
485        for cb in &callbacks {
486            cb(c);
487        }
488    }
489}
490
491impl<T: Clone + 'static> ChartAggregate<T> {
492    /// Wrap `source`, bucketing every series into groups of `bucket_size`
493    /// source points reduced via `aggregate_fn`. `bucket_size` is clamped
494    /// to a minimum of 1.
495    pub fn new(source: ChartModel<T>, bucket_size: usize, aggregate_fn: ChartAggregateFn) -> Self {
496        let series_ids_fn: SeriesIdsFn = {
497            let m = source.clone();
498            Rc::new(move || m.series_ids())
499        };
500        let point_count_fn: PointCountFn = {
501            let m = source.clone();
502            Rc::new(move |series| m.point_count(series))
503        };
504        let with_point_fn: WithPointFn<T> = {
505            let m = source.clone();
506            Rc::new(move |series, idx, f| {
507                m.with_point(series, idx, |d| f(d));
508            })
509        };
510        let with_series_fn: WithSeriesFn = {
511            let m = source.clone();
512            Rc::new(move |series, f| {
513                m.with_series(series, |name, color, visible| f(name, color, visible));
514            })
515        };
516        let observe_fn: ObserveChartFn = {
517            let m = source;
518            Rc::new(move |callback| m.observe_changes(move |change| callback(change)))
519        };
520
521        let inner = Rc::new(RefCell::new(ChartAggregateInner {
522            series_ids_fn,
523            point_count_fn,
524            with_point_fn,
525            with_series_fn,
526            bucket_size: bucket_size.max(1),
527            aggregate_fn,
528            buckets: HashMap::new(),
529            divergence: HashMap::new(),
530            observers: Vec::new(),
531            next_observer_id: 1,
532            _upstream_handle: None,
533        }));
534
535        rebuild_all(&mut inner.borrow_mut());
536
537        let weak = Rc::downgrade(&inner);
538        let upstream_handle = (observe_fn)(Box::new(move |change| {
539            if let Some(strong) = weak.upgrade() {
540                translate_and_notify(&strong, change);
541            }
542        }));
543        inner.borrow_mut()._upstream_handle = Some(upstream_handle);
544
545        Self { inner }
546    }
547
548    /// Change the bucket size, rebuilding every series and emitting
549    /// `ChartChange::Reset`. Clamped to a minimum of 1.
550    pub fn set_bucket_size(&self, bucket_size: usize) {
551        let callbacks = {
552            let mut guard = self.inner.borrow_mut();
553            guard.bucket_size = bucket_size.max(1);
554            rebuild_all(&mut guard);
555            guard.snapshot_callbacks()
556        };
557        for cb in &callbacks {
558            cb(&ChartChange::Reset);
559        }
560    }
561
562    /// Change the aggregate reduction, rebuilding every series and emitting
563    /// `ChartChange::Reset`.
564    pub fn set_aggregate_fn(&self, aggregate_fn: ChartAggregateFn) {
565        let callbacks = {
566            let mut guard = self.inner.borrow_mut();
567            guard.aggregate_fn = aggregate_fn;
568            rebuild_all(&mut guard);
569            guard.snapshot_callbacks()
570        };
571        for cb in &callbacks {
572            cb(&ChartChange::Reset);
573        }
574    }
575}
576
577impl<T: 'static> ChartAggregate<T> {
578    /// The configured bucket size.
579    pub fn bucket_size(&self) -> usize {
580        self.inner.borrow().bucket_size
581    }
582
583    /// Number of series (same set as the source).
584    pub fn series_count(&self) -> usize {
585        (self.inner.borrow().series_ids_fn)().len()
586    }
587
588    /// The series ids, in the source's display order.
589    pub fn series_ids(&self) -> Vec<SeriesId> {
590        (self.inner.borrow().series_ids_fn)()
591    }
592
593    /// Number of buckets currently materialized for `series`.
594    pub fn point_count(&self, series: SeriesId) -> usize {
595        self.inner
596            .borrow()
597            .buckets
598            .get(&series)
599            .map(|v| v.len())
600            .unwrap_or(0)
601    }
602
603    /// Access a series' metadata (delegates straight through to the
604    /// source). Returns `None` if `series` is unknown.
605    pub fn with_series<R>(
606        &self,
607        series: SeriesId,
608        f: impl FnOnce(&str, Option<&ColorProp>, bool) -> R,
609    ) -> Option<R> {
610        let with_series_fn = self.inner.borrow().with_series_fn.clone();
611        let f_cell: Cell<Option<_>> = Cell::new(Some(f));
612        let slot: Cell<Option<R>> = Cell::new(None);
613        (with_series_fn)(series, &|name, color, visible| {
614            if let Some(f) = f_cell.take() {
615                slot.set(Some(f(name, color, visible)));
616            }
617        });
618        slot.into_inner()
619    }
620
621    /// Access the bucket at `index` within `series`. Returns `None` if
622    /// `series` or `index` is unknown.
623    pub fn with_point<R>(
624        &self,
625        series: SeriesId,
626        index: usize,
627        f: impl FnOnce(&ChartDatum<T>) -> R,
628    ) -> Option<R> {
629        let guard = self.inner.borrow();
630        guard.buckets.get(&series).and_then(|v| v.get(index)).map(f)
631    }
632
633    /// Register an observer for translated bucket changes. Returns an
634    /// `ObserverHandle` — dropping it removes the callback.
635    pub fn observe_changes(&self, f: impl Fn(&ChartChange) + 'static) -> ObserverHandle {
636        let mut guard = self.inner.borrow_mut();
637        let id = guard.next_observer_id;
638        guard.next_observer_id += 1;
639        guard.observers.push(ObserverEntry {
640            id,
641            callback: Rc::new(f),
642        });
643        let inner = self.inner.clone();
644        ObserverHandle::new(
645            self.inner.clone(),
646            id,
647            Rc::new(move |observer_id| {
648                inner.borrow_mut().observers.retain(|e| e.id != observer_id);
649            }),
650        )
651    }
652
653    /// First bucket index of `series` whose content may differ since the
654    /// latest translated change. `None` if `series` is unknown or
655    /// unaffected yet.
656    pub fn first_changed_index(&self, series: SeriesId) -> Option<usize> {
657        self.inner.borrow().divergence.get(&series).copied()
658    }
659}
660
661impl<T: 'static> Clone for ChartAggregate<T> {
662    fn clone(&self) -> Self {
663        Self {
664            inner: self.inner.clone(),
665        }
666    }
667}
668
669impl<T: 'static> std::fmt::Debug for ChartAggregate<T> {
670    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
671        let guard = self.inner.borrow();
672        f.debug_struct("ChartAggregate")
673            .field("bucket_size", &guard.bucket_size)
674            .field("series_count", &(guard.series_ids_fn)().len())
675            .finish()
676    }
677}
678
679#[cfg(test)]
680mod tests {
681    use super::*;
682
683    #[test]
684    fn apply_on_empty_slice_is_zero_for_every_built_in_variant() {
685        // Regression: `Min`/`Max` used to leak their fold seed (`±INFINITY`)
686        // on an empty bucket, inconsistent with `Mean`/`Sum`/`First`/`Last`
687        // all returning `0.0`. No internal caller hits this path, but a
688        // direct caller shouldn't see an infinity fall out of a chart value.
689        for f in [
690            ChartAggregateFn::Mean,
691            ChartAggregateFn::Sum,
692            ChartAggregateFn::Min,
693            ChartAggregateFn::Max,
694            ChartAggregateFn::First,
695            ChartAggregateFn::Last,
696        ] {
697            assert_eq!(f.apply(&[]), 0.0, "{f:?} on empty slice");
698        }
699    }
700
701    fn series_with(values: &[f32]) -> (ChartModel<i32>, SeriesId) {
702        let model: ChartModel<i32> = ChartModel::new();
703        let s = model.add_series("s");
704        for (i, &v) in values.iter().enumerate() {
705            model.push_point(s, i as i32, v);
706        }
707        (model, s)
708    }
709
710    fn values<T: 'static>(agg: &ChartAggregate<T>, s: SeriesId) -> Vec<f32> {
711        (0..agg.point_count(s))
712            .map(|i| agg.with_point(s, i, |d| d.value).unwrap())
713            .collect()
714    }
715
716    fn track(agg: &ChartAggregate<i32>) -> (Rc<RefCell<Vec<ChartChange>>>, ObserverHandle) {
717        let log: Rc<RefCell<Vec<ChartChange>>> = Rc::new(RefCell::new(Vec::new()));
718        let l = log.clone();
719        let handle = agg.observe_changes(move |c| l.borrow_mut().push(c.clone()));
720        (log, handle)
721    }
722
723    #[test]
724    fn bucket_size_one_is_identity() {
725        let (model, s) = series_with(&[0.0, 10.0, 20.0, 30.0]);
726        let agg = ChartAggregate::new(model, 1, ChartAggregateFn::Mean);
727        assert_eq!(agg.point_count(s), 4);
728        assert_eq!(values(&agg, s), vec![0.0, 10.0, 20.0, 30.0]);
729    }
730
731    #[test]
732    fn mean_aggregate() {
733        let (model, s) = series_with(&[1.0, 2.0, 3.0, 4.0]);
734        let agg = ChartAggregate::new(model, 2, ChartAggregateFn::Mean);
735        assert_eq!(agg.point_count(s), 2);
736        assert_eq!(values(&agg, s), vec![1.5, 3.5]);
737    }
738
739    #[test]
740    fn max_aggregate() {
741        let (model, s) = series_with(&[1.0, 5.0, 2.0, 9.0]);
742        let agg = ChartAggregate::new(model, 2, ChartAggregateFn::Max);
743        assert_eq!(values(&agg, s), vec![5.0, 9.0]);
744    }
745
746    #[test]
747    fn custom_aggregate() {
748        let (model, s) = series_with(&[1.0, 2.0, 3.0]);
749        let agg = ChartAggregate::new(
750            model,
751            3,
752            ChartAggregateFn::Custom(Rc::new(|vs: &[f32]| vs.iter().product())),
753        );
754        assert_eq!(values(&agg, s), vec![6.0]);
755    }
756
757    #[test]
758    fn trailing_partial_bucket_is_included() {
759        let (model, s) = series_with(&[1.0, 2.0, 3.0]);
760        let agg = ChartAggregate::new(model, 2, ChartAggregateFn::Sum);
761        assert_eq!(agg.point_count(s), 2);
762        assert_eq!(values(&agg, s), vec![3.0, 3.0]); // [1,2] and trailing [3]
763    }
764
765    #[test]
766    fn tail_append_worked_trace_bucket_size_3() {
767        // 7 points -> 3 buckets: [0,1,2], [3,4,5], [6]
768        let (model, s) = series_with(&[0.0, 1.0, 2.0, 3.0, 4.0, 5.0, 6.0]);
769        let agg = ChartAggregate::new(model.clone(), 3, ChartAggregateFn::Sum);
770        assert_eq!(agg.point_count(s), 3);
771        let (log, _h) = track(&agg);
772
773        // 7 -> 8: same bucket count (last bucket [6,7)).
774        model.push_point(s, 7, 7.0);
775        {
776            let entries = log.borrow();
777            assert_eq!(entries.len(), 1);
778            assert_eq!(
779                entries[0],
780                ChartChange::PointUpdated {
781                    series: s,
782                    index: 2
783                }
784            );
785        }
786        assert_eq!(agg.point_count(s), 3);
787        assert_eq!(agg.with_point(s, 2, |d| d.value), Some(13.0)); // 6+7
788        log.borrow_mut().clear();
789
790        // 8 -> 9: still same bucket count (last bucket [6,9)).
791        model.push_point(s, 8, 8.0);
792        {
793            let entries = log.borrow();
794            assert_eq!(entries.len(), 1);
795            assert_eq!(
796                entries[0],
797                ChartChange::PointUpdated {
798                    series: s,
799                    index: 2
800                }
801            );
802        }
803        assert_eq!(agg.with_point(s, 2, |d| d.value), Some(21.0)); // 6+7+8
804        log.borrow_mut().clear();
805
806        // 9 -> 10: new bucket 3 = [9].
807        model.push_point(s, 9, 9.0);
808        {
809            let entries = log.borrow();
810            assert_eq!(entries.len(), 2);
811            assert_eq!(
812                entries[0],
813                ChartChange::PointUpdated {
814                    series: s,
815                    index: 2
816                }
817            );
818            assert_eq!(
819                entries[1],
820                ChartChange::PointsInserted {
821                    series: s,
822                    range: 3..4
823                }
824            );
825        }
826        assert_eq!(agg.point_count(s), 4);
827        assert_eq!(agg.with_point(s, 3, |d| d.value), Some(9.0));
828    }
829
830    #[test]
831    fn mid_series_insert_falls_back_to_replace() {
832        let (model, s) = series_with(&[0.0, 1.0, 2.0]);
833        let agg = ChartAggregate::new(model.clone(), 2, ChartAggregateFn::Sum);
834        let (log, _h) = track(&agg);
835
836        model.insert_point(s, 1, 99, 99.0);
837
838        let entries = log.borrow();
839        assert_eq!(entries.len(), 1);
840        assert_eq!(entries[0], ChartChange::SeriesDataReplaced { series: s });
841    }
842
843    #[test]
844    fn front_removal_falls_back_to_replace() {
845        // Removing anything but the tail shifts every later point into a
846        // different bucket — must still rebuild.
847        let (model, s) = series_with(&[0.0, 1.0, 2.0, 3.0]);
848        let agg = ChartAggregate::new(model.clone(), 2, ChartAggregateFn::Sum);
849        let (log, _h) = track(&agg);
850
851        model.remove_point(s, 0);
852
853        let entries = log.borrow();
854        assert_eq!(entries.len(), 1);
855        assert_eq!(entries[0], ChartChange::SeriesDataReplaced { series: s });
856        assert_eq!(agg.point_count(s), 2); // 3 remaining points / 2
857    }
858
859    #[test]
860    fn tail_removal_worked_trace_bucket_size_3() {
861        // Mirror of `tail_append_worked_trace_bucket_size_3` for removal.
862        // 10 points -> 4 buckets: [0,1,2], [3,4,5], [6,7,8], [9]
863        let (model, s) = series_with(&[0.0, 1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0]);
864        let agg = ChartAggregate::new(model.clone(), 3, ChartAggregateFn::Sum);
865        assert_eq!(agg.point_count(s), 4);
866        let (log, _h) = track(&agg);
867
868        // 10 -> 9: drops the trailing partial bucket [9] entirely; the new
869        // last bucket [6,7,8] is recomputed (unchanged value, but still
870        // written — same "always recompute the boundary bucket" contract
871        // the insert-side tail-append path uses).
872        model.remove_point(s, 9);
873        {
874            let entries = log.borrow();
875            assert_eq!(entries.len(), 2);
876            assert_eq!(
877                entries[0],
878                ChartChange::PointUpdated {
879                    series: s,
880                    index: 2
881                }
882            );
883            assert_eq!(
884                entries[1],
885                ChartChange::PointsRemoved {
886                    series: s,
887                    range: 3..4
888                }
889            );
890        }
891        assert_eq!(agg.point_count(s), 3);
892        assert_eq!(agg.with_point(s, 2, |d| d.value), Some(21.0)); // 6+7+8
893        log.borrow_mut().clear();
894
895        // 9 -> 8: same bucket count (last bucket shrinks to [6,7]).
896        model.remove_point(s, 8);
897        {
898            let entries = log.borrow();
899            assert_eq!(entries.len(), 1);
900            assert_eq!(
901                entries[0],
902                ChartChange::PointUpdated {
903                    series: s,
904                    index: 2
905                }
906            );
907        }
908        assert_eq!(agg.point_count(s), 3);
909        assert_eq!(agg.with_point(s, 2, |d| d.value), Some(13.0)); // 6+7
910    }
911
912    #[test]
913    fn tail_removal_down_to_empty_series_drops_the_last_bucket() {
914        // new_bc == 0: the "recompute the new last bucket" step must be
915        // skipped entirely (there is no last bucket any more), leaving only
916        // the bucket-removal event.
917        let (model, s) = series_with(&[1.0, 2.0]);
918        let agg = ChartAggregate::new(model.clone(), 3, ChartAggregateFn::Sum);
919        assert_eq!(agg.point_count(s), 1); // one trailing partial bucket
920        let (log, _h) = track(&agg);
921
922        model.remove_point(s, 1);
923        model.remove_point(s, 0);
924
925        let entries = log.borrow();
926        assert_eq!(
927            entries.last(),
928            Some(&ChartChange::PointsRemoved {
929                series: s,
930                range: 0..1
931            })
932        );
933        drop(entries);
934        assert_eq!(agg.point_count(s), 0);
935    }
936
937    #[test]
938    fn point_updated_recomputes_its_bucket() {
939        let (model, s) = series_with(&[1.0, 2.0, 3.0, 4.0]);
940        let agg = ChartAggregate::new(model.clone(), 2, ChartAggregateFn::Sum);
941        let (log, _h) = track(&agg);
942
943        model.update_point(s, 2, 2, 30.0); // bucket 1 = [3,4) -> now [30,4)
944
945        let entries = log.borrow();
946        assert_eq!(entries.len(), 1);
947        assert_eq!(
948            entries[0],
949            ChartChange::PointUpdated {
950                series: s,
951                index: 1
952            }
953        );
954        drop(entries);
955        assert_eq!(agg.with_point(s, 1, |d| d.value), Some(34.0));
956    }
957
958    #[test]
959    fn reset_rebuilds_and_passes_through() {
960        let (model, s) = series_with(&[1.0, 2.0, 3.0]);
961        let agg = ChartAggregate::new(model.clone(), 2, ChartAggregateFn::Sum);
962        let (log, _h) = track(&agg);
963
964        model.clear();
965        assert_eq!(log.borrow().last(), Some(&ChartChange::Reset));
966        assert_eq!(agg.point_count(s), 0);
967    }
968
969    #[test]
970    fn dropping_aggregate_unregisters_upstream_observer() {
971        let model: ChartModel<i32> = ChartModel::new();
972        let before = model.observer_count();
973        let agg = ChartAggregate::new(model.clone(), 3, ChartAggregateFn::Sum);
974        assert_eq!(model.observer_count(), before + 1);
975        drop(agg);
976        assert_eq!(model.observer_count(), before);
977    }
978
979    #[test]
980    fn set_bucket_size_rebuilds() {
981        let (model, s) = series_with(&[1.0, 2.0, 3.0, 4.0]);
982        let agg = ChartAggregate::new(model, 2, ChartAggregateFn::Sum);
983        assert_eq!(agg.point_count(s), 2);
984        agg.set_bucket_size(4);
985        assert_eq!(agg.point_count(s), 1);
986        assert_eq!(values(&agg, s), vec![10.0]);
987    }
988
989    #[test]
990    fn set_aggregate_fn_rebuilds() {
991        let (model, s) = series_with(&[1.0, 5.0, 2.0]);
992        let agg = ChartAggregate::new(model, 3, ChartAggregateFn::Sum);
993        assert_eq!(values(&agg, s), vec![8.0]);
994        agg.set_aggregate_fn(ChartAggregateFn::Max);
995        assert_eq!(values(&agg, s), vec![5.0]);
996    }
997}