use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
#[derive(Debug)]
pub(crate) struct PollTracker {
last_poll: parking_lot::Mutex<tokio::time::Instant>,
max_poll_interval: Duration,
exceeded: std::sync::atomic::AtomicBool,
}
impl PollTracker {
pub(crate) fn new(max_poll_interval: Duration) -> Self {
Self {
last_poll: parking_lot::Mutex::new(tokio::time::Instant::now()),
max_poll_interval,
exceeded: std::sync::atomic::AtomicBool::new(false),
}
}
pub(crate) fn note_poll(&self) {
*self.last_poll.lock() = tokio::time::Instant::now();
}
pub(crate) fn elapsed(&self) -> Duration {
self.last_poll.lock().elapsed()
}
pub(crate) fn max_poll_interval(&self) -> Duration {
self.max_poll_interval
}
pub(crate) fn is_expired(&self) -> bool {
self.elapsed() > self.max_poll_interval
}
pub(crate) fn mark_exceeded(&self) -> bool {
!self
.exceeded
.swap(true, std::sync::atomic::Ordering::SeqCst)
}
pub(crate) fn exceeded(&self) -> bool {
self.exceeded.load(std::sync::atomic::Ordering::SeqCst)
}
pub(crate) fn reset(&self) {
self.exceeded
.store(false, std::sync::atomic::Ordering::SeqCst);
self.note_poll();
}
}
#[derive(Debug, Default)]
pub(crate) struct HeartbeatController {
running: AtomicBool,
rebalance_needed: AtomicBool,
member_invalidated: AtomicBool,
}
impl HeartbeatController {
pub(crate) fn is_running(&self) -> bool {
self.running.load(Ordering::Acquire)
}
pub(crate) fn start(&self) {
self.running.store(true, Ordering::Release);
}
pub(crate) fn stop(&self) {
self.running.store(false, Ordering::Release);
}
pub(crate) fn signal_rebalance(&self) {
self.rebalance_needed.store(true, Ordering::Release);
}
pub(crate) fn take_rebalance_needed(&self) -> bool {
self.rebalance_needed.swap(false, Ordering::AcqRel)
}
pub(crate) fn signal_member_invalidated(&self) {
self.member_invalidated.store(true, Ordering::Release);
self.rebalance_needed.store(true, Ordering::Release);
}
pub(crate) fn take_member_invalidated(&self) -> bool {
self.member_invalidated.swap(false, Ordering::AcqRel)
}
}
#[derive(Debug)]
pub(crate) enum HeartbeatCommand {
Stop,
AcknowledgeRevocation,
}