use std::sync::Arc;
use crate::core::{
CaseId, CaseStatus, DeadlineState, Decision, InboundEvent, OnExpiry, RunId, RuntimeError,
StoreError, SweptAction, TaskState, Timestamp,
};
use crate::journal::{Append, JournalStore, RecordKind};
use super::executor::{LEASE_TTL, Runtime};
use super::metrics::{self, Census};
pub const SOURCE_WORKLIST: &str = "agentplane://worklist";
const SWEEP_EPOCH: crate::core::Epoch = 1;
pub(crate) const SWEEP_OUTCOME: &str = "swept";
struct SweepLedger {
run: Option<RunId>,
wrote: bool,
}
impl SweepLedger {
const fn new() -> Self {
Self {
run: None,
wrote: false,
}
}
async fn note(
&mut self,
store: &Arc<dyn JournalStore>,
case: Option<CaseId>,
subject: String,
action: SweptAction,
detail: Option<String>,
) -> Result<(), RuntimeError> {
let run = *self.run.get_or_insert_with(RunId::generate);
let mut entry = crate::journal::Append::new(
run,
RecordKind::Swept {
subject,
action,
detail,
},
);
if let Some(case) = case {
entry = entry.case(case);
}
store
.append(SWEEP_EPOCH, vec![entry])
.await
.map_err(RuntimeError::from_store)?;
self.wrote = true;
Ok(())
}
async fn seal(self, store: &Arc<dyn JournalStore>) -> SweepRecord {
let Some(run) = self.run else {
return SweepRecord::Quiet;
};
if !self.wrote {
return SweepRecord::Quiet;
}
let head = match store.head(run).await {
Ok(head) => head,
Err(e) => {
tracing::error!(%run, error = %e, "a sweep's record could not be read back");
return SweepRecord::EvidenceLost;
}
};
let sealed = crate::journal::Append::new(
run,
RecordKind::RunConcluded {
outcome: SWEEP_OUTCOME.to_owned(),
reason: None,
exhaustion: None,
live_spend: crate::core::Spend::default(),
chain_head: head.hash,
},
);
if let Err(e) = store.append(SWEEP_EPOCH, vec![sealed]).await {
tracing::error!(%run, error = %e, "a sweep's record could not be closed");
return SweepRecord::EvidenceLost;
}
if let Err(e) = store.seal(run, SWEEP_EPOCH, SWEEP_OUTCOME).await {
tracing::error!(%run, error = %e, "a sweep's record could not be sealed");
return SweepRecord::EvidenceLost;
}
SweepRecord::Recorded(run)
}
}
enum SweepRecord {
Quiet,
Recorded(RunId),
EvidenceLost,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct WokenRuns {
pub fired: usize,
pub failed: usize,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
#[allow(clippy::struct_excessive_bools)]
pub struct Saturation {
pub timers: bool,
pub deadlines: bool,
pub tasks: bool,
pub recovery: bool,
}
impl Saturation {
#[must_use]
pub const fn any(self) -> bool {
self.timers || self.deadlines || self.tasks || self.recovery
}
}
const TIMER_BATCH: usize = 128;
const DEADLINE_BATCH: usize = 512;
const TASK_BATCH: usize = 512;
const RECOVERY_BATCH: usize = 32;
const EVENT_BATCH: usize = 128;
#[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 events_redelivered: usize,
pub timers_fired: usize,
pub wake_failures: usize,
pub runs_recovered: usize,
pub recovery_failures: usize,
pub saturated: Saturation,
pub record: Option<RunId>,
pub evidence_lost: bool,
pub census: Census,
pub census_unavailable: bool,
}
impl SweepReport {
#[must_use]
pub fn needs_attention(&self) -> bool {
self.breached > 0
|| self.tasks_expired > 0
|| self.dead_lettered > 0
|| self.saturated.any()
|| self.evidence_lost
|| self.recovery_failures > 0
|| self.wake_failures > 0
|| self.census_unavailable
}
#[must_use]
pub fn is_quiet(&self) -> bool {
!self.saturated.any()
&& self.warned == 0
&& self.breached == 0
&& self.tasks_expired == 0
&& self.tasks_escalated == 0
&& self.dead_lettered == 0
&& self.events_redelivered == 0
&& self.timers_fired == 0
&& self.wake_failures == 0
&& self.runs_recovered == 0
&& self.recovery_failures == 0
}
}
impl Runtime {
pub async fn sweep(
&self,
now: Timestamp,
event_grace: std::time::Duration,
) -> Result<SweepReport, RuntimeError> {
let mut report = SweepReport::default();
let mut ledger = SweepLedger::new();
let phases: Result<(), RuntimeError> = async {
self.recover_abandoned(&mut report, &mut ledger).await?;
report.warned += self.sweep_deadlines(now, &mut report, &mut ledger).await?;
self.sweep_tasks(now, &mut report, &mut ledger).await?;
if self.events().is_some() {
report.events_redelivered = self.redeliver_claimed(EVENT_BATCH).await?;
report.dead_lettered = self.sweep_events(event_grace).await?;
}
if self.timers().is_some() {
let woken = self.fire_timers(now).await?;
report.timers_fired = woken.fired;
report.wake_failures = woken.failed;
if woken.fired + woken.failed >= TIMER_BATCH {
report.saturated.timers = true;
}
}
Ok(())
}
.await;
match ledger.seal(self.store()).await {
SweepRecord::Quiet => {}
SweepRecord::Recorded(run) => report.record = Some(run),
SweepRecord::EvidenceLost => report.evidence_lost = true,
}
phases?;
match self.census(now).await {
Ok(census) => {
report.census = census;
self.meter().census(&report.census);
}
Err(error) => {
report.census_unavailable = true;
tracing::error!(
%error,
"the census could not be read; this report carries the tick's \
counters and no gauges"
);
}
}
Ok(report)
}
async fn recover_abandoned(
&self,
report: &mut SweepReport,
ledger: &mut SweepLedger,
) -> Result<(), RuntimeError> {
let stranded = self
.store()
.abandoned_runs(RECOVERY_BATCH)
.await
.map_err(RuntimeError::from_store)?;
if stranded.len() >= RECOVERY_BATCH {
report.saturated.recovery = true;
}
for run in stranded {
ledger
.note(
self.store(),
None,
run.to_string(),
SweptAction::RunRecovered,
Some(
"its owner's lease lapsed without release; taking the run \
over and resuming it"
.to_owned(),
),
)
.await?;
match self.recover_abandoned_run(run).await {
Ok(outcome) => {
if let Ok(lease) = self.store().acquire(run, self.owner_id(), LEASE_TTL).await {
let _ = self.store().release_lease(run, lease.epoch).await;
}
let status = outcome.status.as_str().to_owned();
tracing::info!(
target: super::telemetry::RUN_RECOVERED,
%run,
outcome = %status,
);
self.meter().count(metrics::RUNS_RECOVERED, &status);
report.runs_recovered += 1;
}
Err(RuntimeError::Store(StoreError::LeaseHeld { .. })) => {}
Err(RuntimeError::Store(StoreError::NotFound(_))) => {
ledger
.note(
self.store(),
None,
run.to_string(),
SweptAction::RunRecovered,
Some(
"its owner died between acquiring the lease and the \
admission append; no run exists, so clearing the \
lease is the whole recovery"
.to_owned(),
),
)
.await?;
if let Ok(lease) = self.store().acquire(run, self.owner_id(), LEASE_TTL).await {
if let Err(error) = self.release_empty_quota_reservation(run).await {
tracing::warn!(%run, %error, "could not release the quota slot of an empty abandoned admission");
}
let _ = self.store().release_lease(run, lease.epoch).await;
}
report.runs_recovered += 1;
}
Err(error) => {
tracing::error!(%run, %error, "an abandoned run could not be resumed");
self.meter().count(metrics::RECOVERY_FAILURES, "");
report.recovery_failures += 1;
}
}
}
Ok(())
}
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<WokenRuns, 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, TIMER_BATCH)
.await
.map_err(RuntimeError::from_store)?;
let mut woken = WokenRuns::default();
for timer in due {
match self.fire_one(&timer).await {
Ok(()) => woken.fired += 1,
Err(error) => {
tracing::error!(
run = %timer.run,
step = %timer.step,
%error,
"a due timer's wake failed; the claim lease retries it, \
and a wake recorded before the failure reaches the \
recovery pass",
);
woken.failed += 1;
}
}
}
Ok(woken)
}
async fn fire_one(&self, timer: &crate::core::Timer) -> Result<(), RuntimeError> {
let lease = self
.store()
.acquire(timer.run, self.owner_id(), LEASE_TTL)
.await
.map_err(RuntimeError::from_store)?;
let already_recorded = self
.store()
.read(timer.run, 1)
.await
.map_err(RuntimeError::from_store)?
.iter()
.any(|record| {
record.effect_key() == Some(timer.effect)
&& matches!(record.kind(), RecordKind::EffectDone { .. })
});
if !already_recorded {
let mut record = Append::new(
timer.run,
RecordKind::EffectDone {
output: serde_json::json!({ "fired_at": timer.fire_at.unix_timestamp() }),
source: None,
spend: crate::core::Spend::default(),
declared: crate::core::DeclaredOutput::trusted(),
},
)
.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)?;
}
self.timers()
.ok_or_else(|| RuntimeError::PlanContract("timer store vanished mid-sweep".into()))?
.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,
);
self.meter().count(metrics::TIMERS_FIRED, "");
match self.resume_holding(timer.run, lease).await {
Ok(_) | Err(RuntimeError::LeaseHeld { .. }) => Ok(()),
Err(e) => Err(e),
}
}
async fn sweep_deadlines(
&self,
now: Timestamp,
report: &mut SweepReport,
ledger: &mut SweepLedger,
) -> Result<usize, RuntimeError> {
let Some(cases) = self.cases() else {
return Ok(0);
};
let mut warned = 0;
let due = cases
.due(now, DEADLINE_BATCH)
.await
.map_err(RuntimeError::from_store)?;
if due.len() >= DEADLINE_BATCH {
report.saturated.deadlines = true;
}
for deadline in due {
if deadline.is_due(now) {
ledger
.note(
self.store(),
Some(deadline.case),
deadline.case.to_string(),
SweptAction::DeadlineBreached,
Some(format!(
"'{}' was due {} and was not met",
deadline.name, deadline.resolved_at
)),
)
.await?;
ledger
.note(
self.store(),
Some(deadline.case),
deadline.case.to_string(),
SweptAction::CaseEscalated,
Some(format!("obligation '{}' was breached", deadline.name)),
)
.await?;
cases
.set_status(deadline.case, CaseStatus::Escalated)
.await
.map_err(RuntimeError::from_store)?;
cases
.set_deadline_state(deadline.case, &deadline.name, DeadlineState::Breached)
.await
.map_err(RuntimeError::from_store)?;
tracing::error!(
target: super::telemetry::DEADLINE_BREACHED,
case = %deadline.case,
obligation = %deadline.name,
due = %deadline.resolved_at,
);
self.meter().count(metrics::DEADLINE_BREACHES, "");
report.breached += 1;
} else if deadline.needs_warning(now) {
ledger
.note(
self.store(),
Some(deadline.case),
deadline.case.to_string(),
SweptAction::DeadlineWarned,
Some(format!(
"'{}' comes due {}",
deadline.name, deadline.resolved_at
)),
)
.await?;
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,
ledger: &mut SweepLedger,
) -> Result<(), RuntimeError> {
let Some(tasks) = self.tasks() else {
return Ok(());
};
let overdue = tasks
.overdue(now, TASK_BATCH)
.await
.map_err(RuntimeError::from_store)?;
if overdue.len() >= TASK_BATCH {
report.saturated.tasks = true;
}
for task in overdue {
match task.on_expiry {
OnExpiry::Escalate => {
ledger
.note(
self.store(),
task.case,
task.id.to_hex(),
SweptAction::TaskEscalated,
Some("nobody answered inside the window".to_owned()),
)
.await?;
tasks
.escalate(task.id)
.await
.map_err(RuntimeError::from_store)?;
report.tasks_escalated += 1;
}
policy => {
let decision = Decision::expired(policy);
ledger
.note(
self.store(),
task.case,
task.id.to_hex(),
SweptAction::TaskExpired,
Some(format!("window closed; applied {policy:?}")),
)
.await?;
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(
SOURCE_WORKLIST,
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?;
let delivery = self.answer_task(id, decision).await?;
tasks
.set_state(id, TaskState::Completed)
.await
.map_err(RuntimeError::from_store)?;
Ok(delivery)
}
}