mecha10-diagnostics 0.6.2

Diagnostics and metrics collection for Mecha10 robotics framework
Documentation
//! Diagnostic message types
//!
//! Each diagnostic category has its own message type that can be published to topics.

use mecha10_core::messages::Message;
use serde::{Deserialize, Serialize};

/// Get current timestamp in microseconds since Unix epoch
pub fn now_micros() -> u64 {
    mecha10_core::prelude::now_micros()
}

/// Generic diagnostic message wrapper
///
/// All diagnostic messages follow this pattern for consistency.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DiagnosticMessage<T> {
    /// Timestamp in microseconds since epoch
    pub timestamp: u64,
    /// Source node/component that published this diagnostic
    pub source: String,
    /// The actual diagnostic payload
    pub payload: T,
}

impl<T> DiagnosticMessage<T> {
    /// Create a new diagnostic message
    pub fn new(source: impl Into<String>, payload: T) -> Self {
        Self {
            timestamp: now_micros(),
            source: source.into(),
            payload,
        }
    }

    /// Create a new diagnostic message with explicit timestamp
    pub fn new_with_timestamp(source: impl Into<String>, timestamp: u64, payload: T) -> Self {
        Self {
            timestamp,
            source: source.into(),
            payload,
        }
    }
}

// Implement Message trait for DiagnosticMessage
// This allows it to be published via Context
impl<T: Send + Sync + 'static> Message for DiagnosticMessage<T> {}

// ===== Streaming Pipeline Metrics =====

/// Streaming pipeline metrics (frame throughput at each stage)
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct StreamingPipelineMetrics {
    /// Total frames received from camera
    pub frames_received: u64,
    /// Total frames successfully encoded
    pub frames_encoded: u64,
    /// Total frames sent via WebRTC
    pub frames_sent: u64,
    /// Total frames dropped (encoding queue full, etc.)
    pub frames_dropped: u64,
    /// Current throughput in bytes/second
    pub bytes_per_second: u64,
    /// Current frame rate (frames/second)
    pub fps: f64,
}

/// Streaming latency metrics (end-to-end latency breakdown)
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct StreamingLatencyMetrics {
    /// Camera capture → bridge receive (microseconds)
    pub camera_to_bridge_us: u64,
    /// Frame queuing time (microseconds)
    pub queue_time_us: u64,
    /// Encoding time (microseconds)
    pub encoding_time_us: u64,
    /// WebRTC send time (microseconds)
    pub send_time_us: u64,
    /// Total end-to-end latency (microseconds)
    pub end_to_end_us: u64,
}

/// Encoding performance metrics
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct EncodingMetrics {
    /// Total frames encoded
    pub total_frames: u64,
    /// Slow frames (exceeded target frame time)
    pub slow_frames: u64,
    /// Average encoding time (milliseconds)
    pub avg_encode_time_ms: f64,
    /// P50 encoding time (milliseconds)
    pub p50_encode_time_ms: f64,
    /// P95 encoding time (milliseconds)
    pub p95_encode_time_ms: f64,
    /// P99 encoding time (milliseconds)
    pub p99_encode_time_ms: f64,
    /// Max encoding time (milliseconds)
    pub max_encode_time_ms: f64,
    /// Encoder queue depth (current)
    pub queue_depth: usize,
}

/// Streaming bandwidth metrics
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct BandwidthMetrics {
    /// Current bitrate (bits per second)
    pub bitrate_bps: u64,
    /// Target bitrate (bits per second)
    pub target_bitrate_bps: u64,
    /// Average frame size (bytes)
    pub avg_frame_size_bytes: u64,
    /// Total bytes sent
    pub total_bytes_sent: u64,
    /// Bandwidth utilization (0.0-1.0)
    pub utilization: f64,
}

// ===== WebRTC Metrics =====

/// WebRTC connection metrics
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WebRTCConnectionMetrics {
    /// Number of active peer connections
    pub active_connections: usize,
    /// Total connections established
    pub total_connections: u64,
    /// Failed connection attempts
    pub failed_connections: u64,
    /// Connections by state
    pub connections_by_state: ConnectionStateBreakdown,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ConnectionStateBreakdown {
    pub new: usize,
    pub connecting: usize,
    pub connected: usize,
    pub disconnected: usize,
    pub failed: usize,
    pub closed: usize,
}

/// WebRTC quality metrics
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WebRTCQualityMetrics {
    /// Round-trip time (milliseconds)
    pub rtt_ms: f64,
    /// Packet loss percentage (0.0-100.0)
    pub packet_loss_percent: f64,
    /// Jitter (milliseconds)
    pub jitter_ms: f64,
    /// Bandwidth estimate (bits per second)
    pub bandwidth_estimate_bps: u64,
}

// ===== WebSocket Metrics =====

/// WebSocket connection metrics
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WebSocketConnectionMetrics {
    /// Active WebSocket connections
    pub active_connections: usize,
    /// Total connections established
    pub total_connections: u64,
    /// Failed connection attempts
    pub failed_connections: u64,
    /// Connection type breakdown
    pub by_type: WebSocketTypeBreakdown,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WebSocketTypeBreakdown {
    pub signaling: usize,
    pub dashboard: usize,
}

/// WebSocket message metrics
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WebSocketMessageMetrics {
    /// Messages received per second
    pub messages_received_per_sec: f64,
    /// Messages sent per second
    pub messages_sent_per_sec: f64,
    /// Average message size (bytes)
    pub avg_message_size_bytes: u64,
    /// Total bytes received
    pub total_bytes_received: u64,
    /// Total bytes sent
    pub total_bytes_sent: u64,
    /// Error count
    pub error_count: u64,
}

// ===== Redis Metrics =====

/// Redis connection pool metrics
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RedisPoolMetrics {
    /// Active connections in pool
    pub active_connections: usize,
    /// Idle connections in pool
    pub idle_connections: usize,
    /// Maximum pool size
    pub max_pool_size: usize,
    /// Connection wait time (milliseconds)
    pub avg_wait_time_ms: f64,
    /// Connection failures
    pub connection_failures: u64,
}

/// Redis operation metrics
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RedisOperationMetrics {
    /// Operations per second
    pub ops_per_second: f64,
    /// Average operation latency (microseconds)
    pub avg_latency_us: f64,
    /// P95 operation latency (microseconds)
    pub p95_latency_us: f64,
    /// P99 operation latency (microseconds)
    pub p99_latency_us: f64,
    /// Failed operations
    pub failed_ops: u64,
    /// Operation type breakdown
    pub by_type: RedisOperationTypeBreakdown,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RedisOperationTypeBreakdown {
    pub publish: u64,
    pub subscribe: u64,
    pub get: u64,
    pub set: u64,
    pub other: u64,
}

/// Redis server info metrics from INFO command
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RedisServerInfoMetrics {
    /// Redis server version
    pub redis_version: String,
    /// Server uptime in seconds
    pub uptime_seconds: u64,
    /// Connected clients
    pub connected_clients: u64,
    /// Used memory in bytes
    pub used_memory: u64,
    /// Used memory RSS in bytes
    pub used_memory_rss: u64,
    /// Peak memory usage in bytes
    pub used_memory_peak: u64,
    /// Total connections received
    pub total_connections_received: u64,
    /// Total commands processed
    pub total_commands_processed: u64,
    /// Instantaneous operations per second
    pub instantaneous_ops_per_sec: u64,
    /// Keyspace hits
    pub keyspace_hits: u64,
    /// Keyspace misses
    pub keyspace_misses: u64,
    /// Number of keys in database 0
    pub db0_keys: u64,
    /// Number of keys with expiry in database 0
    pub db0_expires: u64,
}

// ===== Docker Metrics =====

/// Docker container metrics
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DockerContainerMetrics {
    /// Container ID
    pub container_id: String,
    /// Container name
    pub container_name: String,
    /// CPU usage percentage (0.0-100.0 per core)
    pub cpu_percent: f64,
    /// Memory usage (bytes)
    pub memory_usage_bytes: u64,
    /// Memory limit (bytes)
    pub memory_limit_bytes: u64,
    /// Memory percentage (0.0-100.0)
    pub memory_percent: f64,
    /// Network received (bytes)
    pub network_rx_bytes: u64,
    /// Network transmitted (bytes)
    pub network_tx_bytes: u64,
    /// Block I/O read (bytes)
    pub block_io_read_bytes: u64,
    /// Block I/O write (bytes)
    pub block_io_write_bytes: u64,
}

// ===== Simulation Metrics =====

/// Health of a supervised simulation backend subprocess (e.g. `mujoco-sim`'s
/// `mecha10_rl.mujoco.cli serve` child process).
///
/// Simulation nodes are process supervisors, not WebSocket bridges - there's no
/// "connection" in the network sense. The meaningful health signals are whether the
/// subprocess is currently alive, how many times it's been (re)started (crash-loop
/// detection), and what happened last time it exited.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SimulationConnectionMetrics {
    /// Simulation backend identifier (e.g. "mujoco").
    pub backend: String,
    /// Whether the subprocess is currently running.
    pub running: bool,
    /// Number of times the subprocess has been (re)started, including the initial start.
    pub restart_count: u64,
    /// Timestamp (microseconds since epoch) of the most recent (re)start.
    pub last_start_us: u64,
    /// Reason the subprocess last exited (`None` until it has exited at least once).
    pub last_exit_reason: Option<String>,
}

// ===== System Metrics =====

/// System-wide resource metrics
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SystemResourceMetrics {
    /// CPU usage percentage (0.0-100.0, averaged across cores)
    pub cpu_percent: f64,
    /// Per-core CPU usage
    pub cpu_per_core: Vec<f64>,
    /// Total memory (bytes)
    pub memory_total_bytes: u64,
    /// Used memory (bytes)
    pub memory_used_bytes: u64,
    /// Memory usage percentage (0.0-100.0)
    pub memory_percent: f64,
    /// Total disk space (bytes)
    pub disk_total_bytes: u64,
    /// Used disk space (bytes)
    pub disk_used_bytes: u64,
    /// Disk usage percentage (0.0-100.0)
    pub disk_percent: f64,
    /// Network received (bytes/sec)
    pub network_rx_bytes_per_sec: u64,
    /// Network transmitted (bytes/sec)
    pub network_tx_bytes_per_sec: u64,
}

// ===== Node Health Metrics =====

/// Per-node health metrics
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct NodeHealthMetrics {
    /// Node ID
    pub node_id: String,
    /// Health status
    pub healthy: bool,
    /// Messages processed
    pub messages_processed: u64,
    /// Message processing rate (msgs/sec)
    pub message_rate: f64,
    /// Error count
    pub error_count: u64,
    /// Uptime (seconds)
    pub uptime_seconds: u64,
    /// Custom health message
    pub message: Option<String>,
}