use std::sync::atomic::{AtomicI64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use aion_core::{
AlarmCause, HealthSample, HealthStatus, InvariantAlarm, InvariantSpec, ToleranceSpec,
WorkflowId, WorkloopArming, WorkloopSpec,
};
use aion_store::InMemoryStore;
use aion_store::workloop::WorkloopStore;
use async_trait::async_trait;
use chrono::{DateTime, TimeZone, Utc};
use super::error::WorkloopError;
use super::service::{LoopEventSink, WorkloopService, WorkloopWaker};
use crate::engine_seam::RecordOutcome;
fn base() -> DateTime<Utc> {
Utc.with_ymd_and_hms(2026, 8, 25, 6, 0, 0)
.single()
.unwrap_or_default()
}
#[derive(Clone, Default)]
struct TestClock(Arc<AtomicI64>);
impl TestClock {
fn now_fn(&self) -> impl Fn() -> DateTime<Utc> + Send + Sync + 'static {
let offset = Arc::clone(&self.0);
move || base() + chrono::Duration::seconds(offset.load(Ordering::SeqCst))
}
fn set(&self, offset: i64) {
self.0.store(offset, Ordering::SeqCst);
}
}
#[derive(Default)]
struct FakeSink {
fires: Mutex<Vec<(WorkflowId, u64)>>,
alarms: Mutex<Vec<(WorkflowId, InvariantAlarm)>>,
refuse_cadence_as_terminal: std::sync::atomic::AtomicBool,
refuse_cadence_as_retired: std::sync::atomic::AtomicBool,
}
impl FakeSink {
fn fires(&self) -> Vec<(WorkflowId, u64)> {
self.fires
.lock()
.map(|fires| fires.clone())
.unwrap_or_default()
}
fn alarms(&self) -> Vec<(WorkflowId, InvariantAlarm)> {
self.alarms
.lock()
.map(|alarms| alarms.clone())
.unwrap_or_default()
}
}
#[async_trait]
impl LoopEventSink for FakeSink {
async fn record_cadence_fired(
&self,
loop_id: &WorkflowId,
window_seq: u64,
) -> Result<RecordOutcome, WorkloopError> {
if self.refuse_cadence_as_retired.load(Ordering::SeqCst) {
return Ok(RecordOutcome::RefusedRetired);
}
if self.refuse_cadence_as_terminal.load(Ordering::SeqCst) {
return Ok(RecordOutcome::RefusedTerminal);
}
self.fires
.lock()
.map_err(|_| WorkloopError::Engine {
reason: "fires lock poisoned".to_owned(),
})?
.push((loop_id.clone(), window_seq));
Ok(RecordOutcome::Recorded)
}
async fn record_invariant_unconfirmed(
&self,
loop_id: &WorkflowId,
alarm: InvariantAlarm,
) -> Result<(), WorkloopError> {
self.alarms
.lock()
.map_err(|_| WorkloopError::Engine {
reason: "alarms lock poisoned".to_owned(),
})?
.push((loop_id.clone(), alarm));
Ok(())
}
}
#[derive(Default)]
struct FakeWaker {
woken: Mutex<Vec<WorkflowId>>,
}
impl FakeWaker {
fn woken(&self) -> Vec<WorkflowId> {
self.woken
.lock()
.map(|woken| woken.clone())
.unwrap_or_default()
}
}
#[async_trait]
impl WorkloopWaker for FakeWaker {
async fn wake(&self, loop_id: &WorkflowId) -> Result<(), WorkloopError> {
self.woken
.lock()
.map_err(|_| WorkloopError::Engine {
reason: "woken lock poisoned".to_owned(),
})?
.push(loop_id.clone());
Ok(())
}
}
struct Rig {
service: WorkloopService,
store: Arc<InMemoryStore>,
sink: Arc<FakeSink>,
waker: Arc<FakeWaker>,
clock: TestClock,
}
fn rig() -> Result<Rig, WorkloopError> {
let store = Arc::new(InMemoryStore::default());
let sink = Arc::new(FakeSink::default());
let waker = Arc::new(FakeWaker::default());
let clock = TestClock::default();
let service = WorkloopService::with_clock(
store.clone(),
sink.clone(),
waker.clone(),
Duration::from_secs(1),
clock.now_fn(),
)?;
Ok(Rig {
service,
store,
sink,
waker,
clock,
})
}
fn cadence_spec(
period_secs: u64,
tolerance: ToleranceSpec,
) -> Result<WorkloopSpec, Box<dyn std::error::Error>> {
Ok(WorkloopSpec::new(
WorkloopArming::every(Duration::from_secs(period_secs))?,
vec![InvariantSpec {
name: String::from("serving"),
record_type: String::from("ServeState"),
tolerance,
confirms: vec![String::from("sweep")],
}],
Duration::from_secs(86_400),
)?)
}
fn confirmed_sample(window_seq: Option<u64>) -> HealthSample {
HealthSample {
invariant: String::from("serving"),
status: HealthStatus::Confirmed,
window_seq,
}
}
fn red_sample(window_seq: Option<u64>) -> HealthSample {
HealthSample {
invariant: String::from("serving"),
status: HealthStatus::Unconfirmed,
window_seq,
}
}
#[tokio::test]
async fn a_zero_sweep_interval_is_refused() -> Result<(), Box<dyn std::error::Error>> {
let store = Arc::new(InMemoryStore::default());
let refused = WorkloopService::new(
store.clone(),
Arc::new(FakeSink::default()),
Arc::new(FakeWaker::default()),
Duration::ZERO,
);
assert!(matches!(refused, Err(WorkloopError::ZeroSweepInterval)));
Ok(())
}
#[tokio::test]
async fn a_registered_loop_fires_on_its_window_and_wakes() -> Result<(), Box<dyn std::error::Error>>
{
let rig = rig()?;
let loop_id = WorkflowId::new_v4();
rig.service
.register(
loop_id.clone(),
String::from("default"),
cadence_spec(100, ToleranceSpec::count(3))?,
)
.await?;
rig.clock.set(99);
let report = rig.service.tick().await?;
assert_eq!(report.swept, 0);
assert!(rig.sink.fires().is_empty());
rig.clock.set(100);
let report = rig.service.tick().await?;
assert_eq!(report.swept, 1);
assert_eq!(report.fired, vec![(loop_id.clone(), 1)]);
assert_eq!(rig.sink.fires(), vec![(loop_id.clone(), 1)]);
assert_eq!(rig.waker.woken(), vec![loop_id.clone()]);
let record = rig
.store
.get_workloop(&loop_id)
.await?
.ok_or("record must persist")?;
assert_eq!(record.window_seq, 1);
assert_eq!(
record.next_window_at,
Some(base() + chrono::Duration::seconds(200))
);
Ok(())
}
#[tokio::test]
async fn missed_windows_exceed_count_tolerance_exactly_once_with_window_missed_cause()
-> Result<(), Box<dyn std::error::Error>> {
let rig = rig()?;
let loop_id = WorkflowId::new_v4();
rig.service
.register(
loop_id.clone(),
String::from("default"),
cadence_spec(100, ToleranceSpec::count(2))?,
)
.await?;
for window in 1..=3 {
rig.clock.set(window * 100);
let report = rig.service.tick().await?;
assert!(
report.alarms.is_empty(),
"window {window}: within tolerance must not alarm"
);
}
rig.clock.set(400);
let report = rig.service.tick().await?;
assert_eq!(report.alarms.len(), 1, "exceeding tolerance alarms");
let (alarmed_loop, alarm) = &report.alarms[0];
assert_eq!(alarmed_loop, &loop_id);
assert_eq!(alarm.cause, AlarmCause::WindowMissed);
assert_eq!(alarm.consecutive_unconfirmed, 3);
assert_eq!(alarm.last_confirmed_at, None);
assert_eq!(alarm.window_seq, Some(4));
rig.clock.set(500);
let report = rig.service.tick().await?;
assert!(report.alarms.is_empty(), "a latched alarm must not repeat");
assert_eq!(rig.sink.alarms().len(), 1);
Ok(())
}
#[tokio::test]
async fn a_closing_iteration_confirms_and_resets_the_deadman()
-> Result<(), Box<dyn std::error::Error>> {
let rig = rig()?;
let loop_id = WorkflowId::new_v4();
rig.service
.register(
loop_id.clone(),
String::from("default"),
cadence_spec(100, ToleranceSpec::count(0))?,
)
.await?;
rig.clock.set(100);
rig.service.tick().await?;
let alarms = rig
.service
.note_iteration_closed(&loop_id, &[confirmed_sample(Some(1))])
.await?;
assert!(alarms.is_empty());
rig.clock.set(200);
let report = rig.service.tick().await?;
assert!(report.alarms.is_empty());
let record = rig
.store
.get_workloop(&loop_id)
.await?
.ok_or("record must persist")?;
let health = record
.invariant_health
.get("serving")
.ok_or("health state must exist")?;
assert_eq!(health.consecutive_unconfirmed, 0);
assert_eq!(
health.last_confirmed_at,
Some(base() + chrono::Duration::seconds(100))
);
Ok(())
}
#[tokio::test]
async fn a_red_sample_exceeds_tolerance_zero_with_sample_red_cause()
-> Result<(), Box<dyn std::error::Error>> {
let rig = rig()?;
let loop_id = WorkflowId::new_v4();
rig.service
.register(
loop_id.clone(),
String::from("default"),
cadence_spec(100, ToleranceSpec::count(0))?,
)
.await?;
rig.clock.set(100);
rig.service.tick().await?;
let alarms = rig
.service
.note_iteration_closed(&loop_id, &[red_sample(Some(1))])
.await?;
assert_eq!(alarms.len(), 1);
assert_eq!(alarms[0].cause, AlarmCause::SampleRed);
assert_eq!(alarms[0].consecutive_unconfirmed, 1);
Ok(())
}
#[tokio::test]
async fn recovery_after_an_alarm_rearms_the_latch() -> Result<(), Box<dyn std::error::Error>> {
let rig = rig()?;
let loop_id = WorkflowId::new_v4();
rig.service
.register(
loop_id.clone(),
String::from("default"),
cadence_spec(100, ToleranceSpec::count(0))?,
)
.await?;
rig.clock.set(100);
rig.service.tick().await?;
let first = rig
.service
.note_iteration_closed(&loop_id, &[red_sample(Some(1))])
.await?;
assert_eq!(first.len(), 1);
let none = rig
.service
.note_iteration_closed(&loop_id, &[confirmed_sample(Some(1))])
.await?;
assert!(none.is_empty());
let second = rig
.service
.note_iteration_closed(&loop_id, &[red_sample(Some(1))])
.await?;
assert_eq!(second.len(), 1);
assert_eq!(rig.sink.alarms().len(), 2);
Ok(())
}
#[tokio::test]
async fn a_silent_signal_only_loop_alarms_unconfirmed_unknown_at_its_duration_deadline()
-> Result<(), Box<dyn std::error::Error>> {
let rig = rig()?;
let loop_id = WorkflowId::new_v4();
let spec = WorkloopSpec::new(
WorkloopArming::signal_only(vec![String::from("task_ready")])?,
vec![InvariantSpec {
name: String::from("serving"),
record_type: String::from("ServeState"),
tolerance: ToleranceSpec::duration(Duration::from_secs(300))?,
confirms: vec![String::from("sweep")],
}],
Duration::from_secs(86_400),
)?;
rig.service
.register(loop_id.clone(), String::from("default"), spec)
.await?;
rig.clock.set(300);
let report = rig.service.tick().await?;
assert!(report.alarms.is_empty());
assert!(
rig.sink.fires().is_empty(),
"no windows on a signal-only loop"
);
rig.clock.set(301);
let report = rig.service.tick().await?;
assert_eq!(report.alarms.len(), 1);
let alarm = &report.alarms[0].1;
assert_eq!(alarm.cause, AlarmCause::UnconfirmedUnknown);
assert_eq!(alarm.window_seq, None);
assert_eq!(alarm.consecutive_unconfirmed, 0);
let record = rig
.store
.get_workloop(&loop_id)
.await?
.ok_or("record must persist")?;
assert_eq!(record.next_check_at, None);
rig.clock.set(1000);
let report = rig.service.tick().await?;
assert_eq!(report.swept, 0);
Ok(())
}
#[tokio::test]
async fn a_terminal_loop_is_declared_dead_with_fanned_out_loop_dead_alarms()
-> Result<(), Box<dyn std::error::Error>> {
let rig = rig()?;
let loop_id = WorkflowId::new_v4();
let spec = WorkloopSpec::new(
WorkloopArming::every(Duration::from_secs(100))?,
vec![
InvariantSpec {
name: String::from("serving"),
record_type: String::from("ServeState"),
tolerance: ToleranceSpec::count(3),
confirms: vec![String::from("sweep")],
},
InvariantSpec {
name: String::from("drained"),
record_type: String::from("DrainState"),
tolerance: ToleranceSpec::count(3),
confirms: vec![String::from("drain")],
},
],
Duration::from_secs(86_400),
)?;
rig.service
.register(loop_id.clone(), String::from("default"), spec)
.await?;
rig.sink
.refuse_cadence_as_terminal
.store(true, Ordering::SeqCst);
rig.clock.set(100);
let report = rig.service.tick().await?;
assert_eq!(report.dead, vec![loop_id.clone()]);
assert_eq!(report.alarms.len(), 2);
assert!(
report
.alarms
.iter()
.all(|(_, alarm)| alarm.cause == AlarmCause::LoopDead)
);
assert!(rig.store.get_workloop(&loop_id).await?.is_none());
rig.clock.set(200);
let report = rig.service.tick().await?;
assert_eq!(report.swept, 0);
Ok(())
}
#[tokio::test]
async fn a_retired_loop_is_withdrawn_without_a_single_alarm()
-> Result<(), Box<dyn std::error::Error>> {
let rig = rig()?;
let loop_id = WorkflowId::new_v4();
let spec = WorkloopSpec::new(
WorkloopArming::every(Duration::from_secs(100))?,
vec![
InvariantSpec {
name: String::from("serving"),
record_type: String::from("ServeState"),
tolerance: ToleranceSpec::count(3),
confirms: vec![String::from("sweep")],
},
InvariantSpec {
name: String::from("drained"),
record_type: String::from("DrainState"),
tolerance: ToleranceSpec::count(3),
confirms: vec![String::from("drain")],
},
],
Duration::from_secs(86_400),
)?;
rig.service
.register(loop_id.clone(), String::from("default"), spec)
.await?;
rig.sink
.refuse_cadence_as_retired
.store(true, Ordering::SeqCst);
rig.clock.set(100);
let report = rig.service.tick().await?;
assert!(
report.dead.is_empty(),
"a declared retirement is not a death: {:?}",
report.dead
);
assert!(
report.alarms.is_empty(),
"a declared retirement must raise NO alarms: {:?}",
report.alarms
);
assert_eq!(
rig.sink.alarms(),
Vec::new(),
"and must append none to the loop's history either"
);
assert_eq!(report.retired, vec![loop_id.clone()]);
assert!(rig.store.get_workloop(&loop_id).await?.is_none());
rig.clock.set(200);
assert_eq!(rig.service.tick().await?.swept, 0);
Ok(())
}
#[tokio::test]
async fn boot_reconciliation_withdraws_only_registrations_with_no_history()
-> Result<(), Box<dyn std::error::Error>> {
use aion_core::{ContentType, Event, EventEnvelope, PackageVersion, Payload, RunId};
use aion_store::{EventStore, WriteToken};
let rig = rig()?;
let orphan = WorkflowId::new_v4();
let started = WorkflowId::new_v4();
let spec = WorkloopSpec::new(
WorkloopArming::every(Duration::from_secs(100))?,
vec![InvariantSpec {
name: String::from("serving"),
record_type: String::from("ServeState"),
tolerance: ToleranceSpec::count(3),
confirms: vec![String::from("sweep")],
}],
Duration::from_secs(86_400),
)?;
rig.service
.register(orphan.clone(), String::from("default"), spec.clone())
.await?;
rig.service
.register(started.clone(), String::from("default"), spec)
.await?;
let events: Arc<dyn EventStore> = Arc::clone(&rig.store) as Arc<dyn EventStore>;
events
.append(
WriteToken::recorder(),
&started,
&[Event::WorkflowStarted {
envelope: EventEnvelope {
seq: 1,
recorded_at: base(),
workflow_id: started.clone(),
},
workflow_type: String::from("queue_watch"),
input: Payload::new(ContentType::Json, b"{}".to_vec()),
run_id: RunId::new_v4(),
parent_run_id: None,
parent_workflow_id: None,
package_version: PackageVersion::new("a".repeat(64)),
}],
0,
)
.await?;
let withdrawn = super::service::withdraw_unstarted_registrations(
&(Arc::clone(&rig.store) as Arc<dyn WorkloopStore>),
&events,
)
.await?;
assert_eq!(
withdrawn,
vec![orphan.clone()],
"only the registration whose workflow has no history is withdrawn"
);
assert!(rig.store.get_workloop(&orphan).await?.is_none());
assert!(
rig.store.get_workloop(&started).await?.is_some(),
"a loop that really started must survive boot reconciliation"
);
Ok(())
}
#[tokio::test]
async fn a_late_sweep_fires_once_and_rearms_on_the_grid() -> Result<(), Box<dyn std::error::Error>>
{
let rig = rig()?;
let loop_id = WorkflowId::new_v4();
rig.service
.register(
loop_id.clone(),
String::from("default"),
cadence_spec(100, ToleranceSpec::count(10))?,
)
.await?;
rig.clock.set(350);
let report = rig.service.tick().await?;
assert_eq!(report.fired.len(), 1);
let record = rig
.store
.get_workloop(&loop_id)
.await?
.ok_or("record must persist")?;
assert_eq!(record.window_seq, 1);
assert_eq!(
record.next_window_at,
Some(base() + chrono::Duration::seconds(400))
);
Ok(())
}
#[tokio::test]
async fn duplicate_registration_is_refused_and_deregistration_stops_the_sweep()
-> Result<(), Box<dyn std::error::Error>> {
let rig = rig()?;
let loop_id = WorkflowId::new_v4();
let spec = cadence_spec(100, ToleranceSpec::count(3))?;
rig.service
.register(loop_id.clone(), String::from("default"), spec.clone())
.await?;
let duplicate = rig
.service
.register(loop_id.clone(), String::from("default"), spec)
.await;
assert!(matches!(
duplicate,
Err(WorkloopError::AlreadyRegistered { .. })
));
assert!(rig.service.deregister(&loop_id).await?);
rig.clock.set(100);
let report = rig.service.tick().await?;
assert_eq!(report.swept, 0);
assert!(rig.sink.fires().is_empty());
Ok(())
}
#[tokio::test]
async fn a_sample_for_an_undeclared_invariant_is_refused() -> Result<(), Box<dyn std::error::Error>>
{
let rig = rig()?;
let loop_id = WorkflowId::new_v4();
rig.service
.register(
loop_id.clone(),
String::from("default"),
cadence_spec(100, ToleranceSpec::count(3))?,
)
.await?;
let refused = rig
.service
.note_iteration_closed(
&loop_id,
&[HealthSample {
invariant: String::from("phantom"),
status: HealthStatus::Confirmed,
window_seq: Some(1),
}],
)
.await;
assert!(matches!(
refused,
Err(WorkloopError::UndeclaredInvariant { .. })
));
Ok(())
}