canic-cli 0.111.0

Operator CLI for Canic fleet setup, builds, evidence, catalog, backup, and restore workflows
Documentation
//! Module: metrics::transport
//!
//! Responsibility: collect typed metric observations for current Fleet canisters.
//! Does not own: metric DTOs, report rendering, or Fleet registry authority.
//! Boundary: preserves query causes until projecting per-canister report diagnostics.

use crate::metrics::{
    MetricsCommandError,
    model::{
        MetricEntry, MetricValue, MetricsCanisterReport, MetricsCanisterStatus, MetricsKind,
        MetricsReport,
    },
    options::MetricsOptions,
    parse::metric_page,
};
use canic_contracts::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_contracts::dto::metrics::MetricsKind {
    match kind {
        MetricsKind::Core => canic_contracts::dto::metrics::MetricsKind::Core,
        MetricsKind::Placement => canic_contracts::dto::metrics::MetricsKind::Placement,
        MetricsKind::Platform => canic_contracts::dto::metrics::MetricsKind::Platform,
        MetricsKind::Runtime => canic_contracts::dto::metrics::MetricsKind::Runtime,
        MetricsKind::Security => canic_contracts::dto::metrics::MetricsKind::Security,
        MetricsKind::Storage => canic_contracts::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_contracts::dto::role::MetricsStatusRequest {
        kind: metrics_kind_dto(options.kind),
        page: canic_contracts::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,
        }
    }

    // Ensure method-missing responses do not stretch the table with raw ICP output.
    #[test]
    fn shortens_metrics_unavailable_errors() {
        let error = MetricsQueryError::Observability(FleetObservabilityError::Protocol(
            CanisterProtocolError::Invocation {
                canister: candid::Principal::management_canister(),
                method: canic_contracts::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(&registry_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));
    }

    // Ensure empty successful metric tiers point operators at profile/deployed-Wasm checks.
    #[test]
    fn empty_metrics_reports_carry_profile_hint() {
        let report = metrics_empty_report(&registry_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(&registry_entry(), true);
        assert_eq!(filtered.status, MetricsCanisterStatus::Empty);
        assert_eq!(filtered.error.as_deref(), Some(METRICS_NONZERO_EMPTY_HINT));
    }

    // Ensure zero filtering treats every payload shape consistently.
    #[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 }));
    }

    // Ensure transport preserves the Candid metric kind vocabulary.
    #[test]
    fn maps_metric_kind_to_observability_dto() {
        assert!(matches!(
            metrics_kind_dto(MetricsKind::Security),
            canic_contracts::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_contracts::protocol::CANIC_ROOT_COMMAND,
                source: canic_host::icp::decode_json_result_response::<u64>("{}").unwrap_err(),
            },
        ));
        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::Envelope(_))
            ) {
                preserved = true;
                break;
            }
            source = cause.source();
        }

        assert!(preserved, "typed response cause must remain in the chain");
    }
}