Skip to main content

agent_effects_store/
testkit.rs

1//! Conformance suite for [`EffectStore`] implementations.
2//!
3//! Enabled by the `testkit` feature. Add it as a dev-dependency feature:
4//!
5//! ```toml
6//! [dev-dependencies]
7//! agent-effects-store = { version = "0.1", features = ["testkit"] }
8//! ```
9
10use std::future::Future;
11use std::sync::Arc;
12use std::time::{Duration, SystemTime};
13
14use serde_json::json;
15
16use crate::failure::FailureClass;
17use crate::id::{EffectId, EffectKey, EffectName, LogicalKey, WorkerId};
18use crate::kind::EffectKind;
19use crate::state::{ALL_STATUSES, ALL_TRANSITIONS, EffectStatus, Transition};
20use crate::{
21    EffectEvent, EffectRecord, EffectStore, ErrorRecord, Lease, ListQuery, NewEffect, PruneQuery,
22    StoreError, TransitionRequest,
23};
24
25/// Runs the store conformance suite against stores built by `make_store`.
26///
27/// Every case gets a fresh store. A failing case panics. Some cases spawn
28/// tasks, so call this from inside a Tokio runtime, ideally a multi-threaded
29/// one so races are real.
30///
31/// ```ignore
32/// #[tokio::test(flavor = "multi_thread")]
33/// async fn conformance() {
34///     agent_effects_store::testkit::conformance(|| async { MyStore::open_temp().await }).await;
35/// }
36/// ```
37pub async fn conformance<S, F, Fut>(make_store: F)
38where
39    S: EffectStore,
40    F: Fn() -> Fut,
41    Fut: Future<Output = S>,
42{
43    insert_and_read_back(make_store().await).await;
44    insert_is_idempotent_per_key(make_store().await).await;
45    concurrent_inserts_converge(make_store().await).await;
46    missing_records_are_reported(make_store().await).await;
47    leases_are_exclusive(make_store().await).await;
48    takeover_fences_the_old_owner(make_store().await).await;
49    renewal_extends_and_is_strict(make_store().await).await;
50    release_frees_the_lease(make_store().await).await;
51    transitions_persist_with_events(make_store().await).await;
52    rejected_transitions_change_nothing(make_store().await).await;
53    unleased_transitions_for_operators(make_store().await).await;
54    terminal_records_are_final(make_store().await).await;
55    listing_filters_and_pages(make_store().await).await;
56    doubt_is_persisted_and_guards_failure(make_store().await).await;
57    compensation_is_persisted(make_store().await).await;
58    approval_is_persisted(make_store().await).await;
59    pruning_removes_only_settled_idle_old_records(make_store().await).await;
60}
61
62const TTL: Duration = Duration::from_secs(30);
63
64/// Test time `secs` seconds after a fixed base, at millisecond precision.
65fn t(secs: u64) -> SystemTime {
66    SystemTime::UNIX_EPOCH + Duration::from_secs(1_700_000_000 + secs)
67}
68
69fn key(n: u32) -> EffectKey {
70    EffectKey::new(
71        EffectName::new("testkit.effect").unwrap(),
72        LogicalKey::new(format!("key-{n}")).unwrap(),
73    )
74}
75
76fn new_effect(n: u32) -> NewEffect {
77    NewEffect {
78        input: Some(json!({ "n": n, "nested": { "b": [1, 2], "a": null } })),
79        input_fingerprint: Some(format!("fingerprint-{n}")),
80        created_by: Some("agent:testkit".into()),
81        ..NewEffect::new(key(n), EffectKind::IrreversibleWrite, t(0))
82    }
83}
84
85fn worker(name: &str) -> WorkerId {
86    WorkerId::new(name)
87}
88
89async fn insert<S: EffectStore>(store: &S, n: u32) -> EffectRecord {
90    store.insert_or_get(new_effect(n)).await.unwrap().record
91}
92
93async fn reload<S: EffectStore>(store: &S, id: EffectId) -> EffectRecord {
94    store.get(id).await.unwrap().expect("record exists")
95}
96
97/// Applies `transitions` in order under `lease` at time `now`.
98async fn drive<S: EffectStore>(
99    store: &S,
100    mut record: EffectRecord,
101    lease: Option<&Lease>,
102    transitions: &[Transition],
103    now: SystemTime,
104) -> EffectRecord {
105    for &transition in transitions {
106        record = store
107            .transition(TransitionRequest::new(&record, lease, transition, now))
108            .await
109            .unwrap_or_else(|e| panic!("{transition} from {}: {e}", record.status));
110    }
111    record
112}
113
114async fn insert_and_read_back<S: EffectStore>(store: S) {
115    let new = new_effect(1);
116    let outcome = store.insert_or_get(new.clone()).await.unwrap();
117    assert!(outcome.inserted);
118    let record = outcome.record;
119    assert_eq!(record.id, new.id);
120    assert_eq!(record.key, new.key);
121    assert_eq!(record.kind, new.kind);
122    assert_eq!(record.status, EffectStatus::Pending);
123    assert_eq!(record.input, new.input);
124    assert_eq!(record.input_fingerprint, new.input_fingerprint);
125    assert_eq!(record.created_by, new.created_by);
126    assert_eq!(
127        (record.version, record.lease_epoch, record.attempt_count),
128        (0, 0, 0)
129    );
130    assert_eq!((record.created_at, record.updated_at), (t(0), t(0)));
131    assert_eq!(record.lease_owner, None);
132
133    assert_eq!(store.get(record.id).await.unwrap(), Some(record.clone()));
134    assert_eq!(
135        store.get_by_key(&record.key).await.unwrap(),
136        Some(record.clone())
137    );
138    assert_eq!(
139        store.events(record.id).await.unwrap(),
140        Vec::<EffectEvent>::new()
141    );
142}
143
144async fn insert_is_idempotent_per_key<S: EffectStore>(store: S) {
145    let first = insert(&store, 1).await;
146    let again = NewEffect {
147        input: Some(json!("different")),
148        ..NewEffect::new(key(1), EffectKind::Read, t(5))
149    };
150    let outcome = store.insert_or_get(again).await.unwrap();
151    assert!(
152        !outcome.inserted,
153        "a second insert for the same key must not insert"
154    );
155    assert_eq!(
156        outcome.record, first,
157        "the existing record must come back untouched"
158    );
159
160    let other = insert(&store, 2).await;
161    assert_ne!(other.id, first.id);
162}
163
164async fn concurrent_inserts_converge<S: EffectStore>(store: S) {
165    let store = Arc::new(store);
166    let tasks: Vec<_> = (0..16)
167        .map(|_| {
168            let store = Arc::clone(&store);
169            tokio::spawn(async move { store.insert_or_get(new_effect(7)).await.unwrap() })
170        })
171        .collect();
172    let mut outcomes = Vec::new();
173    for task in tasks {
174        outcomes.push(task.await.unwrap());
175    }
176    let inserted = outcomes.iter().filter(|o| o.inserted).count();
177    assert_eq!(inserted, 1, "exactly one concurrent insert may win");
178    let id = outcomes[0].record.id;
179    assert!(
180        outcomes.iter().all(|o| o.record.id == id),
181        "all callers must see one record"
182    );
183}
184
185async fn missing_records_are_reported<S: EffectStore>(store: S) {
186    let id = EffectId::new();
187    assert_eq!(store.get(id).await.unwrap(), None);
188    assert_eq!(store.get_by_key(&key(404)).await.unwrap(), None);
189    assert!(matches!(
190        store.acquire_lease(id, &worker("a"), t(0), TTL).await,
191        Err(StoreError::NotFound(missing)) if missing == id
192    ));
193    let mut request = TransitionRequest::new(
194        &insert(&store, 1).await,
195        None,
196        Transition::StartAttempt,
197        t(0),
198    );
199    request.id = id;
200    assert!(matches!(
201        store.transition(request).await,
202        Err(StoreError::NotFound(_))
203    ));
204    assert!(matches!(
205        store.events(id).await,
206        Err(StoreError::NotFound(_))
207    ));
208}
209
210async fn leases_are_exclusive<S: EffectStore>(store: S) {
211    let record = insert(&store, 1).await;
212    let lease = store
213        .acquire_lease(record.id, &worker("a"), t(0), TTL)
214        .await
215        .unwrap();
216    assert_eq!((lease.epoch, lease.expires_at), (1, t(30)));
217    assert_eq!(lease.owner, worker("a"));
218
219    for contender in ["b", "a"] {
220        match store
221            .acquire_lease(record.id, &worker(contender), t(1), TTL)
222            .await
223        {
224            Err(StoreError::LeaseHeld { owner, expires_at }) => {
225                assert_eq!(owner, worker("a"));
226                assert_eq!(expires_at, t(30));
227            }
228            other => panic!("{contender} acquired a held lease: {other:?}"),
229        }
230    }
231}
232
233async fn takeover_fences_the_old_owner<S: EffectStore>(store: S) {
234    let record = insert(&store, 1).await;
235    let old = store
236        .acquire_lease(record.id, &worker("a"), t(0), TTL)
237        .await
238        .unwrap();
239    assert!(matches!(
240        store
241            .acquire_lease(record.id, &worker("b"), t(29), TTL)
242            .await,
243        Err(StoreError::LeaseHeld { .. })
244    ));
245    let new = store
246        .acquire_lease(record.id, &worker("b"), t(30), TTL)
247        .await
248        .unwrap();
249    assert_eq!(new.epoch, 2);
250
251    let stale = TransitionRequest::new(&record, Some(&old), Transition::StartAttempt, t(31));
252    assert!(matches!(
253        store.transition(stale).await,
254        Err(StoreError::LeaseLost)
255    ));
256    assert!(matches!(
257        store.renew_lease(&old, t(31), TTL).await,
258        Err(StoreError::LeaseLost)
259    ));
260    store.release_lease(&old).await.unwrap();
261
262    let current = reload(&store, record.id).await;
263    assert_eq!(
264        current.lease_owner,
265        Some(worker("b")),
266        "a stale release must not free b's lease"
267    );
268    assert_eq!((current.lease_epoch, current.version), (2, 0));
269}
270
271async fn renewal_extends_and_is_strict<S: EffectStore>(store: S) {
272    let record = insert(&store, 1).await;
273    let lease = store
274        .acquire_lease(record.id, &worker("a"), t(0), TTL)
275        .await
276        .unwrap();
277    let renewed = store.renew_lease(&lease, t(20), TTL).await.unwrap();
278    assert_eq!((renewed.epoch, renewed.expires_at), (lease.epoch, t(50)));
279    assert!(matches!(
280        store
281            .acquire_lease(record.id, &worker("b"), t(40), TTL)
282            .await,
283        Err(StoreError::LeaseHeld { .. })
284    ));
285    assert!(
286        matches!(
287            store.renew_lease(&renewed, t(50), TTL).await,
288            Err(StoreError::LeaseLost)
289        ),
290        "an expired lease must not be revived"
291    );
292    assert_eq!(
293        reload(&store, record.id).await.version,
294        0,
295        "lease operations must not bump the version"
296    );
297}
298
299async fn release_frees_the_lease<S: EffectStore>(store: S) {
300    let record = insert(&store, 1).await;
301    let lease = store
302        .acquire_lease(record.id, &worker("a"), t(0), TTL)
303        .await
304        .unwrap();
305    store.release_lease(&lease).await.unwrap();
306    store.release_lease(&lease).await.unwrap();
307    let next = store
308        .acquire_lease(record.id, &worker("b"), t(1), TTL)
309        .await
310        .unwrap();
311    assert_eq!(next.epoch, 2);
312}
313
314async fn transitions_persist_with_events<St: EffectStore>(store: St) {
315    use EffectStatus as S;
316    use Transition as T;
317
318    let record = insert(&store, 1).await;
319    let lease = store
320        .acquire_lease(record.id, &worker("a"), t(0), TTL)
321        .await
322        .unwrap();
323    let record = drive(
324        &store,
325        record,
326        Some(&lease),
327        &[Transition::StartAttempt],
328        t(1),
329    )
330    .await;
331    assert_eq!(record.status, EffectStatus::Executing);
332    assert_eq!(record.attempt_started_at, Some(t(1)));
333
334    let error = ErrorRecord {
335        class: Some(FailureClass::RateLimited {
336            retry_after: Some(Duration::from_secs(9)),
337        }),
338        message: "429".into(),
339    };
340    let mut retry = TransitionRequest::new(&record, Some(&lease), Transition::ScheduleRetry, t(2));
341    retry.next_attempt_at = Some(t(11));
342    retry.error = Some(error.clone());
343    retry.actor = Some("worker:a".into());
344    retry.payload = Some(json!({ "delay_ms": 9000 }));
345    let record = store.transition(retry).await.unwrap();
346    assert_eq!(
347        record,
348        reload(&store, record.id).await,
349        "transition must return the stored record"
350    );
351    assert_eq!(record.next_attempt_at, Some(t(11)));
352    assert_eq!(record.attempt_ended_at, Some(t(2)));
353    assert_eq!(record.last_error, Some(error));
354
355    let record = drive(
356        &store,
357        record,
358        Some(&lease),
359        &[Transition::StartAttempt],
360        t(11),
361    )
362    .await;
363    let output = json!({ "payment_id": "pi_123", "amount": 4200 });
364    let mut verify =
365        TransitionRequest::new(&record, Some(&lease), Transition::StartVerification, t(12));
366    verify.output = Some(output.clone());
367    let record = store.transition(verify).await.unwrap();
368    let record = drive(
369        &store,
370        record,
371        Some(&lease),
372        &[Transition::VerificationConfirmed],
373        t(13),
374    )
375    .await;
376
377    let stored = reload(&store, record.id).await;
378    assert_eq!(stored.status, EffectStatus::Committed);
379    assert_eq!(stored.output, Some(output));
380    assert_eq!(stored.committed_at, Some(t(13)));
381    assert_eq!((stored.attempt_count, stored.version), (2, 5));
382    assert_eq!(stored.next_attempt_at, None);
383    assert_eq!(stored.attempt_ended_at, Some(t(12)));
384
385    let events = store.events(record.id).await.unwrap();
386    let trail: Vec<_> = events
387        .iter()
388        .map(|e| (e.sequence, e.transition, e.from, e.to, e.attempt))
389        .collect();
390    assert_eq!(
391        trail,
392        [
393            (1, T::StartAttempt, S::Pending, S::Executing, 1),
394            (2, T::ScheduleRetry, S::Executing, S::Pending, 1),
395            (3, T::StartAttempt, S::Pending, S::Executing, 2),
396            (4, T::StartVerification, S::Executing, S::Verifying, 2),
397            (5, T::VerificationConfirmed, S::Verifying, S::Committed, 2),
398        ]
399    );
400    assert_eq!(events[1].actor.as_deref(), Some("worker:a"));
401    assert_eq!(events[1].payload, Some(json!({ "delay_ms": 9000 })));
402    assert_eq!(events[1].at, t(2));
403}
404
405async fn rejected_transitions_change_nothing<S: EffectStore>(store: S) {
406    let record = insert(&store, 1).await;
407    let lease = store
408        .acquire_lease(record.id, &worker("a"), t(0), TTL)
409        .await
410        .unwrap();
411    let record = reload(&store, record.id).await;
412
413    let illegal = TransitionRequest::new(&record, Some(&lease), Transition::Succeeded, t(1));
414    let stale_version = TransitionRequest {
415        expected_version: 7,
416        ..TransitionRequest::new(&record, Some(&lease), Transition::StartAttempt, t(1))
417    };
418    let unleased = TransitionRequest::new(&record, None, Transition::StartAttempt, t(1));
419    let expired = TransitionRequest::new(&record, Some(&lease), Transition::StartAttempt, t(30));
420
421    assert!(matches!(
422        store.transition(illegal).await,
423        Err(StoreError::InvalidTransition(_))
424    ));
425    assert!(matches!(
426        store.transition(stale_version).await,
427        Err(StoreError::VersionConflict {
428            expected: 7,
429            actual: 0
430        })
431    ));
432    assert!(matches!(
433        store.transition(unleased).await,
434        Err(StoreError::LeaseHeld { .. })
435    ));
436    assert!(matches!(
437        store.transition(expired).await,
438        Err(StoreError::LeaseLost)
439    ));
440
441    assert_eq!(reload(&store, record.id).await, record);
442    assert_eq!(
443        store.events(record.id).await.unwrap(),
444        Vec::<EffectEvent>::new()
445    );
446}
447
448async fn unleased_transitions_for_operators<S: EffectStore>(store: S) {
449    let record = insert(&store, 1).await;
450    let lease = store
451        .acquire_lease(record.id, &worker("a"), t(0), TTL)
452        .await
453        .unwrap();
454    let record = drive(
455        &store,
456        record,
457        Some(&lease),
458        &[Transition::StartAttempt, Transition::OutcomeUnknown],
459        t(1),
460    )
461    .await;
462    store.release_lease(&lease).await.unwrap();
463
464    let record = drive(&store, record, None, &[Transition::Escalate], t(2)).await;
465    assert_eq!(record.status, EffectStatus::NeedsIntervention);
466    let mut resolve = TransitionRequest::new(&record, None, Transition::ResolvedApplied, t(3));
467    resolve.output = Some(json!("confirmed by operator"));
468    resolve.actor = Some("operator:dennis".into());
469    let record = store.transition(resolve).await.unwrap();
470    assert_eq!(record.status, EffectStatus::Committed);
471    assert_eq!(record.output, Some(json!("confirmed by operator")));
472}
473
474async fn terminal_records_are_final<S: EffectStore>(store: S) {
475    let record = insert(&store, 1).await;
476    let lease = store
477        .acquire_lease(record.id, &worker("a"), t(0), TTL)
478        .await
479        .unwrap();
480    let record = drive(
481        &store,
482        record,
483        Some(&lease),
484        &[Transition::StartAttempt, Transition::Succeeded],
485        t(1),
486    )
487    .await;
488    store.release_lease(&lease).await.unwrap();
489    let record = reload(&store, record.id).await;
490    // Committed has one exit, compensation; nothing else may leave it.
491    for transition in ALL_TRANSITIONS
492        .into_iter()
493        .filter(|t| *t != Transition::StartCompensation)
494    {
495        let request = TransitionRequest::new(&record, None, transition, t(2));
496        assert!(
497            matches!(
498                store.transition(request).await,
499                Err(StoreError::InvalidTransition(_))
500            ),
501            "{transition} left a committed record"
502        );
503    }
504    assert_eq!(reload(&store, record.id).await, record);
505
506    // Compensated is terminal.
507    let record = drive(
508        &store,
509        record,
510        None,
511        &[
512            Transition::StartCompensation,
513            Transition::CompensationSucceeded,
514        ],
515        t(3),
516    )
517    .await;
518    for transition in ALL_TRANSITIONS {
519        let request = TransitionRequest::new(&record, None, transition, t(4));
520        assert!(
521            matches!(
522                store.transition(request).await,
523                Err(StoreError::InvalidTransition(_))
524            ),
525            "{transition} left a compensated record"
526        );
527    }
528}
529
530async fn listing_filters_and_pages<S: EffectStore>(store: S) {
531    let now = t(100);
532    let mut ids = Vec::new();
533    for n in 0..5 {
534        ids.push(insert(&store, n).await.id);
535    }
536    // 0: Pending, no lease.
537    // 1: Executing, live lease.
538    let lease = store
539        .acquire_lease(ids[1], &worker("a"), t(90), TTL)
540        .await
541        .unwrap();
542    drive(
543        &store,
544        reload(&store, ids[1]).await,
545        Some(&lease),
546        &[Transition::StartAttempt],
547        t(90),
548    )
549    .await;
550    // 2: Executing, lease expired.
551    let lease = store
552        .acquire_lease(ids[2], &worker("a"), t(0), TTL)
553        .await
554        .unwrap();
555    drive(
556        &store,
557        reload(&store, ids[2]).await,
558        Some(&lease),
559        &[Transition::StartAttempt],
560        t(1),
561    )
562    .await;
563    // 3: Verifying, lease released.
564    let lease = store
565        .acquire_lease(ids[3], &worker("a"), t(0), TTL)
566        .await
567        .unwrap();
568    drive(
569        &store,
570        reload(&store, ids[3]).await,
571        Some(&lease),
572        &[Transition::StartAttempt, Transition::StartVerification],
573        t(1),
574    )
575    .await;
576    store.release_lease(&lease).await.unwrap();
577    // 4: Unknown.
578    let lease = store
579        .acquire_lease(ids[4], &worker("a"), t(0), TTL)
580        .await
581        .unwrap();
582    drive(
583        &store,
584        reload(&store, ids[4]).await,
585        Some(&lease),
586        &[Transition::StartAttempt, Transition::OutcomeUnknown],
587        t(1),
588    )
589    .await;
590
591    let listed = |query: ListQuery| {
592        let store = &store;
593        async move {
594            store
595                .list(query)
596                .await
597                .unwrap()
598                .into_iter()
599                .map(|r| r.id)
600                .collect::<Vec<_>>()
601        }
602    };
603    assert_eq!(
604        listed(ListQuery::statuses([EffectStatus::Pending])).await,
605        [ids[0]]
606    );
607    assert_eq!(
608        listed(ListQuery::expired_leases(now)).await,
609        [ids[2], ids[3]]
610    );
611    assert_eq!(
612        listed(ListQuery::statuses([]).limit(2)).await,
613        [ids[0], ids[1]]
614    );
615    assert_eq!(
616        listed(ListQuery::statuses([]).after(ids[1])).await,
617        [ids[2], ids[3], ids[4]]
618    );
619    assert_eq!(
620        listed(ListQuery::expired_leases(now).after(ids[2])).await,
621        [ids[3]]
622    );
623
624    let records = store
625        .list(ListQuery::statuses([EffectStatus::Unknown]))
626        .await
627        .unwrap();
628    assert_eq!(
629        records,
630        [reload(&store, ids[4]).await],
631        "listed records must be complete"
632    );
633}
634
635async fn doubt_is_persisted_and_guards_failure<S: EffectStore>(store: S) {
636    let record = insert(&store, 1).await;
637    assert!(!record.may_have_applied);
638    let lease = store
639        .acquire_lease(record.id, &worker("a"), t(0), TTL)
640        .await
641        .unwrap();
642    let record = drive(
643        &store,
644        record,
645        Some(&lease),
646        &[
647            Transition::StartAttempt,
648            Transition::OutcomeUnknown,
649            Transition::ScheduleRetry,
650            Transition::StartAttempt,
651        ],
652        t(1),
653    )
654    .await;
655    assert!(
656        reload(&store, record.id).await.may_have_applied,
657        "the store must persist may_have_applied"
658    );
659    let fail = TransitionRequest::new(&record, Some(&lease), Transition::FailedDefinitively, t(2));
660    assert!(
661        matches!(
662            store.transition(fail).await,
663            Err(StoreError::InvalidTransition(_))
664        ),
665        "a failed attempt must not fail an effect that may have applied"
666    );
667    assert_eq!(reload(&store, record.id).await, record);
668}
669
670async fn compensation_is_persisted<S: EffectStore>(store: S) {
671    let record = insert(&store, 1).await;
672    let lease = store
673        .acquire_lease(record.id, &worker("a"), t(0), TTL)
674        .await
675        .unwrap();
676    let record = drive(
677        &store,
678        record,
679        Some(&lease),
680        &[
681            Transition::StartAttempt,
682            Transition::Succeeded,
683            Transition::StartCompensation,
684        ],
685        t(1),
686    )
687    .await;
688    let mut retry = TransitionRequest::new(
689        &record,
690        Some(&lease),
691        Transition::ScheduleCompensationRetry,
692        t(2),
693    );
694    retry.next_attempt_at = Some(t(5));
695    let record = store.transition(retry).await.unwrap();
696    let stored = reload(&store, record.id).await;
697    assert_eq!(stored, record, "transition must return the stored record");
698    assert_eq!(stored.status, EffectStatus::Compensating);
699    assert_eq!(
700        stored.compensation_attempts, 1,
701        "a scheduled retry is not an attempt until it starts"
702    );
703    assert_eq!(stored.next_attempt_at, Some(t(5)));
704
705    let record = drive(
706        &store,
707        record,
708        Some(&lease),
709        &[Transition::StartCompensationRetry],
710        t(5),
711    )
712    .await;
713    let stored = reload(&store, record.id).await;
714    assert_eq!(stored.compensation_attempts, 2);
715    assert_eq!(stored.next_attempt_at, None);
716
717    let record = drive(
718        &store,
719        record,
720        Some(&lease),
721        &[Transition::CompensationFailed],
722        t(6),
723    )
724    .await;
725    store.release_lease(&lease).await.unwrap();
726    let record = drive(
727        &store,
728        reload(&store, record.id).await,
729        None,
730        &[Transition::ResolvedRetry],
731        t(7),
732    )
733    .await;
734    assert_eq!(
735        record.status,
736        EffectStatus::Compensating,
737        "an operator can retry it"
738    );
739    let listed = store
740        .list(ListQuery::statuses([EffectStatus::Compensating]))
741        .await
742        .unwrap();
743    assert_eq!(listed, [record]);
744}
745
746async fn approval_is_persisted<S: EffectStore>(store: S) {
747    let record = insert(&store, 1).await;
748    assert!(!record.approved);
749    let record = drive(&store, record, None, &[Transition::RequestApproval], t(1)).await;
750    assert_eq!(record.status, EffectStatus::AwaitingApproval);
751    let mut approve = TransitionRequest::new(&record, None, Transition::Approve, t(2));
752    approve.actor = Some("operator:dennis".into());
753    let record = store.transition(approve).await.unwrap();
754    let stored = reload(&store, record.id).await;
755    assert_eq!(stored, record, "transition must return the stored record");
756    assert_eq!(stored.status, EffectStatus::Pending);
757    assert!(stored.approved, "the store must persist approval");
758
759    let other = insert(&store, 2).await;
760    let other = drive(
761        &store,
762        other,
763        None,
764        &[Transition::RequestApproval, Transition::Deny],
765        t(3),
766    )
767    .await;
768    assert_eq!(other.status, EffectStatus::Rejected);
769    assert!(!reload(&store, other.id).await.approved);
770}
771
772/// Runs effect `n` through `transitions` at time `at`, then releases its
773/// lease.
774async fn settle<S: EffectStore>(
775    store: &S,
776    n: u32,
777    transitions: &[Transition],
778    at: u64,
779) -> EffectId {
780    let record = insert(store, n).await;
781    let lease = store
782        .acquire_lease(record.id, &worker("a"), t(at), TTL)
783        .await
784        .unwrap();
785    drive(store, record, Some(&lease), transitions, t(at)).await;
786    store.release_lease(&lease).await.unwrap();
787    lease.effect_id
788}
789
790async fn pruning_removes_only_settled_idle_old_records<S: EffectStore>(store: S) {
791    let committed = [Transition::StartAttempt, Transition::Succeeded];
792    let old_a = settle(&store, 0, &committed, 10).await;
793    let failed = settle(
794        &store,
795        1,
796        &[Transition::StartAttempt, Transition::FailedDefinitively],
797        10,
798    )
799    .await;
800    let young = settle(&store, 2, &committed, 50).await;
801    let unknown = settle(
802        &store,
803        3,
804        &[Transition::StartAttempt, Transition::OutcomeUnknown],
805        10,
806    )
807    .await;
808    let escalated = settle(
809        &store,
810        6,
811        &[
812            Transition::StartAttempt,
813            Transition::OutcomeUnknown,
814            Transition::Escalate,
815        ],
816        10,
817    )
818    .await;
819    let leased = settle(&store, 4, &committed, 10).await;
820    let old_b = settle(&store, 5, &committed, 10).await;
821    // A compensation is about to start on `leased`: it holds a live lease.
822    store
823        .acquire_lease(leased, &worker("b"), t(95), TTL)
824        .await
825        .unwrap();
826
827    let now = t(100);
828    let minute = Duration::from_secs(60);
829    let prune = |status, older_than, limit| {
830        let store = &store;
831        async move {
832            store
833                .prune(PruneQuery::new(status, older_than, now).limit(limit))
834                .await
835                .unwrap()
836        }
837    };
838    let exists = |id| {
839        let store = &store;
840        async move { store.get(id).await.unwrap().is_some() }
841    };
842
843    assert_eq!(
844        prune(EffectStatus::Committed, minute, 1).await,
845        1,
846        "the limit caps a batch"
847    );
848    assert!(!exists(old_a).await, "lowest id first");
849    assert!(exists(old_b).await);
850    assert_eq!(
851        prune(EffectStatus::Committed, minute, 100).await,
852        1,
853        "only old_b is left to prune: young is too young, leased is leased"
854    );
855    assert!(!exists(old_b).await);
856    assert!(exists(young).await && exists(leased).await);
857    assert_eq!(prune(EffectStatus::Committed, minute, 100).await, 0);
858
859    // Spelled out, not `is_settled`, so the suite does not share the
860    // store's definition.
861    let settled = [
862        EffectStatus::Committed,
863        EffectStatus::Failed,
864        EffectStatus::Rejected,
865        EffectStatus::Compensated,
866    ];
867    for status in ALL_STATUSES.into_iter().filter(|s| !settled.contains(s)) {
868        assert_eq!(
869            prune(status, Duration::ZERO, 100).await,
870            0,
871            "{status} is not settled and is never pruned"
872        );
873    }
874    assert!(exists(unknown).await && exists(escalated).await);
875    assert_eq!(
876        prune(
877            EffectStatus::Failed,
878            Duration::from_secs(200_000_000_000),
879            100
880        )
881        .await,
882        0,
883        "an age older than any record prunes nothing"
884    );
885    assert_eq!(prune(EffectStatus::Failed, minute, 100).await, 1);
886    assert!(!exists(failed).await);
887
888    pruned_records_leave_nothing_behind(&store, old_a, young).await;
889}
890
891/// After `pruned` (key 0) was pruned and `kept` was not.
892async fn pruned_records_leave_nothing_behind<S: EffectStore>(
893    store: &S,
894    pruned: EffectId,
895    kept: EffectId,
896) {
897    // The audit trail goes with the record; others keep theirs.
898    assert!(matches!(
899        store.events(pruned).await,
900        Err(StoreError::NotFound(id)) if id == pruned
901    ));
902    assert_eq!(store.events(kept).await.unwrap().len(), 2);
903    // The key is free again: the next request is a new effect.
904    assert!(store.get_by_key(&key(0)).await.unwrap().is_none());
905    let again = store.insert_or_get(new_effect(0)).await.unwrap();
906    assert!(again.inserted);
907    assert_ne!(again.record.id, pruned);
908    assert_eq!(again.record.status, EffectStatus::Pending);
909    assert_eq!(store.events(again.record.id).await.unwrap(), Vec::new());
910}