use std::sync::Arc;
use std::time::Duration;
use aion_core::{HealthSample, HealthStatus, InvariantAlarm, WorkflowId, WorkloopSpec};
use aion_store::workloop::{WorkloopRecord, WorkloopStore};
use async_trait::async_trait;
use chrono::{DateTime, Utc};
use super::error::WorkloopError;
use super::health::{
UnconfirmedEvidence, latch_alarm, observe_confirmed, observe_unconfirmed, peek_alarm,
};
use super::windows::{advance_window, initial_window, next_check_at};
use crate::engine_seam::RecordOutcome;
#[async_trait]
pub trait LoopEventSink: Send + Sync {
async fn record_cadence_fired(
&self,
loop_id: &WorkflowId,
window_seq: u64,
) -> Result<RecordOutcome, WorkloopError>;
async fn record_invariant_unconfirmed(
&self,
loop_id: &WorkflowId,
alarm: InvariantAlarm,
) -> Result<(), WorkloopError>;
}
#[async_trait]
pub trait WorkloopWaker: Send + Sync {
async fn wake(&self, loop_id: &WorkflowId) -> Result<(), WorkloopError>;
}
#[derive(Debug, Default)]
pub struct SweepReport {
pub swept: usize,
pub fired: Vec<(WorkflowId, u64)>,
pub alarms: Vec<(WorkflowId, InvariantAlarm)>,
pub dead: Vec<WorkflowId>,
pub retired: Vec<WorkflowId>,
pub faults: Vec<(WorkflowId, String)>,
}
pub struct WorkloopService {
store: Arc<dyn WorkloopStore>,
sink: Arc<dyn LoopEventSink>,
waker: Arc<dyn WorkloopWaker>,
sweep_interval: Duration,
now: Arc<dyn Fn() -> DateTime<Utc> + Send + Sync>,
}
impl WorkloopService {
pub fn new(
store: Arc<dyn WorkloopStore>,
sink: Arc<dyn LoopEventSink>,
waker: Arc<dyn WorkloopWaker>,
sweep_interval: Duration,
) -> Result<Self, WorkloopError> {
Self::with_clock(store, sink, waker, sweep_interval, Utc::now)
}
pub fn with_clock(
store: Arc<dyn WorkloopStore>,
sink: Arc<dyn LoopEventSink>,
waker: Arc<dyn WorkloopWaker>,
sweep_interval: Duration,
now: impl Fn() -> DateTime<Utc> + Send + Sync + 'static,
) -> Result<Self, WorkloopError> {
if sweep_interval.is_zero() {
return Err(WorkloopError::ZeroSweepInterval);
}
Ok(Self {
store,
sink,
waker,
sweep_interval,
now: Arc::new(now),
})
}
#[must_use]
pub const fn sweep_interval(&self) -> Duration {
self.sweep_interval
}
pub async fn register(
&self,
loop_id: WorkflowId,
namespace: String,
spec: WorkloopSpec,
) -> Result<WorkloopRecord, WorkloopError> {
if self.store.get_workloop(&loop_id).await?.is_some() {
return Err(WorkloopError::AlreadyRegistered { loop_id });
}
let now = (self.now)();
let next_window_at = match spec.arming().cadence_period() {
Some(period) => Some(initial_window(now, period)?),
None => None,
};
let mut record = WorkloopRecord {
loop_id,
namespace,
spec,
window_seq: 0,
next_window_at,
next_check_at: None,
last_iteration_closed_window: None,
invariant_health: std::collections::BTreeMap::new(),
registered_at: now,
updated_at: now,
};
record.next_check_at = next_check_at(&record);
self.store.put_workloop(record.clone()).await?;
Ok(record)
}
pub async fn wake_now(&self, loop_id: &WorkflowId) -> Result<(), WorkloopError> {
self.waker.wake(loop_id).await
}
pub async fn deregister(&self, loop_id: &WorkflowId) -> Result<bool, WorkloopError> {
Ok(self.store.remove_workloop(loop_id).await?)
}
pub async fn note_iteration_closed(
&self,
loop_id: &WorkflowId,
samples: &[HealthSample],
) -> Result<Vec<InvariantAlarm>, WorkloopError> {
let mut record = self.store.get_workloop(loop_id).await?.ok_or_else(|| {
WorkloopError::NotRegistered {
loop_id: loop_id.clone(),
}
})?;
for sample in samples {
if !record
.spec
.invariants()
.iter()
.any(|invariant| invariant.name == sample.invariant)
{
return Err(WorkloopError::UndeclaredInvariant {
loop_id: loop_id.clone(),
invariant: sample.invariant.clone(),
});
}
}
let now = (self.now)();
for sample in samples {
let state = record.invariant_health.entry(sample.invariant.clone());
let state = state.or_default();
match sample.status {
HealthStatus::Confirmed => observe_confirmed(state, now),
HealthStatus::Unconfirmed => {
observe_unconfirmed(state, UnconfirmedEvidence::SampleRed);
}
}
}
record.last_iteration_closed_window = Some(record.window_seq);
let mut alarms = Vec::new();
let window_ctx = window_context(&record);
let anchor = record.registered_at;
for invariant in record.spec.invariants().to_vec() {
let Some(state) = record.invariant_health.get_mut(&invariant.name) else {
continue;
};
if let Some(alarm) = peek_alarm(
&invariant.name,
&invariant.tolerance,
state,
anchor,
now,
window_ctx,
) {
self.sink
.record_invariant_unconfirmed(loop_id, alarm.clone())
.await?;
latch_alarm(state);
alarms.push(alarm);
}
}
record.next_check_at = next_check_at(&record);
record.updated_at = now;
self.store.put_workloop(record).await?;
Ok(alarms)
}
pub async fn tick(&self) -> Result<SweepReport, WorkloopError> {
let now = (self.now)();
let mut report = SweepReport::default();
for record in self.store.due_workloops(now).await? {
let loop_id = record.loop_id.clone();
if let Err(error) = self.sweep_loop(record, now, &mut report).await {
report.faults.push((loop_id, error.to_string()));
}
report.swept += 1;
}
Ok(report)
}
async fn sweep_loop(
&self,
mut record: WorkloopRecord,
now: DateTime<Utc>,
report: &mut SweepReport,
) -> Result<(), WorkloopError> {
let loop_id = record.loop_id.clone();
let mut wake_pending = false;
if let (Some(period), Some(window_at)) =
(record.spec.arming().cadence_period(), record.next_window_at)
&& window_at <= now
{
if record.window_seq >= 1
&& record.last_iteration_closed_window != Some(record.window_seq)
{
for invariant in record.spec.invariants().to_vec() {
let state = record
.invariant_health
.entry(invariant.name.clone())
.or_default();
observe_unconfirmed(state, UnconfirmedEvidence::WindowMissed);
}
}
let window_seq = record.window_seq.saturating_add(1);
match self.sink.record_cadence_fired(&loop_id, window_seq).await? {
RecordOutcome::Recorded | RecordOutcome::AlreadyRecorded => {
record.window_seq = window_seq;
record.next_window_at = Some(advance_window(window_at, period, now)?);
report.fired.push((loop_id.clone(), window_seq));
wake_pending = true;
}
RecordOutcome::RefusedTerminal => {
return self.declare_loop_dead(record, report).await;
}
RecordOutcome::RefusedRetired => {
let removed = self.store.remove_workloop(&loop_id).await?;
tracing::info!(
%loop_id,
row_existed = removed,
"sweep found a retired workloop still on the sweep set and withdrew it; \
a declared retirement raises no alarms"
);
report.retired.push(loop_id);
return Ok(());
}
}
}
let window_ctx = window_context(&record);
let anchor = record.registered_at;
for invariant in record.spec.invariants().to_vec() {
let state = record
.invariant_health
.entry(invariant.name.clone())
.or_default();
if let Some(alarm) = peek_alarm(
&invariant.name,
&invariant.tolerance,
state,
anchor,
now,
window_ctx,
) {
self.sink
.record_invariant_unconfirmed(&loop_id, alarm.clone())
.await?;
latch_alarm(state);
report.alarms.push((loop_id.clone(), alarm));
}
}
record.next_check_at = next_check_at(&record);
record.updated_at = now;
self.store.put_workloop(record).await?;
if wake_pending && let Err(error) = self.waker.wake(&loop_id).await {
report
.faults
.push((loop_id, format!("wake failed: {error}")));
}
Ok(())
}
async fn declare_loop_dead(
&self,
record: WorkloopRecord,
report: &mut SweepReport,
) -> Result<(), WorkloopError> {
let loop_id = record.loop_id.clone();
let window_ctx = window_context(&record);
for invariant in record.spec.invariants() {
let state = record
.invariant_health
.get(&invariant.name)
.cloned()
.unwrap_or_default();
let alarm = InvariantAlarm {
invariant: invariant.name.clone(),
cause: aion_core::AlarmCause::LoopDead,
window_seq: window_ctx,
last_confirmed_at: state.last_confirmed_at,
consecutive_unconfirmed: state.consecutive_unconfirmed,
};
self.sink
.record_invariant_unconfirmed(&loop_id, alarm.clone())
.await?;
report.alarms.push((loop_id.clone(), alarm));
}
self.store.remove_workloop(&loop_id).await?;
report.dead.push(loop_id);
Ok(())
}
pub async fn run(self: Arc<Self>, mut shutdown: tokio::sync::watch::Receiver<bool>) {
let mut interval = tokio::time::interval(self.sweep_interval);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
tokio::select! {
_ = interval.tick() => {
match self.tick().await {
Ok(report) => {
for (loop_id, fault) in &report.faults {
tracing::warn!(%loop_id, fault, "workloop sweep fault");
}
}
Err(error) => {
tracing::error!(%error, "workloop sweep pass failed");
}
}
}
changed = shutdown.changed() => {
if changed.is_err() || *shutdown.borrow() {
break;
}
}
}
}
}
}
pub async fn withdraw_unstarted_registrations(
workloop_store: &Arc<dyn WorkloopStore>,
store: &Arc<dyn aion_store::EventStore>,
) -> Result<Vec<WorkflowId>, WorkloopError> {
let listing = workloop_store.list_workloops().await?;
let mut withdrawn = Vec::new();
for record in listing.workloops {
let loop_id = record.loop_id;
if !store.read_history(&loop_id).await?.is_empty() {
continue;
}
let existed = workloop_store.remove_workloop(&loop_id).await?;
tracing::warn!(
%loop_id,
row_existed = existed,
"withdrawing a workloop registration whose workflow has no recorded history: its \
start never landed, so the sweep set carried a row for a loop that never ran and \
the dead-man switch would have alarmed on it forever"
);
withdrawn.push(loop_id);
}
for row in listing.undecodable {
tracing::error!(
loop_id = %row.loop_id,
"workloop registration row does not decode; it is left in place rather than \
withdrawn, because an unreadable row is not evidence that its loop never started"
);
}
Ok(withdrawn)
}
pub(crate) fn window_context(record: &WorkloopRecord) -> Option<u64> {
(record.spec.arming().cadence_period().is_some() && record.window_seq >= 1)
.then_some(record.window_seq)
}