use std::io;
use std::sync::atomic::AtomicUsize;
use std::sync::atomic::Ordering;
use super::TimerFailurePoint;
use crate::ManualMonotonicClock;
use crate::MonotonicClock;
use crate::MonotonicInstant;
use crate::TimeError;
use crate::Timer;
use crate::TimerFuture;
use crate::TimerUnavailableError;
pub struct FaultInjectingTimer {
clock: ManualMonotonicClock,
failure_point: TimerFailurePoint,
error_factory: Box<dyn Fn() -> TimeError + Send + Sync + 'static>,
registration_count: AtomicUsize,
}
impl FaultInjectingTimer {
#[must_use]
pub fn new<F>(failure_point: TimerFailurePoint, error_factory: F) -> Self
where
F: Fn() -> TimeError + Send + Sync + 'static,
{
Self {
clock: ManualMonotonicClock::new(),
failure_point,
error_factory: Box::new(error_factory),
registration_count: AtomicUsize::new(0),
}
}
#[must_use]
pub fn backend_unavailable(failure_point: TimerFailurePoint, backend: &'static str, message: &str) -> Self {
let message = message.to_owned();
Self::new(failure_point, move || TimeError::TimerUnavailable {
source: TimerUnavailableError::BackendUnavailable {
backend,
source: Box::new(io::Error::other(message.clone())),
},
})
}
#[must_use]
pub fn failure_point(&self) -> TimerFailurePoint {
self.failure_point
}
#[must_use]
#[inline(always)]
pub fn registration_count(&self) -> usize {
self.registration_count.load(Ordering::Relaxed)
}
}
impl Timer for FaultInjectingTimer {
#[inline(always)]
fn clock(&self) -> &dyn MonotonicClock {
&self.clock
}
fn at(&self, deadline: MonotonicInstant) -> Result<TimerFuture, TimeError> {
let now = self.clock.now();
deadline.validate_domain(now.domain())?;
if deadline <= now {
return Ok(Box::pin(std::future::ready(Ok(()))));
}
self.registration_count.fetch_add(1, Ordering::Relaxed);
let error = (self.error_factory)();
match self.failure_point {
TimerFailurePoint::Registration => Err(error),
TimerFailurePoint::Completion => Ok(Box::pin(std::future::ready(Err(error)))),
}
}
}