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}