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 (unit, mut points) = series.map_or_else(
164            || (None, Vec::new()),
165            |series| (Some(series.unit), series.slots),
166        );
167        points
168            .retain(|point| point.slot <= slot && slot - point.slot < PUBLIC_HISTORY_SLOTS as u64);
169        let state = if !selected {
170            PublicSnapshotState::Disabled
171        } else if let Some(point) = points.last() {
172            if now_ns.saturating_sub(point.observed_at_ns) > PUBLIC_METRICS_STALE_AFTER_NS {
173                PublicSnapshotState::Stale
174            } else {
175                PublicSnapshotState::Fresh
176            }
177        } else {
178            PublicSnapshotState::Unavailable
179        };
180        let coverage_start_ns = points
181            .first()
182            .map(|point| point.slot * PUBLIC_METRICS_CADENCE_NS);
183        let total = points.len() as u64;
184        let entries = points
185            .iter()
186            .enumerate()
187            .map(|(index, point)| {
188                let delta = index
189                    .checked_sub(1)
190                    .and_then(|previous| counter_delta(&points[previous], point));
191                PublicHistoryPoint {
192                    delta,
193                    slot_start_ns: point.slot * PUBLIC_METRICS_CADENCE_NS,
194                    observed_at_ns: point.observed_at_ns,
195                    value: point.value,
196                    kind: point.kind,
197                }
198            })
199            .skip(usize::try_from(request.page.offset.min(total)).unwrap_or(PUBLIC_HISTORY_SLOTS))
200            .take(
201                usize::try_from(request.page.limit.min(PUBLIC_HISTORY_SLOTS as u64))
202                    .unwrap_or(PUBLIC_HISTORY_SLOTS),
203            )
204            .collect();
205        PublicHistorySnapshot {
206            state,
207            unit,
208            heap_started_at_ns: selected
209                .then(PublicHistoryCache::heap_started_at_ns)
210                .flatten(),
211            canister_version: 0,
212            coverage_start_ns,
213            cadence_ns: PUBLIC_METRICS_CADENCE_NS,
214            retention_ns: PUBLIC_HISTORY_RETENTION_NS,
215            stale_after_ns: PUBLIC_METRICS_STALE_AFTER_NS,
216            truncated: selected && PublicHistoryCache::truncated(),
217            series_limit: MAX_HISTORY_SERIES as u64,
218            byte_limit: MAX_HISTORY_BYTES as u64,
219            reserved_bytes: if selected {
220                PublicHistoryCache::reserved_bytes() as u64
221            } else {
222                0
223            },
224            points: Page { entries, total },
225        }
226    }
227
228    pub fn record_application(metrics: Vec<PublicMetric>) -> Result<(), InternalError> {
229        if !Self::enabled().contains(&PublicMetricFamily::Application) {
230            return Ok(());
231        }
232        let rows = metrics.into_iter().map(|row| PublicMetricSample {
233            name: row.name,
234            canister_id: row.canister_id,
235            value: row.value,
236            unit: row.unit,
237            observed_at_ns: row.observed_at_ns,
238            kind: row.kind,
239        });
240        PublicMetricsCache::replace(PublicMetricFamily::Application, IcOps::now_nanos(), rows)
241    }
242
243    pub fn sample_family(family: PublicMetricFamily, now: u64) -> Result<(), InternalError> {
244        let mut rows = match family {
245            PublicMetricFamily::Application => {
246                if let Some(sample) = APPLICATION_SAMPLER.get() {
247                    let metrics = (sample.collect)().map_err(|_| InternalError::invalid_input())?;
248                    return Self::record_application(metrics);
249                }
250                return Ok(());
251            }
252            PublicMetricFamily::Cycles => vec![PublicMetricSample {
253                name: "balance".into(),
254                canister_id: Some(IcOps::canister_self()),
255                value: IcOps::canister_cycle_balance().to_u128(),
256                unit: "cycles".into(),
257                observed_at_ns: 0,
258                kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
259            }],
260            PublicMetricFamily::Operations => {
261                let mut rows = process::operations()?;
262                rows.extend(operation_metrics()?);
263                rows
264            }
265            PublicMetricFamily::Performance => {
266                let mut rows = process::memory();
267                rows.extend(performance_metrics()?);
268                rows
269            }
270            PublicMetricFamily::ShardOccupancy => shard_metrics(),
271        };
272        for row in &mut rows {
273            row.observed_at_ns = now;
274            // Preserve source-owned timer registrations; aggregate timer events remain gauges.
275            let counter = (family == PublicMetricFamily::Operations
276                && !row.name.starts_with("cycles_funding.icp_refill.")
277                && !row.name.starts_with("timer."))
278                || (family == PublicMetricFamily::Performance
279                    && !row.name.starts_with("perf.timer.")
280                    && !row.name.starts_with("memory."));
281            if counter {
282                row.kind = PublicMetricKind::Counter {
283                    window_id: 0,
284                    saturated: if matches!(row.unit.as_str(), "cycles" | "icp_e8s") {
285                        row.value == u128::MAX
286                    } else {
287                        row.value == u128::from(u64::MAX)
288                    },
289                };
290            }
291        }
292        if family == PublicMetricFamily::Performance {
293            // Allocation rows own their source times, including retained values on failure.
294            rows.splice(0..0, memory::sample(now));
295        }
296        PublicMetricsCache::replace(family, now, rows)
297    }
298}
299
300fn counter_delta(
301    previous: &crate::model::public_metrics::PublicHistorySample,
302    current: &crate::model::public_metrics::PublicHistorySample,
303) -> Option<PublicCounterDelta> {
304    let PublicMetricKind::Counter {
305        window_id,
306        saturated: false,
307    } = previous.kind
308    else {
309        return None;
310    };
311    if current.kind
312        != (PublicMetricKind::Counter {
313            window_id,
314            saturated: false,
315        })
316        || previous.slot.checked_add(1) != Some(current.slot)
317    {
318        return None;
319    }
320    let elapsed_ns = current
321        .observed_at_ns
322        .checked_sub(previous.observed_at_ns)
323        .filter(|elapsed| *elapsed > 0)?;
324    Some(PublicCounterDelta {
325        amount: current.value.checked_sub(previous.value)?,
326        elapsed_ns,
327    })
328}
329
330fn page(rows: Vec<PublicMetric>, request: PageRequest) -> Page<PublicMetric> {
331    let total = u64::try_from(rows.len()).unwrap_or(u64::MAX);
332    let start = usize::try_from(request.offset.min(total)).unwrap_or(rows.len());
333    let limit = usize::try_from(request.limit.min(total)).unwrap_or(rows.len());
334    Page {
335        entries: rows.into_iter().skip(start).take(limit).collect(),
336        total,
337    }
338}
339
340// Validate the complete series identity before allocating its formatted name.
341fn metric_name(labels: &[String], suffix_bytes: usize) -> Result<String, InternalError> {
342    let bytes = labels.iter().try_fold(
343        suffix_bytes + labels.len().saturating_sub(1),
344        |bytes, label| bytes.checked_add(label.len()),
345    );
346    if bytes.is_none_or(|bytes| bytes > MAX_PUBLIC_METRIC_TEXT_BYTES) {
347        return Err(InternalError::invalid_input());
348    }
349    Ok(labels.join("."))
350}
351
352fn operation_metrics() -> Result<Vec<PublicMetricSample>, InternalError> {
353    metrics::bounded_core_entries(MAX_PUBLIC_METRICS + 1)?
354        .into_iter()
355        .take(MAX_PUBLIC_METRICS + 1)
356        .map(|row| {
357            let suffix_bytes = if matches!(&row.value, MetricValue::CountAndU64 { .. }) {
358                6
359            } else {
360                0
361            };
362            let name = metric_name(&row.labels, suffix_bytes)?;
363            let amount_unit = if row.labels.iter().any(|label| label == "amount_e8s") {
364                "icp_e8s"
365            } else {
366                "cycles"
367            };
368            Ok(match row.value {
369                MetricValue::Count(value) => vec![PublicMetricSample {
370                    name,
371                    canister_id: row.principal,
372                    value: u128::from(value),
373                    unit: "count".into(),
374                    observed_at_ns: 0,
375                    kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
376                }],
377                MetricValue::U128(value) => vec![PublicMetricSample {
378                    name,
379                    canister_id: row.principal,
380                    value,
381                    unit: amount_unit.into(),
382                    observed_at_ns: 0,
383                    kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
384                }],
385                MetricValue::CountAndU64 { count, value_u64 } => vec![
386                    PublicMetricSample {
387                        name: format!("{name}.count"),
388                        canister_id: row.principal,
389                        value: u128::from(count),
390                        unit: "count".into(),
391                        observed_at_ns: 0,
392                        kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
393                    },
394                    PublicMetricSample {
395                        name,
396                        canister_id: row.principal,
397                        value: u128::from(value_u64),
398                        unit: "value".into(),
399                        observed_at_ns: 0,
400                        kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
401                    },
402                ],
403            })
404        })
405        .collect::<Result<Vec<_>, InternalError>>()
406        .map(|rows| rows.into_iter().flatten().collect())
407}
408
409fn performance_metrics() -> Result<Vec<PublicMetricSample>, InternalError> {
410    let inventory = ic_timers::timer_inventory().ok();
411    let timers = inventory
412        .as_ref()
413        .map_or(&[][..], |inventory| inventory.timers());
414    metrics::bounded_performance_entries(MAX_PUBLIC_METRICS / 2 + 1, timers)?
415        .into_iter()
416        .take(MAX_PUBLIC_METRICS / 2 + 1)
417        .map(|row| {
418            let MetricValue::CountAndU64 { count, value_u64 } = row.value else {
419                return Ok(Vec::new());
420            };
421            let name = metric_name(&row.labels, 6)?;
422            let registration = timers.iter().find_map(|timer| {
423                let identity = timer.identity();
424                let labels = &row.labels;
425                let matches = labels.len() == 6
426                    && labels[..5].iter().map(String::as_str).eq([
427                        "perf",
428                        "timer",
429                        identity.owner(),
430                        identity.subsystem(),
431                        identity.name(),
432                    ]);
433                matches.then(|| timer.registration_id())
434            });
435            Ok(vec![
436                PublicMetricSample {
437                    name: format!("{name}.calls"),
438                    canister_id: None,
439                    value: u128::from(count),
440                    unit: "count".into(),
441                    observed_at_ns: 0,
442                    kind: timer_measurement_kind(registration, count),
443                },
444                PublicMetricSample {
445                    name,
446                    canister_id: None,
447                    value: u128::from(value_u64),
448                    unit: "instructions".into(),
449                    observed_at_ns: 0,
450                    kind: timer_measurement_kind(registration, value_u64),
451                },
452            ])
453        })
454        .collect::<Result<Vec<_>, InternalError>>()
455        .map(|rows| rows.into_iter().flatten().collect())
456}
457
458fn timer_measurement_kind(
459    registration: Option<ic_timers::TimerRegistrationId>,
460    value: u64,
461) -> PublicMetricKind {
462    registration.map_or(PublicMetricKind::Gauge, |id| {
463        PublicMetricKind::TimerCounter {
464            registration: crate::domain::public_metrics::TimerMetricRegistration {
465                canister_version: id.epoch().canister_version(),
466                started_at_ns: id.epoch().started_at_ns(),
467                sequence: id.sequence(),
468            },
469            saturated: value == u64::MAX,
470        }
471    })
472}
473
474#[cfg(feature = "sharding")]
475fn shard_metrics() -> Vec<PublicMetricSample> {
476    ShardingRegistryOps::bounded_registry_entries(MAX_PUBLIC_METRICS / 2 + 1)
477        .into_iter()
478        .flat_map(|row| {
479            vec![
480                PublicMetricSample {
481                    name: format!("{}.assigned", row.entry.pool),
482                    canister_id: Some(row.pid),
483                    value: u128::from(row.entry.count),
484                    unit: "assignments".into(),
485                    observed_at_ns: 0,
486                    kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
487                },
488                PublicMetricSample {
489                    name: format!("{}.capacity", row.entry.pool),
490                    canister_id: Some(row.pid),
491                    value: u128::from(row.entry.capacity),
492                    unit: "assignments".into(),
493                    observed_at_ns: 0,
494                    kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
495                },
496            ]
497        })
498        .collect()
499}
500#[cfg(not(feature = "sharding"))]
501const fn shard_metrics() -> Vec<PublicMetricSample> {
502    Vec::new()
503}
504
505// -----------------------------------------------------------------------------
506// Tests
507// -----------------------------------------------------------------------------
508#[cfg(test)]
509mod tests {
510    use super::*;
511    #[cfg(feature = "sharding")]
512    use crate::ids::CanisterRole;
513    use crate::model::public_metrics::MAX_PUBLIC_METRICS;
514
515    fn request(family: PublicMetricFamily) -> PublicMetricsRequest {
516        PublicMetricsRequest {
517            family,
518            page: PageRequest {
519                limit: 1_000,
520                offset: 0,
521            },
522        }
523    }
524    fn publish(
525        family: PublicMetricFamily,
526        now: u64,
527        rows: impl IntoIterator<Item = PublicMetricSample>,
528    ) -> Result<(), InternalError> {
529        PublicMetricsCache::replace(
530            family,
531            now,
532            rows.into_iter().map(|mut row| {
533                row.observed_at_ns = now;
534                row
535            }),
536        )
537    }
538
539    fn sample(value: u128) -> PublicMetricSample {
540        PublicMetricSample {
541            name: format!("assigned.{value:04}"),
542            canister_id: None,
543            value,
544            unit: "assignments".into(),
545            observed_at_ns: 0,
546            kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
547        }
548    }
549
550    #[test]
551    fn publication_is_disabled_even_when_a_cached_snapshot_exists() {
552        let family = PublicMetricFamily::Cycles;
553        publish(family, 10, vec![sample(7)]).unwrap();
554        let result = PublicMetricsOps::project(request(family), &BTreeSet::new(), 10);
555        assert_eq!(result.state, PublicSnapshotState::Disabled);
556        assert_eq!(result.sampled_at_ns, None);
557        assert!(result.metrics.entries.is_empty());
558    }
559
560    #[test]
561    fn reads_preserve_sample_time_and_report_staleness_without_refresh() {
562        let family = PublicMetricFamily::Performance;
563        let enabled = BTreeSet::from([family]);
564        let missing = PublicMetricsOps::project(request(family), &enabled, 10);
565        assert_eq!(missing.state, PublicSnapshotState::Unavailable);
566        publish(family, 10, vec![sample(3)]).unwrap();
567        let fresh = PublicMetricsOps::project(request(family), &enabled, 10);
568        assert_eq!(fresh.state, PublicSnapshotState::Fresh);
569        let stale = PublicMetricsOps::project(
570            request(family),
571            &enabled,
572            11 + PUBLIC_METRICS_STALE_AFTER_NS,
573        );
574        assert_eq!(stale.state, PublicSnapshotState::Stale);
575        assert_eq!(stale.sampled_at_ns, Some(10));
576        assert_eq!(stale.metrics.entries, fresh.metrics.entries);
577        assert_eq!(
578            PublicMetricsCache::snapshot(family).unwrap().sampled_at_ns,
579            10
580        );
581    }
582
583    #[test]
584    fn public_timer_measurements_preserve_registration_without_synthesizing_history_rates() {
585        let kind = PublicMetricKind::TimerCounter {
586            registration: crate::domain::public_metrics::TimerMetricRegistration {
587                canister_version: 1,
588                started_at_ns: 5,
589                sequence: 9,
590            },
591            saturated: false,
592        };
593        let family = PublicMetricFamily::Performance;
594        let mut row = sample(7);
595        row.kind = kind;
596        publish(family, 10, vec![row]).unwrap();
597        let snapshot = PublicMetricsOps::project(request(family), &BTreeSet::from([family]), 10);
598        assert_eq!(snapshot.metrics.entries[0].kind, kind);
599        let before = crate::model::public_metrics::PublicHistorySample {
600            slot: 1,
601            observed_at_ns: 10,
602            value: 7,
603            kind,
604        };
605        let after = crate::model::public_metrics::PublicHistorySample {
606            slot: 2,
607            observed_at_ns: 20,
608            value: 12,
609            kind,
610        };
611        assert_eq!(counter_delta(&before, &after), None);
612    }
613
614    #[test]
615    fn publication_selection_is_exact_and_pages_are_bounded() {
616        let family = PublicMetricFamily::ShardOccupancy;
617        let enabled = BTreeSet::from([family]);
618        publish(
619            family,
620            20,
621            (0..=MAX_PUBLIC_METRICS).map(|v| sample(v as u128)),
622        )
623        .unwrap();
624        let all = PublicMetricsOps::project(request(family), &enabled, 20);
625        assert!(all.truncated);
626        assert_eq!(all.metrics.entries.len(), MAX_PUBLIC_METRICS);
627        let mut req = request(family);
628        req.page = PageRequest {
629            limit: 2,
630            offset: 1,
631        };
632        let page = PublicMetricsOps::project(req, &enabled, 20);
633        assert_eq!(page.metrics.entries, all.metrics.entries[1..3]);
634        assert_eq!(
635            PublicMetricsOps::project(request(PublicMetricFamily::Operations), &enabled, 20).state,
636            PublicSnapshotState::Disabled
637        );
638    }
639
640    #[test]
641    #[cfg(feature = "sharding")]
642    fn shard_occupancy_samples_assignments_and_capacity_without_keys() {
643        let shard = crate::cdk::types::Principal::from_slice(&[42; 29]);
644        ShardingRegistryOps::clear_for_test();
645        ShardingRegistryOps::create(shard, "demo", 0, &CanisterRole::new("shard"), 4, 0).unwrap();
646        ShardingRegistryOps::assign("demo", "private-key-a", shard).unwrap();
647        ShardingRegistryOps::assign("demo", "private-key-b", shard).unwrap();
648        let family = PublicMetricFamily::ShardOccupancy;
649        let enabled = BTreeSet::from([family]);
650        PublicMetricsOps::sample_family(family, 10).unwrap();
651        let snapshot = PublicMetricsOps::project(request(family), &enabled, 10);
652        assert_eq!(
653            snapshot.metrics.entries,
654            vec![
655                PublicMetric {
656                    name: "demo.assigned".into(),
657                    canister_id: Some(shard),
658                    value: 2,
659                    unit: "assignments".into(),
660                    observed_at_ns: 10,
661                    kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
662                },
663                PublicMetric {
664                    name: "demo.capacity".into(),
665                    canister_id: Some(shard),
666                    value: 4,
667                    unit: "assignments".into(),
668                    observed_at_ns: 10,
669                    kind: crate::domain::public_metrics::PublicMetricKind::Gauge,
670                },
671            ]
672        );
673        ShardingRegistryOps::release("demo", "private-key-a").unwrap();
674        let cached = PublicMetricsOps::project(request(family), &enabled, 11);
675        assert_eq!(cached.metrics.entries, snapshot.metrics.entries);
676        PublicMetricsOps::sample_family(family, 12).unwrap();
677        let refreshed = PublicMetricsOps::project(request(family), &enabled, 12);
678        assert_eq!(refreshed.metrics.entries[0].value, 1);
679        assert_eq!(refreshed.sampled_at_ns, Some(12));
680        ShardingRegistryOps::clear_for_test();
681    }
682
683    #[test]
684    fn rejected_sample_preserves_previous_snapshot() {
685        let family = PublicMetricFamily::Application;
686        publish(family, 10, vec![sample(1)]).unwrap();
687        let mut invalid = sample(2);
688        invalid.name.clear();
689        assert_eq!(
690            publish(family, 20, vec![invalid]).unwrap_err().code(),
691            crate::diagnostics::codes::REQUEST_INVALID
692        );
693        assert_eq!(
694            PublicMetricsCache::snapshot(family).unwrap().sampled_at_ns,
695            10
696        );
697    }
698    #[test]
699    fn cache_consumes_only_one_bounded_prefix_and_reports_truncation() {
700        let consumed = std::cell::Cell::new(0);
701        publish(
702            PublicMetricFamily::Application,
703            10,
704            (0..).map(|value| {
705                consumed.set(consumed.get() + 1);
706                sample(value)
707            }),
708        )
709        .unwrap();
710        assert_eq!(consumed.get(), MAX_PUBLIC_METRICS + 1);
711        let snapshot = PublicMetricsCache::snapshot(PublicMetricFamily::Application).unwrap();
712        assert!(snapshot.truncated);
713        assert_eq!(snapshot.metrics.len(), MAX_PUBLIC_METRICS);
714    }
715
716    #[test]
717    fn performance_sampling_is_bounded_and_independent_of_recording_order() {
718        let family = PublicMetricFamily::Performance;
719        crate::perf::reset();
720        for value in (0..1024).rev() {
721            crate::perf::record_checkpoint("bounded", &format!("sample_{value:04}"), value);
722        }
723        PublicMetricsOps::sample_family(family, 10).unwrap();
724        let first = PublicMetricsOps::project(request(family), &BTreeSet::from([family]), 10);
725        assert!(first.truncated);
726        assert_eq!(first.metrics.entries.len(), MAX_PUBLIC_METRICS);
727        assert_eq!(
728            first.metrics.entries[0].name,
729            "perf.checkpoint.bounded.sample_0000"
730        );
731        crate::perf::reset();
732        for value in 0..1024 {
733            crate::perf::record_checkpoint("bounded", &format!("sample_{value:04}"), value);
734        }
735        PublicMetricsOps::sample_family(family, 20).unwrap();
736        let second = PublicMetricsOps::project(request(family), &BTreeSet::from([family]), 20);
737        for (first, second) in first.metrics.entries.iter().zip(&second.metrics.entries) {
738            assert_eq!(first.name, second.name);
739            assert_eq!(first.value, second.value);
740            assert_eq!(first.kind, second.kind);
741            assert_eq!(first.observed_at_ns, 10);
742            assert_eq!(second.observed_at_ns, 20);
743        }
744        crate::perf::reset();
745    }
746
747    #[test]
748    #[cfg(feature = "sharding")]
749    fn shard_sampling_bounds_registry_visits_and_retained_rows() {
750        ShardingRegistryOps::clear_for_test();
751        for value in 0_u32..300 {
752            let shard = crate::cdk::types::Principal::from_slice(&value.to_be_bytes());
753            ShardingRegistryOps::create(shard, "bounded", value, &CanisterRole::new("shard"), 4, 0)
754                .unwrap();
755        }
756        assert_eq!(
757            ShardingRegistryOps::bounded_registry_entries(129).len(),
758            129
759        );
760        PublicMetricsOps::sample_family(PublicMetricFamily::ShardOccupancy, 10).unwrap();
761        let snapshot = PublicMetricsCache::snapshot(PublicMetricFamily::ShardOccupancy).unwrap();
762        assert!(snapshot.truncated);
763        assert_eq!(snapshot.metrics.len(), MAX_PUBLIC_METRICS);
764        ShardingRegistryOps::clear_for_test();
765    }
766}
767
768#[cfg(test)]
769mod history_tests {
770    use super::*;
771
772    #[test]
773    fn history_reads_bound_pages_hide_disabled_data_and_expire_without_mutation() {
774        let family = PublicMetricFamily::Cycles;
775        for slot in [1, 2, 5] {
776            PublicMetricsCache::replace(
777                family,
778                slot * PUBLIC_METRICS_CADENCE_NS,
779                [PublicMetricSample {
780                    name: "balance".into(),
781                    canister_id: None,
782                    value: slot.into(),
783                    unit: "cycles".into(),
784                    observed_at_ns: slot * PUBLIC_METRICS_CADENCE_NS,
785                    kind: PublicMetricKind::Gauge,
786                }],
787            )
788            .unwrap();
789        }
790        let request = PublicHistoryRequest {
791            family,
792            name: "balance".into(),
793            canister_id: None,
794            page: PageRequest {
795                offset: 1,
796                limit: u64::MAX,
797            },
798        };
799        let enabled = BTreeSet::from([family]);
800        let view = PublicMetricsOps::project_history(
801            request.clone(),
802            &enabled,
803            5 * PUBLIC_METRICS_CADENCE_NS,
804        );
805        assert_eq!(view.points.total, 3);
806        assert_eq!(
807            view.points
808                .entries
809                .iter()
810                .map(|point| point.value)
811                .collect::<Vec<_>>(),
812            [2, 5]
813        );
814        assert_eq!(view.coverage_start_ns, Some(PUBLIC_METRICS_CADENCE_NS));
815        let disabled = PublicMetricsOps::project_history(
816            request.clone(),
817            &BTreeSet::new(),
818            5 * PUBLIC_METRICS_CADENCE_NS,
819        );
820        assert_eq!(disabled.state, PublicSnapshotState::Disabled);
821        assert!(disabled.points.entries.is_empty());
822        assert_eq!(disabled.reserved_bytes, 0);
823        let expired =
824            PublicMetricsOps::project_history(request, &enabled, 400 * PUBLIC_METRICS_CADENCE_NS);
825        assert_eq!(expired.state, PublicSnapshotState::Unavailable);
826        assert!(expired.points.entries.is_empty());
827        assert_eq!(expired.reserved_bytes, view.reserved_bytes);
828        assert_eq!(
829            PublicMetricsCache::snapshot(family).unwrap().sampled_at_ns,
830            5 * PUBLIC_METRICS_CADENCE_NS
831        );
832    }
833}
834
835#[cfg(test)]
836mod counter_tests {
837    use super::*;
838    use crate::model::public_metrics::PublicHistorySample;
839
840    #[test]
841    fn public_metrics_counter_deltas_require_adjacent_unsaturated_same_window_observations() {
842        let first = PublicHistorySample {
843            slot: 1,
844            observed_at_ns: 10,
845            value: 7,
846            kind: PublicMetricKind::Counter {
847                window_id: 4,
848                saturated: false,
849            },
850        };
851        let second = PublicHistorySample {
852            slot: 2,
853            observed_at_ns: 20,
854            value: 12,
855            ..first
856        };
857        assert_eq!(
858            counter_delta(&first, &second),
859            Some(PublicCounterDelta {
860                amount: 5,
861                elapsed_ns: 10
862            })
863        );
864        for incompatible in [
865            PublicHistorySample {
866                kind: PublicMetricKind::Gauge,
867                ..second
868            },
869            PublicHistorySample {
870                kind: PublicMetricKind::Counter {
871                    window_id: 5,
872                    saturated: false,
873                },
874                ..second
875            },
876            PublicHistorySample {
877                kind: PublicMetricKind::Counter {
878                    window_id: 4,
879                    saturated: true,
880                },
881                ..second
882            },
883            PublicHistorySample { slot: 3, ..second },
884            PublicHistorySample {
885                observed_at_ns: 10,
886                ..second
887            },
888            PublicHistorySample { value: 1, ..second },
889        ] {
890            assert_eq!(counter_delta(&first, &incompatible), None);
891        }
892    }
893}