use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::{Duration, Instant};
pub const POLL_INTERVAL_MS: i32 = 1_000;
pub const RECONNECT_MIN_MS: i32 = 1_000;
pub const RECONNECT_MAX_MS: i32 = 15_000;
pub const HEARTBEAT_IVL_MS: i32 = 5_000;
pub const HEARTBEAT_TIMEOUT_MS: i32 = 10_000;
const MONITORED: i32 = zmq::SocketEvent::CONNECTED as i32
| zmq::SocketEvent::DISCONNECTED as i32
| zmq::SocketEvent::CONNECT_RETRIED as i32
| zmq::SocketEvent::HANDSHAKE_FAILED_NO_DETAIL as i32
| zmq::SocketEvent::HANDSHAKE_FAILED_PROTOCOL as i32
| zmq::SocketEvent::HANDSHAKE_FAILED_AUTH as i32;
static MONITORS: AtomicUsize = AtomicUsize::new(0);
pub(crate) struct Connection {
pub(crate) socket: zmq::Socket,
monitor: zmq::Socket,
}
impl Connection {
pub(crate) fn open(
ctx: &zmq::Context,
url: &str,
) -> Result<Self, zmq::Error> {
let socket = ctx.socket(zmq::SUB)?;
socket.set_reconnect_ivl(RECONNECT_MIN_MS)?;
socket.set_reconnect_ivl_max(RECONNECT_MAX_MS)?;
socket.set_heartbeat_ivl(HEARTBEAT_IVL_MS)?;
socket.set_heartbeat_timeout(HEARTBEAT_TIMEOUT_MS)?;
socket.set_rcvtimeo(POLL_INTERVAL_MS)?;
let endpoint = format!(
"inproc://eddn-monitor-{}",
MONITORS.fetch_add(1, Ordering::Relaxed)
);
socket.monitor(&endpoint, MONITORED)?;
let monitor = ctx.socket(zmq::PAIR)?;
monitor.connect(&endpoint)?;
socket.connect(url)?;
socket.set_subscribe(&[])?;
Ok(Connection { socket, monitor })
}
pub(crate) fn events(&self) -> Vec<zmq::SocketEvent> {
let mut events = Vec::new();
while let Ok(frames) = self.monitor.recv_multipart(zmq::DONTWAIT) {
match frames.first() {
Some(frame) if frame.len() >= 2 => {
let id = u16::from_le_bytes([frame[0], frame[1]]);
events.push(zmq::SocketEvent::from_raw(id));
}
_ => {}
}
}
events
}
}
pub(crate) struct Stall {
timeout: Duration,
quiet_since: Instant,
}
impl Stall {
pub(crate) fn new(timeout: Duration) -> Self {
Stall { timeout, quiet_since: Instant::now() }
}
pub(crate) fn restart(&mut self) {
self.quiet_since = Instant::now();
}
pub(crate) fn overrun(&self) -> Option<Duration> {
let quiet = self.quiet_since.elapsed();
(quiet >= self.timeout).then_some(quiet)
}
}
#[cfg(test)]
mod tests {
use super::*;
fn opened() -> Connection {
let ctx = zmq::Context::new();
Connection::open(&ctx, "tcp://127.0.0.1:9599").unwrap()
}
#[test]
fn quiet_for_longer_than_the_stall_timeout_is_an_overrun() {
let mut stall = Stall::new(Duration::from_secs(30));
assert!(stall.overrun().is_none());
stall.quiet_since = Instant::now() - Duration::from_secs(29);
assert!(stall.overrun().is_none());
stall.quiet_since = Instant::now() - Duration::from_secs(31);
assert_eq!(stall.overrun().map(|d| d.as_secs()), Some(31));
stall.restart();
assert!(stall.overrun().is_none());
}
#[test]
fn a_connection_pings_and_gives_up_on_an_unanswered_one() {
let connection = opened();
let ivl = connection.socket.get_heartbeat_ivl().unwrap();
let timeout = connection.socket.get_heartbeat_timeout().unwrap();
assert_ne!(ivl, 0);
assert_eq!(ivl + timeout, 15_000);
}
#[test]
fn a_connection_retries_from_the_floor_and_no_slower_than_the_ceiling() {
let connection = opened();
let floor = connection.socket.get_reconnect_ivl().unwrap();
let ceiling = connection.socket.get_reconnect_ivl_max().unwrap();
assert!(floor >= 1_000, "the floor is the rate, and it is {}", floor);
assert!(ceiling >= floor);
assert_eq!(connection.socket.get_rcvtimeo(), Ok(POLL_INTERVAL_MS));
}
}