1use std::collections::{HashMap, VecDeque};
5use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
6use std::time::{Duration, SystemTime};
7
8use crate::clock::Clock;
9use crate::effect::EffectFailure;
10use crate::failure::FailureClass;
11use crate::id::IdempotencyKey;
12
13#[derive(Clone, Copy, Debug, PartialEq, Eq)]
15pub enum Behavior {
16 Succeed,
18 Fail(FailureClass),
20 Unreachable,
22 CommitThenDrop,
25 LoseRequest,
28 Hang,
30}
31
32#[derive(Clone)]
47pub struct FakeRemote {
48 clock: Arc<dyn Clock>,
49 state: Arc<Mutex<State>>,
50}
51
52#[derive(Default)]
53struct State {
54 script: VecDeque<Behavior>,
55 lag: Duration,
56 requests: u32,
57 created: HashMap<String, (String, SystemTime)>,
58 applications: HashMap<String, u32>,
59 by_idempotency_key: HashMap<IdempotencyKey, String>,
60 cancellations: HashMap<String, u32>,
61 cancelled_keys: HashMap<IdempotencyKey, String>,
62}
63
64impl FakeRemote {
65 pub fn new(clock: impl Clock) -> Self {
67 Self {
68 clock: Arc::new(clock),
69 state: Arc::default(),
70 }
71 }
72
73 #[must_use]
75 pub fn script(self, behaviors: impl IntoIterator<Item = Behavior>) -> Self {
76 self.lock().script.extend(behaviors);
77 self
78 }
79
80 #[must_use]
82 pub fn lag(self, lag: Duration) -> Self {
83 self.lock().lag = lag;
84 self
85 }
86
87 pub async fn create(
94 &self,
95 resource: &str,
96 idempotency_key: Option<IdempotencyKey>,
97 ) -> Result<String, EffectFailure> {
98 let (behavior, replay) = {
99 let mut state = self.lock();
100 state.requests += 1;
101 let replay = idempotency_key.and_then(|k| state.by_idempotency_key.get(&k).cloned());
102 (
103 state.script.pop_front().unwrap_or(Behavior::Succeed),
104 replay,
105 )
106 };
107 match behavior {
108 Behavior::Succeed => Ok(self.apply(resource, idempotency_key)),
109 Behavior::Fail(class) => match replay {
112 Some(id) => Ok(id),
113 None => Err(EffectFailure::new(class, "remote refused the request")),
114 },
115 Behavior::Unreachable => {
116 Err(EffectFailure::ambiguous("connection refused").request_sent(false))
117 }
118 Behavior::CommitThenDrop => {
119 self.apply(resource, idempotency_key);
120 Err(EffectFailure::ambiguous(
121 "connection reset before the response",
122 ))
123 }
124 Behavior::LoseRequest => {
125 Err(EffectFailure::ambiguous("timed out waiting for a response"))
126 }
127 Behavior::Hang => std::future::pending().await,
128 }
129 }
130
131 #[allow(unknown_lints, clippy::unused_async, clippy::unused_async_trait_impl)]
138 pub async fn find(&self, resource: &str) -> Result<Option<String>, EffectFailure> {
139 let now = self.clock.now();
140 let state = self.lock();
141 Ok(state
142 .created
143 .get(resource)
144 .filter(|(_, at)| *at + state.lag <= now)
145 .map(|(id, _)| id.clone()))
146 }
147
148 pub fn applications(&self, resource: &str) -> u32 {
150 self.lock().applications.get(resource).copied().unwrap_or(0)
151 }
152
153 pub async fn cancel(
163 &self,
164 resource: &str,
165 idempotency_key: Option<IdempotencyKey>,
166 ) -> Result<(), EffectFailure> {
167 let behavior = {
168 let mut state = self.lock();
169 state.requests += 1;
170 state.script.pop_front().unwrap_or(Behavior::Succeed)
171 };
172 match behavior {
173 Behavior::Succeed => {
174 self.apply_cancel(resource, idempotency_key);
175 Ok(())
176 }
177 Behavior::Fail(class) => Err(EffectFailure::new(class, "remote refused the cancel")),
178 Behavior::Unreachable => {
179 Err(EffectFailure::ambiguous("connection refused").request_sent(false))
180 }
181 Behavior::CommitThenDrop => {
182 self.apply_cancel(resource, idempotency_key);
183 Err(EffectFailure::ambiguous(
184 "connection reset before the response",
185 ))
186 }
187 Behavior::LoseRequest => {
188 Err(EffectFailure::ambiguous("timed out waiting for a response"))
189 }
190 Behavior::Hang => std::future::pending().await,
191 }
192 }
193
194 pub fn cancellations(&self, resource: &str) -> u32 {
197 self.lock()
198 .cancellations
199 .get(resource)
200 .copied()
201 .unwrap_or(0)
202 }
203
204 pub fn exists(&self, resource: &str) -> bool {
206 self.lock().created.contains_key(resource)
207 }
208
209 fn apply_cancel(&self, resource: &str, idempotency_key: Option<IdempotencyKey>) {
210 let mut state = self.lock();
211 if idempotency_key.is_some_and(|k| state.cancelled_keys.contains_key(&k)) {
212 return;
213 }
214 if state.created.remove(resource).is_some() {
215 *state.cancellations.entry(resource.to_owned()).or_default() += 1;
216 }
217 if let Some(key) = idempotency_key {
218 state.cancelled_keys.insert(key, resource.to_owned());
219 }
220 }
221
222 pub fn requests(&self) -> u32 {
224 self.lock().requests
225 }
226
227 fn apply(&self, resource: &str, idempotency_key: Option<IdempotencyKey>) -> String {
228 let now = self.clock.now();
229 let mut state = self.lock();
230 if let Some(id) = idempotency_key.and_then(|k| state.by_idempotency_key.get(&k)) {
231 return id.clone();
232 }
233 let count = state.applications.entry(resource.to_owned()).or_default();
234 *count += 1;
235 let id = format!("{resource}#{count}");
236 state.created.insert(resource.to_owned(), (id.clone(), now));
237 if let Some(key) = idempotency_key {
238 state.by_idempotency_key.insert(key, id.clone());
239 }
240 id
241 }
242
243 fn lock(&self) -> MutexGuard<'_, State> {
244 self.state.lock().unwrap_or_else(PoisonError::into_inner)
245 }
246}