use crate::metrics::{
MetricsCommandError,
model::{
MetricEntry, MetricValue, MetricsCanisterReport, MetricsCanisterStatus, MetricsKind,
MetricsReport,
},
options::MetricsOptions,
parse::metric_page,
};
use canic_core::dto::observability::{CanisterObservabilityRequest, CanisterObservabilityResponse};
use canic_host::{
CanisterProtocolError,
fleet_ensure::{CurrentFleetResolution, resolve_current_fleet},
icp::{IcpCli, IcpDiagnostic},
icp_config::resolve_current_canic_icp_root,
observability::{FleetObservabilityError, observe_fleet_canister},
registry::RegistryEntry,
};
use std::{path::Path, sync::Arc, thread};
use thiserror::Error as ThisError;
const METRICS_UNAVAILABLE_HINT: &str =
"role-owned Metrics status unavailable; check deployed Wasm and metrics profile";
const METRICS_EMPTY_HINT: &str =
"no metrics rows; check whether this tier is enabled by the deployed role profile";
const METRICS_NONZERO_EMPTY_HINT: &str =
"no nonzero metrics rows; rerun without --nonzero or check the deployed role profile";
const METRICS_WORKER_PANIC: &str = "metrics query worker panicked";
#[derive(Debug, ThisError)]
enum MetricsQueryError {
#[error(transparent)]
Observability(#[from] FleetObservabilityError),
}
pub(super) fn metrics_report(
options: &MetricsOptions,
) -> Result<MetricsReport, MetricsCommandError> {
let root = resolve_current_canic_icp_root().map_err(MetricsCommandError::IcpRoot)?;
let fleet = load_registry(options)?;
let canisters = collect_metrics_reports(options, &fleet, &root);
Ok(MetricsReport {
fleet: options.fleet.clone(),
environment: options.environment.clone(),
kind: options.kind,
canisters,
})
}
fn load_registry(options: &MetricsOptions) -> Result<CurrentFleetResolution, MetricsCommandError> {
resolve_metrics_fleet(options)
}
fn matches_metrics_filter(options: &MetricsOptions, entry: &RegistryEntry) -> bool {
if let Some(role) = &options.role
&& entry.role.as_deref() != Some(role.as_str())
{
return false;
}
if let Some(canister) = &options.canister
&& entry.pid != *canister
{
return false;
}
true
}
fn collect_metrics_reports(
options: &MetricsOptions,
fleet: &CurrentFleetResolution,
icp_root: &Path,
) -> Vec<MetricsCanisterReport> {
let query = Arc::new(options.clone());
let fleet = Arc::new(fleet.clone());
let mut handles = Vec::new();
for entry in fleet
.registry
.entries
.iter()
.filter(|entry| matches_metrics_filter(options, entry))
{
let entry = entry.clone();
let worker_entry = entry.clone();
let query = Arc::clone(&query);
let fleet = Arc::clone(&fleet);
let icp_root = icp_root.to_path_buf();
handles.push((
worker_entry,
thread::spawn(move || metrics_canister_report(&query, &fleet, &icp_root, &entry)),
));
}
collect_metrics_worker_reports(handles)
}
fn collect_metrics_worker_reports(
handles: Vec<(RegistryEntry, thread::JoinHandle<MetricsCanisterReport>)>,
) -> Vec<MetricsCanisterReport> {
handles
.into_iter()
.map(|(entry, handle)| {
handle
.join()
.unwrap_or_else(|_| metrics_message_error_report(&entry, METRICS_WORKER_PANIC))
})
.collect()
}
fn metrics_canister_report(
options: &MetricsOptions,
fleet: &CurrentFleetResolution,
icp_root: &Path,
entry: &RegistryEntry,
) -> MetricsCanisterReport {
match query_metrics(options, fleet, icp_root, entry) {
Ok(mut entries) => {
if options.nonzero {
entries.retain(|entry| !metric_value_is_zero(&entry.value));
}
if entries.is_empty() {
return metrics_empty_report(entry, options.nonzero);
}
MetricsCanisterReport {
role: entry.role.clone().unwrap_or_else(|| "-".to_string()),
canister_id: entry.pid.clone(),
status: MetricsCanisterStatus::Ok,
entries,
error: None,
}
}
Err(error) => metrics_query_error_report(entry, &error),
}
}
fn metrics_empty_report(entry: &RegistryEntry, nonzero: bool) -> MetricsCanisterReport {
MetricsCanisterReport {
role: entry.role.clone().unwrap_or_else(|| "-".to_string()),
canister_id: entry.pid.clone(),
status: MetricsCanisterStatus::Empty,
entries: Vec::new(),
error: Some(
if nonzero {
METRICS_NONZERO_EMPTY_HINT
} else {
METRICS_EMPTY_HINT
}
.to_string(),
),
}
}
fn metrics_query_error_report(
entry: &RegistryEntry,
error: &MetricsQueryError,
) -> MetricsCanisterReport {
if matches!(
error,
MetricsQueryError::Observability(FleetObservabilityError::Protocol(
CanisterProtocolError::Invocation { source, .. },
)) if matches!(source.diagnostic(), Some(IcpDiagnostic::MethodMissing))
) {
return metrics_failure_report(
entry,
MetricsCanisterStatus::Unavailable,
METRICS_UNAVAILABLE_HINT,
);
}
let error = error.to_string();
metrics_failure_report(
entry,
MetricsCanisterStatus::Error,
error.lines().next().unwrap_or(&error),
)
}
fn metrics_message_error_report(entry: &RegistryEntry, error: &str) -> MetricsCanisterReport {
metrics_failure_report(
entry,
MetricsCanisterStatus::Error,
error.lines().next().unwrap_or(error),
)
}
fn metrics_failure_report(
entry: &RegistryEntry,
status: MetricsCanisterStatus,
error: &str,
) -> MetricsCanisterReport {
MetricsCanisterReport {
role: entry.role.clone().unwrap_or_else(|| "-".to_string()),
canister_id: entry.pid.clone(),
status,
entries: Vec::new(),
error: Some(error.to_string()),
}
}
const fn metric_value_is_zero(value: &MetricValue) -> bool {
match value {
MetricValue::Count { count } => *count == 0,
MetricValue::CountAndU64 { count, value_u64 } => *count == 0 && *value_u64 == 0,
MetricValue::U128 { value } => *value == 0,
}
}
const fn metrics_kind_dto(kind: MetricsKind) -> canic_core::dto::metrics::MetricsKind {
match kind {
MetricsKind::Core => canic_core::dto::metrics::MetricsKind::Core,
MetricsKind::Placement => canic_core::dto::metrics::MetricsKind::Placement,
MetricsKind::Platform => canic_core::dto::metrics::MetricsKind::Platform,
MetricsKind::Runtime => canic_core::dto::metrics::MetricsKind::Runtime,
MetricsKind::Security => canic_core::dto::metrics::MetricsKind::Security,
MetricsKind::Storage => canic_core::dto::metrics::MetricsKind::Storage,
}
}
fn query_metrics(
options: &MetricsOptions,
fleet: &CurrentFleetResolution,
icp_root: &Path,
entry: &RegistryEntry,
) -> Result<Vec<MetricEntry>, MetricsQueryError> {
let mut icp = IcpCli::new(&options.icp, Some(options.environment.clone()));
icp = icp.with_cwd(icp_root);
let request = canic_core::dto::role::MetricsStatusRequest {
kind: metrics_kind_dto(options.kind),
page: canic_core::dto::page::PageRequest {
offset: 0,
limit: options.limit,
},
};
let response = observe_fleet_canister(
&icp,
icp_root,
&options.environment,
fleet,
entry,
CanisterObservabilityRequest::Metrics(request),
)?;
let CanisterObservabilityResponse::Metrics(page) = response else {
unreachable!("Metrics request returned a different observability response");
};
Ok(metric_page(page))
}
fn resolve_metrics_fleet(
options: &MetricsOptions,
) -> Result<CurrentFleetResolution, MetricsCommandError> {
let root = resolve_current_canic_icp_root().map_err(MetricsCommandError::IcpRoot)?;
resolve_current_fleet(&root, &options.environment, &options.fleet)
.map_err(MetricsCommandError::from)
}
#[cfg(test)]
mod tests {
use super::*;
use canic_host::icp::{IcpCommandError, IcpJsonResponseError};
fn registry_entry() -> RegistryEntry {
RegistryEntry {
pid: "aaaaa-aa".to_string(),
role: Some("wasm_store".to_string()),
parent_pid: None,
module_hash: None,
protocol_binding: None,
}
}
#[test]
fn shortens_metrics_unavailable_errors() {
let error = MetricsQueryError::Observability(FleetObservabilityError::Protocol(
CanisterProtocolError::Invocation {
canister: candid::Principal::management_canister(),
method: canic_core::protocol::CANIC_ROOT_COMMAND,
source: IcpCommandError::Failed {
command: "icp canister call".to_string(),
stderr: "Canister has no query method 'canic_observability'.".to_string(),
},
},
));
let report = metrics_query_error_report(®istry_entry(), &error);
assert_eq!(report.status, MetricsCanisterStatus::Unavailable);
assert_eq!(
serde_json::to_value(&report).expect("serialize metrics report")["status"],
"unavailable"
);
assert_eq!(report.error.as_deref(), Some(METRICS_UNAVAILABLE_HINT));
}
#[test]
fn empty_metrics_reports_carry_profile_hint() {
let report = metrics_empty_report(®istry_entry(), false);
assert_eq!(report.status, MetricsCanisterStatus::Empty);
assert_eq!(
serde_json::to_value(&report).expect("serialize metrics report")["status"],
"empty"
);
assert_eq!(report.error.as_deref(), Some(METRICS_EMPTY_HINT));
let filtered = metrics_empty_report(®istry_entry(), true);
assert_eq!(filtered.status, MetricsCanisterStatus::Empty);
assert_eq!(filtered.error.as_deref(), Some(METRICS_NONZERO_EMPTY_HINT));
}
#[test]
fn detects_zero_metric_values() {
assert!(metric_value_is_zero(&MetricValue::Count { count: 0 }));
assert!(metric_value_is_zero(&MetricValue::CountAndU64 {
count: 0,
value_u64: 0
}));
assert!(!metric_value_is_zero(&MetricValue::U128 { value: 1 }));
}
#[test]
fn maps_metric_kind_to_observability_dto() {
assert!(matches!(
metrics_kind_dto(MetricsKind::Security),
canic_core::dto::metrics::MetricsKind::Security
));
}
#[test]
fn panicked_metrics_worker_becomes_an_explicit_canister_error() {
let entry = registry_entry();
let reports = collect_metrics_worker_reports(vec![(
entry.clone(),
thread::spawn(|| panic!("simulated metrics worker panic")),
)]);
assert_eq!(reports.len(), 1);
assert_eq!(reports[0].canister_id, entry.pid);
assert_eq!(reports[0].status, MetricsCanisterStatus::Error);
assert_eq!(reports[0].error.as_deref(), Some(METRICS_WORKER_PANIC));
}
#[test]
fn metrics_response_failure_preserves_typed_cause_until_projection() {
let error = MetricsQueryError::Observability(FleetObservabilityError::Protocol(
CanisterProtocolError::Response {
canister: candid::Principal::management_canister(),
method: canic_core::protocol::CANIC_ROOT_COMMAND,
source: IcpJsonResponseError::MissingResponseBytes,
},
));
let mut source = std::error::Error::source(&error);
let mut preserved = false;
while let Some(cause) = source {
if matches!(
cause.downcast_ref::<IcpJsonResponseError>(),
Some(IcpJsonResponseError::MissingResponseBytes)
) {
preserved = true;
break;
}
source = cause.source();
}
assert!(preserved, "typed response cause must remain in the chain");
}
}