1mod 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#[derive(Clone, Copy)]
45pub struct ApplicationMetricsSampler {
46 collect: fn() -> Result<Vec<PublicMetric>, crate::dto::error::Error>,
47}
48
49impl ApplicationMetricsSampler {
50 #[must_use]
52 pub const fn new(collect: fn() -> Result<Vec<PublicMetric>, crate::dto::error::Error>) -> Self {
53 Self { collect }
54 }
55}
56
57pub struct PublicMetricsOps;
59
60impl PublicMetricsOps {
61 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 pub fn expire_history(now_ns: u64) {
141 PublicHistoryCache::expire(now_ns);
142 }
143
144 #[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 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 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
340fn 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#[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}