use std::collections::{HashMap, VecDeque};
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use std::time::{Duration, SystemTime};
use crate::clock::Clock;
use crate::effect::EffectFailure;
use crate::failure::FailureClass;
use crate::id::IdempotencyKey;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Behavior {
Succeed,
Fail(FailureClass),
Unreachable,
CommitThenDrop,
LoseRequest,
Hang,
}
#[derive(Clone)]
pub struct FakeRemote {
clock: Arc<dyn Clock>,
state: Arc<Mutex<State>>,
}
#[derive(Default)]
struct State {
script: VecDeque<Behavior>,
lag: Duration,
requests: u32,
created: HashMap<String, (String, SystemTime)>,
applications: HashMap<String, u32>,
by_idempotency_key: HashMap<IdempotencyKey, String>,
cancellations: HashMap<String, u32>,
cancelled_keys: HashMap<IdempotencyKey, String>,
}
impl FakeRemote {
pub fn new(clock: impl Clock) -> Self {
Self {
clock: Arc::new(clock),
state: Arc::default(),
}
}
#[must_use]
pub fn script(self, behaviors: impl IntoIterator<Item = Behavior>) -> Self {
self.lock().script.extend(behaviors);
self
}
#[must_use]
pub fn lag(self, lag: Duration) -> Self {
self.lock().lag = lag;
self
}
pub async fn create(
&self,
resource: &str,
idempotency_key: Option<IdempotencyKey>,
) -> Result<String, EffectFailure> {
let (behavior, replay) = {
let mut state = self.lock();
state.requests += 1;
let replay = idempotency_key.and_then(|k| state.by_idempotency_key.get(&k).cloned());
(
state.script.pop_front().unwrap_or(Behavior::Succeed),
replay,
)
};
match behavior {
Behavior::Succeed => Ok(self.apply(resource, idempotency_key)),
Behavior::Fail(class) => match replay {
Some(id) => Ok(id),
None => Err(EffectFailure::new(class, "remote refused the request")),
},
Behavior::Unreachable => {
Err(EffectFailure::ambiguous("connection refused").request_sent(false))
}
Behavior::CommitThenDrop => {
self.apply(resource, idempotency_key);
Err(EffectFailure::ambiguous(
"connection reset before the response",
))
}
Behavior::LoseRequest => {
Err(EffectFailure::ambiguous("timed out waiting for a response"))
}
Behavior::Hang => std::future::pending().await,
}
}
#[allow(unknown_lints, clippy::unused_async, clippy::unused_async_trait_impl)]
pub async fn find(&self, resource: &str) -> Result<Option<String>, EffectFailure> {
let now = self.clock.now();
let state = self.lock();
Ok(state
.created
.get(resource)
.filter(|(_, at)| *at + state.lag <= now)
.map(|(id, _)| id.clone()))
}
pub fn applications(&self, resource: &str) -> u32 {
self.lock().applications.get(resource).copied().unwrap_or(0)
}
pub async fn cancel(
&self,
resource: &str,
idempotency_key: Option<IdempotencyKey>,
) -> Result<(), EffectFailure> {
let behavior = {
let mut state = self.lock();
state.requests += 1;
state.script.pop_front().unwrap_or(Behavior::Succeed)
};
match behavior {
Behavior::Succeed => {
self.apply_cancel(resource, idempotency_key);
Ok(())
}
Behavior::Fail(class) => Err(EffectFailure::new(class, "remote refused the cancel")),
Behavior::Unreachable => {
Err(EffectFailure::ambiguous("connection refused").request_sent(false))
}
Behavior::CommitThenDrop => {
self.apply_cancel(resource, idempotency_key);
Err(EffectFailure::ambiguous(
"connection reset before the response",
))
}
Behavior::LoseRequest => {
Err(EffectFailure::ambiguous("timed out waiting for a response"))
}
Behavior::Hang => std::future::pending().await,
}
}
pub fn cancellations(&self, resource: &str) -> u32 {
self.lock()
.cancellations
.get(resource)
.copied()
.unwrap_or(0)
}
pub fn exists(&self, resource: &str) -> bool {
self.lock().created.contains_key(resource)
}
fn apply_cancel(&self, resource: &str, idempotency_key: Option<IdempotencyKey>) {
let mut state = self.lock();
if idempotency_key.is_some_and(|k| state.cancelled_keys.contains_key(&k)) {
return;
}
if state.created.remove(resource).is_some() {
*state.cancellations.entry(resource.to_owned()).or_default() += 1;
}
if let Some(key) = idempotency_key {
state.cancelled_keys.insert(key, resource.to_owned());
}
}
pub fn requests(&self) -> u32 {
self.lock().requests
}
fn apply(&self, resource: &str, idempotency_key: Option<IdempotencyKey>) -> String {
let now = self.clock.now();
let mut state = self.lock();
if let Some(id) = idempotency_key.and_then(|k| state.by_idempotency_key.get(&k)) {
return id.clone();
}
let count = state.applications.entry(resource.to_owned()).or_default();
*count += 1;
let id = format!("{resource}#{count}");
state.created.insert(resource.to_owned(), (id.clone(), now));
if let Some(key) = idempotency_key {
state.by_idempotency_key.insert(key, id.clone());
}
id
}
fn lock(&self) -> MutexGuard<'_, State> {
self.state.lock().unwrap_or_else(PoisonError::into_inner)
}
}