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.chars().count() > MAX_ERROR_LEN {
215 let prefix: String = msg.chars().take(MAX_ERROR_LEN).collect();
216 format!("{prefix}…")
217 } else {
218 msg.to_string()
219 };
220 if let Ok(mut guard) = self.last_error.lock() {
221 *guard = Some(truncated);
222 }
223 }
224
225 pub fn snapshot(&self) -> ReverseMetricsSnapshot {
227 let last_error = self.last_error.lock().ok().and_then(|guard| guard.clone());
228 let state_time_ms = self
229 .state_time_ms
230 .lock()
231 .map(|guard| *guard)
232 .unwrap_or([0u64; 6]);
233 ReverseMetricsSnapshot {
234 control_connections_active: self.control_connections_active.load(Ordering::Relaxed),
235 control_connections_accepted_total: self
236 .control_connections_accepted_total
237 .load(Ordering::Relaxed),
238 control_connections_rejected_total: self
239 .control_connections_rejected_total
240 .load(Ordering::Relaxed),
241 control_reconnects_total: self.control_reconnects_total.load(Ordering::Relaxed),
242 auth_failures_total: self.auth_failures_total.load(Ordering::Relaxed),
243 heartbeat_failures_total: self.heartbeat_failures_total.load(Ordering::Relaxed),
244 drain_total: self.drain_total.load(Ordering::Relaxed),
245 drain_duration_ms_total: self.drain_duration_ms_total.load(Ordering::Relaxed),
246 streams_opened_total: self.streams_opened_total.load(Ordering::Relaxed),
247 streams_closed_total: self.streams_closed_total.load(Ordering::Relaxed),
248 stream_bytes_total: self.stream_bytes_total.load(Ordering::Relaxed),
249 state_time_ms,
250 last_error,
251 }
252 }
253
254 pub fn render_prometheus(&self) -> String {
256 let s = self.snapshot();
257 let mut out = String::with_capacity(1024);
258
259 out.push_str("# HELP eggress_reverse_control_connections_active Currently active control connections.\n");
260 out.push_str("# TYPE eggress_reverse_control_connections_active gauge\n");
261 out.push_str(&format!(
262 "eggress_reverse_control_connections_active {}\n",
263 s.control_connections_active
264 ));
265
266 out.push_str("# HELP eggress_reverse_control_connections_accepted_total Total accepted control connections.\n");
267 out.push_str("# TYPE eggress_reverse_control_connections_accepted_total counter\n");
268 out.push_str(&format!(
269 "eggress_reverse_control_connections_accepted_total {}\n",
270 s.control_connections_accepted_total
271 ));
272
273 out.push_str("# HELP eggress_reverse_control_connections_rejected_total Total rejected control connections.\n");
274 out.push_str("# TYPE eggress_reverse_control_connections_rejected_total counter\n");
275 out.push_str(&format!(
276 "eggress_reverse_control_connections_rejected_total {}\n",
277 s.control_connections_rejected_total
278 ));
279
280 out.push_str(
281 "# HELP eggress_reverse_control_reconnects_total Total client reconnect attempts.\n",
282 );
283 out.push_str("# TYPE eggress_reverse_control_reconnects_total counter\n");
284 out.push_str(&format!(
285 "eggress_reverse_control_reconnects_total {}\n",
286 s.control_reconnects_total
287 ));
288
289 out.push_str(
290 "# HELP eggress_reverse_auth_failures_total Total control channel auth failures.\n",
291 );
292 out.push_str("# TYPE eggress_reverse_auth_failures_total counter\n");
293 out.push_str(&format!(
294 "eggress_reverse_auth_failures_total {}\n",
295 s.auth_failures_total
296 ));
297
298 out.push_str(
299 "# HELP eggress_reverse_heartbeat_failures_total Total heartbeat / keepalive failures.\n",
300 );
301 out.push_str("# TYPE eggress_reverse_heartbeat_failures_total counter\n");
302 out.push_str(&format!(
303 "eggress_reverse_heartbeat_failures_total {}\n",
304 s.heartbeat_failures_total
305 ));
306
307 out.push_str("# HELP eggress_reverse_drain_total Total drain operations completed.\n");
308 out.push_str("# TYPE eggress_reverse_drain_total counter\n");
309 out.push_str(&format!("eggress_reverse_drain_total {}\n", s.drain_total));
310
311 out.push_str(
312 "# HELP eggress_reverse_drain_duration_ms_total Cumulative drain duration in milliseconds.\n",
313 );
314 out.push_str("# TYPE eggress_reverse_drain_duration_ms_total counter\n");
315 out.push_str(&format!(
316 "eggress_reverse_drain_duration_ms_total {}\n",
317 s.drain_duration_ms_total
318 ));
319
320 out.push_str(
321 "# HELP eggress_reverse_streams_opened_total Total control sessions initiated.\n",
322 );
323 out.push_str("# TYPE eggress_reverse_streams_opened_total counter\n");
324 out.push_str(&format!(
325 "eggress_reverse_streams_opened_total {}\n",
326 s.streams_opened_total
327 ));
328
329 out.push_str("# HELP eggress_reverse_streams_closed_total Total control sessions completed cleanly.\n");
330 out.push_str("# TYPE eggress_reverse_streams_closed_total counter\n");
331 out.push_str(&format!(
332 "eggress_reverse_streams_closed_total {}\n",
333 s.streams_closed_total
334 ));
335
336 out.push_str("# HELP eggress_reverse_stream_bytes_total Total bytes relayed through control channels.\n");
337 out.push_str("# TYPE eggress_reverse_stream_bytes_total counter\n");
338 out.push_str(&format!(
339 "eggress_reverse_stream_bytes_total {}\n",
340 s.stream_bytes_total
341 ));
342
343 out.push_str(
344 "# HELP eggress_reverse_state_time_ms Cumulative time spent in each control state.\n",
345 );
346 out.push_str("# TYPE eggress_reverse_state_time_ms counter\n");
347 let state_labels = [
348 "disconnected",
349 "connecting",
350 "authenticating",
351 "ready",
352 "draining",
353 "closed",
354 ];
355 for (idx, label) in state_labels.iter().enumerate() {
356 out.push_str(&format!(
357 "eggress_reverse_state_time_ms{{state=\"{}\"}} {}\n",
358 label, s.state_time_ms[idx]
359 ));
360 }
361
362 out
363 }
364}
365
366#[cfg(test)]
367mod tests {
368 use super::*;
369
370 #[test]
371 fn new_counters_are_zero() {
372 let m = ReverseMetrics::new();
373 assert_eq!(m.control_connections_active.load(Ordering::Relaxed), 0);
374 assert_eq!(
375 m.control_connections_accepted_total.load(Ordering::Relaxed),
376 0
377 );
378 assert_eq!(
379 m.control_connections_rejected_total.load(Ordering::Relaxed),
380 0
381 );
382 assert_eq!(m.control_reconnects_total.load(Ordering::Relaxed), 0);
383 assert_eq!(m.auth_failures_total.load(Ordering::Relaxed), 0);
384 assert_eq!(m.heartbeat_failures_total.load(Ordering::Relaxed), 0);
385 assert_eq!(m.drain_total.load(Ordering::Relaxed), 0);
386 assert_eq!(m.drain_duration_ms_total.load(Ordering::Relaxed), 0);
387 assert_eq!(m.streams_opened_total.load(Ordering::Relaxed), 0);
388 assert_eq!(m.streams_closed_total.load(Ordering::Relaxed), 0);
389 assert_eq!(m.stream_bytes_total.load(Ordering::Relaxed), 0);
390 assert_eq!(*m.state_time_ms.lock().unwrap(), [0u64; 6]);
391 assert!(m.last_error.lock().unwrap().is_none());
392 }
393
394 fn peer() -> SocketAddr {
395 "127.0.0.1:9999".parse().unwrap()
396 }
397
398 #[test]
399 fn record_control_accepted() {
400 let m = ReverseMetrics::new();
401 m.record_control_accepted(peer());
402 m.record_control_accepted(peer());
403 assert_eq!(m.control_connections_active.load(Ordering::Relaxed), 2);
404 assert_eq!(
405 m.control_connections_accepted_total.load(Ordering::Relaxed),
406 2
407 );
408 }
409
410 #[test]
411 fn record_control_rejected() {
412 let m = ReverseMetrics::new();
413 m.record_control_rejected(peer(), "max_control_connections");
414 assert_eq!(
415 m.control_connections_rejected_total.load(Ordering::Relaxed),
416 1
417 );
418 assert_eq!(m.auth_failures_total.load(Ordering::Relaxed), 0);
419 assert_eq!(
420 m.last_error.lock().unwrap().as_deref(),
421 Some("max_control_connections")
422 );
423 }
424
425 #[test]
426 fn record_auth_failure() {
427 let m = ReverseMetrics::new();
428 m.record_auth_failure(peer(), "bad credentials");
429 assert_eq!(m.auth_failures_total.load(Ordering::Relaxed), 1);
430 assert_eq!(
431 m.last_error.lock().unwrap().as_deref(),
432 Some("bad credentials")
433 );
434 }
435
436 #[test]
437 fn record_heartbeat_failure() {
438 let m = ReverseMetrics::new();
439 m.record_heartbeat_failure();
440 m.record_heartbeat_failure();
441 assert_eq!(m.heartbeat_failures_total.load(Ordering::Relaxed), 2);
442 }
443
444 #[test]
445 fn record_drain() {
446 let m = ReverseMetrics::new();
447 m.record_drain(150);
448 m.record_drain(50);
449 assert_eq!(m.drain_total.load(Ordering::Relaxed), 2);
450 assert_eq!(m.drain_duration_ms_total.load(Ordering::Relaxed), 200);
451 }
452
453 #[test]
454 fn record_reconnect() {
455 let m = ReverseMetrics::new();
456 m.record_reconnect();
457 m.record_reconnect();
458 m.record_reconnect();
459 assert_eq!(m.control_reconnects_total.load(Ordering::Relaxed), 3);
460 }
461
462 #[test]
463 fn record_stream_opened_and_closed() {
464 let m = ReverseMetrics::new();
465 m.record_stream_opened();
466 m.record_stream_opened();
467 assert_eq!(m.streams_opened_total.load(Ordering::Relaxed), 2);
468
469 m.record_stream_closed(1024);
470 assert_eq!(m.streams_closed_total.load(Ordering::Relaxed), 1);
471 assert_eq!(m.stream_bytes_total.load(Ordering::Relaxed), 1024);
472
473 m.record_stream_closed(512);
474 assert_eq!(m.streams_closed_total.load(Ordering::Relaxed), 2);
475 assert_eq!(m.stream_bytes_total.load(Ordering::Relaxed), 1536);
476 }
477
478 #[test]
479 fn record_state_duration() {
480 let m = ReverseMetrics::new();
481 m.record_state_duration(ControlState::Connecting, 100);
482 m.record_state_duration(ControlState::Ready, 500);
483 m.record_state_duration(ControlState::Ready, 250);
484 let snap = m.snapshot();
485 assert_eq!(snap.state_ms(ControlState::Connecting), 100);
486 assert_eq!(snap.state_ms(ControlState::Ready), 750);
487 assert_eq!(snap.state_ms(ControlState::Draining), 0);
488 }
489
490 #[test]
491 fn record_error_truncates_long_messages() {
492 let m = ReverseMetrics::new();
493 let long_msg = "x".repeat(500);
494 m.record_error(&long_msg);
495 let err = m.last_error.lock().unwrap().clone().unwrap();
496 assert_eq!(err.chars().count(), MAX_ERROR_LEN + 1);
498 assert!(err.ends_with('…'));
499 }
500
501 #[test]
502 fn record_error_handles_multibyte_messages() {
503 let m = ReverseMetrics::new();
504 m.record_error(&"é".repeat(MAX_ERROR_LEN + 1));
505 let err = m.last_error.lock().unwrap().clone().unwrap();
506 assert_eq!(err.chars().count(), MAX_ERROR_LEN + 1);
507 assert!(err.ends_with('…'));
508 }
509
510 #[test]
511 fn snapshot_returns_expected_values() {
512 let m = ReverseMetrics::new();
513 m.record_control_accepted(peer());
514 m.record_control_accepted(peer());
515 m.record_stream_opened();
516 m.record_stream_closed(2048);
517 m.record_reconnect();
518 m.record_control_rejected(peer(), "timeout");
519 m.record_auth_failure(peer(), "bad password");
520 m.record_heartbeat_failure();
521 m.record_drain(75);
522 m.record_state_duration(ControlState::Ready, 1000);
523
524 let snap = m.snapshot();
525 assert_eq!(snap.control_connections_active, 2);
526 assert_eq!(snap.control_connections_accepted_total, 2);
527 assert_eq!(snap.control_connections_rejected_total, 2);
528 assert_eq!(snap.auth_failures_total, 1);
529 assert_eq!(snap.heartbeat_failures_total, 1);
530 assert_eq!(snap.drain_total, 1);
531 assert_eq!(snap.drain_duration_ms_total, 75);
532 assert_eq!(snap.streams_opened_total, 1);
533 assert_eq!(snap.streams_closed_total, 1);
534 assert_eq!(snap.stream_bytes_total, 2048);
535 assert_eq!(snap.state_ms(ControlState::Ready), 1000);
536 assert_eq!(snap.last_error.as_deref(), Some("bad password"));
537 }
538
539 #[test]
540 fn snapshot_display_summary() {
541 let m = ReverseMetrics::new();
542 m.record_control_accepted(peer());
543 m.record_stream_opened();
544 m.record_stream_closed(100);
545 m.record_drain(20);
546 let snap = m.snapshot();
547 let display = snap.display_summary();
548 assert!(display.contains("active=1"));
549 assert!(display.contains("accepted=1"));
550 assert!(display.contains("bytes=100"));
551 assert!(display.contains("drain_total=1"));
552 assert!(display.contains("drain_ms=20"));
553 assert!(display.contains("last_error=(none)"));
554 }
555
556 #[test]
557 fn prometheus_output_contains_expected_names() {
558 let m = ReverseMetrics::new();
559 m.record_control_accepted(peer());
560 m.record_stream_opened();
561 m.record_stream_closed(42);
562 m.record_heartbeat_failure();
563 m.record_drain(5);
564 m.record_state_duration(ControlState::Ready, 100);
565
566 let prom = m.render_prometheus();
567 assert!(prom.contains("eggress_reverse_control_connections_active"));
568 assert!(prom.contains("eggress_reverse_control_connections_accepted_total"));
569 assert!(prom.contains("eggress_reverse_control_connections_rejected_total"));
570 assert!(prom.contains("eggress_reverse_control_reconnects_total"));
571 assert!(prom.contains("eggress_reverse_auth_failures_total"));
572 assert!(prom.contains("eggress_reverse_heartbeat_failures_total"));
573 assert!(prom.contains("eggress_reverse_drain_total"));
574 assert!(prom.contains("eggress_reverse_drain_duration_ms_total"));
575 assert!(prom.contains("eggress_reverse_streams_opened_total"));
576 assert!(prom.contains("eggress_reverse_streams_closed_total"));
577 assert!(prom.contains("eggress_reverse_stream_bytes_total"));
578 assert!(prom.contains("eggress_reverse_state_time_ms"));
579 assert!(prom.contains("TYPE eggress_reverse_control_connections_active gauge"));
580 assert!(prom.contains("TYPE eggress_reverse_streams_opened_total counter"));
581 assert!(prom.contains("state=\"ready\""));
582 }
583
584 #[test]
585 fn snapshot_is_clone() {
586 let m = ReverseMetrics::new();
587 m.record_control_accepted(peer());
588 let s1 = m.snapshot();
589 let s2 = s1.clone();
590 assert_eq!(s1.control_connections_active, s2.control_connections_active);
591 }
592
593 #[test]
594 fn snapshot_is_debug() {
595 let m = ReverseMetrics::new();
596 let snap = m.snapshot();
597 let debug_str = format!("{:?}", snap);
598 assert!(debug_str.contains("ReverseMetricsSnapshot"));
599 }
600
601 #[test]
602 fn snapshot_is_serialize() {
603 let m = ReverseMetrics::new();
604 let snap = m.snapshot();
605 let json = serde_json::to_string(&snap).unwrap();
606 assert!(json.contains("control_connections_active"));
607 assert!(json.contains("stream_bytes_total"));
608 assert!(json.contains("auth_failures_total"));
609 assert!(json.contains("state_time_ms"));
610 }
611}