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::ctx::Mode;
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;
const SWEEP_OUTCOME: &str = "swept";
struct SweepLedger {
run: Option<RunId>,
entries: Vec<crate::journal::Append>,
}
impl SweepLedger {
const fn new() -> Self {
Self {
run: None,
entries: Vec::new(),
}
}
fn note(
&mut self,
case: Option<CaseId>,
subject: String,
action: SweptAction,
detail: Option<String>,
) {
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);
}
self.entries.push(entry);
}
async fn seal(self, store: &Arc<dyn JournalStore>) -> SweepRecord {
let Some(run) = self.run else {
return SweepRecord::Quiet;
};
if let Err(e) = store.append(SWEEP_EPOCH, self.entries).await {
tracing::error!(%run, error = %e, "a sweep's own record could not be written");
return SweepRecord::EvidenceLost;
}
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::RunSealed {
outcome: SWEEP_OUTCOME.to_owned(),
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 Saturation {
pub timers: bool,
pub deadlines: bool,
pub tasks: bool,
}
impl Saturation {
#[must_use]
pub const fn any(self) -> bool {
self.timers || self.deadlines || self.tasks
}
}
const TIMER_BATCH: usize = 128;
const DEADLINE_BATCH: usize = 512;
const TASK_BATCH: usize = 512;
#[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 saturated: Saturation,
pub record: Option<RunId>,
pub evidence_lost: bool,
pub census: Census,
}
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
}
#[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.timers_fired == 0
}
}
impl Runtime {
pub async fn sweep(
&self,
now: Timestamp,
event_grace: time::Duration,
) -> Result<SweepReport, RuntimeError> {
let mut report = SweepReport::default();
let mut ledger = SweepLedger::new();
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.dead_lettered = self.sweep_events(event_grace).await?;
}
if self.timers().is_some() {
report.timers_fired = self.fire_timers(now).await?;
if report.timers_fired >= TIMER_BATCH {
report.saturated.timers = true;
}
}
match ledger.seal(self.store()).await {
SweepRecord::Quiet => {}
SweepRecord::Recorded(run) => report.record = Some(run),
SweepRecord::EvidenceLost => report.evidence_lost = true,
}
report.census = self.census(now).await?;
self.meter().census(&report.census);
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, TIMER_BATCH)
.await
.map_err(RuntimeError::from_store)?;
let mut fired = 0;
for timer in due {
let lease = self
.store()
.acquire(timer.run, self.owner_id(), 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() }),
source: None,
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,
);
self.meter().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,
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) {
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)?;
ledger.note(
Some(deadline.case),
deadline.case.to_string(),
SweptAction::DeadlineBreached,
Some(format!(
"'{}' was due {} and was not met",
deadline.name, deadline.resolved_at
)),
);
ledger.note(
Some(deadline.case),
deadline.case.to_string(),
SweptAction::CaseEscalated,
Some(format!("obligation '{}' was breached", deadline.name)),
);
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) {
cases
.set_deadline_state(deadline.case, &deadline.name, DeadlineState::Warned)
.await
.map_err(RuntimeError::from_store)?;
ledger.note(
Some(deadline.case),
deadline.case.to_string(),
SweptAction::DeadlineWarned,
Some(format!(
"'{}' comes due {}",
deadline.name, deadline.resolved_at
)),
);
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 if task.state != TaskState::Escalated => {
tasks
.set_state(task.id, TaskState::Escalated)
.await
.map_err(RuntimeError::from_store)?;
ledger.note(
task.case,
task.id.to_hex(),
SweptAction::TaskEscalated,
Some("nobody answered inside the window".to_owned()),
);
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)?;
ledger.note(
task.case,
task.id.to_hex(),
SweptAction::TaskExpired,
Some(format!("window closed; applied {policy:?}")),
);
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.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)
}
}