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