use std::time::{Duration, Instant};
pub const MAX_HEARTBEAT_INTERVAL_SECS: u64 = 3600;
const TEST_REQUEST_GRACE_DIVISOR: u32 = 5;
const MIN_TEST_REQUEST_GRACE: Duration = Duration::from_millis(250);
#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
#[error("counterparty HeartBtInt (108) of {secs}s exceeds the maximum supported {max}s")]
pub struct HeartbeatIntervalError {
pub secs: u64,
pub max: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TestRequestOutcome {
NonePending,
Confirmed,
SupersededByTraffic,
}
pub fn negotiate_interval(
requested: Duration,
confirmed_secs: u64,
) -> Result<Duration, HeartbeatIntervalError> {
let confirmed = Duration::from_secs(confirmed_secs);
if confirmed == requested {
return Ok(confirmed);
}
if confirmed_secs > MAX_HEARTBEAT_INTERVAL_SECS {
return Err(HeartbeatIntervalError {
secs: confirmed_secs,
max: MAX_HEARTBEAT_INTERVAL_SECS,
});
}
Ok(confirmed)
}
fn grace_for(interval: Duration) -> Duration {
if interval.is_zero() {
return Duration::ZERO;
}
(interval / TEST_REQUEST_GRACE_DIVISOR).max(MIN_TEST_REQUEST_GRACE)
}
#[derive(Debug)]
pub struct HeartbeatManager {
interval: Duration,
grace: Duration,
last_sent: Instant,
last_received: Instant,
test_request_pending: Option<String>,
test_request_sent_at: Option<Instant>,
}
impl HeartbeatManager {
#[must_use]
pub fn new(interval: Duration) -> Self {
let now = Instant::now();
Self {
interval,
grace: grace_for(interval),
last_sent: now,
last_received: now,
test_request_pending: None,
test_request_sent_at: None,
}
}
#[must_use]
pub const fn is_enabled(&self) -> bool {
!self.interval.is_zero()
}
#[inline]
pub fn on_message_sent(&mut self) {
self.last_sent = Instant::now();
}
pub fn on_message_received(
&mut self,
is_heartbeat: bool,
test_req_id: Option<&str>,
) -> TestRequestOutcome {
self.last_received = Instant::now();
let Some(pending) = self.test_request_pending.take() else {
return TestRequestOutcome::NonePending;
};
self.test_request_sent_at = None;
if is_heartbeat && test_req_id == Some(pending.as_str()) {
TestRequestOutcome::Confirmed
} else {
TestRequestOutcome::SupersededByTraffic
}
}
#[must_use]
pub fn should_send_heartbeat(&self) -> bool {
self.is_enabled() && self.last_sent.elapsed() >= self.interval
}
#[must_use]
pub fn should_send_test_request(&self) -> bool {
if !self.is_enabled() || self.test_request_pending.is_some() {
return false;
}
match self.interval.checked_add(self.grace) {
Some(due_after) => self.last_received.elapsed() >= due_after,
None => false,
}
}
#[must_use]
pub fn is_timed_out(&self) -> bool {
if !self.is_enabled() {
return false;
}
match self.test_request_sent_at {
Some(sent_at) => sent_at.elapsed() >= self.interval,
None => false,
}
}
pub fn on_test_request_sent(&mut self, test_req_id: String) {
self.test_request_pending = Some(test_req_id);
self.test_request_sent_at = Some(Instant::now());
self.last_sent = Instant::now();
}
#[must_use]
pub fn pending_test_request(&self) -> Option<&str> {
self.test_request_pending.as_deref()
}
#[must_use]
pub fn time_since_last_received(&self) -> Duration {
self.last_received.elapsed()
}
#[must_use]
pub fn time_since_last_sent(&self) -> Duration {
self.last_sent.elapsed()
}
#[must_use]
pub const fn interval(&self) -> Duration {
self.interval
}
#[must_use]
pub const fn test_request_grace(&self) -> Duration {
self.grace
}
pub fn reset(&mut self) {
let now = Instant::now();
self.last_sent = now;
self.last_received = now;
self.test_request_pending = None;
self.test_request_sent_at = None;
}
}
#[must_use]
pub fn generate_test_req_id() -> String {
use std::time::{SystemTime, UNIX_EPOCH};
let nanos = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_nanos();
format!("TEST{}", nanos)
}
#[cfg(test)]
mod tests {
use super::*;
use std::thread::sleep;
#[test]
fn test_heartbeat_manager_new() {
let mgr = HeartbeatManager::new(Duration::from_secs(30));
assert_eq!(mgr.interval(), Duration::from_secs(30));
assert!(mgr.pending_test_request().is_none());
assert!(mgr.is_enabled());
}
#[test]
fn test_should_send_heartbeat() {
let mgr = HeartbeatManager::new(Duration::from_millis(10));
assert!(!mgr.should_send_heartbeat());
sleep(Duration::from_millis(15));
assert!(mgr.should_send_heartbeat());
}
#[test]
fn test_on_message_sent() {
let mut mgr = HeartbeatManager::new(Duration::from_millis(10));
sleep(Duration::from_millis(15));
assert!(mgr.should_send_heartbeat());
mgr.on_message_sent();
assert!(!mgr.should_send_heartbeat());
}
#[test]
fn test_test_request_pending() {
let mut mgr = HeartbeatManager::new(Duration::from_secs(30));
mgr.on_test_request_sent("TEST123".to_string());
assert_eq!(mgr.pending_test_request(), Some("TEST123"));
let outcome = mgr.on_message_received(true, Some("TEST123"));
assert_eq!(outcome, TestRequestOutcome::Confirmed);
assert!(mgr.pending_test_request().is_none());
}
#[test]
fn test_generate_test_req_id() {
let id1 = generate_test_req_id();
std::thread::sleep(std::time::Duration::from_nanos(1));
let id2 = generate_test_req_id();
assert!(id1.starts_with("TEST"));
assert!(id2.starts_with("TEST"));
assert!(id1.len() > 4);
assert!(id2.len() > 4);
}
#[test]
fn test_zero_interval_reports_disabled() {
let mgr = HeartbeatManager::new(Duration::ZERO);
assert!(!mgr.is_enabled());
assert_eq!(mgr.test_request_grace(), Duration::ZERO);
}
#[test]
fn test_zero_interval_never_sends_heartbeat_or_test_request() {
let mgr = HeartbeatManager::new(Duration::ZERO);
assert!(!mgr.should_send_heartbeat());
assert!(!mgr.should_send_test_request());
sleep(Duration::from_millis(20));
assert!(!mgr.should_send_heartbeat());
assert!(!mgr.should_send_test_request());
}
#[test]
fn test_zero_interval_never_times_out() {
let mut mgr = HeartbeatManager::new(Duration::ZERO);
assert!(!mgr.is_timed_out());
mgr.on_test_request_sent("TEST-ZERO".to_string());
sleep(Duration::from_millis(20));
assert!(!mgr.is_timed_out());
}
#[test]
fn test_test_request_grace_is_one_fifth_of_a_long_interval() {
let mgr = HeartbeatManager::new(Duration::from_secs(30));
assert_eq!(mgr.test_request_grace(), Duration::from_secs(6));
}
#[test]
fn test_test_request_grace_uses_floor_for_short_intervals() {
let mgr = HeartbeatManager::new(Duration::from_millis(100));
assert_eq!(mgr.test_request_grace(), MIN_TEST_REQUEST_GRACE);
}
#[test]
fn test_should_send_test_request_waits_for_interval_plus_grace() {
let interval = Duration::from_millis(2000);
let mgr = HeartbeatManager::new(interval);
assert_eq!(mgr.test_request_grace(), Duration::from_millis(400));
sleep(Duration::from_millis(50));
assert!(!mgr.should_send_test_request());
}
#[test]
fn test_should_send_test_request_boundary() {
let mgr = HeartbeatManager::new(Duration::from_secs(1));
assert_eq!(mgr.test_request_grace(), Duration::from_millis(250));
assert!(
!mgr.should_send_test_request(),
"not due before interval plus grace elapses"
);
sleep(Duration::from_millis(1400));
assert!(
mgr.should_send_test_request(),
"due once interval plus grace has elapsed"
);
}
#[test]
fn test_should_send_test_request_suppressed_while_pending() {
let mut mgr = HeartbeatManager::new(Duration::from_millis(50));
sleep(Duration::from_millis(320));
assert!(mgr.should_send_test_request());
mgr.on_test_request_sent("TEST-PENDING".to_string());
assert!(!mgr.should_send_test_request());
}
#[test]
fn test_should_send_test_request_unreachable_interval_is_never_due() {
let mgr = HeartbeatManager::new(Duration::MAX);
assert!(!mgr.should_send_test_request());
}
#[test]
fn test_predicates_near_duration_max_do_not_panic() {
let mut mgr = HeartbeatManager::new(Duration::MAX);
assert!(!mgr.should_send_heartbeat());
assert!(!mgr.should_send_test_request());
assert!(!mgr.is_timed_out());
mgr.on_test_request_sent("TEST-NEAR-MAX".to_string());
assert!(!mgr.is_timed_out());
assert!(!mgr.should_send_test_request());
}
#[test]
fn test_is_timed_out_false_without_pending_test_request() {
let mgr = HeartbeatManager::new(Duration::from_millis(10));
sleep(Duration::from_millis(50));
assert!(!mgr.is_timed_out());
}
#[test]
fn test_is_timed_out_false_before_the_interval_elapses() {
let mut mgr = HeartbeatManager::new(Duration::from_millis(200));
mgr.on_test_request_sent("TEST-EARLY".to_string());
sleep(Duration::from_millis(20));
assert!(!mgr.is_timed_out());
}
#[test]
fn test_is_timed_out_true_after_silence_since_the_test_request() {
let mut mgr = HeartbeatManager::new(Duration::from_millis(50));
mgr.on_test_request_sent("TEST-SILENT".to_string());
sleep(Duration::from_millis(80));
assert!(mgr.is_timed_out());
}
#[test]
fn test_application_traffic_clears_pending_test_request() {
let mut mgr = HeartbeatManager::new(Duration::from_millis(50));
mgr.on_test_request_sent("TEST-TRAFFIC".to_string());
let outcome = mgr.on_message_received(false, None);
assert_eq!(outcome, TestRequestOutcome::SupersededByTraffic);
assert!(mgr.pending_test_request().is_none());
sleep(Duration::from_millis(80));
assert!(
!mgr.is_timed_out(),
"traffic after the TestRequest proves the peer is alive"
);
}
#[test]
fn test_heartbeat_without_test_req_id_clears_pending_test_request() {
let mut mgr = HeartbeatManager::new(Duration::from_millis(50));
mgr.on_test_request_sent("TEST-NO-112".to_string());
let outcome = mgr.on_message_received(true, None);
assert_eq!(outcome, TestRequestOutcome::SupersededByTraffic);
assert!(mgr.pending_test_request().is_none());
assert!(!mgr.is_timed_out());
}
#[test]
fn test_heartbeat_with_wrong_test_req_id_clears_pending_test_request() {
let mut mgr = HeartbeatManager::new(Duration::from_millis(50));
mgr.on_test_request_sent("TEST-WANTED".to_string());
let outcome = mgr.on_message_received(true, Some("TEST-OTHER"));
assert_eq!(outcome, TestRequestOutcome::SupersededByTraffic);
assert!(mgr.pending_test_request().is_none());
assert!(!mgr.is_timed_out());
}
#[test]
fn test_non_heartbeat_with_matching_id_is_not_confirmation() {
let mut mgr = HeartbeatManager::new(Duration::from_millis(50));
mgr.on_test_request_sent("TEST-ECHO".to_string());
let outcome = mgr.on_message_received(false, Some("TEST-ECHO"));
assert_eq!(outcome, TestRequestOutcome::SupersededByTraffic);
assert!(mgr.pending_test_request().is_none());
}
#[test]
fn test_on_message_received_without_pending_reports_none_pending() {
let mut mgr = HeartbeatManager::new(Duration::from_secs(30));
let outcome = mgr.on_message_received(true, Some("TEST-UNSOLICITED"));
assert_eq!(outcome, TestRequestOutcome::NonePending);
}
#[test]
fn test_reset_clears_pending_test_request_and_timers() {
let mut mgr = HeartbeatManager::new(Duration::from_millis(50));
mgr.on_test_request_sent("TEST-RESET".to_string());
sleep(Duration::from_millis(80));
assert!(mgr.is_timed_out());
mgr.reset();
assert!(mgr.pending_test_request().is_none());
assert!(!mgr.is_timed_out());
assert!(!mgr.should_send_test_request());
assert_eq!(mgr.interval(), Duration::from_millis(50));
assert_eq!(mgr.test_request_grace(), MIN_TEST_REQUEST_GRACE);
}
#[test]
fn test_negotiate_interval_accepts_the_echoed_value() {
let requested = Duration::from_secs(30);
assert_eq!(negotiate_interval(requested, 30), Ok(requested));
}
#[test]
fn test_negotiate_interval_accepts_zero() {
assert_eq!(
negotiate_interval(Duration::from_secs(30), 0),
Ok(Duration::ZERO)
);
}
#[test]
fn test_negotiate_interval_accepts_a_differing_value_within_bounds() {
assert_eq!(
negotiate_interval(Duration::from_secs(30), 60),
Ok(Duration::from_secs(60))
);
}
#[test]
fn test_negotiate_interval_accepts_the_maximum() {
assert_eq!(
negotiate_interval(Duration::from_secs(30), MAX_HEARTBEAT_INTERVAL_SECS),
Ok(Duration::from_secs(MAX_HEARTBEAT_INTERVAL_SECS))
);
}
#[test]
fn test_negotiate_interval_refuses_above_the_maximum() {
let err = negotiate_interval(Duration::from_secs(30), MAX_HEARTBEAT_INTERVAL_SECS + 1);
assert_eq!(
err,
Err(HeartbeatIntervalError {
secs: MAX_HEARTBEAT_INTERVAL_SECS + 1,
max: MAX_HEARTBEAT_INTERVAL_SECS,
})
);
}
#[test]
fn test_negotiate_interval_refuses_an_absurd_value() {
assert!(negotiate_interval(Duration::from_secs(30), u64::MAX).is_err());
}
#[test]
fn test_negotiate_interval_accepts_an_out_of_bounds_echo_of_our_own_request() {
let requested = Duration::from_secs(MAX_HEARTBEAT_INTERVAL_SECS + 100);
assert_eq!(
negotiate_interval(requested, MAX_HEARTBEAT_INTERVAL_SECS + 100),
Ok(requested)
);
}
}