Skip to main content

agent_effects/
testkit.rs

1//! Test doubles for code built on `agent-effects`. Enabled by the `testkit`
2//! feature.
3
4use 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/// What the [`FakeRemote`] does with one request.
14#[derive(Clone, Copy, Debug, PartialEq, Eq)]
15pub enum Behavior {
16    /// Apply the request and answer.
17    Succeed,
18    /// Refuse the request without applying it, with this failure class.
19    Fail(FailureClass),
20    /// Refuse the connection: nothing was sent, nothing applied.
21    Unreachable,
22    /// Apply the request, then drop the connection before answering. The
23    /// caller sees an ambiguous failure for an effect that happened.
24    CommitThenDrop,
25    /// Lose the request: nothing applied, and the caller sees an ambiguous
26    /// failure. Indistinguishable from `CommitThenDrop` for the caller.
27    LoseRequest,
28    /// Never answer.
29    Hang,
30}
31
32/// A scripted remote system for exercising failure handling.
33///
34/// It creates resources by name and counts how often each was really
35/// created: [`FakeRemote::applications`] above 1 is a duplicated side effect.
36/// Each request consumes the next scripted [`Behavior`]; once the script is
37/// empty, requests succeed.
38///
39/// A request with an idempotency key the remote has already applied returns
40/// the earlier result without applying again, like a provider that honours
41/// `Idempotency-Key`. That includes a request scripted to
42/// [`Behavior::Fail`]: the remote answers a replay before it would evaluate
43/// the request. Network-level behaviors (`Unreachable`, `LoseRequest`,
44/// `CommitThenDrop`'s dropped answer, `Hang`) still happen. [`FakeRemote::find`] sees a resource only `lag` after
45/// it was created, like an eventually consistent search API.
46#[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    /// A remote that reads time from `clock` (for the lookup lag).
66    pub fn new(clock: impl Clock) -> Self {
67        Self {
68            clock: Arc::new(clock),
69            state: Arc::default(),
70        }
71    }
72
73    /// Queues behaviors for the next requests, in order.
74    #[must_use]
75    pub fn script(self, behaviors: impl IntoIterator<Item = Behavior>) -> Self {
76        self.lock().script.extend(behaviors);
77        self
78    }
79
80    /// Makes [`Self::find`] lag `lag` behind creation.
81    #[must_use]
82    pub fn lag(self, lag: Duration) -> Self {
83        self.lock().lag = lag;
84        self
85    }
86
87    /// Creates `resource`, deduplicating on `idempotency_key` if given.
88    /// Returns the resource's id.
89    ///
90    /// # Errors
91    ///
92    /// Whatever the script says.
93    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            // A remote that deduplicates answers a replayed key with the
110            // original result before it would evaluate the request again.
111            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    /// Looks `resource` up, seeing it only once the lookup lag has passed.
132    ///
133    /// # Errors
134    ///
135    /// Never; the signature matches a real lookup.
136    // `async` to match a real lookup's signature.
137    #[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    /// How often `resource` was really created. More than 1 is a duplicate.
149    pub fn applications(&self, resource: &str) -> u32 {
150        self.lock().applications.get(resource).copied().unwrap_or(0)
151    }
152
153    /// Undoes `resource`, deduplicating on `idempotency_key` if given. Each
154    /// request consumes the next scripted [`Behavior`], like `create`:
155    /// `CommitThenDrop` cancels and then reports an ambiguous failure.
156    /// Cancelling a resource that does not exist succeeds and changes
157    /// nothing, as idempotent deletes do.
158    ///
159    /// # Errors
160    ///
161    /// Whatever the script says.
162    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    /// How often `resource` was really cancelled. More than 1 means an undo
195    /// ran twice.
196    pub fn cancellations(&self, resource: &str) -> u32 {
197        self.lock()
198            .cancellations
199            .get(resource)
200            .copied()
201            .unwrap_or(0)
202    }
203
204    /// Whether `resource` exists: created and not cancelled.
205    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    /// How many create requests arrived.
223    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}