Skip to main content

mecha10_diagnostics/
types.rs

1//! Diagnostic message types
2//!
3//! Each diagnostic category has its own message type that can be published to topics.
4
5use mecha10_core::messages::Message;
6use serde::{Deserialize, Serialize};
7
8/// Get current timestamp in microseconds since Unix epoch
9pub fn now_micros() -> u64 {
10    mecha10_core::prelude::now_micros()
11}
12
13/// Generic diagnostic message wrapper
14///
15/// All diagnostic messages follow this pattern for consistency.
16#[derive(Debug, Clone, Serialize, Deserialize)]
17pub struct DiagnosticMessage<T> {
18    /// Timestamp in microseconds since epoch
19    pub timestamp: u64,
20    /// Source node/component that published this diagnostic
21    pub source: String,
22    /// The actual diagnostic payload
23    pub payload: T,
24}
25
26impl<T> DiagnosticMessage<T> {
27    /// Create a new diagnostic message
28    pub fn new(source: impl Into<String>, payload: T) -> Self {
29        Self {
30            timestamp: now_micros(),
31            source: source.into(),
32            payload,
33        }
34    }
35
36    /// Create a new diagnostic message with explicit timestamp
37    pub fn new_with_timestamp(source: impl Into<String>, timestamp: u64, payload: T) -> Self {
38        Self {
39            timestamp,
40            source: source.into(),
41            payload,
42        }
43    }
44}
45
46// Implement Message trait for DiagnosticMessage
47// This allows it to be published via Context
48impl<T: Send + Sync + 'static> Message for DiagnosticMessage<T> {}
49
50// ===== Streaming Pipeline Metrics =====
51
52/// Streaming pipeline metrics (frame throughput at each stage)
53#[derive(Debug, Clone, Serialize, Deserialize)]
54pub struct StreamingPipelineMetrics {
55    /// Total frames received from camera
56    pub frames_received: u64,
57    /// Total frames successfully encoded
58    pub frames_encoded: u64,
59    /// Total frames sent via WebRTC
60    pub frames_sent: u64,
61    /// Total frames dropped (encoding queue full, etc.)
62    pub frames_dropped: u64,
63    /// Current throughput in bytes/second
64    pub bytes_per_second: u64,
65    /// Current frame rate (frames/second)
66    pub fps: f64,
67}
68
69/// Streaming latency metrics (end-to-end latency breakdown)
70#[derive(Debug, Clone, Serialize, Deserialize)]
71pub struct StreamingLatencyMetrics {
72    /// Camera capture → bridge receive (microseconds)
73    pub camera_to_bridge_us: u64,
74    /// Frame queuing time (microseconds)
75    pub queue_time_us: u64,
76    /// Encoding time (microseconds)
77    pub encoding_time_us: u64,
78    /// WebRTC send time (microseconds)
79    pub send_time_us: u64,
80    /// Total end-to-end latency (microseconds)
81    pub end_to_end_us: u64,
82}
83
84/// Encoding performance metrics
85#[derive(Debug, Clone, Serialize, Deserialize)]
86pub struct EncodingMetrics {
87    /// Total frames encoded
88    pub total_frames: u64,
89    /// Slow frames (exceeded target frame time)
90    pub slow_frames: u64,
91    /// Average encoding time (milliseconds)
92    pub avg_encode_time_ms: f64,
93    /// P50 encoding time (milliseconds)
94    pub p50_encode_time_ms: f64,
95    /// P95 encoding time (milliseconds)
96    pub p95_encode_time_ms: f64,
97    /// P99 encoding time (milliseconds)
98    pub p99_encode_time_ms: f64,
99    /// Max encoding time (milliseconds)
100    pub max_encode_time_ms: f64,
101    /// Encoder queue depth (current)
102    pub queue_depth: usize,
103}
104
105/// Streaming bandwidth metrics
106#[derive(Debug, Clone, Serialize, Deserialize)]
107pub struct BandwidthMetrics {
108    /// Current bitrate (bits per second)
109    pub bitrate_bps: u64,
110    /// Target bitrate (bits per second)
111    pub target_bitrate_bps: u64,
112    /// Average frame size (bytes)
113    pub avg_frame_size_bytes: u64,
114    /// Total bytes sent
115    pub total_bytes_sent: u64,
116    /// Bandwidth utilization (0.0-1.0)
117    pub utilization: f64,
118}
119
120// ===== WebRTC Metrics =====
121
122/// WebRTC connection metrics
123#[derive(Debug, Clone, Serialize, Deserialize)]
124pub struct WebRTCConnectionMetrics {
125    /// Number of active peer connections
126    pub active_connections: usize,
127    /// Total connections established
128    pub total_connections: u64,
129    /// Failed connection attempts
130    pub failed_connections: u64,
131    /// Connections by state
132    pub connections_by_state: ConnectionStateBreakdown,
133}
134
135#[derive(Debug, Clone, Serialize, Deserialize)]
136pub struct ConnectionStateBreakdown {
137    pub new: usize,
138    pub connecting: usize,
139    pub connected: usize,
140    pub disconnected: usize,
141    pub failed: usize,
142    pub closed: usize,
143}
144
145/// WebRTC quality metrics
146#[derive(Debug, Clone, Serialize, Deserialize)]
147pub struct WebRTCQualityMetrics {
148    /// Round-trip time (milliseconds)
149    pub rtt_ms: f64,
150    /// Packet loss percentage (0.0-100.0)
151    pub packet_loss_percent: f64,
152    /// Jitter (milliseconds)
153    pub jitter_ms: f64,
154    /// Bandwidth estimate (bits per second)
155    pub bandwidth_estimate_bps: u64,
156}
157
158// ===== WebSocket Metrics =====
159
160/// WebSocket connection metrics
161#[derive(Debug, Clone, Serialize, Deserialize)]
162pub struct WebSocketConnectionMetrics {
163    /// Active WebSocket connections
164    pub active_connections: usize,
165    /// Total connections established
166    pub total_connections: u64,
167    /// Failed connection attempts
168    pub failed_connections: u64,
169    /// Connection type breakdown
170    pub by_type: WebSocketTypeBreakdown,
171}
172
173#[derive(Debug, Clone, Serialize, Deserialize)]
174pub struct WebSocketTypeBreakdown {
175    pub signaling: usize,
176    pub dashboard: usize,
177}
178
179/// WebSocket message metrics
180#[derive(Debug, Clone, Serialize, Deserialize)]
181pub struct WebSocketMessageMetrics {
182    /// Messages received per second
183    pub messages_received_per_sec: f64,
184    /// Messages sent per second
185    pub messages_sent_per_sec: f64,
186    /// Average message size (bytes)
187    pub avg_message_size_bytes: u64,
188    /// Total bytes received
189    pub total_bytes_received: u64,
190    /// Total bytes sent
191    pub total_bytes_sent: u64,
192    /// Error count
193    pub error_count: u64,
194}
195
196// ===== Redis Metrics =====
197
198/// Redis connection pool metrics
199#[derive(Debug, Clone, Serialize, Deserialize)]
200pub struct RedisPoolMetrics {
201    /// Active connections in pool
202    pub active_connections: usize,
203    /// Idle connections in pool
204    pub idle_connections: usize,
205    /// Maximum pool size
206    pub max_pool_size: usize,
207    /// Connection wait time (milliseconds)
208    pub avg_wait_time_ms: f64,
209    /// Connection failures
210    pub connection_failures: u64,
211}
212
213/// Redis operation metrics
214#[derive(Debug, Clone, Serialize, Deserialize)]
215pub struct RedisOperationMetrics {
216    /// Operations per second
217    pub ops_per_second: f64,
218    /// Average operation latency (microseconds)
219    pub avg_latency_us: f64,
220    /// P95 operation latency (microseconds)
221    pub p95_latency_us: f64,
222    /// P99 operation latency (microseconds)
223    pub p99_latency_us: f64,
224    /// Failed operations
225    pub failed_ops: u64,
226    /// Operation type breakdown
227    pub by_type: RedisOperationTypeBreakdown,
228}
229
230#[derive(Debug, Clone, Serialize, Deserialize)]
231pub struct RedisOperationTypeBreakdown {
232    pub publish: u64,
233    pub subscribe: u64,
234    pub get: u64,
235    pub set: u64,
236    pub other: u64,
237}
238
239/// Redis server info metrics from INFO command
240#[derive(Debug, Clone, Serialize, Deserialize)]
241pub struct RedisServerInfoMetrics {
242    /// Redis server version
243    pub redis_version: String,
244    /// Server uptime in seconds
245    pub uptime_seconds: u64,
246    /// Connected clients
247    pub connected_clients: u64,
248    /// Used memory in bytes
249    pub used_memory: u64,
250    /// Used memory RSS in bytes
251    pub used_memory_rss: u64,
252    /// Peak memory usage in bytes
253    pub used_memory_peak: u64,
254    /// Total connections received
255    pub total_connections_received: u64,
256    /// Total commands processed
257    pub total_commands_processed: u64,
258    /// Instantaneous operations per second
259    pub instantaneous_ops_per_sec: u64,
260    /// Keyspace hits
261    pub keyspace_hits: u64,
262    /// Keyspace misses
263    pub keyspace_misses: u64,
264    /// Number of keys in database 0
265    pub db0_keys: u64,
266    /// Number of keys with expiry in database 0
267    pub db0_expires: u64,
268}
269
270// ===== Docker Metrics =====
271
272/// Docker container metrics
273#[derive(Debug, Clone, Serialize, Deserialize)]
274pub struct DockerContainerMetrics {
275    /// Container ID
276    pub container_id: String,
277    /// Container name
278    pub container_name: String,
279    /// CPU usage percentage (0.0-100.0 per core)
280    pub cpu_percent: f64,
281    /// Memory usage (bytes)
282    pub memory_usage_bytes: u64,
283    /// Memory limit (bytes)
284    pub memory_limit_bytes: u64,
285    /// Memory percentage (0.0-100.0)
286    pub memory_percent: f64,
287    /// Network received (bytes)
288    pub network_rx_bytes: u64,
289    /// Network transmitted (bytes)
290    pub network_tx_bytes: u64,
291    /// Block I/O read (bytes)
292    pub block_io_read_bytes: u64,
293    /// Block I/O write (bytes)
294    pub block_io_write_bytes: u64,
295}
296
297// ===== Simulation Metrics =====
298
299/// Health of a supervised simulation backend subprocess (e.g. `simulator`'s
300/// `mecha10_rl.mujoco.cli serve` child process).
301///
302/// Simulation nodes are process supervisors, not WebSocket bridges - there's no
303/// "connection" in the network sense. The meaningful health signals are whether the
304/// subprocess is currently alive, how many times it's been (re)started (crash-loop
305/// detection), and what happened last time it exited.
306#[derive(Debug, Clone, Serialize, Deserialize)]
307pub struct SimulationConnectionMetrics {
308    /// Simulation backend identifier (e.g. "mujoco").
309    pub backend: String,
310    /// Whether the subprocess is currently running.
311    pub running: bool,
312    /// Number of times the subprocess has been (re)started, including the initial start.
313    pub restart_count: u64,
314    /// Timestamp (microseconds since epoch) of the most recent (re)start.
315    pub last_start_us: u64,
316    /// Reason the subprocess last exited (`None` until it has exited at least once).
317    pub last_exit_reason: Option<String>,
318}
319
320// ===== System Metrics =====
321
322/// System-wide resource metrics
323#[derive(Debug, Clone, Serialize, Deserialize)]
324pub struct SystemResourceMetrics {
325    /// CPU usage percentage (0.0-100.0, averaged across cores)
326    pub cpu_percent: f64,
327    /// Per-core CPU usage
328    pub cpu_per_core: Vec<f64>,
329    /// Total memory (bytes)
330    pub memory_total_bytes: u64,
331    /// Used memory (bytes)
332    pub memory_used_bytes: u64,
333    /// Memory usage percentage (0.0-100.0)
334    pub memory_percent: f64,
335    /// Total disk space (bytes)
336    pub disk_total_bytes: u64,
337    /// Used disk space (bytes)
338    pub disk_used_bytes: u64,
339    /// Disk usage percentage (0.0-100.0)
340    pub disk_percent: f64,
341    /// Network received (bytes/sec)
342    pub network_rx_bytes_per_sec: u64,
343    /// Network transmitted (bytes/sec)
344    pub network_tx_bytes_per_sec: u64,
345}
346
347// ===== Node Health Metrics =====
348
349/// Per-node health metrics
350#[derive(Debug, Clone, Serialize, Deserialize)]
351pub struct NodeHealthMetrics {
352    /// Node ID
353    pub node_id: String,
354    /// Health status
355    pub healthy: bool,
356    /// Messages processed
357    pub messages_processed: u64,
358    /// Message processing rate (msgs/sec)
359    pub message_rate: f64,
360    /// Error count
361    pub error_count: u64,
362    /// Uptime (seconds)
363    pub uptime_seconds: u64,
364    /// Custom health message
365    pub message: Option<String>,
366}