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