1use 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
25pub 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
65fn 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
98async 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 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 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 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 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 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 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 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
813async 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 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 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
932async fn pruned_records_leave_nothing_behind<S: EffectStore>(
934 store: &S,
935 pruned: EffectId,
936 kept: EffectId,
937) {
938 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 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}