use std::sync::Arc;
use crate::core::{
CaseId, DeadlineState, Decision, InboundEvent, OnExpiry, RunId, RuntimeError, StoreError,
SweptAction, TaskState, Timestamp,
};
use crate::journal::{Append, JournalStore, RecordKind};
use super::executor::{LEASE_TTL, RunStatus, 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)]
pub struct Redelivered {
pub finished: usize,
pub examined: 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,
pub redeliveries: bool,
}
impl Saturation {
#[must_use]
pub const fn any(self) -> bool {
self.timers || self.deadlines || self.tasks || self.recovery || self.redeliveries
}
}
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 cosignatures: usize,
pub witness_shortfall: usize,
pub witness_integrity: usize,
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
|| self.witness_shortfall > 0
|| self.witness_integrity > 0
}
#[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
&& self.cosignatures == 0
&& self.witness_shortfall == 0
&& self.witness_integrity == 0
}
}
impl Runtime {
async fn cosign_checkpoint(&self, report: &mut SweepReport) -> Result<(), RuntimeError> {
use std::sync::atomic::Ordering;
let Some(witnessing) = self.witnessing() else {
return Ok(());
};
let checkpoint = self.store().checkpoint().await?;
if checkpoint.size == witnessing.submitted.load(Ordering::Relaxed) {
return Ok(());
}
let outcome = crate::journal::cosign_quorum(
self.store().as_ref(),
&checkpoint,
&witnessing.witnesses,
witnessing.quorum,
)
.await?;
report.cosignatures = outcome.cosignatures.len();
report.witness_shortfall = outcome.shortfall();
report.witness_integrity = outcome.integrity.len();
for (index, refusal) in &outcome.integrity {
tracing::error!(
target: crate::runtime::telemetry::WITNESS_INTEGRITY,
witness = index,
origin = %checkpoint.origin,
size = checkpoint.size,
refusal = %refusal,
"a witness refused this plane's checkpoint on integrity grounds"
);
}
if outcome.met() {
witnessing
.submitted
.store(checkpoint.size, Ordering::Relaxed);
}
Ok(())
}
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() {
let again = self.redeliver_claimed(EVENT_BATCH).await?;
report.events_redelivered = again.finished;
report.saturated.redeliveries = again.examined >= EVENT_BATCH;
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;
}
}
self.cosign_checkpoint(&mut report).await?;
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::LeaseHeld { .. } | RuntimeError::Fenced { .. }) => {}
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;
c.unaccounted_breaches = cc.breached;
}
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?;
}
c.quarantined_runs = self
.journal()
.count_by_outcome(RunStatus::Quarantined(String::new()).as_str())
.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,
by: 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::Fenced { .. }) => 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 owed = cases
.breaches_to_note(DEADLINE_BATCH)
.await
.map_err(RuntimeError::from_store)?;
if owed.len() >= DEADLINE_BATCH {
report.saturated.deadlines = true;
}
for deadline in owed {
self.account_for_breach(&deadline, report, ledger).await?;
}
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) {
let applied = cases
.breach_deadline(deadline.case, &deadline.name, now)
.await
.map_err(RuntimeError::from_store)?;
if applied {
self.account_for_breach(&deadline, report, ledger).await?;
}
} 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?;
match cases
.set_deadline_state(deadline.case, &deadline.name, DeadlineState::Warned)
.await
{
Ok(()) => warned += 1,
Err(StoreError::DeadlineFinal { .. }) => {
ledger
.note(
self.store(),
Some(deadline.case),
deadline.case.to_string(),
SweptAction::NotApplied,
Some(format!(
"warning for '{}' not applied: the obligation was settled first",
deadline.name
)),
)
.await?;
}
Err(e) => return Err(RuntimeError::from_store(e)),
}
}
}
Ok(warned)
}
async fn account_for_breach(
&self,
deadline: &crate::core::Deadline,
report: &mut SweepReport,
ledger: &mut SweepLedger,
) -> Result<(), RuntimeError> {
let Some(cases) = self.cases() else {
return Ok(());
};
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
.mark_breach_noted(deadline.case, &deadline.name)
.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;
Ok(())
}
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?;
let delivery = self.answer_task(task.id, &decision).await?;
if delivery == crate::core::Delivery::Duplicate {
match self.settle_answered(task.id).await? {
Some(crate::case::Minter::Operator(by)) => {
ledger
.note(
self.store(),
task.case,
task.id.to_hex(),
SweptAction::NotApplied,
Some(format!(
"expiry not applied: '{}' answered first",
by.actor()
)),
)
.await?;
}
Some(crate::case::Minter::Nobody) => report.tasks_expired += 1,
None => {}
}
continue;
}
if tasks
.set_state(task.id, TaskState::Expired)
.await
.map_err(RuntimeError::from_store)?
{
report.tasks_expired += 1;
}
}
}
}
Ok(())
}
async fn settle_answered(
&self,
id: crate::core::TaskId,
) -> Result<Option<crate::case::Minter>, RuntimeError> {
let (Some(tasks), Some(events)) = (self.tasks(), self.events()) else {
return Ok(None);
};
let minter = events
.minter(SOURCE_WORKLIST, &decision_event_id(id))
.await
.map_err(RuntimeError::from_store)?;
let state = match &minter {
Some(crate::case::Minter::Operator(_)) => TaskState::Completed,
Some(crate::case::Minter::Nobody) => TaskState::Expired,
None => return Ok(None),
};
if tasks
.set_state(id, state)
.await
.map_err(RuntimeError::from_store)?
{
Ok(minter)
} else {
Ok(None)
}
}
pub(crate) async fn answer_task(
&self,
id: crate::core::TaskId,
decision: &Decision,
) -> Result<crate::core::Delivery, RuntimeError> {
let event = InboundEvent::new(
SOURCE_WORKLIST,
decision_event_id(id),
super::ctx::TASK_DECIDED,
serde_json::to_value(decision)?,
)
.correlate(crate::core::CorrelationKey::new("task", id.to_hex()));
let event = match decision.decided.operator() {
Some(by) => event.minted_by(by.clone()),
None => event,
};
self.deliver_minted(&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()))?;
let Some(by) = decision.decided.operator() else {
return Err(RuntimeError::PlanContract(
"an expired window is the sweeper's to apply: this door records a person's \
answer, and an expiry has no actor to check a claim or an exclusion against"
.to_owned(),
));
};
if decision.approved
&& let Some(task) = tasks.task(id).await.map_err(RuntimeError::from_store)?
&& let Some(reason) = task.withheld
{
return Err(RuntimeError::ProposalWithheld {
task: id.to_string(),
reason,
});
}
let claimed = tasks.claim(id, by.actor(), roles).await?;
let mut decision = decision.clone();
decision.reviewed = Some(claimed.justification.digest());
let delivery = self.answer_task(id, &decision).await?;
if matches!(delivery, crate::core::Delivery::Duplicate) {
let events = self.events().ok_or_else(|| {
RuntimeError::PlanContract("this runtime has no event store".into())
})?;
let minter = events
.minter(SOURCE_WORKLIST, &decision_event_id(id))
.await
.map_err(RuntimeError::from_store)?;
let own = matches!(
&minter,
Some(crate::case::Minter::Operator(first)) if first.actor() == by.actor()
);
self.settle_answered(id).await?;
if own {
return Ok(delivery);
}
return Err(RuntimeError::TaskClaim(
crate::core::ClaimError::AlreadyAnswered { task: id },
));
}
tasks
.set_state(id, TaskState::Completed)
.await
.map_err(RuntimeError::from_store)?;
Ok(delivery)
}
}
fn decision_event_id(id: crate::core::TaskId) -> String {
format!("task-decision:{}", id.to_hex())
}