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