use std::sync::Arc;
use aion_core::{Event, EventEnvelope, TimerCancelCause, TimerId, WorkflowId};
use aion_store::{ReadableEventStore, StoreError};
use chrono::{DateTime, Utc};
use dashmap::DashSet;
use crate::engine_seam::{
EngineHandle, EngineSeamError, RecordOutcome, TimerWheelEntry, WorkflowMailboxMessage,
WorkflowResidency,
};
use crate::time::deadline::{DeadlineHandler, deadline_run_id, is_deadline_timer};
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>,
) -> Result<(), TimerServiceError> {
self.store
.schedule_timer(&workflow_id, &timer_id, fire_at)
.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> {
if !self.timer_is_live(&workflow_id, &timer_id).await? {
return Ok(());
}
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,
cause,
};
self.engine.record_workflow_event(&workflow_id, event)?;
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> {
if !self.timer_is_live(&workflow_id, &timer_id).await? {
return Ok(());
}
if is_deadline_timer(&timer_id) {
return self.fire_deadline(workflow_id, timer_id).await;
}
let event = Event::TimerFired {
envelope: self.next_envelope(&workflow_id).await?,
timer_id: timer_id.clone(),
};
if self.engine.record_workflow_event(&workflow_id, event)? == RecordOutcome::RefusedTerminal
{
return Ok(());
}
if let WorkflowResidency::Resident(process) = self.engine.resolve_workflow(&workflow_id)? {
self.engine.deliver_workflow_message(
process,
WorkflowMailboxMessage::TimerFired { timer_id, fire_at },
)?;
}
Ok(())
}
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()))
}
async fn timer_is_live(
&self,
workflow_id: &WorkflowId,
timer_id: &TimerId,
) -> Result<bool, StoreError> {
let history = self.store.read_history(workflow_id).await?;
Ok(live_timers_in_active_segment(&history).contains(timer_id))
}
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 live_timers_in_active_segment(history: &[Event]) -> Vec<TimerId> {
let segment_start = history
.iter()
.rposition(|event| matches!(event, Event::WorkflowStarted { .. }))
.unwrap_or(0);
let mut live: Vec<TimerId> = Vec::new();
for event in &history[segment_start..] {
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
}
#[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::{TimerService, TimerServiceError, live_timers_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,
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]);
}
#[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)
.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)
.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)
.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_and_delivers_once() -> 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
);
assert_eq!(engine.delivered_messages()?.len(), 1);
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(())
}
}