#![allow(clippy::missing_const_for_fn)]
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use chrono::{DateTime, Duration as ChronoDuration, Utc};
use crate::entropy::{Entropy, SeededEntropy};
use crate::time::{ClockSource, TickingClock};
pub(crate) const CHAOS_STREAM_SALT: u64 = 0xC7A0_5EED_C7A0_5EED;
const CHAOS_SKEW_SALT: u64 = 0x5C0F_F5E7_5C0F_F5E7;
#[cfg(feature = "mail")]
const CHAOS_MAIL_SALT: u64 = 0x3A11_FA17_3A11_FA17;
#[non_exhaustive]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ChaosHook {
DbCheckout,
JobDelivery,
MailSend,
}
#[cfg(feature = "mail")]
#[non_exhaustive]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MailFault {
Fail,
Timeout,
}
#[non_exhaustive]
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ChaosEvent {
pub hook: ChaosHook,
pub seq: u64,
pub fired: bool,
}
#[non_exhaustive]
#[derive(Default, Debug, Clone)]
pub struct Chaos {
db_transient_error_prob: f64,
job_duplicate_prob: f64,
clock_skew: Option<std::time::Duration>,
#[cfg(feature = "mail")]
mail_fault_schedule: std::collections::BTreeMap<usize, MailFault>,
#[cfg(feature = "mail")]
mail_transient_error_prob: f64,
}
impl Chaos {
#[must_use]
pub fn db_transient_errors(mut self, p: f64) -> Self {
self.db_transient_error_prob = clamp_prob(p);
self
}
#[must_use]
pub fn job_duplicate_delivery(mut self, p: f64) -> Self {
self.job_duplicate_prob = clamp_prob(p);
self
}
#[must_use]
pub fn clock_skew(mut self, dur: std::time::Duration) -> Self {
self.clock_skew = Some(dur);
self
}
#[cfg(feature = "mail")]
#[must_use]
pub fn smtp_faults(mut self, schedule: impl IntoIterator<Item = (usize, MailFault)>) -> Self {
for (index, fault) in schedule {
if index >= 1 {
self.mail_fault_schedule.insert(index, fault);
}
}
self
}
#[cfg(feature = "mail")]
#[must_use]
pub fn smtp_transient_errors(mut self, p: f64) -> Self {
self.mail_transient_error_prob = clamp_prob(p);
self
}
#[cfg(feature = "mail")]
fn mail_chaos_active(&self) -> bool {
!self.mail_fault_schedule.is_empty() || self.mail_transient_error_prob > 0.0
}
pub(crate) fn is_active(&self) -> bool {
if self.db_transient_error_prob > 0.0
|| self.job_duplicate_prob > 0.0
|| self.clock_skew.is_some()
{
return true;
}
#[cfg(feature = "mail")]
if self.mail_chaos_active() {
return true;
}
false
}
}
fn clamp_prob(p: f64) -> f64 {
if p.is_nan() { 0.0 } else { p.clamp(0.0, 1.0) }
}
#[allow(clippy::cast_precision_loss)] fn unit_from_draw(draw: u64) -> f64 {
(draw >> 11) as f64 / (1u64 << 53) as f64
}
pub(crate) struct ChaosState {
stream: Arc<dyn Entropy>,
db_transient_error_prob: f64,
job_duplicate_prob: f64,
checkout_seq: AtomicU64,
job_seq: AtomicU64,
suppress_next_duplicate: AtomicBool,
#[cfg(feature = "mail")]
mail_stream: Arc<dyn Entropy>,
#[cfg(feature = "mail")]
mail_fault_schedule: std::collections::BTreeMap<usize, MailFault>,
#[cfg(feature = "mail")]
mail_transient_error_prob: f64,
#[cfg(feature = "mail")]
mail_seq: AtomicU64,
events: Mutex<Vec<ChaosEvent>>,
}
impl ChaosState {
pub(crate) fn new(seed: u64, chaos: &Chaos) -> Arc<Self> {
Arc::new(Self {
stream: SeededEntropy::shared(seed ^ CHAOS_STREAM_SALT),
db_transient_error_prob: chaos.db_transient_error_prob,
job_duplicate_prob: chaos.job_duplicate_prob,
checkout_seq: AtomicU64::new(0),
job_seq: AtomicU64::new(0),
suppress_next_duplicate: AtomicBool::new(false),
#[cfg(feature = "mail")]
mail_stream: SeededEntropy::shared(seed ^ CHAOS_MAIL_SALT),
#[cfg(feature = "mail")]
mail_fault_schedule: chaos.mail_fault_schedule.clone(),
#[cfg(feature = "mail")]
mail_transient_error_prob: chaos.mail_transient_error_prob,
#[cfg(feature = "mail")]
mail_seq: AtomicU64::new(0),
events: Mutex::new(Vec::new()),
})
}
fn decide(&self, prob: f64) -> bool {
unit_from_draw(self.stream.next_u64()) < prob
}
#[cfg(feature = "mail")]
fn resolve_mail_fault(&self, send_index: usize) -> Option<MailFault> {
if let Some(&fault) = self.mail_fault_schedule.get(&send_index) {
return Some(fault);
}
if self.mail_transient_error_prob > 0.0
&& unit_from_draw(self.mail_stream.next_u64()) < self.mail_transient_error_prob
{
return Some(MailFault::Fail);
}
None
}
fn record(&self, hook: ChaosHook, seq: u64, fired: bool) {
self.events
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(ChaosEvent { hook, seq, fired });
}
pub(crate) fn events(&self) -> Vec<ChaosEvent> {
self.events
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
}
struct SkewClock {
inner: TickingClock,
offset: ChronoDuration,
}
impl ClockSource for SkewClock {
fn now(&self) -> DateTime<Utc> {
self.inner.now() + self.offset
}
fn monotonic(&self) -> crate::time::MonotonicInstant {
self.inner.monotonic()
}
}
fn deterministic_skew(seed: u64, dur: std::time::Duration) -> ChronoDuration {
let max_nanos = dur.as_nanos();
if max_nanos == 0 {
return ChronoDuration::zero();
}
let stream = SeededEntropy::shared(seed ^ CHAOS_SKEW_SALT);
let draw = u128::from(stream.next_u64());
let picked = draw % (max_nanos.saturating_add(1));
let nanos = i64::try_from(picked).unwrap_or(i64::MAX);
ChronoDuration::nanoseconds(nanos)
}
pub(crate) fn install(
app: crate::test::TestApp,
chaos: &Chaos,
seed: u64,
ticking: TickingClock,
state: Arc<ChaosState>,
) -> crate::test::TestApp {
let mut app = match chaos.clock_skew {
Some(dur) => app.with_clock(SkewClock {
inner: ticking,
offset: deterministic_skew(seed, dur),
}),
None => app.with_clock(ticking),
};
app = app.with_job_interceptor(ChaosJobInterceptor {
state: Arc::clone(&state),
});
#[cfg(feature = "db")]
{
app = app.with_db_interceptor(ChaosDbInterceptor {
state: Arc::clone(&state),
});
}
#[cfg(feature = "mail")]
if !state.mail_fault_schedule.is_empty() || state.mail_transient_error_prob > 0.0 {
app = app.with_mail_interceptor(ChaosMailInterceptor {
state: Arc::clone(&state),
});
}
drop(state);
app
}
struct ChaosJobInterceptor {
state: Arc<ChaosState>,
}
impl crate::interceptor::JobInterceptor for ChaosJobInterceptor {
fn intercept_enqueue<'a>(
&'a self,
name: &'a str,
payload: &'a serde_json::Value,
next: std::pin::Pin<
Box<dyn std::future::Future<Output = crate::AutumnResult<()>> + Send + 'a>,
>,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = crate::AutumnResult<()>> + Send + 'a>>
{
Box::pin(async move {
if self
.state
.suppress_next_duplicate
.swap(false, Ordering::SeqCst)
{
return next.await;
}
let seq = self.state.job_seq.fetch_add(1, Ordering::SeqCst);
let fired = self.state.decide(self.state.job_duplicate_prob);
self.state.record(ChaosHook::JobDelivery, seq, fired);
let res = next.await;
if fired && res.is_ok() {
self.state
.suppress_next_duplicate
.store(true, Ordering::SeqCst);
if crate::job::enqueue(name, payload.clone()).await.is_err() {
self.state
.suppress_next_duplicate
.store(false, Ordering::SeqCst);
}
}
res
})
}
fn intercept_execute<'a>(
&'a self,
_name: &'a str,
_payload: &'a serde_json::Value,
next: std::pin::Pin<
Box<dyn std::future::Future<Output = crate::AutumnResult<()>> + Send + 'a>,
>,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = crate::AutumnResult<()>> + Send + 'a>>
{
next
}
}
#[cfg(feature = "db")]
struct ChaosDbInterceptor {
state: Arc<ChaosState>,
}
#[cfg(feature = "db")]
impl crate::interceptor::DbConnectionInterceptor for ChaosDbInterceptor {
fn intercept_checkout<'a>(
&'a self,
_ctx: crate::interceptor::DbCheckoutContext,
next: std::pin::Pin<
Box<
dyn std::future::Future<
Output = Result<crate::db::PooledConnection, crate::AutumnError>,
> + Send
+ 'a,
>,
>,
) -> std::pin::Pin<
Box<
dyn std::future::Future<
Output = Result<crate::db::PooledConnection, crate::AutumnError>,
> + Send
+ 'a,
>,
> {
Box::pin(async move {
let seq = self.state.checkout_seq.fetch_add(1, Ordering::SeqCst);
let fired = self.state.decide(self.state.db_transient_error_prob);
self.state.record(ChaosHook::DbCheckout, seq, fired);
if fired {
Err(crate::AutumnError::service_unavailable_msg(
"chaos: injected transient database checkout error",
))
} else {
next.await
}
})
}
}
#[cfg(feature = "mail")]
struct ChaosMailInterceptor {
state: Arc<ChaosState>,
}
#[cfg(feature = "mail")]
impl crate::interceptor::MailInterceptor for ChaosMailInterceptor {
fn intercept<'a>(
&'a self,
_mail: &'a crate::mail::Mail,
next: std::pin::Pin<
Box<dyn std::future::Future<Output = Result<(), crate::mail::MailError>> + Send + 'a>,
>,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = Result<(), crate::mail::MailError>> + Send + 'a>,
> {
Box::pin(async move {
let send_index = self.state.mail_seq.fetch_add(1, Ordering::SeqCst) + 1;
let fault = self
.state
.resolve_mail_fault(usize::try_from(send_index).unwrap_or(usize::MAX));
self.state
.record(ChaosHook::MailSend, send_index - 1, fault.is_some());
match fault {
Some(MailFault::Fail) => Err(crate::mail::MailError::RuntimeUnavailable(
"chaos: injected SMTP send failure".to_owned(),
)),
Some(MailFault::Timeout) => Err(crate::mail::MailError::Io(std::io::Error::new(
std::io::ErrorKind::TimedOut,
"chaos: injected SMTP send timeout",
))),
None => next.await,
}
})
}
}
#[cfg(test)]
mod tests {
use super::{Chaos, ChaosHook, ChaosState, clamp_prob, deterministic_skew};
#[test]
fn default_chaos_is_inactive() {
assert!(!Chaos::default().is_active());
}
#[test]
fn any_configured_fault_activates() {
assert!(Chaos::default().db_transient_errors(0.01).is_active());
assert!(Chaos::default().job_duplicate_delivery(0.01).is_active());
assert!(
Chaos::default()
.clock_skew(std::time::Duration::from_secs(1))
.is_active()
);
}
#[test]
fn probabilities_are_clamped() {
assert!((clamp_prob(-1.0) - 0.0).abs() < f64::EPSILON);
assert!((clamp_prob(2.0) - 1.0).abs() < f64::EPSILON);
assert!((clamp_prob(f64::NAN) - 0.0).abs() < f64::EPSILON);
assert!(!Chaos::default().db_transient_errors(0.0).is_active());
}
#[test]
fn decision_stream_is_seed_deterministic() {
let chaos = Chaos::default().db_transient_errors(0.5);
let draw = |seed: u64| {
let state = ChaosState::new(seed, &chaos);
(0..64).map(|_| state.decide(0.5)).collect::<Vec<bool>>()
};
assert_eq!(
draw(42),
draw(42),
"same seed must replay the same schedule"
);
assert_ne!(
draw(42),
draw(43),
"different seeds should (overwhelmingly likely) diverge"
);
}
#[test]
fn decide_extremes_are_exact() {
let state = ChaosState::new(7, &Chaos::default());
assert!((0..32).all(|_| state.decide(1.0)), "p=1.0 always fires");
assert!((0..32).all(|_| !state.decide(0.0)), "p=0.0 never fires");
}
#[test]
fn recorded_events_snapshot_in_order() {
let state = ChaosState::new(1, &Chaos::default());
state.record(ChaosHook::DbCheckout, 0, true);
state.record(ChaosHook::JobDelivery, 0, false);
let events = state.events();
assert_eq!(events.len(), 2);
assert_eq!(events[0].hook, ChaosHook::DbCheckout);
assert!(events[0].fired);
assert_eq!(events[1].hook, ChaosHook::JobDelivery);
assert!(!events[1].fired);
}
#[test]
fn clock_skew_offset_is_deterministic_and_bounded() {
let dur = std::time::Duration::from_secs(5);
let a = deterministic_skew(99, dur);
let b = deterministic_skew(99, dur);
assert_eq!(a, b, "same seed ⇒ same skew offset");
assert!(a >= chrono::Duration::zero());
assert!(
a <= chrono::Duration::from_std(dur).unwrap(),
"offset stays within [0, dur]"
);
assert_eq!(
deterministic_skew(99, std::time::Duration::ZERO),
chrono::Duration::zero()
);
}
#[cfg(feature = "mail")]
#[test]
fn mail_faults_activate_chaos() {
use super::MailFault;
assert!(
Chaos::default()
.smtp_faults([(7, MailFault::Fail), (8, MailFault::Timeout)])
.is_active(),
"an explicit schedule activates chaos"
);
assert!(
Chaos::default().smtp_transient_errors(0.01).is_active(),
"a probabilistic rate activates chaos"
);
assert!(!Chaos::default().smtp_faults([]).is_active());
assert!(!Chaos::default().smtp_transient_errors(0.0).is_active());
}
#[cfg(feature = "mail")]
#[test]
fn explicit_schedule_is_1_based_and_seed_independent() {
use super::MailFault;
let chaos = Chaos::default().smtp_faults([(7, MailFault::Fail), (8, MailFault::Timeout)]);
for seed in [1u64, 42, 0x5EED] {
let state = ChaosState::new(seed, &chaos);
assert_eq!(state.resolve_mail_fault(6), None, "send #6 is unscheduled");
assert_eq!(state.resolve_mail_fault(7), Some(MailFault::Fail));
assert_eq!(state.resolve_mail_fault(8), Some(MailFault::Timeout));
assert_eq!(state.resolve_mail_fault(9), None, "send #9 is unscheduled");
}
let zeroed = Chaos::default().smtp_faults([(0, MailFault::Fail)]);
assert!(!zeroed.is_active());
}
#[cfg(feature = "mail")]
#[test]
fn schedule_wins_over_probabilistic_and_prob_is_seed_deterministic() {
use super::MailFault;
let chaos = Chaos::default()
.smtp_faults([(8, MailFault::Timeout)])
.smtp_transient_errors(1.0);
let state = ChaosState::new(3, &chaos);
assert_eq!(
state.resolve_mail_fault(1),
Some(MailFault::Fail),
"p=1.0 faults an unscheduled send transiently"
);
let state2 = ChaosState::new(3, &chaos);
assert_eq!(
state2.resolve_mail_fault(8),
Some(MailFault::Timeout),
"an explicit entry wins over the probabilistic rate"
);
let prob_only = Chaos::default().smtp_transient_errors(0.5);
let draw = |seed: u64| {
let state = ChaosState::new(seed, &prob_only);
(1..=32)
.map(|i| state.resolve_mail_fault(i).is_some())
.collect::<Vec<bool>>()
};
assert_eq!(
draw(77),
draw(77),
"same seed replays the same mail schedule"
);
assert_ne!(draw(77), draw(78), "different seeds diverge");
}
}