use crate::connection::{RECONNECT_MAX_MS, RECONNECT_MIN_MS};
use std::time::{Duration, Instant};
use tracing::Level;
pub(crate) struct Note {
pub(crate) level: Level,
pub(crate) message: String,
}
const WARN_INTERVAL: Duration = Duration::from_secs(30);
#[derive(Default)]
pub(crate) struct Reporter {
lost_at: Option<Instant>,
loss_reported: bool,
said_at: Option<Instant>,
unwarned: u64,
}
impl Reporter {
pub(crate) fn observe(&mut self, events: &[zmq::SocketEvent]) -> Vec<Note> {
let mut said = Vec::new();
for event in events {
match event {
zmq::SocketEvent::DISCONNECTED if self.lost_at.is_none() => {
self.lost_at = Some(Instant::now());
self.loss_reported = self.warn(
&format!(
"Connection lost, retrying every {} to {} seconds",
RECONNECT_MIN_MS / 1_000,
RECONNECT_MAX_MS / 1_000,
),
&mut said,
);
}
zmq::SocketEvent::CONNECTED => {
if let Some(at) = self.lost_at.take() {
if self.loss_reported {
said.push(Note {
level: Level::INFO,
message: format!(
"Connected again, {}s without one",
at.elapsed().as_secs()
),
});
}
}
}
zmq::SocketEvent::HANDSHAKE_FAILED_NO_DETAIL
| zmq::SocketEvent::HANDSHAKE_FAILED_PROTOCOL
| zmq::SocketEvent::HANDSHAKE_FAILED_AUTH => {
self.warn(
"Connected to something that is not EDDN",
&mut said,
);
}
_ => {}
}
}
said
}
fn warn(&mut self, message: &str, said: &mut Vec<Note>) -> bool {
if self.said_at.is_some_and(|at| at.elapsed() < WARN_INTERVAL) {
self.unwarned += 1;
return false;
}
said.push(Note {
level: Level::WARN,
message: match self.unwarned {
0 => message.to_string(),
n => format!("{} ({} more went unreported)", message, n),
},
});
self.unwarned = 0;
self.said_at = Some(Instant::now());
true
}
pub(crate) fn replaced(&mut self) {
self.lost_at = None;
self.loss_reported = false;
}
}
#[cfg(test)]
mod tests {
use super::*;
use zmq::SocketEvent::*;
fn age(reporter: &mut Reporter, by: Duration) {
reporter.said_at = reporter.said_at.map(|at| at - by);
}
#[test]
fn losing_a_connection_over_and_over_is_reported_once_an_interval() {
let mut reporter = Reporter::default();
let said = reporter.observe(&[DISCONNECTED, CONNECTED]);
assert_eq!(said.len(), 2);
assert_eq!(said[0].level, Level::WARN);
assert!(said[0].message.starts_with("Connection lost"));
assert_eq!(said[1].level, Level::INFO);
assert!(said[1].message.starts_with("Connected again"));
for _ in 0..50 {
assert!(reporter.observe(&[DISCONNECTED, CONNECTED]).is_empty());
}
age(&mut reporter, WARN_INTERVAL);
let said = reporter.observe(&[DISCONNECTED]);
assert_eq!(said.len(), 1);
assert!(
said[0].message.contains("50 more went unreported"),
"{}",
said[0].message
);
}
#[test]
fn a_connection_coming_back_is_reported_only_if_its_loss_was() {
let mut reporter = Reporter::default();
assert_eq!(reporter.observe(&[DISCONNECTED]).len(), 1);
assert_eq!(reporter.observe(&[CONNECTED]).len(), 1);
assert!(reporter.observe(&[DISCONNECTED]).is_empty());
assert!(reporter.observe(&[CONNECTED]).is_empty());
}
#[test]
fn a_connection_replaced_by_hand_is_not_waiting_on_a_recovery() {
let mut reporter = Reporter::default();
assert_eq!(reporter.observe(&[DISCONNECTED]).len(), 1);
reporter.replaced();
assert!(reporter.observe(&[CONNECTED]).is_empty());
}
}