use crate::core::{
CaseStatus, DeadlineState, Decision, InboundEvent, OnExpiry, RuntimeError, StoreError,
TaskState, Timestamp,
};
use crate::journal::{Append, RecordKind};
use super::ctx::Mode;
use super::executor::{LEASE_TTL, Runtime};
use super::metrics::{self, Census};
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct SweepReport {
pub warned: usize,
pub breached: usize,
pub tasks_expired: usize,
pub tasks_escalated: usize,
pub dead_lettered: usize,
pub timers_fired: usize,
pub census: Census,
}
impl SweepReport {
#[must_use]
pub fn needs_attention(&self) -> bool {
self.breached > 0 || self.tasks_expired > 0 || self.dead_lettered > 0
}
#[must_use]
pub fn is_quiet(&self) -> bool {
self.warned == 0
&& self.breached == 0
&& self.tasks_expired == 0
&& self.tasks_escalated == 0
&& self.dead_lettered == 0
&& self.timers_fired == 0
}
}
impl Runtime {
pub async fn sweep(
&self,
now: Timestamp,
event_grace: time::Duration,
) -> Result<SweepReport, RuntimeError> {
let mut report = SweepReport::default();
report.warned += self.sweep_deadlines(now, &mut report).await?;
self.sweep_tasks(now, &mut report).await?;
if self.events().is_some() {
report.dead_lettered = self.sweep_events(event_grace).await?;
}
if self.timers().is_some() {
report.timers_fired = self.fire_timers(now).await?;
}
report.census = self.census(now).await?;
report.census.emit();
Ok(report)
}
pub async fn census(&self, now: Timestamp) -> Result<Census, StoreError> {
let mut c = Census::default();
if let Some(cases) = self.cases() {
let cc = cases.census(now).await?;
c.open_cases = cc.open;
c.oldest_case_age_secs = cc.oldest_age_secs;
c.due_deadlines = cc.due;
}
if let Some(timers) = self.timers() {
c.pending_timers = timers.pending_count().await?;
}
if let Some(tasks) = self.tasks() {
c.open_tasks = tasks.open_count().await?;
}
Ok(c)
}
pub async fn fire_timers(&self, now: Timestamp) -> Result<usize, RuntimeError> {
let timers = self.timers().ok_or_else(|| {
RuntimeError::PlanContract(
"this runtime has no timer store — build it with `.timers(store)`".into(),
)
})?;
let due = timers
.claim_due(now, 128)
.await
.map_err(RuntimeError::from_store)?;
let mut fired = 0;
for timer in due {
let lease = self
.store()
.acquire(timer.run, self.owner(), LEASE_TTL)
.await
.map_err(RuntimeError::from_store)?;
let mut record = Append::new(
timer.run,
RecordKind::EffectDone {
output: serde_json::json!({ "fired_at": timer.fire_at.unix_timestamp() }),
spend: crate::core::Spend::default(),
},
)
.effect(timer.effect)
.step(timer.step)
.phase(timer.phase);
if let Some(c) = timer.case {
record = record.case(c);
}
self.store()
.append(lease.epoch, vec![record])
.await
.map_err(RuntimeError::from_store)?;
timers
.disarm(timer.run, timer.effect)
.await
.map_err(RuntimeError::from_store)?;
tracing::info!(
target: super::telemetry::TIMER_FIRED,
run = %timer.run,
step = %timer.step,
due = %timer.fire_at,
);
metrics::count(metrics::TIMERS_FIRED, "");
self.replay(timer.run, Mode::Resume).await?;
fired += 1;
}
Ok(fired)
}
async fn sweep_deadlines(
&self,
now: Timestamp,
report: &mut SweepReport,
) -> Result<usize, RuntimeError> {
let Some(cases) = self.cases() else {
return Ok(0);
};
let mut warned = 0;
for deadline in cases
.due(now, 512)
.await
.map_err(RuntimeError::from_store)?
{
if deadline.is_due(now) {
cases
.set_deadline_state(deadline.case, &deadline.name, DeadlineState::Breached)
.await
.map_err(RuntimeError::from_store)?;
cases
.set_status(deadline.case, CaseStatus::Escalated)
.await
.map_err(RuntimeError::from_store)?;
tracing::error!(
target: super::telemetry::DEADLINE_BREACHED,
case = %deadline.case,
obligation = %deadline.name,
due = %deadline.resolved_at,
);
metrics::count(metrics::DEADLINE_BREACHES, "");
report.breached += 1;
} else if deadline.needs_warning(now) {
cases
.set_deadline_state(deadline.case, &deadline.name, DeadlineState::Warned)
.await
.map_err(RuntimeError::from_store)?;
warned += 1;
}
}
Ok(warned)
}
async fn sweep_tasks(
&self,
now: Timestamp,
report: &mut SweepReport,
) -> Result<(), RuntimeError> {
let Some(tasks) = self.tasks() else {
return Ok(());
};
for task in tasks
.overdue(now, 512)
.await
.map_err(RuntimeError::from_store)?
{
match task.on_expiry {
OnExpiry::Escalate if task.state != TaskState::Escalated => {
tasks
.set_state(task.id, TaskState::Escalated)
.await
.map_err(RuntimeError::from_store)?;
report.tasks_escalated += 1;
}
OnExpiry::Escalate => {}
policy => {
let decision = Decision::expired(policy);
self.answer_task(task.id, &decision).await?;
tasks
.set_state(task.id, TaskState::Expired)
.await
.map_err(RuntimeError::from_store)?;
report.tasks_expired += 1;
}
}
}
Ok(())
}
pub async fn answer_task(
&self,
id: crate::core::TaskId,
decision: &Decision,
) -> Result<crate::core::Delivery, RuntimeError> {
let event = InboundEvent::new(
format!("task-decision:{}", id.to_hex()),
super::ctx::TASK_DECIDED,
serde_json::to_value(decision)?,
)
.correlate(crate::core::CorrelationKey::new("task", id.to_hex()));
self.deliver(&event).await
}
pub async fn decide_task(
&self,
id: crate::core::TaskId,
decision: &Decision,
roles: &[String],
) -> Result<crate::core::Delivery, RuntimeError> {
let tasks = self
.tasks()
.ok_or_else(|| RuntimeError::PlanContract("this runtime has no task store".into()))?;
tasks.claim(id, &decision.actor, roles).await.map_err(|e| {
RuntimeError::PolicyDenied(crate::core::PolicyError::Denied {
principal: decision.actor.clone(),
action: "task/decide".into(),
resource: format!("{id}: {e}"),
})
})?;
let delivery = self.answer_task(id, decision).await?;
tasks
.set_state(id, TaskState::Completed)
.await
.map_err(RuntimeError::from_store)?;
Ok(delivery)
}
}