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 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
64fn 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
97async 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 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 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 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 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 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 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
772async 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 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 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
891async fn pruned_records_leave_nothing_behind<S: EffectStore>(
893 store: &S,
894 pruned: EffectId,
895 kept: EffectId,
896) {
897 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 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}