1use std::net::SocketAddr;
2use std::sync::atomic::{AtomicU64, Ordering};
3use std::sync::Mutex;
4
5use serde::Serialize;
6
7use crate::ControlState;
8
9#[derive(Debug)]
15pub struct ReverseMetrics {
16 pub control_connections_active: AtomicU64,
18 pub control_connections_accepted_total: AtomicU64,
20 pub control_connections_rejected_total: AtomicU64,
22 pub control_reconnects_total: AtomicU64,
24 pub auth_failures_total: AtomicU64,
27 pub heartbeat_failures_total: AtomicU64,
29 pub drain_total: AtomicU64,
31 pub drain_duration_ms_total: AtomicU64,
33 pub streams_opened_total: AtomicU64,
35 pub streams_closed_total: AtomicU64,
37 pub stream_bytes_total: AtomicU64,
39 state_time_ms: Mutex<[u64; 6]>,
41 pub last_error: Mutex<Option<String>>,
43}
44
45const STATE_DISCONNECTED: usize = 0;
46const STATE_CONNECTING: usize = 1;
47const STATE_AUTHENTICATING: usize = 2;
48const STATE_READY: usize = 3;
49const STATE_DRAINING: usize = 4;
50const STATE_CLOSED: usize = 5;
51
52#[derive(Debug, Clone, Serialize)]
54pub struct ReverseMetricsSnapshot {
55 pub control_connections_active: u64,
56 pub control_connections_accepted_total: u64,
57 pub control_connections_rejected_total: u64,
58 pub control_reconnects_total: u64,
59 pub auth_failures_total: u64,
60 pub heartbeat_failures_total: u64,
61 pub drain_total: u64,
62 pub drain_duration_ms_total: u64,
63 pub streams_opened_total: u64,
64 pub streams_closed_total: u64,
65 pub stream_bytes_total: u64,
66 pub state_time_ms: [u64; 6],
69 pub last_error: Option<String>,
70}
71
72impl ReverseMetricsSnapshot {
73 pub fn display_summary(&self) -> String {
75 format!(
76 "reverse: active={} accepted={} rejected={} reconnects={} \
77 auth_failures={} heartbeat_failures={} drain_total={} drain_ms={} \
78 streams_open={} streams_closed={} bytes={} last_error={}",
79 self.control_connections_active,
80 self.control_connections_accepted_total,
81 self.control_connections_rejected_total,
82 self.control_reconnects_total,
83 self.auth_failures_total,
84 self.heartbeat_failures_total,
85 self.drain_total,
86 self.drain_duration_ms_total,
87 self.streams_opened_total,
88 self.streams_closed_total,
89 self.stream_bytes_total,
90 self.last_error.as_deref().unwrap_or("(none)"),
91 )
92 }
93
94 pub fn state_ms(&self, state: ControlState) -> u64 {
96 self.state_time_ms[state_index(state)]
97 }
98}
99
100fn state_index(state: ControlState) -> usize {
101 match state {
102 ControlState::Disconnected => STATE_DISCONNECTED,
103 ControlState::Connecting => STATE_CONNECTING,
104 ControlState::Authenticating => STATE_AUTHENTICATING,
105 ControlState::Ready => STATE_READY,
106 ControlState::Draining => STATE_DRAINING,
107 ControlState::Closed => STATE_CLOSED,
108 }
109}
110
111const MAX_ERROR_LEN: usize = 256;
112
113impl Default for ReverseMetrics {
114 fn default() -> Self {
115 Self::new()
116 }
117}
118
119impl ReverseMetrics {
120 pub fn new() -> Self {
121 Self {
122 control_connections_active: AtomicU64::new(0),
123 control_connections_accepted_total: AtomicU64::new(0),
124 control_connections_rejected_total: AtomicU64::new(0),
125 control_reconnects_total: AtomicU64::new(0),
126 auth_failures_total: AtomicU64::new(0),
127 heartbeat_failures_total: AtomicU64::new(0),
128 drain_total: AtomicU64::new(0),
129 drain_duration_ms_total: AtomicU64::new(0),
130 streams_opened_total: AtomicU64::new(0),
131 streams_closed_total: AtomicU64::new(0),
132 stream_bytes_total: AtomicU64::new(0),
133 state_time_ms: Mutex::new([0u64; 6]),
134 last_error: Mutex::new(None),
135 }
136 }
137
138 pub fn record_control_accepted(&self, _peer: SocketAddr) {
140 self.control_connections_active
141 .fetch_add(1, Ordering::Relaxed);
142 self.control_connections_accepted_total
143 .fetch_add(1, Ordering::Relaxed);
144 }
145
146 pub fn record_control_closed(&self) {
148 let _ = self.control_connections_active.fetch_update(
152 Ordering::Relaxed,
153 Ordering::Relaxed,
154 |active| active.checked_sub(1),
155 );
156 }
157
158 pub fn record_control_rejected(&self, _peer: SocketAddr, reason: &str) {
161 self.control_connections_rejected_total
162 .fetch_add(1, Ordering::Relaxed);
163 self.record_error(reason);
164 }
165
166 pub fn record_auth_failure(&self, _peer: SocketAddr, reason: &str) {
169 self.control_connections_rejected_total
170 .fetch_add(1, Ordering::Relaxed);
171 self.auth_failures_total.fetch_add(1, Ordering::Relaxed);
172 self.record_error(reason);
173 }
174
175 pub fn record_heartbeat_failure(&self) {
177 self.heartbeat_failures_total
178 .fetch_add(1, Ordering::Relaxed);
179 }
180
181 pub fn record_drain(&self, duration_ms: u64) {
183 self.drain_total.fetch_add(1, Ordering::Relaxed);
184 self.drain_duration_ms_total
185 .fetch_add(duration_ms, Ordering::Relaxed);
186 }
187
188 pub fn record_reconnect(&self) {
190 self.control_reconnects_total
191 .fetch_add(1, Ordering::Relaxed);
192 }
193
194 pub fn record_stream_opened(&self) {
196 self.streams_opened_total.fetch_add(1, Ordering::Relaxed);
197 }
198
199 pub fn record_stream_closed(&self, bytes: u64) {
201 self.streams_closed_total.fetch_add(1, Ordering::Relaxed);
202 self.stream_bytes_total.fetch_add(bytes, Ordering::Relaxed);
203 }
204
205 pub fn record_state_duration(&self, state: ControlState, duration_ms: u64) {
207 if let Ok(mut guard) = self.state_time_ms.lock() {
208 guard[state_index(state)] = guard[state_index(state)].saturating_add(duration_ms);
209 }
210 }
211
212 pub fn record_error(&self, msg: &str) {
214 let truncated = if msg.char_indices().nth(MAX_ERROR_LEN).is_some() {
217 let prefix: String = msg.chars().take(MAX_ERROR_LEN).collect();
218 format!("{prefix}…")
219 } else {
220 msg.to_string()
221 };
222 if let Ok(mut guard) = self.last_error.lock() {
223 *guard = Some(truncated);
224 }
225 }
226
227 pub fn snapshot(&self) -> ReverseMetricsSnapshot {
229 let last_error = self.last_error.lock().ok().and_then(|guard| guard.clone());
230 let state_time_ms = self
231 .state_time_ms
232 .lock()
233 .map(|guard| *guard)
234 .unwrap_or([0u64; 6]);
235 ReverseMetricsSnapshot {
236 control_connections_active: self.control_connections_active.load(Ordering::Relaxed),
237 control_connections_accepted_total: self
238 .control_connections_accepted_total
239 .load(Ordering::Relaxed),
240 control_connections_rejected_total: self
241 .control_connections_rejected_total
242 .load(Ordering::Relaxed),
243 control_reconnects_total: self.control_reconnects_total.load(Ordering::Relaxed),
244 auth_failures_total: self.auth_failures_total.load(Ordering::Relaxed),
245 heartbeat_failures_total: self.heartbeat_failures_total.load(Ordering::Relaxed),
246 drain_total: self.drain_total.load(Ordering::Relaxed),
247 drain_duration_ms_total: self.drain_duration_ms_total.load(Ordering::Relaxed),
248 streams_opened_total: self.streams_opened_total.load(Ordering::Relaxed),
249 streams_closed_total: self.streams_closed_total.load(Ordering::Relaxed),
250 stream_bytes_total: self.stream_bytes_total.load(Ordering::Relaxed),
251 state_time_ms,
252 last_error,
253 }
254 }
255
256 pub fn render_prometheus(&self) -> String {
258 let s = self.snapshot();
259 let mut out = String::with_capacity(1024);
260
261 out.push_str("# HELP eggress_reverse_control_connections_active Currently active control connections.\n");
262 out.push_str("# TYPE eggress_reverse_control_connections_active gauge\n");
263 out.push_str(&format!(
264 "eggress_reverse_control_connections_active {}\n",
265 s.control_connections_active
266 ));
267
268 out.push_str("# HELP eggress_reverse_control_connections_accepted_total Total accepted control connections.\n");
269 out.push_str("# TYPE eggress_reverse_control_connections_accepted_total counter\n");
270 out.push_str(&format!(
271 "eggress_reverse_control_connections_accepted_total {}\n",
272 s.control_connections_accepted_total
273 ));
274
275 out.push_str("# HELP eggress_reverse_control_connections_rejected_total Total rejected control connections.\n");
276 out.push_str("# TYPE eggress_reverse_control_connections_rejected_total counter\n");
277 out.push_str(&format!(
278 "eggress_reverse_control_connections_rejected_total {}\n",
279 s.control_connections_rejected_total
280 ));
281
282 out.push_str(
283 "# HELP eggress_reverse_control_reconnects_total Total client reconnect attempts.\n",
284 );
285 out.push_str("# TYPE eggress_reverse_control_reconnects_total counter\n");
286 out.push_str(&format!(
287 "eggress_reverse_control_reconnects_total {}\n",
288 s.control_reconnects_total
289 ));
290
291 out.push_str(
292 "# HELP eggress_reverse_auth_failures_total Total control channel auth failures.\n",
293 );
294 out.push_str("# TYPE eggress_reverse_auth_failures_total counter\n");
295 out.push_str(&format!(
296 "eggress_reverse_auth_failures_total {}\n",
297 s.auth_failures_total
298 ));
299
300 out.push_str(
301 "# HELP eggress_reverse_heartbeat_failures_total Total heartbeat / keepalive failures.\n",
302 );
303 out.push_str("# TYPE eggress_reverse_heartbeat_failures_total counter\n");
304 out.push_str(&format!(
305 "eggress_reverse_heartbeat_failures_total {}\n",
306 s.heartbeat_failures_total
307 ));
308
309 out.push_str("# HELP eggress_reverse_drain_total Total drain operations completed.\n");
310 out.push_str("# TYPE eggress_reverse_drain_total counter\n");
311 out.push_str(&format!("eggress_reverse_drain_total {}\n", s.drain_total));
312
313 out.push_str(
314 "# HELP eggress_reverse_drain_duration_ms_total Cumulative drain duration in milliseconds.\n",
315 );
316 out.push_str("# TYPE eggress_reverse_drain_duration_ms_total counter\n");
317 out.push_str(&format!(
318 "eggress_reverse_drain_duration_ms_total {}\n",
319 s.drain_duration_ms_total
320 ));
321
322 out.push_str(
323 "# HELP eggress_reverse_streams_opened_total Total control sessions initiated.\n",
324 );
325 out.push_str("# TYPE eggress_reverse_streams_opened_total counter\n");
326 out.push_str(&format!(
327 "eggress_reverse_streams_opened_total {}\n",
328 s.streams_opened_total
329 ));
330
331 out.push_str("# HELP eggress_reverse_streams_closed_total Total control sessions completed cleanly.\n");
332 out.push_str("# TYPE eggress_reverse_streams_closed_total counter\n");
333 out.push_str(&format!(
334 "eggress_reverse_streams_closed_total {}\n",
335 s.streams_closed_total
336 ));
337
338 out.push_str("# HELP eggress_reverse_stream_bytes_total Total bytes relayed through control channels.\n");
339 out.push_str("# TYPE eggress_reverse_stream_bytes_total counter\n");
340 out.push_str(&format!(
341 "eggress_reverse_stream_bytes_total {}\n",
342 s.stream_bytes_total
343 ));
344
345 out.push_str(
346 "# HELP eggress_reverse_state_time_ms Cumulative time spent in each control state.\n",
347 );
348 out.push_str("# TYPE eggress_reverse_state_time_ms counter\n");
349 let state_labels = [
350 "disconnected",
351 "connecting",
352 "authenticating",
353 "ready",
354 "draining",
355 "closed",
356 ];
357 for (idx, label) in state_labels.iter().enumerate() {
358 out.push_str(&format!(
359 "eggress_reverse_state_time_ms{{state=\"{}\"}} {}\n",
360 label, s.state_time_ms[idx]
361 ));
362 }
363
364 out
365 }
366}
367
368#[cfg(test)]
369mod tests {
370 use super::*;
371
372 #[test]
373 fn new_counters_are_zero() {
374 let m = ReverseMetrics::new();
375 assert_eq!(m.control_connections_active.load(Ordering::Relaxed), 0);
376 assert_eq!(
377 m.control_connections_accepted_total.load(Ordering::Relaxed),
378 0
379 );
380 assert_eq!(
381 m.control_connections_rejected_total.load(Ordering::Relaxed),
382 0
383 );
384 assert_eq!(m.control_reconnects_total.load(Ordering::Relaxed), 0);
385 assert_eq!(m.auth_failures_total.load(Ordering::Relaxed), 0);
386 assert_eq!(m.heartbeat_failures_total.load(Ordering::Relaxed), 0);
387 assert_eq!(m.drain_total.load(Ordering::Relaxed), 0);
388 assert_eq!(m.drain_duration_ms_total.load(Ordering::Relaxed), 0);
389 assert_eq!(m.streams_opened_total.load(Ordering::Relaxed), 0);
390 assert_eq!(m.streams_closed_total.load(Ordering::Relaxed), 0);
391 assert_eq!(m.stream_bytes_total.load(Ordering::Relaxed), 0);
392 assert_eq!(*m.state_time_ms.lock().unwrap(), [0u64; 6]);
393 assert!(m.last_error.lock().unwrap().is_none());
394 }
395
396 fn peer() -> SocketAddr {
397 "127.0.0.1:9999".parse().unwrap()
398 }
399
400 #[test]
401 fn record_control_accepted() {
402 let m = ReverseMetrics::new();
403 m.record_control_accepted(peer());
404 m.record_control_accepted(peer());
405 assert_eq!(m.control_connections_active.load(Ordering::Relaxed), 2);
406 assert_eq!(
407 m.control_connections_accepted_total.load(Ordering::Relaxed),
408 2
409 );
410 }
411
412 #[test]
413 fn record_control_rejected() {
414 let m = ReverseMetrics::new();
415 m.record_control_rejected(peer(), "max_control_connections");
416 assert_eq!(
417 m.control_connections_rejected_total.load(Ordering::Relaxed),
418 1
419 );
420 assert_eq!(m.auth_failures_total.load(Ordering::Relaxed), 0);
421 assert_eq!(
422 m.last_error.lock().unwrap().as_deref(),
423 Some("max_control_connections")
424 );
425 }
426
427 #[test]
428 fn record_auth_failure() {
429 let m = ReverseMetrics::new();
430 m.record_auth_failure(peer(), "bad credentials");
431 assert_eq!(m.auth_failures_total.load(Ordering::Relaxed), 1);
432 assert_eq!(
433 m.last_error.lock().unwrap().as_deref(),
434 Some("bad credentials")
435 );
436 }
437
438 #[test]
439 fn record_heartbeat_failure() {
440 let m = ReverseMetrics::new();
441 m.record_heartbeat_failure();
442 m.record_heartbeat_failure();
443 assert_eq!(m.heartbeat_failures_total.load(Ordering::Relaxed), 2);
444 }
445
446 #[test]
447 fn record_drain() {
448 let m = ReverseMetrics::new();
449 m.record_drain(150);
450 m.record_drain(50);
451 assert_eq!(m.drain_total.load(Ordering::Relaxed), 2);
452 assert_eq!(m.drain_duration_ms_total.load(Ordering::Relaxed), 200);
453 }
454
455 #[test]
456 fn record_reconnect() {
457 let m = ReverseMetrics::new();
458 m.record_reconnect();
459 m.record_reconnect();
460 m.record_reconnect();
461 assert_eq!(m.control_reconnects_total.load(Ordering::Relaxed), 3);
462 }
463
464 #[test]
465 fn record_stream_opened_and_closed() {
466 let m = ReverseMetrics::new();
467 m.record_stream_opened();
468 m.record_stream_opened();
469 assert_eq!(m.streams_opened_total.load(Ordering::Relaxed), 2);
470
471 m.record_stream_closed(1024);
472 assert_eq!(m.streams_closed_total.load(Ordering::Relaxed), 1);
473 assert_eq!(m.stream_bytes_total.load(Ordering::Relaxed), 1024);
474
475 m.record_stream_closed(512);
476 assert_eq!(m.streams_closed_total.load(Ordering::Relaxed), 2);
477 assert_eq!(m.stream_bytes_total.load(Ordering::Relaxed), 1536);
478 }
479
480 #[test]
481 fn record_state_duration() {
482 let m = ReverseMetrics::new();
483 m.record_state_duration(ControlState::Connecting, 100);
484 m.record_state_duration(ControlState::Ready, 500);
485 m.record_state_duration(ControlState::Ready, 250);
486 let snap = m.snapshot();
487 assert_eq!(snap.state_ms(ControlState::Connecting), 100);
488 assert_eq!(snap.state_ms(ControlState::Ready), 750);
489 assert_eq!(snap.state_ms(ControlState::Draining), 0);
490 }
491
492 #[test]
493 fn record_error_truncates_long_messages() {
494 let m = ReverseMetrics::new();
495 let long_msg = "x".repeat(500);
496 m.record_error(&long_msg);
497 let err = m.last_error.lock().unwrap().clone().unwrap();
498 assert_eq!(err.chars().count(), MAX_ERROR_LEN + 1);
500 assert!(err.ends_with('…'));
501 }
502
503 #[test]
504 fn record_error_handles_multibyte_messages() {
505 let m = ReverseMetrics::new();
506 m.record_error(&"é".repeat(MAX_ERROR_LEN + 1));
507 let err = m.last_error.lock().unwrap().clone().unwrap();
508 assert_eq!(err.chars().count(), MAX_ERROR_LEN + 1);
509 assert!(err.ends_with('…'));
510 }
511
512 #[test]
513 fn snapshot_returns_expected_values() {
514 let m = ReverseMetrics::new();
515 m.record_control_accepted(peer());
516 m.record_control_accepted(peer());
517 m.record_stream_opened();
518 m.record_stream_closed(2048);
519 m.record_reconnect();
520 m.record_control_rejected(peer(), "timeout");
521 m.record_auth_failure(peer(), "bad password");
522 m.record_heartbeat_failure();
523 m.record_drain(75);
524 m.record_state_duration(ControlState::Ready, 1000);
525
526 let snap = m.snapshot();
527 assert_eq!(snap.control_connections_active, 2);
528 assert_eq!(snap.control_connections_accepted_total, 2);
529 assert_eq!(snap.control_connections_rejected_total, 2);
530 assert_eq!(snap.auth_failures_total, 1);
531 assert_eq!(snap.heartbeat_failures_total, 1);
532 assert_eq!(snap.drain_total, 1);
533 assert_eq!(snap.drain_duration_ms_total, 75);
534 assert_eq!(snap.streams_opened_total, 1);
535 assert_eq!(snap.streams_closed_total, 1);
536 assert_eq!(snap.stream_bytes_total, 2048);
537 assert_eq!(snap.state_ms(ControlState::Ready), 1000);
538 assert_eq!(snap.last_error.as_deref(), Some("bad password"));
539 }
540
541 #[test]
542 fn snapshot_display_summary() {
543 let m = ReverseMetrics::new();
544 m.record_control_accepted(peer());
545 m.record_stream_opened();
546 m.record_stream_closed(100);
547 m.record_drain(20);
548 let snap = m.snapshot();
549 let display = snap.display_summary();
550 assert!(display.contains("active=1"));
551 assert!(display.contains("accepted=1"));
552 assert!(display.contains("bytes=100"));
553 assert!(display.contains("drain_total=1"));
554 assert!(display.contains("drain_ms=20"));
555 assert!(display.contains("last_error=(none)"));
556 }
557
558 #[test]
559 fn prometheus_output_contains_expected_names() {
560 let m = ReverseMetrics::new();
561 m.record_control_accepted(peer());
562 m.record_stream_opened();
563 m.record_stream_closed(42);
564 m.record_heartbeat_failure();
565 m.record_drain(5);
566 m.record_state_duration(ControlState::Ready, 100);
567
568 let prom = m.render_prometheus();
569 assert!(prom.contains("eggress_reverse_control_connections_active"));
570 assert!(prom.contains("eggress_reverse_control_connections_accepted_total"));
571 assert!(prom.contains("eggress_reverse_control_connections_rejected_total"));
572 assert!(prom.contains("eggress_reverse_control_reconnects_total"));
573 assert!(prom.contains("eggress_reverse_auth_failures_total"));
574 assert!(prom.contains("eggress_reverse_heartbeat_failures_total"));
575 assert!(prom.contains("eggress_reverse_drain_total"));
576 assert!(prom.contains("eggress_reverse_drain_duration_ms_total"));
577 assert!(prom.contains("eggress_reverse_streams_opened_total"));
578 assert!(prom.contains("eggress_reverse_streams_closed_total"));
579 assert!(prom.contains("eggress_reverse_stream_bytes_total"));
580 assert!(prom.contains("eggress_reverse_state_time_ms"));
581 assert!(prom.contains("TYPE eggress_reverse_control_connections_active gauge"));
582 assert!(prom.contains("TYPE eggress_reverse_streams_opened_total counter"));
583 assert!(prom.contains("state=\"ready\""));
584 }
585
586 #[test]
587 fn snapshot_is_clone() {
588 let m = ReverseMetrics::new();
589 m.record_control_accepted(peer());
590 let s1 = m.snapshot();
591 let s2 = s1.clone();
592 assert_eq!(s1.control_connections_active, s2.control_connections_active);
593 }
594
595 #[test]
596 fn snapshot_is_debug() {
597 let m = ReverseMetrics::new();
598 let snap = m.snapshot();
599 let debug_str = format!("{:?}", snap);
600 assert!(debug_str.contains("ReverseMetricsSnapshot"));
601 }
602
603 #[test]
604 fn snapshot_is_serialize() {
605 let m = ReverseMetrics::new();
606 let snap = m.snapshot();
607 let json = serde_json::to_string(&snap).unwrap();
608 assert!(json.contains("control_connections_active"));
609 assert!(json.contains("stream_bytes_total"));
610 assert!(json.contains("auth_failures_total"));
611 assert!(json.contains("state_time_ms"));
612 }
613}