use std::time::Duration;
use tokio::time::Instant;
#[must_use]
#[inline]
fn deadline(from: Instant, budget: Duration) -> Instant {
from.checked_add(budget).unwrap_or(from)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum LivenessAction {
Idle,
SendHeartbeat,
InboundStalled {
budget: Duration,
},
}
#[derive(Debug)]
pub(crate) struct Liveness {
slack: Duration,
inbound_budget: Option<Duration>,
last_inbound: Instant,
heartbeat_every: Option<Duration>,
last_outbound: Instant,
}
impl Liveness {
#[must_use]
pub(crate) fn new(slack: Duration, heartbeat_every: Option<Duration>, now: Instant) -> Self {
Self {
slack,
inbound_budget: None,
last_inbound: now,
heartbeat_every,
last_outbound: now,
}
}
pub(crate) fn on_stream_opened(&mut self, open_timeout: Duration, now: Instant) {
self.inbound_budget = Some(open_timeout);
self.last_inbound = now;
}
pub(crate) fn on_bound(&mut self, keep_alive: Duration, now: Instant) {
self.inbound_budget = Some(keep_alive.checked_add(self.slack).unwrap_or(keep_alive));
self.last_inbound = now;
}
pub(crate) fn on_stream_closed(&mut self) {
self.inbound_budget = None;
}
pub(crate) fn on_inbound(&mut self, now: Instant) {
self.last_inbound = now;
}
pub(crate) fn on_outbound(&mut self, now: Instant) {
self.last_outbound = now;
}
#[must_use]
pub(crate) fn next_deadline(&self) -> Option<Instant> {
let inbound = self
.inbound_budget
.map(|budget| deadline(self.last_inbound, budget));
let outbound = self
.heartbeat_every
.map(|every| deadline(self.last_outbound, every));
match (inbound, outbound) {
(Some(a), Some(b)) => Some(a.min(b)),
(Some(a), None) => Some(a),
(None, Some(b)) => Some(b),
(None, None) => None,
}
}
#[must_use]
pub(crate) fn due(&self, now: Instant) -> LivenessAction {
if let Some(budget) = self.inbound_budget
&& now >= deadline(self.last_inbound, budget)
{
return LivenessAction::InboundStalled { budget };
}
if let Some(every) = self.heartbeat_every
&& now >= deadline(self.last_outbound, every)
{
return LivenessAction::SendHeartbeat;
}
LivenessAction::Idle
}
}
#[cfg(test)]
mod tests {
use super::*;
const KEEP_ALIVE: Duration = Duration::from_millis(5000);
const SLACK: Duration = Duration::from_millis(3000);
#[tokio::test(start_paused = true)]
async fn test_liveness_disarmed_before_open_has_no_deadline() {
let liveness = Liveness::new(SLACK, None, Instant::now());
assert_eq!(liveness.next_deadline(), None);
assert_eq!(liveness.due(Instant::now()), LivenessAction::Idle);
}
#[tokio::test(start_paused = true)]
async fn test_liveness_open_timeout_arms_the_inbound_clock() {
let start = Instant::now();
let mut liveness = Liveness::new(SLACK, None, start);
liveness.on_stream_opened(Duration::from_secs(10), start);
assert_eq!(liveness.due(start), LivenessAction::Idle);
assert_eq!(
liveness.due(deadline(start, Duration::from_secs(10))),
LivenessAction::InboundStalled {
budget: Duration::from_secs(10)
}
);
}
#[tokio::test(start_paused = true)]
async fn test_liveness_budget_is_keepalive_plus_slack_after_conok() {
let start = Instant::now();
let mut liveness = Liveness::new(SLACK, None, start);
liveness.on_bound(KEEP_ALIVE, start);
assert_eq!(
liveness.due(deadline(start, KEEP_ALIVE)),
LivenessAction::Idle
);
assert_eq!(
liveness.due(deadline(start, Duration::from_millis(8000))),
LivenessAction::InboundStalled {
budget: Duration::from_millis(8000)
}
);
}
#[tokio::test(start_paused = true)]
async fn test_liveness_any_inbound_line_resets_the_budget() {
let start = Instant::now();
let mut liveness = Liveness::new(SLACK, None, start);
liveness.on_bound(KEEP_ALIVE, start);
let later = deadline(start, Duration::from_millis(7000));
liveness.on_inbound(later);
assert_eq!(
liveness.due(deadline(start, Duration::from_millis(8000))),
LivenessAction::Idle
);
assert_eq!(
liveness.due(deadline(later, Duration::from_millis(8000))),
LivenessAction::InboundStalled {
budget: Duration::from_millis(8000)
}
);
}
#[tokio::test(start_paused = true)]
async fn test_liveness_disarm_makes_silence_meaningless_again() {
let start = Instant::now();
let mut liveness = Liveness::new(SLACK, None, start);
liveness.on_bound(KEEP_ALIVE, start);
liveness.on_stream_closed();
assert_eq!(
liveness.due(deadline(start, Duration::from_secs(600))),
LivenessAction::Idle
);
assert_eq!(liveness.next_deadline(), None);
}
#[tokio::test(start_paused = true)]
async fn test_liveness_heartbeat_is_due_after_the_configured_interval() {
let start = Instant::now();
let mut liveness = Liveness::new(SLACK, Some(Duration::from_secs(2)), start);
liveness.on_bound(KEEP_ALIVE, start);
assert_eq!(liveness.due(start), LivenessAction::Idle);
assert_eq!(
liveness.due(deadline(start, Duration::from_secs(2))),
LivenessAction::SendHeartbeat
);
}
#[tokio::test(start_paused = true)]
async fn test_liveness_any_outbound_request_defers_the_heartbeat() {
let start = Instant::now();
let mut liveness = Liveness::new(SLACK, Some(Duration::from_secs(2)), start);
liveness.on_bound(KEEP_ALIVE, start);
let sent = deadline(start, Duration::from_millis(1500));
liveness.on_outbound(sent);
assert_eq!(
liveness.due(deadline(start, Duration::from_secs(2))),
LivenessAction::Idle
);
assert_eq!(
liveness.due(deadline(sent, Duration::from_secs(2))),
LivenessAction::SendHeartbeat
);
}
#[tokio::test(start_paused = true)]
async fn test_liveness_stall_outranks_heartbeat() {
let start = Instant::now();
let mut liveness = Liveness::new(SLACK, Some(Duration::from_secs(2)), start);
liveness.on_bound(KEEP_ALIVE, start);
assert_eq!(
liveness.due(deadline(start, Duration::from_secs(60))),
LivenessAction::InboundStalled {
budget: Duration::from_millis(8000)
}
);
}
#[tokio::test(start_paused = true)]
async fn test_liveness_next_deadline_is_the_earlier_clock() {
let start = Instant::now();
let mut liveness = Liveness::new(SLACK, Some(Duration::from_secs(2)), start);
liveness.on_bound(KEEP_ALIVE, start);
assert_eq!(
liveness.next_deadline(),
Some(deadline(start, Duration::from_secs(2)))
);
}
}