use std::sync::Arc;
use aion_core::{Event, EventEnvelope, TimerCancelCause, TimerId, WorkflowId};
use aion_store::{ReadableEventStore, StoreError, TimerRetirement};
use chrono::{DateTime, Utc};
use dashmap::DashSet;
use crate::engine_seam::{
EngineHandle, EngineSeamError, RecordOutcome, RedeliveredFire, TimerWheelEntry,
WorkflowMailboxMessage, WorkflowResidency,
};
use crate::time::deadline::{DeadlineHandler, deadline_run_id, is_deadline_timer};
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum RetireAttempt {
Retired,
Superseded,
Failed,
}
pub struct TimerService {
engine: Arc<dyn EngineHandle>,
store: Arc<dyn ReadableEventStore>,
recorded_at: fn() -> DateTime<Utc>,
terminal_updates: Arc<DashSet<(WorkflowId, TimerId)>>,
deadline_handler: Option<Arc<dyn DeadlineHandler>>,
}
struct TerminalUpdateSlot<'a> {
terminal_updates: &'a DashSet<(WorkflowId, TimerId)>,
key: (WorkflowId, TimerId),
}
impl Drop for TerminalUpdateSlot<'_> {
fn drop(&mut self) {
self.terminal_updates.remove(&self.key);
}
}
#[derive(thiserror::Error, Debug, Clone, PartialEq, Eq)]
pub enum TimerServiceError {
#[error("timer store operation failed: {0}")]
Store(#[from] StoreError),
#[error("timer engine operation failed: {0}")]
Engine(#[from] EngineSeamError),
#[error("deadline timer routing failed: {0}")]
Deadline(String),
}
impl TimerService {
#[must_use]
pub fn new(engine: Arc<dyn EngineHandle>, store: Arc<dyn ReadableEventStore>) -> Self {
Self::with_recorded_at(engine, store, Utc::now)
}
#[must_use]
pub fn with_recorded_at(
engine: Arc<dyn EngineHandle>,
store: Arc<dyn ReadableEventStore>,
recorded_at: fn() -> DateTime<Utc>,
) -> Self {
Self {
engine,
store,
recorded_at,
terminal_updates: Arc::new(DashSet::new()),
deadline_handler: None,
}
}
#[must_use]
pub fn with_terminal_updates(
mut self,
terminal_updates: Arc<DashSet<(WorkflowId, TimerId)>>,
) -> Self {
self.terminal_updates = terminal_updates;
self
}
#[must_use]
pub fn with_deadline_handler(mut self, handler: Arc<dyn DeadlineHandler>) -> Self {
self.deadline_handler = Some(handler);
self
}
pub async fn schedule(
&self,
workflow_id: WorkflowId,
timer_id: TimerId,
fire_at: DateTime<Utc>,
armed_seq: u64,
) -> Result<(), TimerServiceError> {
self.store
.schedule_timer(&workflow_id, &timer_id, fire_at, armed_seq)
.await?;
if let WorkflowResidency::Resident(process) = self.engine.resolve_workflow(&workflow_id)? {
self.engine.arm_timer(TimerWheelEntry {
process,
timer_id,
fire_at,
})?;
}
Ok(())
}
pub async fn cancel(
&self,
workflow_id: WorkflowId,
timer_id: TimerId,
cause: TimerCancelCause,
) -> Result<(), TimerServiceError> {
let key = (workflow_id.clone(), timer_id.clone());
let terminal_update_slot = self.wait_for_terminal_update_slot(key).await;
let result = self.cancel_guarded(workflow_id, timer_id, cause).await;
drop(terminal_update_slot);
result
}
async fn cancel_guarded(
&self,
workflow_id: WorkflowId,
timer_id: TimerId,
cause: TimerCancelCause,
) -> Result<(), TimerServiceError> {
let history = self.store.read_history(&workflow_id).await?;
if !matches!(
timer_disposition_in_active_segment(&history, &timer_id),
TimerDisposition::Live
) {
return Ok(());
}
let arming = last_recorded_arming(&history, &timer_id);
if let WorkflowResidency::Resident(process) = self.engine.resolve_workflow(&workflow_id)? {
self.engine.disarm_timer(process, &timer_id)?;
}
let event = Event::TimerCancelled {
envelope: self.next_envelope(&workflow_id).await?,
timer_id: timer_id.clone(),
cause,
};
self.engine.record_workflow_event(&workflow_id, event)?;
if let Some((fire_at, armed_seq)) = arming {
self.retire_consumed_row(&workflow_id, &timer_id, fire_at, armed_seq)
.await;
}
Ok(())
}
pub async fn fire_timer(
&self,
workflow_id: WorkflowId,
timer_id: TimerId,
fire_at: DateTime<Utc>,
) -> Result<(), TimerServiceError> {
let key = (workflow_id.clone(), timer_id.clone());
let terminal_update_slot = self.wait_for_terminal_update_slot(key).await;
let result = self
.fire_timer_guarded(workflow_id, timer_id, fire_at)
.await;
drop(terminal_update_slot);
result
}
async fn wait_for_terminal_update_slot(
&self,
key: (WorkflowId, TimerId),
) -> TerminalUpdateSlot<'_> {
loop {
if self.terminal_updates.insert(key.clone()) {
return TerminalUpdateSlot {
terminal_updates: self.terminal_updates.as_ref(),
key,
};
}
tokio::task::yield_now().await;
}
}
async fn fire_timer_guarded(
&self,
workflow_id: WorkflowId,
timer_id: TimerId,
fire_at: DateTime<Utc>,
) -> Result<(), TimerServiceError> {
let history = self.store.read_history(&workflow_id).await?;
let armed_seq = last_recorded_arming(&history, &timer_id).map_or(0, |(_, seq)| seq);
match timer_disposition_in_active_segment(&history, &timer_id) {
TimerDisposition::Live => {}
TimerDisposition::Fired if !is_deadline_timer(&timer_id) => {
return self
.redeliver_owed_wake(workflow_id, timer_id, fire_at, armed_seq)
.await
.map(|_| ());
}
TimerDisposition::Fired | TimerDisposition::Cancelled | TimerDisposition::Absent => {
self.retire_consumed_row(&workflow_id, &timer_id, fire_at, armed_seq)
.await;
return Ok(());
}
}
if is_deadline_timer(&timer_id) {
self.fire_deadline(workflow_id.clone(), timer_id.clone())
.await?;
self.retire_consumed_row(&workflow_id, &timer_id, fire_at, armed_seq)
.await;
return Ok(());
}
let event = Event::TimerFired {
envelope: self.next_envelope(&workflow_id).await?,
timer_id: timer_id.clone(),
};
match self.engine.record_workflow_event(&workflow_id, event)? {
RecordOutcome::RefusedTerminal | RecordOutcome::RefusedRetired => {
self.retire_consumed_row(&workflow_id, &timer_id, fire_at, armed_seq)
.await;
return Ok(());
}
RecordOutcome::Recorded | RecordOutcome::AlreadyRecorded => {}
}
if let WorkflowResidency::Resident(process) = self.engine.resolve_workflow(&workflow_id)? {
self.engine.deliver_workflow_message(
process,
WorkflowMailboxMessage::TimerFired {
timer_id: timer_id.clone(),
fire_at,
},
)?;
}
self.retire_consumed_row(&workflow_id, &timer_id, fire_at, armed_seq)
.await;
Ok(())
}
pub(crate) async fn redeliver_owed_wake(
&self,
workflow_id: WorkflowId,
timer_id: TimerId,
fire_at: DateTime<Utc>,
armed_seq: u64,
) -> Result<(bool, RetireAttempt), TimerServiceError> {
let WorkflowResidency::Resident(process) = self.engine.resolve_workflow(&workflow_id)?
else {
let row = self
.retire_consumed_row(&workflow_id, &timer_id, fire_at, armed_seq)
.await;
return Ok((false, row));
};
match self
.engine
.record_redelivered_timer_fire(&workflow_id, &timer_id)?
{
RedeliveredFire::RefusedTerminal | RedeliveredFire::NotOwed => {
let row = self
.retire_consumed_row(&workflow_id, &timer_id, fire_at, armed_seq)
.await;
Ok((false, row))
}
RedeliveredFire::WakeOwed => {
self.engine.deliver_workflow_message(
process,
WorkflowMailboxMessage::TimerFired {
timer_id: timer_id.clone(),
fire_at,
},
)?;
tracing::info!(
%workflow_id,
%timer_id,
"timer fire was already durably recorded; delivered the owed mailbox wake"
);
let row = self
.retire_consumed_row(&workflow_id, &timer_id, fire_at, armed_seq)
.await;
Ok((true, row))
}
}
}
async fn fire_deadline(
&self,
workflow_id: WorkflowId,
timer_id: TimerId,
) -> Result<(), TimerServiceError> {
let handler = self.deadline_handler.as_ref().ok_or_else(|| {
TimerServiceError::Deadline(format!(
"no deadline handler registered for {timer_id} on workflow {workflow_id}"
))
})?;
let run_id = deadline_run_id(&timer_id).ok_or_else(|| {
TimerServiceError::Deadline(format!(
"malformed deadline timer {timer_id} on workflow {workflow_id}"
))
})?;
handler
.on_deadline_elapsed(workflow_id, run_id)
.await
.map_err(|error| TimerServiceError::Deadline(error.to_string()))
}
pub(crate) async fn retire_consumed_row(
&self,
workflow_id: &WorkflowId,
timer_id: &TimerId,
fire_at: DateTime<Utc>,
armed_seq: u64,
) -> RetireAttempt {
match self
.store
.retire_timer(workflow_id, timer_id, fire_at, armed_seq)
.await
{
Ok(TimerRetirement::Retired) => RetireAttempt::Retired,
Ok(TimerRetirement::Superseded) => RetireAttempt::Superseded,
Err(error) => {
tracing::warn!(
%workflow_id,
%timer_id,
%fire_at,
%error,
"consumed timer row could not be retired; the row survives until a \
later fire or boot/adoption sweep retires it"
);
RetireAttempt::Failed
}
}
}
async fn next_envelope(&self, workflow_id: &WorkflowId) -> Result<EventEnvelope, StoreError> {
let history = self.store.read_history(workflow_id).await?;
let seq = history.iter().map(Event::seq).max().unwrap_or_default() + 1;
Ok(EventEnvelope {
seq,
recorded_at: (self.recorded_at)(),
workflow_id: workflow_id.clone(),
})
}
}
pub(crate) fn armed_fire_at_in_active_segment(
history: &[Event],
timer_id: &TimerId,
) -> Option<DateTime<Utc>> {
let mut armed = None;
for event in active_segment(history) {
match event {
Event::TimerStarted {
timer_id: id,
fire_at,
..
} if id == timer_id => {
armed = Some(*fire_at);
}
Event::TimerFired { timer_id: id, .. } | Event::TimerCancelled { timer_id: id, .. }
if id == timer_id =>
{
armed = None;
}
_ => {}
}
}
armed
}
fn last_recorded_arming(history: &[Event], timer_id: &TimerId) -> Option<(DateTime<Utc>, u64)> {
history.iter().rev().find_map(|event| match event {
Event::TimerStarted {
envelope,
timer_id: id,
fire_at,
} if id == timer_id => Some((*fire_at, envelope.seq)),
_ => None,
})
}
pub(crate) fn live_timers_in_active_segment(history: &[Event]) -> Vec<TimerId> {
let mut live: Vec<TimerId> = Vec::new();
for event in active_segment(history) {
match event {
Event::TimerStarted { timer_id, .. } if !live.contains(timer_id) => {
live.push(timer_id.clone());
}
Event::TimerFired { timer_id, .. } | Event::TimerCancelled { timer_id, .. } => {
live.retain(|id| id != timer_id);
}
_ => {}
}
}
live
}
fn active_segment(history: &[Event]) -> &[Event] {
let segment_start = history
.iter()
.rposition(|event| matches!(event, Event::WorkflowStarted { .. }))
.unwrap_or(0);
&history[segment_start..]
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum TimerDisposition {
Live,
Fired,
Cancelled,
Absent,
}
pub(crate) fn timer_disposition_in_active_segment(
history: &[Event],
timer_id: &TimerId,
) -> TimerDisposition {
let mut disposition = TimerDisposition::Absent;
for event in active_segment(history) {
match event {
Event::TimerStarted { timer_id: id, .. } if id == timer_id => {
disposition = TimerDisposition::Live;
}
Event::TimerFired { timer_id: id, .. } if id == timer_id => {
disposition = TimerDisposition::Fired;
}
Event::TimerCancelled { timer_id: id, .. } if id == timer_id => {
disposition = TimerDisposition::Cancelled;
}
_ => {}
}
}
disposition
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use aion_core::{Event, EventEnvelope, RunId, TimerCancelCause, TimerId, WorkflowId};
use aion_store::{InMemoryStore, ReadableEventStore, StoreError, WritableEventStore};
use chrono::{DateTime, Utc};
use super::{
TimerDisposition, TimerService, TimerServiceError, live_timers_in_active_segment,
timer_disposition_in_active_segment,
};
use crate::engine_seam::test_support::{
DeliveredWorkflowMessage, FakeEngineHandle, FakeEngineOperation,
};
use crate::engine_seam::{
EngineHandle, TimerWheelEntry, WorkflowProcessHandle, WorkflowResidency,
};
use crate::time::deadline::{DeadlineHandler, DeadlineHandlerError, deadline_timer_id};
fn instant(offset_seconds: i64) -> DateTime<Utc> {
DateTime::from_timestamp(1_700_000_000 + offset_seconds, 0).unwrap_or_default()
}
fn workflow_id() -> WorkflowId {
WorkflowId::new_v4()
}
fn timer_id() -> TimerId {
TimerId::anonymous(7)
}
fn service() -> (Arc<InMemoryStore>, Arc<FakeEngineHandle>, TimerService) {
let concrete_store = Arc::new(InMemoryStore::default());
let recorder_store: Arc<dyn WritableEventStore> = concrete_store.clone();
let readable_store: Arc<dyn ReadableEventStore> = concrete_store.clone();
let engine = Arc::new(FakeEngineHandle::recording_to(recorder_store));
let service = TimerService::with_recorded_at(engine.clone(), readable_store, recorded_at);
(concrete_store, engine, service)
}
fn recorded_at() -> DateTime<Utc> {
instant(1)
}
async fn history(
store: &InMemoryStore,
workflow_id: &WorkflowId,
) -> Result<Vec<Event>, StoreError> {
store.read_history(workflow_id).await
}
fn count_timer_fired(events: &[Event], timer_id: &TimerId) -> usize {
events
.iter()
.filter(|event| {
matches!(event, Event::TimerFired { timer_id: recorded, .. } if recorded == timer_id)
})
.count()
}
fn timer_started_event(workflow_id: &WorkflowId, timer_id: &TimerId, seq: u64) -> Event {
Event::TimerStarted {
envelope: EventEnvelope {
seq,
recorded_at: instant(0),
workflow_id: workflow_id.clone(),
},
timer_id: timer_id.clone(),
fire_at: instant(5),
}
}
fn workflow_started_event(workflow_id: &WorkflowId, seq: u64) -> Event {
Event::WorkflowStarted {
envelope: EventEnvelope {
seq,
recorded_at: instant(0),
workflow_id: workflow_id.clone(),
},
workflow_type: "fixture".to_owned(),
input: aion_core::Payload::new(aion_core::ContentType::Json, b"null".to_vec()),
run_id: aion_core::RunId::new_v4(),
parent_run_id: None,
parent_workflow_id: None,
package_version: aion_core::PackageVersion::new("a".repeat(64)),
}
}
fn timer_fired_event(workflow_id: &WorkflowId, timer_id: &TimerId, seq: u64) -> Event {
Event::TimerFired {
envelope: EventEnvelope {
seq,
recorded_at: instant(0),
workflow_id: workflow_id.clone(),
},
timer_id: timer_id.clone(),
}
}
fn timer_cancelled_event(workflow_id: &WorkflowId, timer_id: &TimerId, seq: u64) -> Event {
Event::TimerCancelled {
cause: TimerCancelCause::WorkflowIntent,
envelope: EventEnvelope {
seq,
recorded_at: instant(0),
workflow_id: workflow_id.clone(),
},
timer_id: timer_id.clone(),
}
}
fn make_named(name: &str) -> TimerId {
TimerId::named(name).unwrap_or_else(|_| TimerId::anonymous(0))
}
fn named_timer_id() -> TimerId {
make_named("review-deadline")
}
#[test]
fn started_timer_is_live() {
let workflow_id = workflow_id();
let timer_id = named_timer_id();
let history = vec![
workflow_started_event(&workflow_id, 0),
timer_started_event(&workflow_id, &timer_id, 1),
];
assert_eq!(live_timers_in_active_segment(&history), vec![timer_id]);
}
#[test]
fn started_then_fired_timer_is_dead() {
let workflow_id = workflow_id();
let timer_id = named_timer_id();
let history = vec![
workflow_started_event(&workflow_id, 0),
timer_started_event(&workflow_id, &timer_id, 1),
timer_fired_event(&workflow_id, &timer_id, 2),
];
assert!(live_timers_in_active_segment(&history).is_empty());
}
#[test]
fn started_then_cancelled_timer_is_dead() {
let workflow_id = workflow_id();
let timer_id = named_timer_id();
let history = vec![
workflow_started_event(&workflow_id, 0),
timer_started_event(&workflow_id, &timer_id, 1),
timer_cancelled_event(&workflow_id, &timer_id, 2),
];
assert!(live_timers_in_active_segment(&history).is_empty());
}
#[test]
fn restarted_named_timer_after_fire_is_live() {
let workflow_id = workflow_id();
let timer_id = named_timer_id();
let history = vec![
workflow_started_event(&workflow_id, 0),
timer_started_event(&workflow_id, &timer_id, 1),
timer_fired_event(&workflow_id, &timer_id, 2),
timer_started_event(&workflow_id, &timer_id, 3),
];
assert_eq!(
live_timers_in_active_segment(&history),
vec![timer_id],
"a re-armed named timer is live again"
);
}
#[test]
fn restarted_named_timer_after_cancel_is_live() {
let workflow_id = workflow_id();
let timer_id = named_timer_id();
let history = vec![
workflow_started_event(&workflow_id, 0),
timer_started_event(&workflow_id, &timer_id, 1),
timer_cancelled_event(&workflow_id, &timer_id, 2),
timer_started_event(&workflow_id, &timer_id, 3),
];
assert_eq!(live_timers_in_active_segment(&history), vec![timer_id]);
}
#[test]
fn prior_run_segment_timer_is_not_live() {
let workflow_id = workflow_id();
let prior = named_timer_id();
let current = make_named("current-deadline");
let history = vec![
workflow_started_event(&workflow_id, 0),
timer_started_event(&workflow_id, &prior, 1),
workflow_started_event(&workflow_id, 2),
timer_started_event(&workflow_id, ¤t, 3),
];
assert_eq!(live_timers_in_active_segment(&history), vec![current]);
}
#[test]
fn disposition_tracks_the_last_event_for_the_id_and_agrees_with_liveness() {
let workflow_id = workflow_id();
let timer_id = named_timer_id();
let other = make_named("unrelated");
let assert_agrees = |history: &[Event], expected: TimerDisposition| {
assert_eq!(
timer_disposition_in_active_segment(history, &timer_id),
expected
);
assert_eq!(
live_timers_in_active_segment(history).contains(&timer_id),
expected == TimerDisposition::Live,
"the per-id disposition and the enumerating model disagree on liveness"
);
};
let mut history = vec![
workflow_started_event(&workflow_id, 0),
timer_started_event(&workflow_id, &timer_id, 1),
];
assert_agrees(&history, TimerDisposition::Live);
history.push(timer_fired_event(&workflow_id, &timer_id, 2));
assert_agrees(&history, TimerDisposition::Fired);
history.push(timer_fired_event(&workflow_id, &other, 3));
assert_agrees(&history, TimerDisposition::Fired);
history.push(timer_started_event(&workflow_id, &timer_id, 4));
assert_agrees(&history, TimerDisposition::Live);
history.push(timer_cancelled_event(&workflow_id, &timer_id, 5));
assert_agrees(&history, TimerDisposition::Cancelled);
history.push(workflow_started_event(&workflow_id, 6));
assert_agrees(&history, TimerDisposition::Absent);
assert_agrees(&[], TimerDisposition::Absent);
}
#[tokio::test]
async fn re_armed_named_timer_fires_again() -> Result<(), TimerServiceError> {
let process = WorkflowProcessHandle::new(42);
let (store, engine, service) = service();
let workflow_id = workflow_id();
let timer_id = named_timer_id();
let fire_at = instant(110);
engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &timer_id, 1),
)?;
engine
.record_workflow_event(&workflow_id, timer_fired_event(&workflow_id, &timer_id, 2))?;
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &timer_id, 3),
)?;
service
.fire_timer(workflow_id.clone(), timer_id.clone(), fire_at)
.await?;
assert_eq!(
count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
2,
"the re-armed timer fires again, recording a second TimerFired"
);
assert_eq!(engine.delivered_messages()?.len(), 1);
Ok(())
}
#[tokio::test]
async fn schedule_records_timer_row_without_timer_started_event()
-> Result<(), TimerServiceError> {
let (store, _engine, service) = service();
let workflow_id = workflow_id();
let timer_id = timer_id();
let fire_at = instant(10);
service
.schedule(workflow_id.clone(), timer_id.clone(), fire_at, 1)
.await?;
let expired = store.expired_timers(fire_at).await?;
assert_eq!(expired.len(), 1);
assert_eq!(expired[0].workflow_id, workflow_id);
assert_eq!(expired[0].timer_id, timer_id);
assert_eq!(expired[0].fire_at, fire_at);
assert!(history(&store, &workflow_id).await?.is_empty());
Ok(())
}
#[tokio::test]
async fn schedule_arms_wheel_for_resident_workflow() -> Result<(), TimerServiceError> {
let process = WorkflowProcessHandle::new(42);
let (_store, engine, service) = service();
let workflow_id = workflow_id();
let timer_id = timer_id();
let fire_at = instant(20);
engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
service
.schedule(workflow_id, timer_id.clone(), fire_at, 1)
.await?;
assert_eq!(
engine.armed_timers()?,
vec![TimerWheelEntry {
process,
timer_id,
fire_at
}]
);
Ok(())
}
#[tokio::test]
async fn schedule_for_nonresident_records_without_arming() -> Result<(), TimerServiceError> {
let (store, engine, service) = service();
let workflow_id = workflow_id();
let timer_id = timer_id();
let fire_at = instant(30);
engine.set_residency(workflow_id.clone(), WorkflowResidency::NonResident)?;
service
.schedule(workflow_id.clone(), timer_id, fire_at, 1)
.await?;
assert!(engine.armed_timers()?.is_empty());
assert!(history(&store, &workflow_id).await?.is_empty());
Ok(())
}
#[tokio::test]
async fn fire_records_timer_fired_then_delivers_mailbox_message()
-> Result<(), TimerServiceError> {
let process = WorkflowProcessHandle::new(42);
let (store, engine, service) = service();
let workflow_id = workflow_id();
let timer_id = timer_id();
let fire_at = instant(40);
engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &timer_id, 1),
)?;
service
.fire_timer(workflow_id.clone(), timer_id.clone(), fire_at)
.await?;
assert_eq!(
count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
1
);
assert_eq!(
engine.delivered_messages()?,
vec![(
process,
DeliveredWorkflowMessage::TimerFired {
timer_id: timer_id.clone(),
fire_at
}
)]
);
assert!(matches!(
engine.operations()?.as_slice(),
[
FakeEngineOperation::EventRecorded {
event: Event::TimerStarted { .. },
..
},
FakeEngineOperation::EventRecorded {
workflow_id: recorded_workflow_id,
event: Event::TimerFired { timer_id: recorded_timer_id, .. },
},
FakeEngineOperation::Delivered {
process: delivered_process,
message: DeliveredWorkflowMessage::TimerFired { timer_id: delivered_timer_id, .. },
}
] if recorded_workflow_id == &workflow_id
&& recorded_timer_id == &timer_id
&& delivered_process == &process
&& delivered_timer_id == &timer_id
));
Ok(())
}
#[tokio::test]
async fn fire_records_without_delivery_when_workflow_becomes_nonresident()
-> Result<(), TimerServiceError> {
let (store, engine, service) = service();
let workflow_id = workflow_id();
let timer_id = timer_id();
let fire_at = instant(50);
engine.set_residency(workflow_id.clone(), WorkflowResidency::NonResident)?;
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &timer_id, 1),
)?;
service
.fire_timer(workflow_id.clone(), timer_id.clone(), fire_at)
.await?;
assert_eq!(
count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
1
);
assert!(engine.delivered_messages()?.is_empty());
Ok(())
}
#[tokio::test]
async fn firing_same_timer_twice_records_once_and_redelivers_the_wake()
-> Result<(), TimerServiceError> {
let process = WorkflowProcessHandle::new(42);
let (store, engine, service) = service();
let workflow_id = workflow_id();
let timer_id = timer_id();
let fire_at = instant(60);
engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &timer_id, 1),
)?;
service
.fire_timer(workflow_id.clone(), timer_id.clone(), fire_at)
.await?;
service
.fire_timer(workflow_id.clone(), timer_id.clone(), fire_at)
.await?;
assert_eq!(
count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
1,
"the recorded fire must never be appended a second time"
);
assert_eq!(
engine.delivered_messages()?.len(),
2,
"the second fire re-delivers the owed wake instead of silently no-opping"
);
Ok(())
}
#[tokio::test]
async fn already_recorded_fire_delivers_owed_wake_without_second_append()
-> Result<(), TimerServiceError> {
let process = WorkflowProcessHandle::new(42);
let (store, engine, service) = service();
let workflow_id = workflow_id();
let timer_id = timer_id();
let fire_at = instant(140);
engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &timer_id, 1),
)?;
engine
.record_workflow_event(&workflow_id, timer_fired_event(&workflow_id, &timer_id, 2))?;
service
.fire_timer(workflow_id.clone(), timer_id.clone(), fire_at)
.await?;
assert_eq!(
count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
1,
"the already-recorded fire must not be appended again"
);
assert_eq!(
engine.delivered_messages()?,
vec![(
process,
DeliveredWorkflowMessage::TimerFired { timer_id, fire_at }
)],
"the owed mailbox wake must be delivered"
);
Ok(())
}
#[tokio::test]
async fn already_recorded_fire_for_nonresident_workflow_wakes_nothing()
-> Result<(), TimerServiceError> {
let (store, engine, service) = service();
let workflow_id = workflow_id();
let timer_id = timer_id();
engine.set_residency(workflow_id.clone(), WorkflowResidency::NonResident)?;
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &timer_id, 1),
)?;
engine
.record_workflow_event(&workflow_id, timer_fired_event(&workflow_id, &timer_id, 2))?;
service
.fire_timer(workflow_id.clone(), timer_id.clone(), instant(150))
.await?;
assert_eq!(
count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
1,
"a non-resident redelivery must not re-enter the recorder seam"
);
assert!(
engine.delivered_messages()?.is_empty(),
"no wake is attempted for a non-resident workflow"
);
Ok(())
}
#[tokio::test]
async fn already_recorded_fire_after_run_terminal_delivers_no_wake()
-> Result<(), TimerServiceError> {
let process = WorkflowProcessHandle::new(42);
let (store, engine, service) = service();
let workflow_id = workflow_id();
let timer_id = timer_id();
engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &timer_id, 1),
)?;
engine
.record_workflow_event(&workflow_id, timer_fired_event(&workflow_id, &timer_id, 2))?;
engine.refuse_next_record_as_terminal()?;
service
.fire_timer(workflow_id.clone(), timer_id.clone(), instant(160))
.await?;
assert_eq!(
count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
1
);
assert!(
engine.delivered_messages()?.is_empty(),
"a post-terminal redelivery must not wake the terminated run"
);
Ok(())
}
#[tokio::test]
async fn firing_cancelled_timer_is_noop() -> Result<(), TimerServiceError> {
let process = WorkflowProcessHandle::new(42);
let (store, engine, service) = service();
let workflow_id = workflow_id();
let timer_id = timer_id();
let fire_at = instant(70);
engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &timer_id, 1),
)?;
let cancelled = Event::TimerCancelled {
cause: TimerCancelCause::WorkflowIntent,
envelope: EventEnvelope {
seq: 2,
recorded_at: instant(69),
workflow_id: workflow_id.clone(),
},
timer_id: timer_id.clone(),
};
engine.record_workflow_event(&workflow_id, cancelled)?;
service
.fire_timer(workflow_id.clone(), timer_id.clone(), fire_at)
.await?;
let history = history(&store, &workflow_id).await?;
assert_eq!(count_timer_fired(&history, &timer_id), 0);
assert!(engine.delivered_messages()?.is_empty());
Ok(())
}
#[tokio::test]
async fn fire_resolves_residency_at_fire_time() -> Result<(), TimerServiceError> {
let process = WorkflowProcessHandle::new(42);
let (store, engine, service) = service();
let workflow_id = workflow_id();
let timer_id = timer_id();
let fire_at = instant(80);
engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
engine.set_residency(workflow_id.clone(), WorkflowResidency::NonResident)?;
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &timer_id, 1),
)?;
service
.fire_timer(workflow_id.clone(), timer_id.clone(), fire_at)
.await?;
assert_eq!(
count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
1
);
assert!(engine.delivered_messages()?.is_empty());
Ok(())
}
#[tokio::test]
async fn firing_unstarted_timer_records_nothing() -> Result<(), TimerServiceError> {
let process = WorkflowProcessHandle::new(42);
let (store, engine, service) = service();
let workflow_id = workflow_id();
let timer_id = timer_id();
engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
service
.fire_timer(workflow_id.clone(), timer_id.clone(), instant(90))
.await?;
assert!(history(&store, &workflow_id).await?.is_empty());
assert!(engine.delivered_messages()?.is_empty());
Ok(())
}
struct RecordingDeadlineHandler {
calls: std::sync::Mutex<Vec<(WorkflowId, RunId)>>,
fail: bool,
}
impl RecordingDeadlineHandler {
fn new(fail: bool) -> Self {
Self {
calls: std::sync::Mutex::new(Vec::new()),
fail,
}
}
fn calls(&self) -> Result<Vec<(WorkflowId, RunId)>, TimerServiceError> {
self.calls
.lock()
.map(|calls| calls.clone())
.map_err(|error| TimerServiceError::Deadline(error.to_string()))
}
}
#[async_trait::async_trait]
impl DeadlineHandler for RecordingDeadlineHandler {
async fn on_deadline_elapsed(
&self,
workflow_id: WorkflowId,
run_id: RunId,
) -> Result<(), DeadlineHandlerError> {
self.calls
.lock()
.map_err(|error| DeadlineHandlerError(error.to_string()))?
.push((workflow_id, run_id));
if self.fail {
Err(DeadlineHandlerError(
"deliberate handler failure".to_owned(),
))
} else {
Ok(())
}
}
}
fn service_with_handler(
handler: Arc<dyn DeadlineHandler>,
) -> (Arc<InMemoryStore>, Arc<FakeEngineHandle>, TimerService) {
let concrete_store = Arc::new(InMemoryStore::default());
let recorder_store: Arc<dyn WritableEventStore> = concrete_store.clone();
let readable_store: Arc<dyn ReadableEventStore> = concrete_store.clone();
let engine = Arc::new(FakeEngineHandle::recording_to(recorder_store));
let service = TimerService::with_recorded_at(engine.clone(), readable_store, recorded_at)
.with_deadline_handler(handler);
(concrete_store, engine, service)
}
#[tokio::test]
async fn deadline_fire_routes_to_handler_and_records_no_timer_fired()
-> Result<(), TimerServiceError> {
let run_id = RunId::new_v4();
let deadline_id = deadline_timer_id(&run_id)
.map_err(|error| TimerServiceError::Deadline(error.to_string()))?;
let handler = Arc::new(RecordingDeadlineHandler::new(false));
let (store, engine, service) = service_with_handler(handler.clone());
let workflow_id = workflow_id();
let fire_at = instant(120);
engine.set_residency(
workflow_id.clone(),
WorkflowResidency::Resident(WorkflowProcessHandle::new(9)),
)?;
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &deadline_id, 1),
)?;
service
.fire_timer(workflow_id.clone(), deadline_id.clone(), fire_at)
.await?;
assert_eq!(handler.calls()?, vec![(workflow_id.clone(), run_id)]);
assert_eq!(
count_timer_fired(&history(&store, &workflow_id).await?, &deadline_id),
0,
"a deadline fire never records TimerFired"
);
assert!(engine.delivered_messages()?.is_empty());
Ok(())
}
#[tokio::test]
async fn deadline_fire_without_handler_is_typed_error() -> Result<(), TimerServiceError> {
let run_id = RunId::new_v4();
let deadline_id = deadline_timer_id(&run_id)
.map_err(|error| TimerServiceError::Deadline(error.to_string()))?;
let (store, engine, service) = service();
let workflow_id = workflow_id();
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &deadline_id, 1),
)?;
let result = service
.fire_timer(workflow_id.clone(), deadline_id.clone(), instant(120))
.await;
assert!(
matches!(result, Err(TimerServiceError::Deadline(_))),
"unhandled deadline fire must be a typed error, got {result:?}"
);
assert_eq!(
count_timer_fired(&history(&store, &workflow_id).await?, &deadline_id),
0
);
Ok(())
}
#[tokio::test]
async fn deadline_handler_failure_surfaces_as_typed_error() -> Result<(), TimerServiceError> {
let run_id = RunId::new_v4();
let deadline_id = deadline_timer_id(&run_id)
.map_err(|error| TimerServiceError::Deadline(error.to_string()))?;
let handler = Arc::new(RecordingDeadlineHandler::new(true));
let (_store, engine, service) = service_with_handler(handler);
let workflow_id = workflow_id();
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &deadline_id, 1),
)?;
let result = service
.fire_timer(workflow_id, deadline_id, instant(120))
.await;
assert!(matches!(result, Err(TimerServiceError::Deadline(_))));
Ok(())
}
#[tokio::test]
async fn refused_terminal_fire_records_nothing_and_delivers_no_wake()
-> Result<(), TimerServiceError> {
let process = WorkflowProcessHandle::new(42);
let (store, engine, service) = service();
let workflow_id = workflow_id();
let timer_id = timer_id();
engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &timer_id, 1),
)?;
engine.refuse_next_record_as_terminal()?;
service
.fire_timer(workflow_id.clone(), timer_id.clone(), instant(130))
.await?;
assert_eq!(
count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
0,
"a refused fire records no TimerFired"
);
assert!(
engine.delivered_messages()?.is_empty(),
"a refused fire delivers no wake"
);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn shared_coordinator_serializes_cancel_and_fire_across_services()
-> Result<(), TimerServiceError> {
use dashmap::DashSet;
use tokio::sync::Barrier;
for _ in 0..20 {
let process = WorkflowProcessHandle::new(42);
let concrete_store = Arc::new(InMemoryStore::default());
let recorder_store: Arc<dyn WritableEventStore> = concrete_store.clone();
let readable: Arc<dyn ReadableEventStore> = concrete_store.clone();
let engine = Arc::new(FakeEngineHandle::recording_to(recorder_store));
let coordinator = Arc::new(DashSet::new());
let service_a =
TimerService::with_recorded_at(engine.clone(), readable.clone(), recorded_at)
.with_terminal_updates(Arc::clone(&coordinator));
let service_b =
TimerService::with_recorded_at(engine.clone(), readable.clone(), recorded_at)
.with_terminal_updates(Arc::clone(&coordinator));
let workflow_id = workflow_id();
let timer_id = timer_id();
let fire_at = instant(200);
engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &timer_id, 1),
)?;
let gate = Arc::new(Barrier::new(2));
let (cancel_gate, fire_gate) = (Arc::clone(&gate), gate);
let (cancel_wf, cancel_timer) = (workflow_id.clone(), timer_id.clone());
let cancel = async move {
cancel_gate.wait().await;
service_a
.cancel(cancel_wf, cancel_timer, TimerCancelCause::WorkflowIntent)
.await
};
let (fire_wf, fire_timer) = (workflow_id.clone(), timer_id.clone());
let fire = async move {
fire_gate.wait().await;
service_b.fire_timer(fire_wf, fire_timer, fire_at).await
};
let (cancel_result, fire_result) = tokio::join!(cancel, fire);
cancel_result?;
fire_result?;
let history = history(&concrete_store, &workflow_id).await?;
let terminal_timer_events = history
.iter()
.filter(|event| {
matches!(
event,
Event::TimerFired { timer_id: recorded, .. }
| Event::TimerCancelled { timer_id: recorded, .. }
if recorded == &timer_id
)
})
.count();
assert_eq!(
terminal_timer_events, 1,
"first-recorded wins across shared services: {history:#?}"
);
}
Ok(())
}
#[tokio::test]
async fn firing_prior_run_timer_after_continue_as_new_is_noop() -> Result<(), TimerServiceError>
{
let process = WorkflowProcessHandle::new(42);
let (store, engine, service) = service();
let workflow_id = workflow_id();
let timer_id = timer_id();
engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &timer_id, 1),
)?;
engine.record_workflow_event(&workflow_id, workflow_started_event(&workflow_id, 2))?;
service
.fire_timer(workflow_id.clone(), timer_id.clone(), instant(100))
.await?;
assert_eq!(
count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
0
);
assert!(engine.delivered_messages()?.is_empty());
Ok(())
}
async fn outstanding_rows(
store: &InMemoryStore,
as_of: DateTime<Utc>,
) -> Result<usize, StoreError> {
Ok(store.expired_timers(as_of).await?.len())
}
#[tokio::test]
async fn a_recorded_fire_retires_the_consumed_row() -> Result<(), TimerServiceError> {
let process = WorkflowProcessHandle::new(42);
let (store, engine, service) = service();
let workflow_id = workflow_id();
let timer_id = timer_id();
let fire_at = instant(40);
engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &timer_id, 1),
)?;
service
.schedule(workflow_id.clone(), timer_id.clone(), fire_at, 1)
.await?;
assert_eq!(
outstanding_rows(&store, instant(1_000)).await?,
1,
"precondition: the arming's row is durable before the fire"
);
service
.fire_timer(workflow_id.clone(), timer_id.clone(), fire_at)
.await?;
assert_eq!(
count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
1,
"the fire itself must still record"
);
assert_eq!(
outstanding_rows(&store, instant(1_000)).await?,
0,
"a recorded fire must retire the consumed arming's row"
);
Ok(())
}
#[tokio::test]
async fn a_recorded_cancel_retires_the_consumed_row() -> Result<(), TimerServiceError> {
let process = WorkflowProcessHandle::new(42);
let (store, engine, service) = service();
let workflow_id = workflow_id();
let timer_id = named_timer_id();
engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &timer_id, 1),
)?;
service
.schedule(workflow_id.clone(), timer_id.clone(), instant(5), 1)
.await?;
assert_eq!(outstanding_rows(&store, instant(1_000)).await?, 1);
service
.cancel(
workflow_id.clone(),
timer_id.clone(),
TimerCancelCause::WorkflowIntent,
)
.await?;
assert_eq!(
outstanding_rows(&store, instant(1_000)).await?,
0,
"a recorded cancel must retire the cancelled arming's row"
);
Ok(())
}
#[tokio::test]
async fn a_stale_fire_leaves_a_re_armed_timers_row() -> Result<(), TimerServiceError> {
let process = WorkflowProcessHandle::new(42);
let (store, engine, service) = service();
let workflow_id = workflow_id();
let timer_id = named_timer_id();
let old_fire_at = instant(5);
let new_fire_at = instant(500);
engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &timer_id, 1),
)?;
service
.schedule(workflow_id.clone(), timer_id.clone(), new_fire_at, 2)
.await?;
service
.fire_timer(workflow_id.clone(), timer_id.clone(), old_fire_at)
.await?;
assert_eq!(
outstanding_rows(&store, instant(1_000)).await?,
1,
"the re-armed row is the replacement arming's only durable claim \
to a recovery fire; a stale retire must leave it standing"
);
Ok(())
}
#[tokio::test]
async fn a_terminal_refused_fire_retires_the_row() -> Result<(), TimerServiceError> {
let process = WorkflowProcessHandle::new(42);
let (store, engine, service) = service();
let workflow_id = workflow_id();
let timer_id = timer_id();
let fire_at = instant(40);
engine.set_residency(workflow_id.clone(), WorkflowResidency::Resident(process))?;
engine.record_workflow_event(
&workflow_id,
timer_started_event(&workflow_id, &timer_id, 1),
)?;
service
.schedule(workflow_id.clone(), timer_id.clone(), fire_at, 1)
.await?;
engine.refuse_next_record_as_terminal()?;
service
.fire_timer(workflow_id.clone(), timer_id.clone(), fire_at)
.await?;
assert_eq!(
count_timer_fired(&history(&store, &workflow_id).await?, &timer_id),
0,
"the refused fire must record nothing"
);
assert!(
engine.delivered_messages()?.is_empty(),
"a refused fire must wake nothing"
);
assert_eq!(
outstanding_rows(&store, instant(1_000)).await?,
0,
"a terminal-refused fire's arming is moot forever; its row retires"
);
Ok(())
}
}