Skip to main content

canic_core/ops/runtime/public_metrics/
mod.rs

1//! Module: ops::runtime::public_metrics
2//!
3//! Responsibility: collect local aggregate counters and project the publication cache.
4//! Does not own: timers, endpoint authorization, or application metric semantics.
5//! Boundary: public query projection reads cached values only.
6
7mod process;
8
9use crate::{
10    InternalError,
11    config::{Config, RoleRuntimeConfig},
12    domain::public_metrics::PublicMetricFamily,
13    dto::{
14        metrics::MetricValue,
15        page::{Page, PageRequest},
16        public_status::{
17            PublicCounterDelta, PublicHealth, PublicHealthStatus, PublicHistoryPoint,
18            PublicHistoryRequest, PublicHistorySnapshot, PublicMetric, PublicMetricKind,
19            PublicMetricsRequest, PublicMetricsSnapshot, PublicSnapshotState,
20        },
21    },
22    model::public_metrics::{
23        MAX_HISTORY_BYTES, MAX_HISTORY_SERIES, MAX_PUBLIC_METRIC_TEXT_BYTES, MAX_PUBLIC_METRICS,
24        PUBLIC_HISTORY_RETENTION_NS, PUBLIC_HISTORY_SLOTS, PUBLIC_METRICS_CADENCE_NS,
25        PUBLIC_METRICS_STALE_AFTER_NS, PublicHistoryCache, PublicMetricSample, PublicMetricsCache,
26    },
27    ops::{
28        ic::IcOps,
29        runtime::{env::EnvOps, metrics},
30    },
31};
32use std::{cell::Cell, collections::BTreeSet};
33
34thread_local! {
35    static APPLICATION_SAMPLER: Cell<Option<ApplicationMetricsSampler>> = const { Cell::new(None) };
36}
37
38#[cfg(feature = "sharding")]
39use crate::ops::storage::placement::sharding::ShardingRegistryOps;
40
41/// One synchronous aggregate provider composed by application lifecycle code.
42/// The provider owns bounded source collection; Canic owns the timer and history.
43#[derive(Clone, Copy)]
44pub struct ApplicationMetricsSampler {
45    collect: fn() -> Result<Vec<PublicMetric>, crate::dto::error::Error>,
46}
47
48impl ApplicationMetricsSampler {
49    /// Wrap a provider that reads bounded counters and preserves source time and reset windows.
50    #[must_use]
51    pub const fn new(collect: fn() -> Result<Vec<PublicMetric>, crate::dto::error::Error>) -> Self {
52        Self { collect }
53    }
54}
55
56/// Local sampling and public snapshot projection under immutable publication configuration.
57pub struct PublicMetricsOps;
58
59impl PublicMetricsOps {
60    /// Install the sole synchronous composition callback, with no timer or database ownership.
61    pub fn set_application_sampler(sample: Option<ApplicationMetricsSampler>) {
62        APPLICATION_SAMPLER.set(sample);
63    }
64
65    #[must_use]
66    pub fn enabled() -> BTreeSet<PublicMetricFamily> {
67        RoleRuntimeConfig::try_get()
68            .map(|config| config.public_metrics.clone())
69            .or_else(|| {
70                Config::get()
71                    .ok()
72                    .map(|config| config.public_metrics.clone())
73            })
74            .unwrap_or_default()
75    }
76
77    #[must_use]
78    pub fn health() -> PublicHealth {
79        let now = IcOps::now_nanos();
80        PublicHealth {
81            canister_id: IcOps::canister_self(),
82            role: EnvOps::canister_role().ok().map(|role| role.to_string()),
83            health: PublicHealthStatus::Responding,
84            observed_at_ns: now,
85        }
86    }
87
88    #[must_use]
89    pub fn read(request: PublicMetricsRequest) -> PublicMetricsSnapshot {
90        Self::project(request, &Self::enabled(), IcOps::now_nanos())
91    }
92
93    fn project(
94        request: PublicMetricsRequest,
95        enabled: &BTreeSet<PublicMetricFamily>,
96        now_ns: u64,
97    ) -> PublicMetricsSnapshot {
98        let snapshot = enabled
99            .contains(&request.family)
100            .then(|| PublicMetricsCache::snapshot(request.family))
101            .flatten();
102        let state = if !enabled.contains(&request.family) {
103            PublicSnapshotState::Disabled
104        } else if let Some(snapshot) = &snapshot {
105            if now_ns.saturating_sub(snapshot.sampled_at_ns) > PUBLIC_METRICS_STALE_AFTER_NS {
106                PublicSnapshotState::Stale
107            } else {
108                PublicSnapshotState::Fresh
109            }
110        } else {
111            PublicSnapshotState::Unavailable
112        };
113        let sampled_at_ns = snapshot.as_ref().map(|s| s.sampled_at_ns);
114        let truncated = snapshot.as_ref().is_some_and(|s| s.truncated);
115        let rows = snapshot.map_or_else(Vec::new, |s| {
116            s.metrics
117                .into_iter()
118                .map(|row| PublicMetric {
119                    name: row.name,
120                    canister_id: row.canister_id,
121                    value: row.value,
122                    unit: row.unit,
123                    observed_at_ns: row.observed_at_ns,
124                    kind: row.kind,
125                })
126                .collect()
127        });
128        PublicMetricsSnapshot {
129            family: request.family,
130            state,
131            sampled_at_ns,
132            stale_after_ns: PUBLIC_METRICS_STALE_AFTER_NS,
133            truncated,
134            metrics: page(rows, request.page),
135        }
136    }
137
138    /// Expire heap history from the update-side timer, never from a public query.
139    pub fn expire_history(now_ns: u64) {
140        PublicHistoryCache::expire(now_ns);
141    }
142
143    /// Read one bounded cached series without invoking any producer.
144    #[must_use]
145    pub fn history(request: PublicHistoryRequest) -> PublicHistorySnapshot {
146        let mut snapshot = Self::project_history(request, &Self::enabled(), IcOps::now_nanos());
147        snapshot.canister_version = ic_cdk::api::canister_version();
148        snapshot
149    }
150
151    fn project_history(
152        request: PublicHistoryRequest,
153        enabled: &BTreeSet<PublicMetricFamily>,
154        now_ns: u64,
155    ) -> PublicHistorySnapshot {
156        let selected = enabled.contains(&request.family);
157        let valid_name = request.name.len() <= MAX_PUBLIC_METRIC_TEXT_BYTES;
158        let series = (selected && valid_name)
159            .then(|| PublicHistoryCache::series(request.family, request.name, request.canister_id))
160            .flatten();
161        let slot = now_ns / PUBLIC_METRICS_CADENCE_NS;
162        let mut points: Vec<_> = series.as_ref().map_or_else(Vec::new, |series| {
163            series
164                .slots
165                .iter()
166                .filter(|point| {
167                    point.slot <= slot && slot - point.slot < PUBLIC_HISTORY_SLOTS as u64
168                })
169                .copied()
170                .collect()
171        });
172        points.sort_by_key(|point| point.slot);
173        let state = if !selected {
174            PublicSnapshotState::Disabled
175        } else if let Some(point) = points.last() {
176            if now_ns.saturating_sub(point.observed_at_ns) > PUBLIC_METRICS_STALE_AFTER_NS {
177                PublicSnapshotState::Stale
178            } else {
179                PublicSnapshotState::Fresh
180            }
181        } else {
182            PublicSnapshotState::Unavailable
183        };
184        let coverage_start_ns = points
185            .first()
186            .map(|point| point.slot * PUBLIC_METRICS_CADENCE_NS);
187        let total = points.len() as u64;
188        let entries = points
189            .iter()
190            .enumerate()
191            .map(|(index, point)| {
192                let delta = index
193                    .checked_sub(1)
194                    .and_then(|previous| counter_delta(&points[previous], point));
195                PublicHistoryPoint {
196                    delta,
197                    slot_start_ns: point.slot * PUBLIC_METRICS_CADENCE_NS,
198                    observed_at_ns: point.observed_at_ns,
199                    value: point.value,
200                    kind: point.kind,
201                }
202            })
203            .skip(usize::try_from(request.page.offset.min(total)).unwrap_or(PUBLIC_HISTORY_SLOTS))
204            .take(
205                usize::try_from(request.page.limit.min(PUBLIC_HISTORY_SLOTS as u64))
206                    .unwrap_or(PUBLIC_HISTORY_SLOTS),
207            )
208            .collect();
209        PublicHistorySnapshot {
210            state,
211            unit: series.map(|series| series.unit),
212            heap_started_at_ns: selected
213                .then(PublicHistoryCache::heap_started_at_ns)
214                .flatten(),
215            canister_version: 0,
216            coverage_start_ns,
217            cadence_ns: PUBLIC_METRICS_CADENCE_NS,
218            retention_ns: PUBLIC_HISTORY_RETENTION_NS,
219            stale_after_ns: PUBLIC_METRICS_STALE_AFTER_NS,
220            truncated: selected && PublicHistoryCache::truncated(),
221            series_limit: MAX_HISTORY_SERIES as u64,
222            byte_limit: MAX_HISTORY_BYTES as u64,
223            reserved_bytes: if selected {
224                PublicHistoryCache::reserved_bytes() as u64
225            } else {
226                0
227            },
228            points: Page { entries, total },
229        }
230    }
231
232    pub fn record_application(metrics: Vec<PublicMetric>) -> Result<(), InternalError> {
233        if !Self::enabled().contains(&PublicMetricFamily::Application) {
234            return Ok(());
235        }
236        let rows = metrics.into_iter().map(|row| PublicMetricSample {
237            name: row.name,
238            canister_id: row.canister_id,
239            value: row.value,
240            unit: row.unit,
241            observed_at_ns: row.observed_at_ns,
242            kind: row.kind,
243        });
244        PublicMetricsCache::replace(PublicMetricFamily::Application, IcOps::now_nanos(), rows)
245    }
246
247    pub fn sample_family(family: PublicMetricFamily, now: u64) -> Result<(), InternalError> {
248        let mut rows = match family {
249            PublicMetricFamily::Application => {
250                if let Some(sample) = APPLICATION_SAMPLER.get() {
251                    let metrics = (sample.collect)().map_err(|_| InternalError::invalid_input())?;
252                    return Self::record_application(metrics);
253                }
254                return Ok(());
255            }
256            PublicMetricFamily::Cycles => vec![PublicMetricSample {
257                name: "balance".into(),
258                canister_id: Some(IcOps::canister_self()),
259                value: IcOps::canister_cycle_balance().to_u128(),
260                unit: "cycles".into(),
261                observed_at_ns: 0,
262                kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
263            }],
264            PublicMetricFamily::Operations => {
265                let mut rows = process::operations()?;
266                rows.extend(operation_metrics()?);
267                rows
268            }
269            PublicMetricFamily::Performance => {
270                let mut rows = process::memory();
271                rows.extend(performance_metrics()?);
272                rows
273            }
274            PublicMetricFamily::ShardOccupancy => shard_metrics(),
275        };
276        for row in &mut rows {
277            row.observed_at_ns = now;
278            // Timers expose lifetime summaries without a per-registration reset identity.
279            // Keep these raw observations out of counter delta/rate calculations.
280            let counter = (family == PublicMetricFamily::Operations
281                && !row.name.starts_with("cycles_funding.icp_refill.")
282                && !row.name.starts_with("timer."))
283                || (family == PublicMetricFamily::Performance
284                    && !row.name.starts_with("perf.timer.")
285                    && !row.name.starts_with("memory."));
286            if counter {
287                row.kind = PublicMetricKind::Counter {
288                    window_id: 0,
289                    saturated: if matches!(row.unit.as_str(), "cycles" | "icp_e8s") {
290                        row.value == u128::MAX
291                    } else {
292                        row.value == u128::from(u64::MAX)
293                    },
294                };
295            }
296        }
297        PublicMetricsCache::replace(family, now, rows)
298    }
299}
300
301fn counter_delta(
302    previous: &crate::model::public_metrics::PublicHistorySample,
303    current: &crate::model::public_metrics::PublicHistorySample,
304) -> Option<PublicCounterDelta> {
305    let PublicMetricKind::Counter {
306        window_id,
307        saturated: false,
308    } = previous.kind
309    else {
310        return None;
311    };
312    if current.kind
313        != (PublicMetricKind::Counter {
314            window_id,
315            saturated: false,
316        })
317        || previous.slot.checked_add(1) != Some(current.slot)
318    {
319        return None;
320    }
321    let elapsed_ns = current
322        .observed_at_ns
323        .checked_sub(previous.observed_at_ns)
324        .filter(|elapsed| *elapsed > 0)?;
325    Some(PublicCounterDelta {
326        amount: current.value.checked_sub(previous.value)?,
327        elapsed_ns,
328    })
329}
330
331fn page(rows: Vec<PublicMetric>, request: PageRequest) -> Page<PublicMetric> {
332    let total = u64::try_from(rows.len()).unwrap_or(u64::MAX);
333    let start = usize::try_from(request.offset.min(total)).unwrap_or(rows.len());
334    let limit = usize::try_from(request.limit.min(total)).unwrap_or(rows.len());
335    Page {
336        entries: rows.into_iter().skip(start).take(limit).collect(),
337        total,
338    }
339}
340
341// Validate the complete series identity before allocating its formatted name.
342fn metric_name(labels: &[String], suffix_bytes: usize) -> Result<String, InternalError> {
343    let bytes = labels.iter().try_fold(
344        suffix_bytes + labels.len().saturating_sub(1),
345        |bytes, label| bytes.checked_add(label.len()),
346    );
347    if bytes.is_none_or(|bytes| bytes > MAX_PUBLIC_METRIC_TEXT_BYTES) {
348        return Err(InternalError::invalid_input());
349    }
350    Ok(labels.join("."))
351}
352
353fn operation_metrics() -> Result<Vec<PublicMetricSample>, InternalError> {
354    metrics::bounded_core_entries(MAX_PUBLIC_METRICS + 1)?
355        .into_iter()
356        .take(MAX_PUBLIC_METRICS + 1)
357        .map(|row| {
358            let suffix_bytes = if matches!(&row.value, MetricValue::CountAndU64 { .. }) {
359                6
360            } else {
361                0
362            };
363            let name = metric_name(&row.labels, suffix_bytes)?;
364            let amount_unit = if row.labels.iter().any(|label| label == "amount_e8s") {
365                "icp_e8s"
366            } else {
367                "cycles"
368            };
369            Ok(match row.value {
370                MetricValue::Count(value) => vec![PublicMetricSample {
371                    name,
372                    canister_id: row.principal,
373                    value: u128::from(value),
374                    unit: "count".into(),
375                    observed_at_ns: 0,
376                    kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
377                }],
378                MetricValue::U128(value) => vec![PublicMetricSample {
379                    name,
380                    canister_id: row.principal,
381                    value,
382                    unit: amount_unit.into(),
383                    observed_at_ns: 0,
384                    kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
385                }],
386                MetricValue::CountAndU64 { count, value_u64 } => vec![
387                    PublicMetricSample {
388                        name: format!("{name}.count"),
389                        canister_id: row.principal,
390                        value: u128::from(count),
391                        unit: "count".into(),
392                        observed_at_ns: 0,
393                        kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
394                    },
395                    PublicMetricSample {
396                        name,
397                        canister_id: row.principal,
398                        value: u128::from(value_u64),
399                        unit: "value".into(),
400                        observed_at_ns: 0,
401                        kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
402                    },
403                ],
404            })
405        })
406        .collect::<Result<Vec<_>, InternalError>>()
407        .map(|rows| rows.into_iter().flatten().collect())
408}
409
410fn performance_metrics() -> Result<Vec<PublicMetricSample>, InternalError> {
411    metrics::bounded_performance_entries(MAX_PUBLIC_METRICS / 2 + 1)?
412        .into_iter()
413        .take(MAX_PUBLIC_METRICS / 2 + 1)
414        .map(|row| {
415            let MetricValue::CountAndU64 { count, value_u64 } = row.value else {
416                return Ok(Vec::new());
417            };
418            let name = metric_name(&row.labels, 6)?;
419            Ok(vec![
420                PublicMetricSample {
421                    name: format!("{name}.calls"),
422                    canister_id: None,
423                    value: u128::from(count),
424                    unit: "count".into(),
425                    observed_at_ns: 0,
426                    kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
427                },
428                PublicMetricSample {
429                    name,
430                    canister_id: None,
431                    value: u128::from(value_u64),
432                    unit: "instructions".into(),
433                    observed_at_ns: 0,
434                    kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
435                },
436            ])
437        })
438        .collect::<Result<Vec<_>, InternalError>>()
439        .map(|rows| rows.into_iter().flatten().collect())
440}
441
442#[cfg(feature = "sharding")]
443fn shard_metrics() -> Vec<PublicMetricSample> {
444    ShardingRegistryOps::bounded_registry_entries(MAX_PUBLIC_METRICS / 2 + 1)
445        .into_iter()
446        .flat_map(|row| {
447            vec![
448                PublicMetricSample {
449                    name: format!("{}.assigned", row.entry.pool),
450                    canister_id: Some(row.pid),
451                    value: u128::from(row.entry.count),
452                    unit: "assignments".into(),
453                    observed_at_ns: 0,
454                    kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
455                },
456                PublicMetricSample {
457                    name: format!("{}.capacity", row.entry.pool),
458                    canister_id: Some(row.pid),
459                    value: u128::from(row.entry.capacity),
460                    unit: "assignments".into(),
461                    observed_at_ns: 0,
462                    kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
463                },
464            ]
465        })
466        .collect()
467}
468#[cfg(not(feature = "sharding"))]
469const fn shard_metrics() -> Vec<PublicMetricSample> {
470    Vec::new()
471}
472
473// -----------------------------------------------------------------------------
474// Tests
475// -----------------------------------------------------------------------------
476#[cfg(test)]
477mod tests {
478    use super::*;
479    #[cfg(feature = "sharding")]
480    use crate::ids::CanisterRole;
481    use crate::model::public_metrics::MAX_PUBLIC_METRICS;
482
483    fn request(family: PublicMetricFamily) -> PublicMetricsRequest {
484        PublicMetricsRequest {
485            family,
486            page: PageRequest {
487                limit: 1_000,
488                offset: 0,
489            },
490        }
491    }
492    fn publish(
493        family: PublicMetricFamily,
494        now: u64,
495        rows: impl IntoIterator<Item = PublicMetricSample>,
496    ) -> Result<(), InternalError> {
497        PublicMetricsCache::replace(
498            family,
499            now,
500            rows.into_iter().map(|mut row| {
501                row.observed_at_ns = now;
502                row
503            }),
504        )
505    }
506
507    fn sample(value: u128) -> PublicMetricSample {
508        PublicMetricSample {
509            name: format!("assigned.{value:04}"),
510            canister_id: None,
511            value,
512            unit: "assignments".into(),
513            observed_at_ns: 0,
514            kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
515        }
516    }
517
518    #[test]
519    fn publication_is_disabled_even_when_a_cached_snapshot_exists() {
520        let family = PublicMetricFamily::Cycles;
521        publish(family, 10, vec![sample(7)]).unwrap();
522        let result = PublicMetricsOps::project(request(family), &BTreeSet::new(), 10);
523        assert_eq!(result.state, PublicSnapshotState::Disabled);
524        assert_eq!(result.sampled_at_ns, None);
525        assert!(result.metrics.entries.is_empty());
526    }
527
528    #[test]
529    fn reads_preserve_sample_time_and_report_staleness_without_refresh() {
530        let family = PublicMetricFamily::Performance;
531        let enabled = BTreeSet::from([family]);
532        let missing = PublicMetricsOps::project(request(family), &enabled, 10);
533        assert_eq!(missing.state, PublicSnapshotState::Unavailable);
534        publish(family, 10, vec![sample(3)]).unwrap();
535        let fresh = PublicMetricsOps::project(request(family), &enabled, 10);
536        assert_eq!(fresh.state, PublicSnapshotState::Fresh);
537        let stale = PublicMetricsOps::project(
538            request(family),
539            &enabled,
540            11 + PUBLIC_METRICS_STALE_AFTER_NS,
541        );
542        assert_eq!(stale.state, PublicSnapshotState::Stale);
543        assert_eq!(stale.sampled_at_ns, Some(10));
544        assert_eq!(stale.metrics.entries, fresh.metrics.entries);
545        assert_eq!(
546            PublicMetricsCache::snapshot(family).unwrap().sampled_at_ns,
547            10
548        );
549    }
550
551    #[test]
552    fn publication_selection_is_exact_and_pages_are_bounded() {
553        let family = PublicMetricFamily::ShardOccupancy;
554        let enabled = BTreeSet::from([family]);
555        publish(
556            family,
557            20,
558            (0..=MAX_PUBLIC_METRICS).map(|v| sample(v as u128)),
559        )
560        .unwrap();
561        let all = PublicMetricsOps::project(request(family), &enabled, 20);
562        assert!(all.truncated);
563        assert_eq!(all.metrics.entries.len(), MAX_PUBLIC_METRICS);
564        let mut req = request(family);
565        req.page = PageRequest {
566            limit: 2,
567            offset: 1,
568        };
569        let page = PublicMetricsOps::project(req, &enabled, 20);
570        assert_eq!(page.metrics.entries, all.metrics.entries[1..3]);
571        assert_eq!(
572            PublicMetricsOps::project(request(PublicMetricFamily::Operations), &enabled, 20).state,
573            PublicSnapshotState::Disabled
574        );
575    }
576
577    #[test]
578    #[cfg(feature = "sharding")]
579    fn shard_occupancy_samples_assignments_and_capacity_without_keys() {
580        let shard = crate::cdk::types::Principal::from_slice(&[42; 29]);
581        ShardingRegistryOps::clear_for_test();
582        ShardingRegistryOps::create(shard, "demo", 0, &CanisterRole::new("shard"), 4, 0).unwrap();
583        ShardingRegistryOps::assign("demo", "private-key-a", shard).unwrap();
584        ShardingRegistryOps::assign("demo", "private-key-b", shard).unwrap();
585        let family = PublicMetricFamily::ShardOccupancy;
586        let enabled = BTreeSet::from([family]);
587        PublicMetricsOps::sample_family(family, 10).unwrap();
588        let snapshot = PublicMetricsOps::project(request(family), &enabled, 10);
589        assert_eq!(
590            snapshot.metrics.entries,
591            vec![
592                PublicMetric {
593                    name: "demo.assigned".into(),
594                    canister_id: Some(shard),
595                    value: 2,
596                    unit: "assignments".into(),
597                    observed_at_ns: 10,
598                    kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
599                },
600                PublicMetric {
601                    name: "demo.capacity".into(),
602                    canister_id: Some(shard),
603                    value: 4,
604                    unit: "assignments".into(),
605                    observed_at_ns: 10,
606                    kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
607                },
608            ]
609        );
610        ShardingRegistryOps::release("demo", "private-key-a").unwrap();
611        let cached = PublicMetricsOps::project(request(family), &enabled, 11);
612        assert_eq!(cached.metrics.entries, snapshot.metrics.entries);
613        PublicMetricsOps::sample_family(family, 12).unwrap();
614        let refreshed = PublicMetricsOps::project(request(family), &enabled, 12);
615        assert_eq!(refreshed.metrics.entries[0].value, 1);
616        assert_eq!(refreshed.sampled_at_ns, Some(12));
617        ShardingRegistryOps::clear_for_test();
618    }
619
620    #[test]
621    fn rejected_sample_preserves_previous_snapshot() {
622        let family = PublicMetricFamily::Application;
623        publish(family, 10, vec![sample(1)]).unwrap();
624        let mut invalid = sample(2);
625        invalid.name.clear();
626        assert_eq!(
627            publish(family, 20, vec![invalid]).unwrap_err().code(),
628            crate::diagnostics::codes::REQUEST_INVALID
629        );
630        assert_eq!(
631            PublicMetricsCache::snapshot(family).unwrap().sampled_at_ns,
632            10
633        );
634    }
635    #[test]
636    fn cache_consumes_only_one_bounded_prefix_and_reports_truncation() {
637        let consumed = std::cell::Cell::new(0);
638        publish(
639            PublicMetricFamily::Application,
640            10,
641            (0..).map(|value| {
642                consumed.set(consumed.get() + 1);
643                sample(value)
644            }),
645        )
646        .unwrap();
647        assert_eq!(consumed.get(), MAX_PUBLIC_METRICS + 1);
648        let snapshot = PublicMetricsCache::snapshot(PublicMetricFamily::Application).unwrap();
649        assert!(snapshot.truncated);
650        assert_eq!(snapshot.metrics.len(), MAX_PUBLIC_METRICS);
651    }
652
653    #[test]
654    fn performance_sampling_is_bounded_and_independent_of_recording_order() {
655        let family = PublicMetricFamily::Performance;
656        crate::perf::reset();
657        for value in (0..1024).rev() {
658            crate::perf::record_checkpoint("bounded", &format!("sample_{value:04}"), value);
659        }
660        PublicMetricsOps::sample_family(family, 10).unwrap();
661        let first = PublicMetricsOps::project(request(family), &BTreeSet::from([family]), 10);
662        assert!(first.truncated);
663        assert_eq!(first.metrics.entries.len(), MAX_PUBLIC_METRICS);
664        assert_eq!(
665            first.metrics.entries[0].name,
666            "perf.checkpoint.bounded.sample_0000"
667        );
668        crate::perf::reset();
669        for value in 0..1024 {
670            crate::perf::record_checkpoint("bounded", &format!("sample_{value:04}"), value);
671        }
672        PublicMetricsOps::sample_family(family, 20).unwrap();
673        let second = PublicMetricsOps::project(request(family), &BTreeSet::from([family]), 20);
674        for (first, second) in first.metrics.entries.iter().zip(&second.metrics.entries) {
675            assert_eq!(first.name, second.name);
676            assert_eq!(first.value, second.value);
677            assert_eq!(first.kind, second.kind);
678            assert_eq!(first.observed_at_ns, 10);
679            assert_eq!(second.observed_at_ns, 20);
680        }
681        crate::perf::reset();
682    }
683
684    #[test]
685    #[cfg(feature = "sharding")]
686    fn shard_sampling_bounds_registry_visits_and_retained_rows() {
687        ShardingRegistryOps::clear_for_test();
688        for value in 0_u32..300 {
689            let shard = crate::cdk::types::Principal::from_slice(&value.to_be_bytes());
690            ShardingRegistryOps::create(shard, "bounded", value, &CanisterRole::new("shard"), 4, 0)
691                .unwrap();
692        }
693        assert_eq!(
694            ShardingRegistryOps::bounded_registry_entries(129).len(),
695            129
696        );
697        PublicMetricsOps::sample_family(PublicMetricFamily::ShardOccupancy, 10).unwrap();
698        let snapshot = PublicMetricsCache::snapshot(PublicMetricFamily::ShardOccupancy).unwrap();
699        assert!(snapshot.truncated);
700        assert_eq!(snapshot.metrics.len(), MAX_PUBLIC_METRICS);
701        ShardingRegistryOps::clear_for_test();
702    }
703}
704
705#[cfg(test)]
706mod history_tests {
707    use super::*;
708
709    #[test]
710    fn history_reads_bound_pages_hide_disabled_data_and_expire_without_mutation() {
711        let family = PublicMetricFamily::Cycles;
712        for slot in [1, 2, 5] {
713            PublicMetricsCache::replace(
714                family,
715                slot * PUBLIC_METRICS_CADENCE_NS,
716                [PublicMetricSample {
717                    name: "balance".into(),
718                    canister_id: None,
719                    value: slot.into(),
720                    unit: "cycles".into(),
721                    observed_at_ns: slot * PUBLIC_METRICS_CADENCE_NS,
722                    kind: PublicMetricKind::Gauge,
723                }],
724            )
725            .unwrap();
726        }
727        let request = PublicHistoryRequest {
728            family,
729            name: "balance".into(),
730            canister_id: None,
731            page: PageRequest {
732                offset: 1,
733                limit: u64::MAX,
734            },
735        };
736        let enabled = BTreeSet::from([family]);
737        let view = PublicMetricsOps::project_history(
738            request.clone(),
739            &enabled,
740            5 * PUBLIC_METRICS_CADENCE_NS,
741        );
742        assert_eq!(view.points.total, 3);
743        assert_eq!(
744            view.points
745                .entries
746                .iter()
747                .map(|point| point.value)
748                .collect::<Vec<_>>(),
749            [2, 5]
750        );
751        assert_eq!(view.coverage_start_ns, Some(PUBLIC_METRICS_CADENCE_NS));
752        let disabled = PublicMetricsOps::project_history(
753            request.clone(),
754            &BTreeSet::new(),
755            5 * PUBLIC_METRICS_CADENCE_NS,
756        );
757        assert_eq!(disabled.state, PublicSnapshotState::Disabled);
758        assert!(disabled.points.entries.is_empty());
759        assert_eq!(disabled.reserved_bytes, 0);
760        let expired =
761            PublicMetricsOps::project_history(request, &enabled, 400 * PUBLIC_METRICS_CADENCE_NS);
762        assert_eq!(expired.state, PublicSnapshotState::Unavailable);
763        assert!(expired.points.entries.is_empty());
764        assert_eq!(expired.reserved_bytes, view.reserved_bytes);
765        assert_eq!(
766            PublicMetricsCache::snapshot(family).unwrap().sampled_at_ns,
767            5 * PUBLIC_METRICS_CADENCE_NS
768        );
769    }
770}
771
772#[cfg(test)]
773mod counter_tests {
774    use super::*;
775    use crate::model::public_metrics::PublicHistorySample;
776
777    #[test]
778    fn public_metrics_counter_deltas_require_adjacent_unsaturated_same_window_observations() {
779        let first = PublicHistorySample {
780            slot: 1,
781            observed_at_ns: 10,
782            value: 7,
783            kind: PublicMetricKind::Counter {
784                window_id: 4,
785                saturated: false,
786            },
787        };
788        let second = PublicHistorySample {
789            slot: 2,
790            observed_at_ns: 20,
791            value: 12,
792            ..first
793        };
794        assert_eq!(
795            counter_delta(&first, &second),
796            Some(PublicCounterDelta {
797                amount: 5,
798                elapsed_ns: 10
799            })
800        );
801        for incompatible in [
802            PublicHistorySample {
803                kind: PublicMetricKind::Gauge,
804                ..second
805            },
806            PublicHistorySample {
807                kind: PublicMetricKind::Counter {
808                    window_id: 5,
809                    saturated: false,
810                },
811                ..second
812            },
813            PublicHistorySample {
814                kind: PublicMetricKind::Counter {
815                    window_id: 4,
816                    saturated: true,
817                },
818                ..second
819            },
820            PublicHistorySample { slot: 3, ..second },
821            PublicHistorySample {
822                observed_at_ns: 10,
823                ..second
824            },
825            PublicHistorySample { value: 1, ..second },
826        ] {
827            assert_eq!(counter_delta(&first, &incompatible), None);
828        }
829    }
830}