1use aion_core::{Event, TimerId, WorkflowId, status_from_events};
37use std::collections::HashMap;
38use std::sync::Arc;
39
40use aion_store::{ReadableEventStore, StoreError, TimerEntry};
41use chrono::{DateTime, Utc};
42
43use crate::engine_seam::EngineSeamError;
44use crate::time::timer_service::{
45 RetireAttempt, TimerDisposition, armed_fire_at_in_active_segment,
46 timer_disposition_in_active_segment,
47};
48use crate::time::{TimerService, TimerServiceError, is_deadline_timer};
49
50pub struct TimerRecovery {
52 store: Arc<dyn ReadableEventStore>,
53 timer_service: Arc<TimerService>,
54 warned_orphans: std::sync::Mutex<std::collections::HashSet<WorkflowId>>,
60}
61
62#[derive(thiserror::Error, Debug, Clone, PartialEq, Eq)]
64pub enum TimerRecoveryError {
65 #[error("timer recovery store operation failed: {0}")]
67 Store(#[from] StoreError),
68
69 #[error("timer recovery fire operation failed: {0}")]
71 Timer(#[from] TimerServiceError),
72}
73
74impl TimerRecovery {
75 #[must_use]
81 pub fn new(store: Arc<dyn ReadableEventStore>, timer_service: Arc<TimerService>) -> Self {
82 Self {
83 store,
84 timer_service,
85 warned_orphans: std::sync::Mutex::new(std::collections::HashSet::new()),
86 }
87 }
88
89 pub async fn recover_on_startup(
107 &self,
108 now: DateTime<Utc>,
109 ) -> Result<SweepCounts, TimerRecoveryError> {
110 let sweep_started = std::time::Instant::now();
111 let due_by_workflow = group_by_workflow(self.store.expired_timers(now).await?);
112 let expired_rows: usize = due_by_workflow.values().map(Vec::len).sum();
113 tracing::info!(
114 expired_rows,
115 due_workflows = due_by_workflow.len(),
116 "timer recovery startup sweep started"
117 );
118 let mut counts = SweepCounts::default();
119 let mut rearmed = 0usize;
120 let result = self
121 .startup_passes(now, due_by_workflow, &mut counts, &mut rearmed)
122 .await;
123 let elapsed_ms = u64::try_from(sweep_started.elapsed().as_millis()).unwrap_or(u64::MAX);
127 match &result {
128 Ok(()) => tracing::info!(
129 fired = counts.fired,
130 redelivered = counts.redelivered,
131 retired = counts.retired,
132 superseded = counts.superseded,
133 retire_failures = counts.retire_failures,
134 skipped_orphans = counts.skipped_orphans,
135 rearmed,
136 elapsed_ms,
137 "timer recovery startup sweep finished"
138 ),
139 Err(error) => tracing::warn!(
140 fired = counts.fired,
141 redelivered = counts.redelivered,
142 retired = counts.retired,
143 superseded = counts.superseded,
144 retire_failures = counts.retire_failures,
145 skipped_orphans = counts.skipped_orphans,
146 rearmed,
147 elapsed_ms,
148 %error,
149 "timer recovery startup sweep ABORTED; counters cover the work \
150 completed before the failure"
151 ),
152 }
153 result.map(|()| counts)
154 }
155
156 async fn startup_passes(
160 &self,
161 now: DateTime<Utc>,
162 mut due_by_workflow: HashMap<WorkflowId, Vec<TimerEntry>>,
163 counts: &mut SweepCounts,
164 rearmed: &mut usize,
165 ) -> Result<(), TimerRecoveryError> {
166 for workflow_id in self.store.list_active().await? {
167 let history = self.store.read_history(&workflow_id).await?;
168 for (timer_id, fire_at, armed_seq) in outstanding_future_timers(&history, now) {
169 self.timer_service
170 .schedule(workflow_id.clone(), timer_id, fire_at, armed_seq)
171 .await?;
172 *rearmed += 1;
173 }
174 if let Some(entries) = due_by_workflow.remove(&workflow_id) {
175 self.dispose_due_rows(&workflow_id, &history, entries, counts)
176 .await?;
177 }
178 }
179 for (workflow_id, entries) in due_by_workflow {
180 let history = self.store.read_history(&workflow_id).await?;
181 self.dispose_due_rows(&workflow_id, &history, entries, counts)
182 .await?;
183 }
184 Ok(())
185 }
186
187 fn first_orphan_sighting(&self, workflow_id: &WorkflowId) -> bool {
193 self.warned_orphans
194 .lock()
195 .map_or(true, |mut warned| warned.insert(workflow_id.clone()))
196 }
197
198 async fn dispose_due_rows(
216 &self,
217 workflow_id: &WorkflowId,
218 history: &[Event],
219 entries: Vec<TimerEntry>,
220 counts: &mut SweepCounts,
221 ) -> Result<(), TimerRecoveryError> {
222 if history.is_empty() {
223 counts.skipped_orphans += entries.len();
224 if self.first_orphan_sighting(workflow_id) {
225 tracing::warn!(
226 %workflow_id,
227 rows = entries.len(),
228 "skipping due timer rows for a workflow with no recorded history \
229 (orphaned rows); nothing is fired and nothing is retired — \
230 counted in every sweep's skipped_orphans, warned once"
231 );
232 }
233 return Ok(());
234 }
235 if status_from_events(history).is_terminal() {
236 for entry in entries {
237 tracing::debug!(
238 %workflow_id,
239 timer_id = %entry.timer_id,
240 fire_at = %entry.fire_at,
241 "retiring due timer row of a terminal workflow"
242 );
243 let attempt = self
244 .timer_service
245 .retire_consumed_row(
246 workflow_id,
247 &entry.timer_id,
248 entry.fire_at,
249 entry.armed_seq,
250 )
251 .await;
252 counts.record_row_retirement(attempt);
253 }
254 return Ok(());
255 }
256 for entry in entries {
257 match timer_disposition_in_active_segment(history, &entry.timer_id) {
258 TimerDisposition::Live
272 if armed_fire_at_in_active_segment(history, &entry.timer_id)
273 == Some(entry.fire_at) =>
274 {
275 self.fire_due(workflow_id, &entry, counts).await?;
276 }
277 TimerDisposition::Live => {
278 tracing::warn!(
279 %workflow_id,
280 timer_id = %entry.timer_id,
281 row_fire_at = %entry.fire_at,
282 armed_fire_at = ?armed_fire_at_in_active_segment(history, &entry.timer_id),
283 "due timer row disagrees with the recorded arming; retiring the \
284 stale row instead of firing it"
285 );
286 let attempt = self
287 .timer_service
288 .retire_consumed_row(
289 workflow_id,
290 &entry.timer_id,
291 entry.fire_at,
292 entry.armed_seq,
293 )
294 .await;
295 counts.record_row_retirement(attempt);
296 }
297 TimerDisposition::Fired if !is_deadline_timer(&entry.timer_id) => {
304 self.redeliver_surviving_fired_row(workflow_id, &entry, counts)
305 .await?;
306 }
307 TimerDisposition::Fired
308 | TimerDisposition::Cancelled
309 | TimerDisposition::Absent => {
310 tracing::debug!(
311 %workflow_id,
312 timer_id = %entry.timer_id,
313 fire_at = %entry.fire_at,
314 "retiring consumed timer row without entering the fire path"
315 );
316 let attempt = self
317 .timer_service
318 .retire_consumed_row(
319 workflow_id,
320 &entry.timer_id,
321 entry.fire_at,
322 entry.armed_seq,
323 )
324 .await;
325 counts.record_row_retirement(attempt);
326 }
327 }
328 }
329 Ok(())
330 }
331
332 async fn redeliver_surviving_fired_row(
336 &self,
337 workflow_id: &WorkflowId,
338 entry: &TimerEntry,
339 counts: &mut SweepCounts,
340 ) -> Result<(), TimerRecoveryError> {
341 match self
342 .timer_service
343 .redeliver_owed_wake(
344 workflow_id.clone(),
345 entry.timer_id.clone(),
346 entry.fire_at,
347 entry.armed_seq,
348 )
349 .await
350 {
351 Ok((true, _)) => counts.redelivered += 1,
356 Ok((false, attempt)) => counts.record_row_retirement(attempt),
357 Err(TimerServiceError::Engine(EngineSeamError::UnknownWorkflow {
366 workflow_id: gone,
367 })) => {
368 tracing::info!(
369 workflow_id = %gone,
370 timer_id = %entry.timer_id,
371 "workflow left residency mid-redelivery; retiring the \
372 consumed row without a wake"
373 );
374 let attempt = self
375 .timer_service
376 .retire_consumed_row(
377 workflow_id,
378 &entry.timer_id,
379 entry.fire_at,
380 entry.armed_seq,
381 )
382 .await;
383 counts.record_row_retirement(attempt);
384 }
385 Err(other) => return Err(other.into()),
386 }
387 Ok(())
388 }
389
390 async fn fire_due(
392 &self,
393 workflow_id: &WorkflowId,
394 entry: &TimerEntry,
395 counts: &mut SweepCounts,
396 ) -> Result<(), TimerRecoveryError> {
397 match self
398 .timer_service
399 .fire_timer(workflow_id.clone(), entry.timer_id.clone(), entry.fire_at)
400 .await
401 {
402 Ok(()) => counts.fired += 1,
403 Err(TimerServiceError::Engine(EngineSeamError::UnknownWorkflow { workflow_id })) => {
410 if self.first_orphan_sighting(&workflow_id) {
411 tracing::warn!(
412 %workflow_id,
413 timer_id = %entry.timer_id,
414 "skipping recovered timer for unknown workflow (orphaned timer); \
415 the workflow no longer exists — counted in every sweep's \
416 skipped_orphans, warned once"
417 );
418 }
419 counts.skipped_orphans += 1;
420 }
421 Err(other) => return Err(other.into()),
422 }
423 Ok(())
424 }
425}
426
427#[derive(Clone, Copy, Debug, Default)]
437pub struct SweepCounts {
438 pub fired: usize,
440 pub redelivered: usize,
443 pub retired: usize,
445 pub superseded: usize,
448 pub retire_failures: usize,
452 pub skipped_orphans: usize,
455}
456
457impl SweepCounts {
458 fn record_row_retirement(&mut self, attempt: RetireAttempt) {
460 match attempt {
461 RetireAttempt::Retired => self.retired += 1,
462 RetireAttempt::Superseded => self.superseded += 1,
463 RetireAttempt::Failed => self.retire_failures += 1,
464 }
465 }
466}
467
468fn group_by_workflow(entries: Vec<TimerEntry>) -> HashMap<WorkflowId, Vec<TimerEntry>> {
471 let mut by_workflow: HashMap<WorkflowId, Vec<TimerEntry>> = HashMap::new();
472 for entry in entries {
473 by_workflow
474 .entry(entry.workflow_id.clone())
475 .or_default()
476 .push(entry);
477 }
478 by_workflow
479}
480
481fn outstanding_future_timers(
482 history: &[Event],
483 now: DateTime<Utc>,
484) -> Vec<(TimerId, DateTime<Utc>, u64)> {
485 let mut outstanding: HashMap<TimerId, (DateTime<Utc>, u64)> = HashMap::new();
486 for event in history {
487 match event {
488 Event::TimerStarted {
489 envelope,
490 timer_id,
491 fire_at,
492 } => {
493 outstanding.insert(timer_id.clone(), (*fire_at, envelope.seq));
494 }
495 Event::TimerFired { timer_id, .. } | Event::TimerCancelled { timer_id, .. } => {
496 outstanding.remove(timer_id);
497 }
498 _ => {}
499 }
500 }
501 outstanding
502 .into_iter()
503 .filter(|(_, (fire_at, _))| *fire_at > now)
504 .map(|(timer_id, (fire_at, armed_seq))| (timer_id, fire_at, armed_seq))
505 .collect()
506}
507
508#[cfg(test)]
509mod tests {
510 use std::sync::Arc;
511
512 use aion_core::{Event, EventEnvelope, RunId, TimerCancelCause, TimerId, WorkflowId};
513 use aion_store::{
514 InMemoryStore, ReadableEventStore, StoreError, WritableEventStore, WriteToken,
515 };
516 use chrono::{DateTime, Utc};
517
518 use super::{TimerRecovery, TimerRecoveryError, outstanding_future_timers};
519 use crate::engine_seam::test_support::{DeliveredWorkflowMessage, FakeEngineHandle};
520 use crate::engine_seam::{
521 EngineHandle, EngineSeamError, WorkflowProcessHandle, WorkflowResidency,
522 };
523 use crate::time::TimerService;
524 use crate::time::deadline_timer_id;
525
526 #[derive(Debug, thiserror::Error)]
527 enum TestError {
528 #[error(transparent)]
529 Recovery(#[from] TimerRecoveryError),
530
531 #[error(transparent)]
532 Store(#[from] StoreError),
533
534 #[error(transparent)]
535 Engine(#[from] EngineSeamError),
536 }
537
538 fn instant(offset_seconds: i64) -> DateTime<Utc> {
539 DateTime::from_timestamp(1_700_000_000 + offset_seconds, 0).unwrap_or_default()
540 }
541
542 fn recorded_at() -> DateTime<Utc> {
543 instant(1)
544 }
545
546 fn workflow_id() -> WorkflowId {
547 WorkflowId::new_v4()
548 }
549
550 fn timer_id(sequence: u64) -> TimerId {
551 TimerId::anonymous(sequence)
552 }
553
554 fn recovery() -> (Arc<InMemoryStore>, Arc<FakeEngineHandle>, TimerRecovery) {
555 let concrete_store = Arc::new(InMemoryStore::default());
556 let writable: Arc<dyn WritableEventStore> = concrete_store.clone();
557 let readable: Arc<dyn ReadableEventStore> = concrete_store.clone();
558 let engine = Arc::new(FakeEngineHandle::recording_to(writable));
559 let timer_service = Arc::new(TimerService::with_recorded_at(
560 engine.clone(),
561 readable.clone(),
562 recorded_at,
563 ));
564 let recovery = TimerRecovery::new(readable, timer_service);
565 (concrete_store, engine, recovery)
566 }
567
568 async fn history(
569 store: &InMemoryStore,
570 workflow_id: &WorkflowId,
571 ) -> Result<Vec<Event>, StoreError> {
572 store.read_history(workflow_id).await
573 }
574
575 fn timer_started_event(
579 workflow_id: &WorkflowId,
580 timer_id: &TimerId,
581 seq: u64,
582 fire_at: DateTime<Utc>,
583 ) -> Event {
584 Event::TimerStarted {
585 envelope: EventEnvelope {
586 seq,
587 recorded_at: instant(0),
588 workflow_id: workflow_id.clone(),
589 },
590 timer_id: timer_id.clone(),
591 fire_at,
592 }
593 }
594
595 fn workflow_started_event(workflow_id: &WorkflowId, seq: u64) -> Event {
596 Event::WorkflowStarted {
597 envelope: EventEnvelope {
598 seq,
599 recorded_at: instant(0),
600 workflow_id: workflow_id.clone(),
601 },
602 workflow_type: "fixture".to_owned(),
603 input: aion_core::Payload::new(aion_core::ContentType::Json, b"null".to_vec()),
604 run_id: RunId::new_v4(),
605 parent_run_id: None,
606 parent_workflow_id: None,
607 package_version: aion_core::PackageVersion::new("a".repeat(64)),
608 }
609 }
610
611 fn count_timer_fired(events: &[Event], timer_id: &TimerId) -> usize {
612 events
613 .iter()
614 .filter(|event| {
615 matches!(event, Event::TimerFired { timer_id: recorded, .. } if recorded == timer_id)
616 })
617 .count()
618 }
619
620 #[tokio::test]
621 async fn startup_sweep_fires_past_timer_and_delivers() -> Result<(), TestError> {
622 let process = WorkflowProcessHandle::new(42);
623 let (store, engine, recovery) = recovery();
624 let workflow_id = workflow_id();
625 let timer_id = timer_id(1);
626 let fire_at = instant(10);
627 engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
628 engine.record_workflow_event(
629 &workflow_id,
630 timer_started_event(&workflow_id, &timer_id, 1, fire_at),
631 )?;
632 store
633 .schedule_timer(&workflow_id, &timer_id, fire_at, 1)
634 .await?;
635
636 let recovered = recovery.recover_on_startup(instant(20)).await?;
637
638 assert_eq!(recovered.fired, 1);
639 assert_eq!(
640 count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
641 1
642 );
643 assert_eq!(
644 engine.delivered_messages()?,
645 vec![(
646 process,
647 DeliveredWorkflowMessage::TimerFired {
648 timer_id: timer_id.clone(),
649 fire_at
650 }
651 )]
652 );
653 Ok(())
654 }
655
656 #[tokio::test]
657 async fn startup_sweep_does_not_fire_future_timer() -> Result<(), TestError> {
658 let process = WorkflowProcessHandle::new(42);
659 let (store, engine, recovery) = recovery();
660 let workflow_id = workflow_id();
661 let timer_id = timer_id(2);
662 engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
663 store
664 .schedule_timer(&workflow_id, &timer_id, instant(30), 1)
665 .await?;
666
667 let recovered = recovery.recover_on_startup(instant(20)).await?;
668
669 assert_eq!(recovered.fired, 0);
670 assert_eq!(
671 count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
672 0
673 );
674 assert!(engine.delivered_messages()?.is_empty());
675 Ok(())
676 }
677
678 #[tokio::test]
694 async fn a_stale_past_due_row_is_retired_not_fired_against_a_rearmed_future_arming()
695 -> Result<(), TestError> {
696 let process = WorkflowProcessHandle::new(42);
697 let (store, engine, recovery) = recovery();
698 let workflow_id = workflow_id();
699 let timer_id = timer_id(11);
700 let stale_row_fire_at = instant(5); let armed_fire_at = instant(300); engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
703 engine.record_workflow_event(&workflow_id, workflow_started_event(&workflow_id, 1))?;
705 engine.record_workflow_event(
707 &workflow_id,
708 timer_started_event(&workflow_id, &timer_id, 2, armed_fire_at),
709 )?;
710 store
713 .schedule_timer(&workflow_id, &timer_id, stale_row_fire_at, 1)
714 .await?;
715
716 let recovered = recovery.recover_on_startup(instant(20)).await?;
717
718 assert_eq!(recovered.fired, 0, "a stale row must never fire");
719 assert_eq!(
720 count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
721 0,
722 "no premature TimerFired may be fabricated for the future arming"
723 );
724 assert!(
725 engine.delivered_messages()?.is_empty(),
726 "no wake may be delivered for a stale row"
727 );
728 assert_eq!(
732 recovered.superseded, 1,
733 "the stale row lands in the superseded bin: the re-arm pass's \
734 replacement row won the key"
735 );
736 assert_eq!(
737 recovered.fired
738 + recovered.redelivered
739 + recovered.retired
740 + recovered.superseded
741 + recovered.retire_failures
742 + recovered.skipped_orphans,
743 1,
744 "the one due row lands in exactly one counter"
745 );
746 assert!(
749 store.expired_timers(instant(20)).await?.is_empty(),
750 "the stale past-due row is gone"
751 );
752 let rows = store.expired_timers(armed_fire_at).await?;
753 assert_eq!(rows.len(), 1, "the armed row survives: {rows:?}");
754 assert_eq!(
755 rows[0].fire_at, armed_fire_at,
756 "the row converged to the recorded arming's fire_at"
757 );
758 Ok(())
759 }
760
761 #[tokio::test]
770 async fn a_refused_retirement_lands_in_retire_failures_and_heals_at_the_next_sweep()
771 -> Result<(), TestError> {
772 use crate::store_faults::FlakyStore;
773
774 let flaky = Arc::new(FlakyStore::new());
775 let writable: Arc<dyn WritableEventStore> = flaky.clone();
776 let readable: Arc<dyn ReadableEventStore> = flaky.clone();
777 let engine = Arc::new(FakeEngineHandle::recording_to(writable));
778 let timer_service = Arc::new(TimerService::with_recorded_at(
779 engine.clone(),
780 readable.clone(),
781 recorded_at,
782 ));
783 let recovery = TimerRecovery::new(readable, timer_service);
784
785 let workflow_id = workflow_id();
788 let timer_id = timer_id(12);
789 engine.record_workflow_event(
790 &workflow_id,
791 timer_started_event(&workflow_id, &timer_id, 1, instant(5)),
792 )?;
793 engine.record_workflow_event(
794 &workflow_id,
795 Event::TimerCancelled {
796 cause: TimerCancelCause::WorkflowIntent,
797 envelope: EventEnvelope {
798 seq: 2,
799 recorded_at: instant(6),
800 workflow_id: workflow_id.clone(),
801 },
802 timer_id: timer_id.clone(),
803 },
804 )?;
805 flaky
806 .schedule_timer(&workflow_id, &timer_id, instant(5), 1)
807 .await?;
808
809 flaky.fail_next_retirements(1);
810 let refused = recovery.recover_on_startup(instant(20)).await?;
811
812 assert_eq!(
813 refused.retire_failures, 1,
814 "the refused retirement must be counted as a failure, not a deletion"
815 );
816 assert_eq!(refused.retired, 0);
817 assert_eq!(
818 refused.fired
819 + refused.redelivered
820 + refused.retired
821 + refused.superseded
822 + refused.retire_failures
823 + refused.skipped_orphans,
824 1,
825 "the refused row still lands in exactly one counter"
826 );
827 assert_eq!(
828 flaky.expired_timers(instant(20)).await?.len(),
829 1,
830 "the row survives the refusal — nothing is silently dropped"
831 );
832
833 let healed = recovery.recover_on_startup(instant(20)).await?;
835 assert_eq!(healed.retired, 1, "the next sweep retires the survivor");
836 assert_eq!(healed.retire_failures, 0);
837 assert!(
838 flaky.expired_timers(instant(20)).await?.is_empty(),
839 "the row is gone once the store accepts the retirement"
840 );
841 Ok(())
842 }
843
844 #[tokio::test]
854 async fn a_second_boot_sweep_after_a_fired_row_retires_finds_nothing_due()
855 -> Result<(), TestError> {
856 let process = WorkflowProcessHandle::new(42);
857 let (store, engine, recovery) = recovery();
858 let workflow_id = workflow_id();
859 let timer_id = timer_id(3);
860 let fire_at = instant(25);
861 engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
862 engine.record_workflow_event(
863 &workflow_id,
864 timer_started_event(&workflow_id, &timer_id, 1, fire_at),
865 )?;
866 store
867 .schedule_timer(&workflow_id, &timer_id, fire_at, 1)
868 .await?;
869
870 assert_eq!(recovery.recover_on_startup(instant(30)).await?.fired, 1);
871 assert_eq!(recovery.recover_on_startup(instant(30)).await?.fired, 0);
872
873 assert_eq!(
874 count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
875 1,
876 "the durable TimerFired is recorded exactly once across repeated sweeps"
877 );
878 assert_eq!(engine.delivered_messages()?.len(), 1);
879 Ok(())
880 }
881
882 #[tokio::test]
883 async fn running_startup_sweep_twice_records_due_timer_once_total() -> Result<(), TestError> {
884 let process = WorkflowProcessHandle::new(42);
885 let (store, engine, recovery) = recovery();
886 let workflow_id = workflow_id();
887 let timer_id = timer_id(4);
888 let fire_at = instant(10);
889 engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
890 engine.record_workflow_event(
891 &workflow_id,
892 timer_started_event(&workflow_id, &timer_id, 1, fire_at),
893 )?;
894 store
895 .schedule_timer(&workflow_id, &timer_id, fire_at, 1)
896 .await?;
897
898 recovery.recover_on_startup(instant(20)).await?;
899 recovery.recover_on_startup(instant(20)).await?;
900
901 assert_eq!(
902 count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
903 1,
904 "the durable TimerFired is recorded exactly once across repeated sweeps"
905 );
906 assert_eq!(engine.delivered_messages()?.len(), 1);
910 Ok(())
911 }
912
913 #[tokio::test]
914 async fn cancelled_timer_is_never_fired_by_recovery() -> Result<(), TestError> {
915 let process = WorkflowProcessHandle::new(42);
916 let (store, engine, recovery) = recovery();
917 let workflow_id = workflow_id();
918 let timer_id = timer_id(5);
919 let fire_at = instant(10);
920 engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
921 store
922 .schedule_timer(&workflow_id, &timer_id, fire_at, 1)
923 .await?;
924 engine.record_workflow_event(
925 &workflow_id,
926 Event::TimerCancelled {
927 cause: aion_core::TimerCancelCause::WorkflowIntent,
928 envelope: EventEnvelope {
929 seq: 1,
930 recorded_at: instant(9),
931 workflow_id: workflow_id.clone(),
932 },
933 timer_id: timer_id.clone(),
934 },
935 )?;
936
937 recovery.recover_on_startup(instant(20)).await?;
938
939 assert_eq!(
940 count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
941 0
942 );
943 assert!(engine.delivered_messages()?.is_empty());
944 Ok(())
945 }
946
947 #[test]
954 fn cancelled_predecessor_deadline_is_not_rearmed_after_continue_as_new() {
955 let workflow_id = workflow_id();
956 let predecessor_run = RunId::new_v4();
957 let deadline = deadline_timer_id(&predecessor_run).unwrap_or_else(|_| timer_id(0));
958 let now = instant(0);
959 let fire_at = instant(120); let started = |seq: u64, run: &RunId| Event::WorkflowStarted {
962 envelope: EventEnvelope {
963 seq,
964 recorded_at: instant(0),
965 workflow_id: workflow_id.clone(),
966 },
967 workflow_type: "sleeper".to_owned(),
968 input: aion_core::Payload::new(aion_core::ContentType::Json, b"null".to_vec()),
969 run_id: run.clone(),
970 parent_run_id: None,
971 parent_workflow_id: None,
972 package_version: aion_core::PackageVersion::new("a".repeat(64)),
973 };
974 let deadline_started = Event::TimerStarted {
975 envelope: EventEnvelope {
976 seq: 2,
977 recorded_at: instant(0),
978 workflow_id: workflow_id.clone(),
979 },
980 timer_id: deadline.clone(),
981 fire_at,
982 };
983 let continued = Event::WorkflowContinuedAsNew {
984 envelope: EventEnvelope {
985 seq: 3,
986 recorded_at: instant(1),
987 workflow_id: workflow_id.clone(),
988 },
989 input: aion_core::Payload::new(aion_core::ContentType::Json, b"null".to_vec()),
990 workflow_type: None,
991 parent_run_id: predecessor_run.clone(),
992 };
993
994 let uncancelled = vec![
997 started(1, &predecessor_run),
998 deadline_started.clone(),
999 continued.clone(),
1000 ];
1001 assert!(
1002 outstanding_future_timers(&uncancelled, now)
1003 .into_iter()
1004 .any(|(timer_id, _, _)| timer_id == deadline),
1005 "an uncancelled predecessor deadline WOULD be re-armed after failover"
1006 );
1007
1008 let cancelled = vec![
1010 started(1, &predecessor_run),
1011 deadline_started,
1012 continued,
1013 Event::TimerCancelled {
1014 envelope: EventEnvelope {
1015 seq: 4,
1016 recorded_at: instant(1),
1017 workflow_id: workflow_id.clone(),
1018 },
1019 timer_id: deadline.clone(),
1020 cause: TimerCancelCause::WorkflowIntent,
1021 },
1022 ];
1023 assert!(
1024 !outstanding_future_timers(&cancelled, now)
1025 .into_iter()
1026 .any(|(timer_id, _, _)| timer_id == deadline),
1027 "the WorkflowIntent cancel closes the whole-history re-arm hole"
1028 );
1029 }
1030
1031 #[tokio::test]
1032 async fn orphaned_timer_for_unknown_workflow_is_skipped_not_fatal() -> Result<(), TestError> {
1033 let (store, engine, recovery) = recovery();
1039 let workflow_id = workflow_id();
1040 let timer_id = timer_id(6);
1041 let fire_at = instant(10);
1042
1043 store
1047 .schedule_timer(&workflow_id, &timer_id, fire_at, 1)
1048 .await?;
1049 engine.record_workflow_event(
1050 &workflow_id,
1051 timer_started_event(&workflow_id, &timer_id, 1, fire_at),
1052 )?;
1053 engine.push_record_response(Err(EngineSeamError::UnknownWorkflow {
1057 workflow_id: workflow_id.clone(),
1058 }))?;
1059
1060 let recovered = recovery.recover_on_startup(instant(20)).await?;
1062
1063 assert_eq!(
1064 recovered.fired, 0,
1065 "the orphaned timer is skipped, not fired"
1066 );
1067 assert_eq!(
1068 recovered.skipped_orphans, 1,
1069 "the engine-side orphan is counted in skipped_orphans (round-3 N5)"
1070 );
1071 assert_eq!(
1072 count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
1073 0,
1074 "no TimerFired is recorded for an unknown workflow"
1075 );
1076 assert!(
1077 engine.delivered_messages()?.is_empty(),
1078 "nothing is delivered for an unknown workflow"
1079 );
1080 Ok(())
1081 }
1082
1083 struct CountingStore {
1087 inner: Arc<InMemoryStore>,
1088 history_reads: std::sync::atomic::AtomicUsize,
1089 }
1090
1091 impl CountingStore {
1092 fn new(inner: Arc<InMemoryStore>) -> Self {
1093 Self {
1094 inner,
1095 history_reads: std::sync::atomic::AtomicUsize::new(0),
1096 }
1097 }
1098
1099 fn history_reads(&self) -> usize {
1100 self.history_reads.load(std::sync::atomic::Ordering::SeqCst)
1101 }
1102 }
1103
1104 #[async_trait::async_trait]
1105 impl ReadableEventStore for CountingStore {
1106 async fn read_history(&self, workflow_id: &WorkflowId) -> Result<Vec<Event>, StoreError> {
1107 self.history_reads
1108 .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1109 self.inner.read_history(workflow_id).await
1110 }
1111
1112 async fn read_history_from(
1113 &self,
1114 workflow_id: &WorkflowId,
1115 from_seq: u64,
1116 ) -> Result<Vec<Event>, StoreError> {
1117 self.inner.read_history_from(workflow_id, from_seq).await
1118 }
1119
1120 async fn read_run_chain(
1121 &self,
1122 workflow_id: &WorkflowId,
1123 ) -> Result<Vec<aion_store::RunSummary>, StoreError> {
1124 self.inner.read_run_chain(workflow_id).await
1125 }
1126
1127 async fn list_workflow_ids(&self) -> Result<Vec<WorkflowId>, StoreError> {
1128 self.inner.list_workflow_ids().await
1129 }
1130
1131 async fn list_active(&self) -> Result<Vec<WorkflowId>, StoreError> {
1132 self.inner.list_active().await
1133 }
1134
1135 async fn list_paused(&self) -> Result<Vec<WorkflowId>, StoreError> {
1136 self.inner.list_paused().await
1137 }
1138 async fn stream_heads(
1139 &self,
1140 ) -> Result<Vec<aion_store::visibility::StreamHead>, StoreError> {
1141 self.inner.stream_heads().await
1142 }
1143
1144 async fn query(
1145 &self,
1146 filter: &aion_core::WorkflowFilter,
1147 ) -> Result<Vec<aion_core::WorkflowSummary>, StoreError> {
1148 self.inner.query(filter).await
1149 }
1150
1151 async fn schedule_timer(
1152 &self,
1153 workflow_id: &WorkflowId,
1154 timer_id: &TimerId,
1155 fire_at: DateTime<Utc>,
1156 armed_seq: u64,
1157 ) -> Result<(), StoreError> {
1158 self.inner
1159 .schedule_timer(workflow_id, timer_id, fire_at, armed_seq)
1160 .await
1161 }
1162
1163 async fn expired_timers(
1164 &self,
1165 as_of: DateTime<Utc>,
1166 ) -> Result<Vec<aion_store::TimerEntry>, StoreError> {
1167 self.inner.expired_timers(as_of).await
1168 }
1169
1170 async fn retire_timer(
1171 &self,
1172 workflow_id: &WorkflowId,
1173 timer_id: &TimerId,
1174 fire_at: DateTime<Utc>,
1175 armed_seq: u64,
1176 ) -> Result<aion_store::TimerRetirement, StoreError> {
1177 self.inner
1178 .retire_timer(workflow_id, timer_id, fire_at, armed_seq)
1179 .await
1180 }
1181 }
1182
1183 fn counting_recovery() -> (
1186 Arc<InMemoryStore>,
1187 Arc<CountingStore>,
1188 Arc<FakeEngineHandle>,
1189 TimerRecovery,
1190 ) {
1191 let inner = Arc::new(InMemoryStore::default());
1192 let counting = Arc::new(CountingStore::new(inner.clone()));
1193 let writable: Arc<dyn WritableEventStore> = inner.clone();
1194 let readable: Arc<dyn ReadableEventStore> = counting.clone();
1195 let engine = Arc::new(FakeEngineHandle::recording_to(writable));
1196 let timer_service = Arc::new(TimerService::with_recorded_at(
1197 engine.clone(),
1198 readable.clone(),
1199 recorded_at,
1200 ));
1201 let recovery = TimerRecovery::new(readable, timer_service);
1202 (inner, counting, engine, recovery)
1203 }
1204
1205 fn timer_fired_event(workflow_id: &WorkflowId, timer_id: &TimerId, seq: u64) -> Event {
1206 Event::TimerFired {
1207 envelope: EventEnvelope {
1208 seq,
1209 recorded_at: instant(0),
1210 workflow_id: workflow_id.clone(),
1211 },
1212 timer_id: timer_id.clone(),
1213 }
1214 }
1215
1216 async fn seed_consumed_armings(
1220 engine: &FakeEngineHandle,
1221 store: &InMemoryStore,
1222 workflow_id: &WorkflowId,
1223 count: u64,
1224 ) -> Result<u64, TestError> {
1225 let mut seq = 0;
1226 for i in 0..count {
1227 let consumed = timer_id(i);
1228 seq += 1;
1229 let armed_seq = seq;
1230 engine.record_workflow_event(
1231 workflow_id,
1232 timer_started_event(workflow_id, &consumed, armed_seq, instant(5)),
1233 )?;
1234 seq += 1;
1235 engine.record_workflow_event(
1236 workflow_id,
1237 timer_fired_event(workflow_id, &consumed, seq),
1238 )?;
1239 store
1240 .schedule_timer(workflow_id, &consumed, instant(5), armed_seq)
1241 .await?;
1242 }
1243 Ok(seq)
1244 }
1245
1246 #[tokio::test]
1259 async fn a_grouped_sweep_reads_once_per_workflow_and_retires_the_backlog()
1260 -> Result<(), TestError> {
1261 const CONSUMED_PER_WORKFLOW: u64 = 300;
1262 let process = WorkflowProcessHandle::new(42);
1263 let (inner, counting, engine, recovery) = counting_recovery();
1264
1265 let running = workflow_id();
1268 engine.set_residency(running.clone(), WorkflowResidency::Resident(process))?;
1269 let mut seq =
1270 seed_consumed_armings(&engine, &inner, &running, CONSUMED_PER_WORKFLOW).await?;
1271 let live = timer_id(9_999);
1272 seq += 1;
1273 engine.record_workflow_event(
1274 &running,
1275 timer_started_event(&running, &live, seq, instant(5)),
1276 )?;
1277 inner
1278 .schedule_timer(&running, &live, instant(5), seq)
1279 .await?;
1280
1281 let finished = workflow_id();
1284 let seq = seed_consumed_armings(&engine, &inner, &finished, CONSUMED_PER_WORKFLOW).await?;
1285 engine.record_workflow_event(
1286 &finished,
1287 Event::WorkflowCompleted {
1288 envelope: EventEnvelope {
1289 seq: seq + 1,
1290 recorded_at: instant(6),
1291 workflow_id: finished.clone(),
1292 },
1293 result: aion_core::Payload::new(aion_core::ContentType::Json, b"null".to_vec()),
1294 },
1295 )?;
1296
1297 let recovered = recovery.recover_on_startup(instant(20)).await?;
1298
1299 assert_eq!(recovered.fired, 1, "exactly the one live arming fires");
1300 assert_eq!(
1306 recovered.redelivered, 300,
1307 "one redelivery per surviving fired row"
1308 );
1309 assert_eq!(
1310 recovered.retired, 300,
1311 "the terminal workflow's rows bulk-retire"
1312 );
1313 assert_eq!(recovered.superseded, 0);
1314 assert_eq!(recovered.retire_failures, 0);
1315 assert_eq!(recovered.skipped_orphans, 0);
1316 assert_eq!(
1317 recovered.fired
1318 + recovered.redelivered
1319 + recovered.retired
1320 + recovered.superseded
1321 + recovered.retire_failures
1322 + recovered.skipped_orphans,
1323 601,
1324 "the counters partition the due-row population exactly"
1325 );
1326 assert_eq!(
1327 count_timer_fired(&history(&inner, &running).await?, &live),
1328 1,
1329 "the live arming records its fire exactly once"
1330 );
1331 assert_eq!(
1337 engine.delivered_messages()?.len(),
1338 301,
1339 "one live-fire wake plus one owed-wake redelivery per surviving \
1340 fired row of the resident workflow"
1341 );
1342 assert_eq!(
1343 count_timer_fired(&history(&inner, &running).await?, &timer_id(0)),
1344 1,
1345 "a redelivered wake appends nothing: the recorder seam answers \
1346 AlreadyRecorded and the original fire stays the only record"
1347 );
1348 assert_eq!(
1359 counting.history_reads(),
1360 4,
1361 "the sweep's own history reads must scale with workflows (2) plus \
1362 the fire path's own reads (2 for 1 fire), never with the 601 rows"
1363 );
1364 assert!(
1365 inner.expired_timers(instant(20)).await?.is_empty(),
1366 "the sweep must leave the expired index EMPTY: every consumed row \
1367 retired in bulk and the fired arming retired by the fire path"
1368 );
1369 Ok(())
1370 }
1371
1372 #[tokio::test]
1376 async fn rows_with_no_history_are_skipped_and_survive() -> Result<(), TestError> {
1377 let (inner, _counting, engine, recovery) = counting_recovery();
1378 let orphaned = workflow_id();
1379 for i in 0..3 {
1380 inner
1381 .schedule_timer(&orphaned, &timer_id(i), instant(5), 1)
1382 .await?;
1383 }
1384
1385 let recovered = recovery.recover_on_startup(instant(20)).await?;
1386
1387 assert_eq!(
1388 recovered.fired, 0,
1389 "nothing fires for a workflow with no history"
1390 );
1391 assert_eq!(
1392 recovered.skipped_orphans, 3,
1393 "every orphaned row is COUNTED — the summary line's standing gauge \
1394 of the population the sweep cannot explain (round-3 N5)"
1395 );
1396 assert_eq!(
1397 inner.expired_timers(instant(20)).await?.len(),
1398 3,
1399 "orphaned rows survive the sweep — skip and count, never retire on \
1400 a history that answers nothing"
1401 );
1402 assert!(engine.delivered_messages()?.is_empty());
1403 Ok(())
1404 }
1405
1406 #[tokio::test]
1412 async fn a_surviving_recorded_fire_row_is_redelivered_once_then_retires()
1413 -> Result<(), TestError> {
1414 let process = WorkflowProcessHandle::new(42);
1415 let (store, engine, recovery) = recovery();
1416 let workflow_id = workflow_id();
1417 let timer_id = timer_id(7);
1418 engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
1419 engine.record_workflow_event(
1422 &workflow_id,
1423 timer_started_event(&workflow_id, &timer_id, 1, instant(5)),
1424 )?;
1425 engine
1426 .record_workflow_event(&workflow_id, timer_fired_event(&workflow_id, &timer_id, 2))?;
1427 store
1428 .schedule_timer(&workflow_id, &timer_id, instant(5), 1)
1429 .await?;
1430
1431 assert_eq!(
1432 recovery.recover_on_startup(instant(20)).await?.fired,
1433 0,
1434 "a redelivery is not a fire"
1435 );
1436 assert_eq!(
1437 engine.delivered_messages()?.len(),
1438 1,
1439 "the owed wake is delivered exactly once"
1440 );
1441 assert_eq!(
1442 count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
1443 1,
1444 "redelivery appends nothing — the recorder seam answers AlreadyRecorded"
1445 );
1446 assert!(
1447 store.expired_timers(instant(20)).await?.is_empty(),
1448 "the redelivered row retires — the wake's acknowledgement is durable now"
1449 );
1450
1451 assert_eq!(recovery.recover_on_startup(instant(20)).await?.fired, 0);
1453 assert_eq!(engine.delivered_messages()?.len(), 1);
1454 Ok(())
1455 }
1456
1457 #[tokio::test]
1462 async fn a_surviving_recorded_fire_row_of_a_nonresident_workflow_retires_silently()
1463 -> Result<(), TestError> {
1464 let (store, engine, recovery) = recovery();
1465 let workflow_id = workflow_id();
1466 let timer_id = timer_id(8);
1467 engine.record_workflow_event(
1468 &workflow_id,
1469 timer_started_event(&workflow_id, &timer_id, 1, instant(5)),
1470 )?;
1471 engine
1472 .record_workflow_event(&workflow_id, timer_fired_event(&workflow_id, &timer_id, 2))?;
1473 store
1474 .schedule_timer(&workflow_id, &timer_id, instant(5), 1)
1475 .await?;
1476
1477 assert_eq!(recovery.recover_on_startup(instant(20)).await?.fired, 0);
1478
1479 assert!(
1480 engine.delivered_messages()?.is_empty(),
1481 "no wake is owed to a non-resident workflow"
1482 );
1483 assert_eq!(
1484 count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
1485 1,
1486 "nothing is appended for a non-resident redelivery"
1487 );
1488 assert!(
1489 store.expired_timers(instant(20)).await?.is_empty(),
1490 "the consumed row retires instead of surviving every boot"
1491 );
1492 Ok(())
1493 }
1494
1495 #[tokio::test]
1503 async fn a_redelivery_for_a_rearmed_timer_appends_nothing_and_wakes_nobody()
1504 -> Result<(), TestError> {
1505 let concrete_store = Arc::new(InMemoryStore::default());
1506 let readable: Arc<dyn ReadableEventStore> = concrete_store.clone();
1507 let engine = Arc::new(FakeEngineHandle::new());
1511 let timer_service = Arc::new(TimerService::with_recorded_at(
1512 engine.clone(),
1513 readable.clone(),
1514 recorded_at,
1515 ));
1516 let recovery = TimerRecovery::new(readable, timer_service);
1517
1518 let workflow_id = workflow_id();
1519 let timer_id = timer_id(3);
1520 engine.set_residency(
1521 workflow_id.clone(),
1522 WorkflowResidency::Resident(WorkflowProcessHandle::new(7)),
1523 )?;
1524 engine.record_workflow_event(
1526 &workflow_id,
1527 timer_started_event(&workflow_id, &timer_id, 1, instant(5)),
1528 )?;
1529 engine
1530 .record_workflow_event(&workflow_id, timer_fired_event(&workflow_id, &timer_id, 2))?;
1531 engine.record_workflow_event(
1532 &workflow_id,
1533 timer_started_event(&workflow_id, &timer_id, 3, instant(5)),
1534 )?;
1535 concrete_store
1538 .append(
1539 WriteToken::recorder(),
1540 &workflow_id,
1541 &[
1542 timer_started_event(&workflow_id, &timer_id, 1, instant(5)),
1543 timer_fired_event(&workflow_id, &timer_id, 2),
1544 ],
1545 0,
1546 )
1547 .await?;
1548 concrete_store
1549 .schedule_timer(&workflow_id, &timer_id, instant(5), 1)
1550 .await?;
1551
1552 let recovered = recovery.recover_on_startup(instant(20)).await?;
1553
1554 assert_eq!(recovered.fired, 0, "a stale row is never a live fire");
1555 assert!(
1556 engine.delivered_messages()?.is_empty(),
1557 "no wake may reach a workflow that already ran past the recorded fire"
1558 );
1559 let recorded: Vec<Event> = engine
1560 .recorded_events()?
1561 .into_iter()
1562 .map(|(_, event)| event)
1563 .collect();
1564 assert_eq!(
1565 count_timer_fired(&recorded, &timer_id),
1566 1,
1567 "the redelivery must NOT mint a premature `TimerFired` for the \
1568 re-armed timer — the original fire stays the only record"
1569 );
1570 assert!(
1571 concrete_store.expired_timers(instant(20)).await?.is_empty(),
1572 "the OLD arming's row retires; nothing re-walks it"
1573 );
1574 Ok(())
1575 }
1576}