use std::collections::{BTreeMap, BTreeSet, HashMap};
use std::sync::Arc;
use std::time::Duration;
use serde_json::Value;
use crate::case::{CaseStore, EventStore, TaskStore, TimerStore};
use crate::core::{
ArgSource, Budget, Calendar, Capability, CorrelationKey, Delivery, Digest, InboundEvent,
Ledger, Outcome, Phase, PlanIR, PlanNode, RunId, RuntimeError, Skill, Spend, StepId, Tainted,
WallClock,
};
use crate::journal::{Append, JournalStore, Record, RecordKind, ReplayCursor, StepCursor};
use super::ctx::{CaseContext, Mode, StepCtx};
use super::metrics;
use super::telemetry;
use tracing::Instrument;
pub(super) const LEASE_TTL: Duration = Duration::from_secs(30);
#[derive(Debug, Clone)]
pub struct RunOutcome {
pub run_id: RunId,
pub status: RunStatus,
pub spend: Spend,
pub chain_head: Digest,
pub output: Option<Value>,
}
#[derive(Debug, Clone, PartialEq)]
pub enum RunStatus {
Succeeded,
Failed(String),
Suspended(crate::core::SuspendReason),
Exhausted(crate::core::BudgetExceeded),
Quarantined(String),
Replanning(String),
Cancelled {
actor: String,
reason: String,
},
}
impl RunStatus {
#[must_use]
pub fn is_suspended(&self) -> bool {
matches!(self, Self::Suspended(_))
}
#[must_use]
pub fn is_quarantined(&self) -> bool {
matches!(self, Self::Quarantined(_))
}
#[must_use]
pub fn as_str(&self) -> &'static str {
match self {
Self::Succeeded => "succeeded",
Self::Failed(_) => "failed",
Self::Suspended(_) => "suspended",
Self::Exhausted(_) => "exhausted",
Self::Quarantined(_) => "quarantined",
Self::Replanning(_) => "replanning",
Self::Cancelled { .. } => "cancelled",
}
}
#[must_use]
pub fn is_cancelled(&self) -> bool {
matches!(self, Self::Cancelled { .. })
}
}
#[derive(Debug, Clone)]
pub struct Runtime {
store: Arc<dyn JournalStore>,
skills: HashMap<String, Arc<dyn Skill>>,
by_capability: HashMap<Capability, String>,
owner: String,
budget: Budget,
cases: Option<Arc<dyn CaseStore>>,
events: Option<Arc<dyn EventStore>>,
tasks: Option<Arc<dyn TaskStore>>,
timers: Option<Arc<dyn TimerStore>>,
batches: Option<Arc<dyn crate::batch::BatchStore>>,
policy: Option<Arc<dyn crate::core::PolicyEngine>>,
identity: Option<crate::core::Delegation>,
replanner: Option<Arc<dyn crate::plan::Replanner>>,
calendar: Arc<dyn Calendar>,
signer: Option<Arc<dyn crate::core::Signer>>,
}
impl Runtime {
#[must_use]
pub fn builder(store: Arc<dyn JournalStore>) -> RuntimeBuilder {
RuntimeBuilder {
store,
signer: None,
skills: Vec::new(),
owner: None,
budget: Budget::unlimited(),
cases: None,
events: None,
tasks: None,
timers: None,
batches: None,
policy: None,
identity: None,
replanner: None,
calendar: None,
}
}
#[must_use]
pub fn cases(&self) -> Option<&Arc<dyn CaseStore>> {
self.cases.as_ref()
}
#[must_use]
pub fn tasks(&self) -> Option<&Arc<dyn TaskStore>> {
self.tasks.as_ref()
}
#[must_use]
pub fn events(&self) -> Option<&Arc<dyn EventStore>> {
self.events.as_ref()
}
#[must_use]
pub fn timers(&self) -> Option<&Arc<dyn TimerStore>> {
self.timers.as_ref()
}
#[must_use]
pub fn batches(&self) -> Option<&Arc<dyn crate::batch::BatchStore>> {
self.batches.as_ref()
}
pub async fn request_cancel(
&self,
run: RunId,
actor: &str,
reason: &str,
) -> Result<bool, RuntimeError> {
if self
.store
.head(run)
.await
.map_err(RuntimeError::from_store)?
.seq
== 0
{
return Err(RuntimeError::Store(crate::core::StoreError::NotFound(
run.to_string(),
)));
}
let fresh = self
.store
.request_cancel(run, actor, reason)
.await
.map_err(RuntimeError::from_store)?;
if fresh {
self.replay(run, Mode::Resume).await?;
}
Ok(fresh)
}
pub async fn cancellation(
&self,
run: RunId,
) -> Result<Option<crate::journal::Cancellation>, RuntimeError> {
self.store
.cancellation(run)
.await
.map_err(RuntimeError::from_store)
}
#[must_use]
pub fn policy(&self) -> Option<&Arc<dyn crate::core::PolicyEngine>> {
self.policy.as_ref()
}
#[must_use]
pub fn journal(&self) -> &Arc<dyn JournalStore> {
&self.store
}
#[must_use]
pub fn store(&self) -> &Arc<dyn JournalStore> {
&self.store
}
#[must_use]
pub(super) fn owner(&self) -> &str {
&self.owner
}
fn resolve(&self, target: &str) -> Result<Arc<dyn Skill>, RuntimeError> {
if let Some(s) = self.skills.get(target) {
return Ok(Arc::clone(s));
}
let cap = Capability::new(target);
if let Some(name) = self.by_capability.get(&cap)
&& let Some(s) = self.skills.get(name)
{
return Ok(Arc::clone(s));
}
Err(RuntimeError::NoProvider(target.to_owned()))
}
pub async fn run(&self, target: &str, input: Value) -> Result<RunOutcome, RuntimeError> {
self.admit(target, input, None).await
}
pub async fn run_in_case(
&self,
target: &str,
input: Value,
case_kind: &str,
keys: &[CorrelationKey],
) -> Result<RunOutcome, RuntimeError> {
self.admit(target, input, Some((case_kind.to_owned(), keys.to_vec())))
.await
}
pub async fn run_plan(&self, plan: PlanIR, input: Value) -> Result<RunOutcome, RuntimeError> {
self.admit_plan(plan, input, None).await
}
pub async fn run_plan_in_case(
&self,
plan: PlanIR,
input: Value,
case_kind: &str,
keys: &[CorrelationKey],
) -> Result<RunOutcome, RuntimeError> {
self.admit_plan(plan, input, Some((case_kind.to_owned(), keys.to_vec())))
.await
}
pub(crate) fn contract(&self) -> crate::plan::Contract {
crate::plan::Contract::new(self.by_capability.keys().cloned())
}
async fn admit(
&self,
target: &str,
input: Value,
case: Option<(String, Vec<CorrelationKey>)>,
) -> Result<RunOutcome, RuntimeError> {
let skill = self.resolve(target)?;
let capability = skill
.descriptor()
.provides
.into_iter()
.next()
.unwrap_or_else(|| Capability::new(skill.descriptor().name));
self.admit_plan(PlanIR::single(capability), input, case)
.await
}
async fn admit_plan(
&self,
plan: PlanIR,
input: Value,
case: Option<(String, Vec<CorrelationKey>)>,
) -> Result<RunOutcome, RuntimeError> {
self.admit_plan_as(RunId::generate(), plan, input, case)
.await
}
fn bind_identity(&self, run: RunId, records: &mut Vec<Append>) {
if let Some(chain) = self.identity.as_ref() {
records.push(Append::new(
run,
RecordKind::IdentityBound {
chain: chain.links().cloned().collect(),
},
));
}
}
fn authorize_scope(&self, plan: &PlanIR) -> Result<(), RuntimeError> {
let Some(chain) = self.identity.as_ref() else {
return Ok(());
};
let scope = chain.effective_scope();
for node in &plan.nodes {
if !scope.permits(&node.capability) {
return Err(RuntimeError::PolicyDenied(
crate::core::PolicyError::Denied {
principal: chain.subject().id.clone(),
action: crate::core::ACTION_ADMIT.to_owned(),
resource: node.capability.to_string(),
},
));
}
}
Ok(())
}
fn authorize_admission(&self, agent: &str, input: &Value) -> Result<(), RuntimeError> {
let Some(engine) = self.policy.as_ref() else {
return Ok(());
};
let mut context = serde_json::json!({ "input": input });
super::ctx::merge_identity(&mut context, self.identity.as_ref());
let request = crate::core::PolicyRequest {
principal: agent,
action: crate::core::ACTION_ADMIT,
resource: agent,
context: &context,
};
let crate::core::PolicyDecision::Deny { reason } = engine.authorize(&request) else {
return Ok(());
};
tracing::error!(
target: telemetry::POLICY_DENIED,
action = crate::core::ACTION_ADMIT,
resource = %agent,
%reason,
);
metrics::count(metrics::POLICY_DENIALS, crate::core::ACTION_ADMIT);
Err(RuntimeError::PolicyDenied(
crate::core::PolicyError::Denied {
principal: agent.to_owned(),
action: crate::core::ACTION_ADMIT.to_owned(),
resource: agent.to_owned(),
},
))
}
pub(crate) async fn admit_plan_as(
&self,
run: RunId,
plan: PlanIR,
input: Value,
case: Option<(String, Vec<CorrelationKey>)>,
) -> Result<RunOutcome, RuntimeError> {
crate::plan::validate(&plan, &self.contract())
.map_err(|e| RuntimeError::PlanContract(e.to_string()))?;
let agent = plan
.nodes
.first()
.map_or_else(|| "plan".to_owned(), |n| n.capability.to_string());
self.authorize_scope(&plan)?;
self.authorize_admission(&agent, &input)?;
let lease = self
.store
.acquire(run, &self.owner, LEASE_TTL)
.await
.map_err(RuntimeError::from_store)?;
let mut records = vec![
Append::new(
run,
RecordKind::RunAdmitted {
agent: agent.clone(),
input: input.clone(),
policy: self.policy.as_ref().map(|p| p.digest()),
},
),
Append::new(
run,
RecordKind::PlanFrozen {
digest: plan.digest(),
steps: plan
.nodes
.iter()
.map(|n| n.capability.to_string())
.collect(),
plan: serde_json::to_value(&plan)?,
},
),
];
self.bind_identity(run, &mut records);
let case_ctx = match (case, self.cases.as_ref()) {
(Some((kind, keys)), Some(cases)) => {
let correlation = cases
.correlate_or_open(&kind, &keys, now_for_admission())
.await
.map_err(RuntimeError::from_store)?;
let case_id = correlation.case_id();
cases
.attach_run(case_id, run)
.await
.map_err(RuntimeError::from_store)?;
for r in &mut records {
r.case = Some(case_id);
}
records.push(
Append::new(
run,
RecordKind::CaseBound {
case_kind: kind,
opened: correlation.is_new(),
},
)
.case(case_id),
);
Some(CaseContext {
cases: Arc::clone(cases),
tasks: self.tasks.clone(),
events: self.events.clone(),
calendar: Arc::clone(&self.calendar),
case_id,
})
}
(Some(_), None) => {
return Err(RuntimeError::PlanContract(
"this run was admitted with correlation keys but the runtime has no case \
store — build it with `.cases(store)`"
.into(),
));
}
(None, _) => None,
};
self.store
.append(lease.epoch, records)
.await
.map_err(RuntimeError::from_store)?;
let mut cursor = ReplayCursor::default();
self.execute(
Execution {
run,
epoch: lease.epoch,
plan: &plan,
input,
mode: Mode::Live,
case: case_ctx,
budget: self.budget,
agent,
refusal: None,
successors: Vec::new(),
},
&mut cursor,
)
.await
}
pub async fn replay(&self, run: RunId, mode: Mode) -> Result<RunOutcome, RuntimeError> {
let records = self
.store
.read(run, 1)
.await
.map_err(RuntimeError::from_store)?;
if records.is_empty() {
return Err(RuntimeError::Store(crate::core::StoreError::NotFound(
run.to_string(),
)));
}
Record::verify_chain(&records, Digest::ZERO).map_err(RuntimeError::from_store)?;
let input = records
.iter()
.find_map(|r| match r.kind() {
RecordKind::RunAdmitted { input, .. } => Some(input.clone()),
_ => None,
})
.ok_or_else(|| {
RuntimeError::PlanContract("journal has no RunAdmitted record".into())
})?;
let plan: PlanIR = records
.iter()
.find_map(|r| match r.kind() {
RecordKind::PlanFrozen { plan, .. } => Some(plan.clone()),
_ => None,
})
.ok_or_else(|| RuntimeError::PlanContract("journal has no PlanFrozen record".into()))
.and_then(|v| serde_json::from_value(v).map_err(RuntimeError::Encoding))?;
if mode == Mode::Resume
&& let Some(recorded) = resume_is_closed(&records)
{
let head = self
.store
.head(run)
.await
.map_err(RuntimeError::from_store)?;
return Ok(RunOutcome {
run_id: run,
status: recorded,
chain_head: head.hash,
output: None,
spend: Spend::default(),
});
}
let case_ctx = records
.iter()
.find_map(|r| match r.kind() {
RecordKind::CaseBound { .. } => r.body.case,
_ => None,
})
.and_then(|case_id| {
self.cases.as_ref().map(|cases| CaseContext {
cases: Arc::clone(cases),
tasks: self.tasks.clone(),
events: self.events.clone(),
calendar: Arc::clone(&self.calendar),
case_id,
})
});
let mut cursor = ReplayCursor::from_records(&records);
let epoch = if mode == Mode::Strict {
records.last().map_or(1, |r| r.body.epoch)
} else {
self.store
.acquire(run, &self.owner, LEASE_TTL)
.await
.map_err(RuntimeError::from_store)?
.epoch
};
self.execute(
Execution {
run,
epoch,
plan: &plan,
input,
mode,
case: case_ctx,
budget: self.budget,
agent: recorded_agent(&records),
refusal: recorded_step_refusal(&records),
successors: records
.iter()
.filter_map(|r| match r.kind() {
RecordKind::PlanFrozen { plan, .. } => {
serde_json::from_value::<PlanIR>(plan.clone()).ok()
}
_ => None,
})
.skip(1)
.collect(),
},
&mut cursor,
)
.await
}
async fn execute(
&self,
plan: Execution<'_>,
cursor: &mut ReplayCursor,
) -> Result<RunOutcome, RuntimeError> {
let span = tracing::info_span!(
telemetry::RUN_SPAN,
{ telemetry::GEN_AI_OPERATION } = telemetry::GEN_AI_INVOKE_AGENT,
{ telemetry::RUN_ID } = tracing::field::display(plan.run),
{ telemetry::MODE } = telemetry::mode_str(plan.mode),
{ telemetry::CASE_ID } = plan
.case
.as_ref()
.map(|c| super::ctx::CaseContext::id(c).to_string()),
{ telemetry::OUTCOME } = tracing::field::Empty,
semconv = telemetry::SEMCONV_VERSION,
);
self.execute_inner(plan, cursor).instrument(span).await
}
#[allow(clippy::too_many_lines)]
async fn execute_inner(
&self,
plan: Execution<'_>,
cursor: &mut ReplayCursor,
) -> Result<RunOutcome, RuntimeError> {
let Execution {
run,
epoch,
plan: ir,
input,
mode,
case,
budget,
agent,
refusal: recorded_refusal,
successors,
} = plan;
let writing = !matches!(mode, Mode::Strict);
let case_id = case.as_ref().map(super::ctx::CaseContext::id);
let stamp = |a: Append| match case_id {
Some(c) => a.case(c),
None => a,
};
let mut current: PlanIR = ir.clone();
let mut replans: u32 = 0;
let recorded_successors = successors;
let ledger = Arc::new(std::sync::Mutex::new(Ledger::new(budget)));
let mut done: BTreeSet<StepId> = BTreeSet::new();
let mut completed: Vec<(StepId, Capability)> = Vec::new();
let mut outputs: BTreeMap<StepId, Tainted<Value>> = BTreeMap::new();
loop {
if writing
&& let Some(c) = self
.store
.cancellation(run)
.await
.map_err(RuntimeError::from_store)?
{
let status = RunStatus::Cancelled {
actor: c.actor.clone(),
reason: c.reason.clone(),
};
self.store
.append(
epoch,
vec![stamp(Append::new(
run,
RecordKind::RunCancelled {
actor: c.actor,
reason: c.reason,
},
))],
)
.await
.map_err(RuntimeError::from_store)?;
return self
.stop(
Unwind {
agent: &agent,
run,
epoch,
ir: ¤t,
mode,
case: case.clone(),
ledger: &ledger,
writing,
stamp: &stamp,
},
status,
&completed,
&outputs,
cursor,
case_id,
)
.await;
}
let ready = current.ready(&done);
if ready.is_empty() {
break;
}
let (admitted, refused) = self
.admit_ready(
&ready,
&ledger,
mode,
recorded_refusal.as_ref(),
Journalling {
run,
epoch,
writing: writing && recorded_refusal.is_none(),
stamp: &stamp,
},
)
.await?;
let dispatched = self
.dispatch(
&admitted,
cursor,
Batch {
agent: &agent,
run,
epoch,
ir: ¤t,
mode,
case: &case,
ledger: &ledger,
writing,
stamp: &stamp,
input: &input,
outputs: &outputs,
},
)
.await;
let outcomes = collect(dispatched, &ready, cursor)?;
let mut stopped = apply(¤t, outcomes, &mut done, &mut completed, &mut outputs)
.or(refused.map(RunStatus::Exhausted));
if let Some(RunStatus::Replanning(reason)) = &stopped {
let recorded = recorded_successors.get(replans as usize);
let journal = Journalling {
run,
epoch,
writing: writing && recorded.is_none(),
stamp: &stamp,
};
let cx = Replan {
current: ¤t,
reason,
already_replanned: replans,
max_replans: budget.max_replans,
recorded,
};
match self
.adopt_successor(cx, journal, &outputs, &completed)
.await?
{
Ok(next) => {
replans += 1;
current = next;
continue;
}
Err(refusal) => stopped = Some(refusal),
}
}
if let Some(status) = stopped {
return self
.stop(
Unwind {
agent: &agent,
run,
epoch,
ir: ¤t,
mode,
case: case.clone(),
ledger: &ledger,
writing,
stamp: &stamp,
},
status,
&completed,
&outputs,
cursor,
case_id,
)
.await;
}
}
let spend = ledger.lock().expect("budget mutex").consumed().spend;
self.conclude(
run,
epoch,
completion(¤t, &done),
run_output(¤t, &outputs),
writing,
case_id,
spend,
)
.await
}
async fn dispatch(
&self,
admitted: &[StepId],
cursor: &mut ReplayCursor,
batch: Batch<'_>,
) -> Vec<Dispatched> {
let slices: Vec<(StepId, StepCursor)> = admitted
.iter()
.map(|&s| (s, cursor.take(s, Phase::Forward)))
.collect();
futures_util::future::join_all(slices.into_iter().map(|(step, slice)| async move {
let node = batch
.ir
.node(step)
.ok_or_else(|| RuntimeError::PlanContract(format!("no node {step}")))?;
let (status, out, slice) = self
.run_step(
StepRun {
agent: batch.agent,
run: batch.run,
epoch: batch.epoch,
node,
phase: Phase::Forward,
mode: batch.mode,
case: batch.case.clone(),
ledger: batch.ledger,
writing: batch.writing,
stamp: batch.stamp,
},
batch.input,
batch.outputs,
slice,
)
.await?;
Ok((step, status, out, slice))
}))
.await
}
async fn adopt_successor(
&self,
cx: Replan<'_>,
journal: Journalling<'_>,
outputs: &BTreeMap<StepId, Tainted<Value>>,
completed: &[(StepId, Capability)],
) -> Result<Result<PlanIR, RunStatus>, RuntimeError> {
let next = match self.successor(cx, outputs, completed).await {
Ok(next) => next,
Err(refusal) => return Ok(Err(refusal)),
};
if journal.writing {
self.freeze(journal.run, journal.epoch, &next, journal.stamp)
.await?;
}
announce_replan(journal.run, &next, next.reason.as_deref().unwrap_or(""));
Ok(Ok(next))
}
async fn freeze(
&self,
run: RunId,
epoch: u64,
plan: &PlanIR,
stamp: &(dyn Fn(Append) -> Append + Send + Sync),
) -> Result<(), RuntimeError> {
self.store
.append(
epoch,
vec![stamp(Append::new(
run,
RecordKind::PlanFrozen {
digest: plan.digest(),
steps: plan.nodes.iter().map(|n| n.capability.0.clone()).collect(),
plan: serde_json::to_value(plan)?,
},
))],
)
.await
.map_err(RuntimeError::from_store)?;
Ok(())
}
async fn successor(
&self,
cx: Replan<'_>,
outputs: &BTreeMap<StepId, Tainted<Value>>,
completed: &[(StepId, Capability)],
) -> Result<PlanIR, RunStatus> {
if let Some(source) = untrusted_in(outputs) {
return Err(RunStatus::Failed(format!(
"replanning refused: untrusted data from {source} is already in working memory, and the plan is an authorization graph — letting it change now would let that data choose what runs next ({})",
cx.reason
)));
}
let spent = cx.already_replanned;
if let Some(max) = cx.max_replans
&& spent >= max
{
return Err(RunStatus::Exhausted(crate::core::BudgetExceeded::Replans {
allowed: max,
}));
}
if let Some(recorded) = cx.recorded {
return Ok(recorded.clone());
}
let replanner = self.replanner.as_ref().ok_or_else(|| {
RunStatus::Failed(format!(
"a step asked to replan and this runtime has no planner — build it with `.replanner(..)` ({})",
cx.reason
))
})?;
let next = replanner
.replan(cx.current, cx.reason, completed)
.await
.map_err(|e| RunStatus::Failed(format!("replanning failed: {e}")))?;
crate::plan::validate(&next, &self.contract())
.map_err(|e| RunStatus::Failed(format!("the successor plan is invalid: {e}")))?;
for (step, ran) in completed {
if let Some(node) = next.node(*step)
&& node.capability != *ran
{
return Err(RunStatus::Failed(format!(
"the successor plan reuses step {step} — which already ran \
as '{}' — for '{}'. Keep a completed step's capability or \
leave the step out; effect keys are derived from the step \
id, so new work at a used id cannot be replayed",
ran.0, node.capability.0
)));
}
}
if next.derived_from != Some(cx.current.digest()) {
return Err(RunStatus::Failed(
"the successor plan does not name its predecessor — use `PlanIR::succeed_with`, or the audit trail has a hole where the lineage should be"
.into(),
));
}
Ok(next)
}
async fn admit_ready(
&self,
ready: &[StepId],
ledger: &Arc<std::sync::Mutex<Ledger>>,
mode: Mode,
recorded: Option<&(StepId, String, String)>,
journal: Journalling<'_>,
) -> Result<(Vec<StepId>, Option<crate::core::BudgetExceeded>), RuntimeError> {
let mut admitted = Vec::new();
for &step in ready {
let verdict = if let Some((at, limit, used)) = recorded {
if *at == step {
Err(crate::core::BudgetExceeded::Recorded {
limit: limit.clone(),
used: used.clone(),
})
} else {
Ok(())
}
} else if mode.is_replaying() {
Ok(())
} else {
ledger.lock().expect("budget mutex").admit_step()
};
let Err(exceeded) = verdict else {
admitted.push(step);
continue;
};
if journal.writing {
let used = format!("{:?}", ledger.lock().expect("budget mutex").consumed());
self.store
.append(
journal.epoch,
vec![(journal.stamp)(
Append::new(
journal.run,
RecordKind::BudgetRefused {
limit: exceeded.to_string(),
used,
},
)
.step(step),
)],
)
.await
.map_err(RuntimeError::from_store)?;
}
return Ok((admitted, Some(exceeded)));
}
Ok((admitted, None))
}
#[allow(clippy::too_many_arguments)]
async fn stop(
&self,
cx: Unwind<'_>,
status: RunStatus,
completed: &[(StepId, Capability)],
outputs: &BTreeMap<StepId, Tainted<Value>>,
cursor: &mut ReplayCursor,
case_id: Option<crate::core::CaseId>,
) -> Result<RunOutcome, RuntimeError> {
let (run, epoch, writing, ir) = (cx.run, cx.epoch, cx.writing, cx.ir);
let output = run_output(ir, outputs);
let ledger = cx.ledger.clone();
let unwound = self
.maybe_unwind(cx, status, completed, outputs, cursor)
.await?;
let spend = ledger.lock().expect("budget mutex").consumed().spend;
self.conclude(run, epoch, unwound, output, writing, case_id, spend)
.await
}
async fn maybe_unwind(
&self,
cx: Unwind<'_>,
status: RunStatus,
completed: &[(StepId, Capability)],
outputs: &BTreeMap<StepId, Tainted<Value>>,
cursor: &mut ReplayCursor,
) -> Result<RunStatus, RuntimeError> {
match status {
RunStatus::Failed(_) | RunStatus::Exhausted(_) | RunStatus::Cancelled { .. } => {}
other => return Ok(other),
}
let (mutated, already_undone) = self.unwind_evidence(cx.run).await?;
let extended;
let completed = if status.is_cancelled() {
match self
.stop_list(cx.run, completed, cx.ir, &mutated, &already_undone)
.await?
{
Ok(list) => {
extended = list;
&extended[..]
}
Err(quarantine) => return Ok(quarantine),
}
} else {
completed
};
for (step, capability) in completed.iter().rev().cloned() {
let skill = self.resolve(&capability.0)?;
let declared = skill.compensation();
match declared {
crate::core::Compensation::Pivot => break,
crate::core::Compensation::Unnecessary => continue,
crate::core::Compensation::Undeclared => {
if !mutated.contains(&step) {
continue;
}
return Ok(RunStatus::Quarantined(format!(
"step {step} ('{}') changed external state and declares no \
compensation, so the run cannot be safely unwound — \
declare Compensation on it, or resolve this by hand",
capability.0
)));
}
crate::core::Compensation::Compensatable => {}
}
let result = self
.run_compensation(&cx, step, skill.as_ref(), outputs, cursor)
.await;
if let Err(crate::core::SkillError::Step(crate::core::StepError::Suspended(reason))) =
&result
{
return Ok(RunStatus::Suspended(reason.clone()));
}
let outcome = match &result {
Ok(()) => "compensated".to_owned(),
Err(e) => e.to_string(),
};
if result.is_ok() {
tracing::info!(target: telemetry::COMPENSATED, run = %cx.run, %step);
metrics::count(metrics::COMPENSATIONS, "done");
}
if cx.writing && !already_undone.contains(&step) {
self.store
.append(
cx.epoch,
vec![(cx.stamp)(
Append::new(
cx.run,
RecordKind::StepCompensated {
compensation: declared,
outcome: outcome.clone(),
},
)
.step(step)
.phase(Phase::Compensating),
)],
)
.await
.map_err(RuntimeError::from_store)?;
}
if result.is_err() {
tracing::error!(
target: telemetry::COMPENSATION_FAILED,
run = %cx.run,
%step,
detail = %outcome,
);
metrics::count(metrics::COMPENSATIONS, "failed");
return Ok(RunStatus::Quarantined(format!(
"compensation failed for step {step} ('{}'): {outcome} — the run is \
partially unwound and needs an operator",
capability.0
)));
}
}
Ok(status)
}
async fn unwind_evidence(
&self,
run: RunId,
) -> Result<(BTreeSet<StepId>, BTreeSet<StepId>), RuntimeError> {
let records = self
.store
.read(run, 1)
.await
.map_err(RuntimeError::from_store)?;
let mut mutated = BTreeSet::new();
let mut undone = BTreeSet::new();
for r in &records {
let Some(step) = r.body.step else { continue };
match r.kind() {
RecordKind::EffectStarted { mutates: true, .. } if r.body.phase.is_forward() => {
mutated.insert(step);
}
RecordKind::StepCompensated { .. } => {
undone.insert(step);
}
_ => {}
}
}
Ok((mutated, undone))
}
async fn run_compensation(
&self,
cx: &Unwind<'_>,
step: StepId,
skill: &dyn Skill,
outputs: &BTreeMap<StepId, Tainted<Value>>,
cursor: &mut ReplayCursor,
) -> Result<(), crate::core::SkillError> {
let output = outputs
.get(&step)
.cloned()
.unwrap_or_else(|| Tainted::trusted(Value::Null));
let mut ctx = StepCtx::new(
&self.store,
cursor.take(step, Phase::Compensating),
super::ctx::Frame {
run: cx.run,
epoch: cx.epoch,
step,
phase: Phase::Compensating,
mode: cx.mode,
case: cx.case.clone(),
timers: self.timers.clone(),
ledger: Arc::clone(cx.ledger),
policy: self.policy.clone(),
identity: self.identity.clone(),
agent: cx.agent.to_owned(),
signer: self.signer.clone(),
},
);
let result = skill.compensate(&mut ctx, &output).await;
cursor.restore(step, Phase::Compensating, ctx.into_cursor());
result
}
async fn stop_list(
&self,
run: RunId,
completed: &[(StepId, Capability)],
ir: &PlanIR,
mutated: &BTreeSet<StepId>,
undone: &BTreeSet<StepId>,
) -> Result<Result<Vec<(StepId, Capability)>, RunStatus>, RuntimeError> {
if let Some(step) = self.undecided_effect(run).await? {
return Ok(Err(RunStatus::Quarantined(format!(
"step {step} announced a mutating effect that never concluded, so \
the run cannot be unwound — its outcome is unknown, and \
compensating around it would undo everything except the one thing \
nobody can account for"
))));
}
Ok(Ok(Self::with_interrupted_steps(
completed, ir, mutated, undone,
)))
}
fn with_interrupted_steps(
completed: &[(StepId, Capability)],
ir: &PlanIR,
mutated: &BTreeSet<StepId>,
undone: &BTreeSet<StepId>,
) -> Vec<(StepId, Capability)> {
let mut out = completed.to_vec();
let done: BTreeSet<StepId> = out.iter().map(|(s, _)| *s).collect();
for step in mutated
.iter()
.filter(|s| !done.contains(s) && !undone.contains(s))
{
if let Some(node) = ir.node(*step) {
out.push((*step, node.capability.clone()));
}
}
out
}
async fn undecided_effect(&self, run: RunId) -> Result<Option<StepId>, RuntimeError> {
let records = self
.store
.read(run, 1)
.await
.map_err(RuntimeError::from_store)?;
let mut open: BTreeMap<crate::core::EffectKey, StepId> = BTreeMap::new();
for r in &records {
let Some(key) = r.effect_key() else { continue };
match r.kind() {
RecordKind::EffectStarted { mutates: true, .. } => {
if let Some(step) = r.body.step {
open.insert(key, step);
}
}
RecordKind::EffectDone { .. }
| RecordKind::EffectFailed { .. }
| RecordKind::EffectReconciled { .. } => {
open.remove(&key);
}
_ => {}
}
}
Ok(open.values().min().copied())
}
async fn run_step(
&self,
ctx: StepRun<'_>,
run_input: &Value,
outputs: &BTreeMap<StepId, Tainted<Value>>,
cursor: crate::journal::StepCursor,
) -> Result<
(
RunStatus,
Option<Tainted<Value>>,
crate::journal::StepCursor,
),
RuntimeError,
> {
let span = tracing::info_span!(
telemetry::STEP_SPAN,
{ telemetry::STEP } = tracing::field::display(ctx.node.id),
{ telemetry::CAPABILITY } = tracing::field::display(&ctx.node.capability.0),
{ telemetry::PHASE } = if ctx.phase.is_forward() {
"forward"
} else {
"compensating"
},
{ telemetry::MODE } = telemetry::mode_str(ctx.mode),
{ telemetry::OUTCOME } = tracing::field::Empty,
);
self.run_step_inner(ctx, run_input, outputs, cursor)
.instrument(span)
.await
}
async fn run_step_inner(
&self,
ctx: StepRun<'_>,
run_input: &Value,
outputs: &BTreeMap<StepId, Tainted<Value>>,
cursor: crate::journal::StepCursor,
) -> Result<
(
RunStatus,
Option<Tainted<Value>>,
crate::journal::StepCursor,
),
RuntimeError,
> {
let StepRun {
run,
epoch,
node,
phase,
mode,
case,
ledger,
writing,
stamp,
agent,
} = ctx;
let step = node.id;
let skill = self.resolve(&node.capability.0)?;
let step_input = assemble(node, run_input, outputs)?;
if mode == Mode::Live {
self.store
.append(
epoch,
vec![stamp(
Append::new(
run,
RecordKind::StepStarted {
skill: skill.descriptor().name,
},
)
.step(step),
)],
)
.await
.map_err(RuntimeError::from_store)?;
}
let mut cx = StepCtx::new(
&self.store,
cursor,
super::ctx::Frame {
run,
epoch,
step,
phase,
mode,
case,
timers: self.timers.clone(),
ledger: Arc::clone(ledger),
policy: self.policy.clone(),
identity: self.identity.clone(),
agent: agent.to_owned(),
signer: self.signer.clone(),
},
);
let result = skill.invoke(&mut cx, step_input).await;
let cursor = cx.into_cursor();
ledger.lock().expect("budget mutex").record_step();
let (status, output) = classify(result);
tracing::Span::current().record(telemetry::OUTCOME, status.as_str());
if let RunStatus::Quarantined(why) = &status {
tracing::error!(target: telemetry::QUARANTINED, %step, reason = %why);
}
if writing {
let record = match &status {
RunStatus::Suspended(reason) => RecordKind::RunSuspended {
reason: reason.clone(),
},
other => RecordKind::StepFinished {
outcome: other.as_str().to_owned(),
},
};
self.store
.append(epoch, vec![stamp(Append::new(run, record).step(step))])
.await
.map_err(RuntimeError::from_store)?;
}
Ok((status, output, cursor))
}
#[allow(clippy::too_many_arguments)]
async fn conclude(
&self,
run: RunId,
epoch: u64,
status: RunStatus,
output: Option<Tainted<Value>>,
writing: bool,
case: Option<crate::core::CaseId>,
spend: Spend,
) -> Result<RunOutcome, RuntimeError> {
let chain_head = if writing && !status.is_suspended() {
let before = self
.store
.head(run)
.await
.map_err(RuntimeError::from_store)?;
let mut sealed = Append::new(
run,
RecordKind::RunSealed {
outcome: status.as_str().to_owned(),
chain_head: before.hash,
},
);
if let Some(c) = case {
sealed = sealed.case(c);
}
self.store
.append(epoch, vec![sealed])
.await
.map_err(RuntimeError::from_store)?;
self.store
.seal(run, epoch, status.as_str())
.await
.map_err(RuntimeError::from_store)?
} else {
self.store
.head(run)
.await
.map_err(RuntimeError::from_store)?
.hash
};
announce(run, &status);
Ok(RunOutcome {
run_id: run,
status,
chain_head,
spend,
output: output.map(|t| t.peek().clone()),
})
}
}
fn run_output(ir: &PlanIR, outputs: &BTreeMap<StepId, Tainted<Value>>) -> Option<Tainted<Value>> {
ir.nodes
.iter()
.filter(|n| n.terminal)
.map(|n| n.id)
.min()
.and_then(|id| outputs.get(&id).cloned())
.or_else(|| outputs.iter().next_back().map(|(_, v)| v.clone()))
}
fn completion(ir: &PlanIR, done: &BTreeSet<StepId>) -> RunStatus {
if ir.is_complete(done) {
return RunStatus::Succeeded;
}
let missing: Vec<String> = ir
.nodes
.iter()
.filter(|n| n.terminal && !done.contains(&n.id))
.map(|n| n.id.to_string())
.collect();
RunStatus::Failed(format!(
"plan did not complete: terminal step(s) {} never ran",
missing.join(", ")
))
}
fn announce_replan(run: RunId, next: &PlanIR, reason: &str) {
tracing::info!(
target: telemetry::REPLANNED,
%run,
from = next.derived_from.map(Digest::to_hex),
version = next.version,
%reason,
);
metrics::count(metrics::REPLANS, "");
}
fn announce(run: RunId, status: &RunStatus) {
metrics::count(metrics::RUNS, status.as_str());
match status {
RunStatus::Quarantined(why) => {
tracing::error!(target: telemetry::QUARANTINED, %run, reason = %why);
metrics::count(metrics::QUARANTINES, "");
}
RunStatus::Exhausted(limit) => {
tracing::warn!(target: telemetry::BUDGET_REFUSED, %run, %limit);
}
_ => {}
}
tracing::Span::current().record(telemetry::OUTCOME, status.as_str());
}
fn apply(
plan: &PlanIR,
outcomes: Vec<StepOutcome>,
done: &mut BTreeSet<StepId>,
completed: &mut Vec<(StepId, Capability)>,
outputs: &mut BTreeMap<StepId, Tainted<Value>>,
) -> Option<RunStatus> {
let mut stopped: Option<RunStatus> = None;
for (step, status, output) in outcomes {
let RunStatus::Succeeded = status else {
if stopped
.as_ref()
.is_none_or(|held| severity(&status) > severity(held))
{
stopped = Some(status);
}
continue;
};
if let Some(v) = output {
outputs.insert(step, v);
}
done.insert(step);
if let Some(node) = plan.node(step) {
completed.push((step, node.capability.clone()));
}
}
stopped
}
fn severity(status: &RunStatus) -> u8 {
match status {
RunStatus::Quarantined(_) => 3,
RunStatus::Cancelled { .. } | RunStatus::Failed(_) => 2,
RunStatus::Exhausted(_) => 1,
RunStatus::Replanning(_) | RunStatus::Suspended(_) | RunStatus::Succeeded => 0,
}
}
type Dispatched = Result<(StepId, RunStatus, Option<Tainted<Value>>, StepCursor), RuntimeError>;
type StepOutcome = (StepId, RunStatus, Option<Tainted<Value>>);
fn collect(
dispatched: Vec<Dispatched>,
ready: &[StepId],
cursor: &mut ReplayCursor,
) -> Result<Vec<StepOutcome>, RuntimeError> {
let mut outcomes = Vec::with_capacity(dispatched.len());
for result in dispatched {
let (step, status, out, slice) = result?;
cursor.restore(step, Phase::Forward, slice);
outcomes.push((step, status, out));
}
outcomes.sort_by_key(|(step, _, _)| ready.iter().position(|r| r == step));
Ok(outcomes)
}
struct Batch<'a> {
agent: &'a str,
run: RunId,
epoch: u64,
ir: &'a PlanIR,
mode: Mode,
case: &'a Option<CaseContext>,
ledger: &'a Arc<std::sync::Mutex<Ledger>>,
writing: bool,
stamp: &'a (dyn Fn(Append) -> Append + Send + Sync),
input: &'a Value,
outputs: &'a BTreeMap<StepId, Tainted<Value>>,
}
struct Replan<'a> {
current: &'a PlanIR,
reason: &'a str,
already_replanned: u32,
max_replans: Option<u32>,
recorded: Option<&'a PlanIR>,
}
fn untrusted_in(outputs: &BTreeMap<StepId, Tainted<Value>>) -> Option<String> {
outputs.values().find_map(|v| {
let label = v.label();
label.is_untrusted().then(|| {
label
.provenance
.first()
.map_or_else(|| "an untrusted source".to_owned(), ToString::to_string)
})
})
}
struct Journalling<'a> {
run: RunId,
epoch: u64,
writing: bool,
stamp: &'a (dyn Fn(Append) -> Append + Send + Sync),
}
struct Unwind<'a> {
run: RunId,
epoch: u64,
ir: &'a PlanIR,
mode: Mode,
case: Option<CaseContext>,
ledger: &'a Arc<std::sync::Mutex<Ledger>>,
writing: bool,
stamp: &'a (dyn Fn(Append) -> Append + Send + Sync),
agent: &'a str,
}
struct StepRun<'a> {
run: RunId,
epoch: u64,
node: &'a PlanNode,
phase: Phase,
mode: Mode,
case: Option<CaseContext>,
ledger: &'a Arc<std::sync::Mutex<Ledger>>,
writing: bool,
stamp: &'a (dyn Fn(Append) -> Append + Send + Sync),
agent: &'a str,
}
fn recorded_step_refusal(records: &[Record]) -> Option<(StepId, String, String)> {
records.iter().find_map(|r| match r.kind() {
RecordKind::BudgetRefused { limit, used } if r.effect_key().is_none() => {
r.body.step.map(|s| (s, limit.clone(), used.clone()))
}
_ => None,
})
}
fn recorded_agent(records: &[Record]) -> String {
records
.iter()
.find_map(|r| match r.kind() {
RecordKind::RunAdmitted { agent, .. } => Some(agent.clone()),
_ => None,
})
.unwrap_or_default()
}
fn assemble(
node: &PlanNode,
run_input: &Value,
outputs: &BTreeMap<StepId, Tainted<Value>>,
) -> Result<Tainted<Value>, RuntimeError> {
if node.args.len() == 1
&& let Some((_, only)) = node.args.iter().next()
{
return resolve_arg(node, only, run_input, outputs);
}
let mut label = crate::core::Label::trusted();
let mut map = serde_json::Map::new();
for (name, source) in &node.args {
let v = resolve_arg(node, source, run_input, outputs)?;
label = label.join(v.label());
map.insert(name.clone(), v.peek().clone());
}
Ok(Tainted::with_label(Value::Object(map), label))
}
fn resolve_arg(
node: &PlanNode,
source: &ArgSource,
run_input: &Value,
outputs: &BTreeMap<StepId, Tainted<Value>>,
) -> Result<Tainted<Value>, RuntimeError> {
let pick = |v: &Value, field: &Option<String>| match field {
Some(f) => v.get(f).cloned().unwrap_or(Value::Null),
None => v.clone(),
};
Ok(match source {
ArgSource::RunInput { field } => Tainted::trusted(pick(run_input, field)),
ArgSource::Const { value } => Tainted::trusted(value.clone()),
ArgSource::Node { step, field } => {
let upstream = outputs.get(step).ok_or_else(|| {
RuntimeError::PlanContract(format!(
"step {} read step {step}, which has not produced a value",
node.id
))
})?;
let picked = pick(upstream.peek(), field);
Tainted::with_label(picked, upstream.label().clone())
}
})
}
fn classify(
result: Result<Outcome, crate::core::SkillError>,
) -> (RunStatus, Option<Tainted<Value>>) {
use crate::core::{SkillError, StepError};
match result {
Ok(Outcome::Done(v)) => (RunStatus::Succeeded, Some(v)),
Ok(Outcome::Fail { reason }) => (RunStatus::Failed(reason), None),
Ok(Outcome::Replan { reason }) => (RunStatus::Replanning(reason), None),
Ok(Outcome::Delegate { target, .. }) => (
RunStatus::Failed(format!(
"delegation to '{target}' requires a collaborative topology, which this \
build does not provide"
)),
None,
),
Err(SkillError::Step(StepError::Suspended(reason))) => (RunStatus::Suspended(reason), None),
Err(SkillError::Step(StepError::Budget(exceeded))) => {
(RunStatus::Exhausted(exceeded), None)
}
Err(e) => {
let msg = e.to_string();
match &e {
SkillError::Step(StepError::NonDeterminism {
seq,
expected,
actual,
}) => {
tracing::error!(
target: telemetry::NONDETERMINISM,
%seq, %expected, %actual,
);
metrics::count(metrics::DIVERGENCES, "");
}
SkillError::Step(StepError::ReplayOverrun { actual }) => {
tracing::error!(target: telemetry::NONDETERMINISM, %actual, overrun = true);
metrics::count(metrics::DIVERGENCES, "");
}
SkillError::Step(StepError::Undecidable { key, detail, .. }) => {
tracing::error!(target: telemetry::UNDECIDABLE, %key, %detail);
metrics::count(metrics::UNDECIDABLE, "");
}
_ => {}
}
let untrustworthy = matches!(
e,
SkillError::Step(
StepError::NonDeterminism { .. }
| StepError::ReplayOverrun { .. }
| StepError::Undecidable { .. }
)
);
if untrustworthy {
(RunStatus::Quarantined(msg), None)
} else {
(RunStatus::Failed(msg), None)
}
}
}
}
struct Execution<'a> {
refusal: Option<(StepId, String, String)>,
successors: Vec<PlanIR>,
run: RunId,
epoch: u64,
plan: &'a PlanIR,
input: Value,
mode: Mode,
case: Option<CaseContext>,
budget: Budget,
agent: String,
}
fn resume_is_closed(records: &[Record]) -> Option<RunStatus> {
let outcome = records.iter().rev().find_map(|r| match r.kind() {
RecordKind::RunSealed { outcome, .. } => Some(outcome.as_str()),
_ => None,
})?;
match outcome {
"succeeded" => Some(RunStatus::Succeeded),
"quarantined" => Some(RunStatus::Quarantined(
"recorded as quarantined; a human must resolve it before it can run again".into(),
)),
"cancelled" => Some(RunStatus::Cancelled {
actor: recorded_canceller(records).unwrap_or_else(|| "unknown".into()),
reason: "recorded as cancelled; an operator stopped this run".into(),
}),
_ => None,
}
}
fn recorded_canceller(records: &[Record]) -> Option<String> {
records.iter().rev().find_map(|r| match r.kind() {
RecordKind::RunCancelled { actor, .. } => Some(actor.clone()),
_ => None,
})
}
#[allow(clippy::disallowed_methods)]
fn now_for_admission() -> crate::core::Timestamp {
crate::core::Timestamp::now_utc()
}
#[derive(Debug)]
pub struct RuntimeBuilder {
store: Arc<dyn JournalStore>,
signer: Option<Arc<dyn crate::core::Signer>>,
skills: Vec<Arc<dyn Skill>>,
owner: Option<String>,
budget: Budget,
cases: Option<Arc<dyn CaseStore>>,
events: Option<Arc<dyn EventStore>>,
tasks: Option<Arc<dyn TaskStore>>,
timers: Option<Arc<dyn TimerStore>>,
batches: Option<Arc<dyn crate::batch::BatchStore>>,
policy: Option<Arc<dyn crate::core::PolicyEngine>>,
identity: Option<crate::core::Delegation>,
replanner: Option<Arc<dyn crate::plan::Replanner>>,
calendar: Option<Arc<dyn Calendar>>,
}
impl RuntimeBuilder {
#[must_use]
pub fn skill(mut self, s: impl Skill) -> Self {
self.skills.push(Arc::new(s));
self
}
#[must_use]
pub fn signing_as(mut self, signer: Arc<dyn crate::core::Signer>) -> Self {
self.signer = Some(signer);
self
}
#[must_use]
pub fn owner(mut self, o: impl Into<String>) -> Self {
self.owner = Some(o.into());
self
}
#[must_use]
pub fn budget(mut self, budget: Budget) -> Self {
self.budget = budget;
self
}
#[must_use]
pub fn cases(mut self, cases: Arc<dyn CaseStore>) -> Self {
self.cases = Some(cases);
self
}
#[must_use]
pub fn events(mut self, events: Arc<dyn EventStore>) -> Self {
self.events = Some(events);
self
}
#[must_use]
pub fn tasks(mut self, tasks: Arc<dyn TaskStore>) -> Self {
self.tasks = Some(tasks);
self
}
#[must_use]
pub fn replanner(mut self, r: Arc<dyn crate::plan::Replanner>) -> Self {
self.replanner = Some(r);
self
}
#[must_use]
pub fn timers(mut self, timers: Arc<dyn TimerStore>) -> Self {
self.timers = Some(timers);
self
}
#[must_use]
pub fn batches(mut self, batches: Arc<dyn crate::batch::BatchStore>) -> Self {
self.batches = Some(batches);
self
}
#[must_use]
pub fn policy(mut self, policy: Arc<dyn crate::core::PolicyEngine>) -> Self {
self.policy = Some(policy);
self
}
#[must_use]
pub fn acting_as(mut self, chain: crate::core::Delegation) -> Self {
self.identity = Some(chain);
self
}
#[must_use]
pub fn calendar(mut self, calendar: Arc<dyn Calendar>) -> Self {
self.calendar = Some(calendar);
self
}
#[must_use]
pub fn build(self) -> Runtime {
let mut skills = HashMap::new();
let mut by_capability = HashMap::new();
for s in self.skills {
let d = s.descriptor();
for cap in d.provides {
by_capability.insert(cap, d.name.clone());
}
skills.insert(d.name, s);
}
Runtime {
signer: self.signer,
store: self.store,
skills,
by_capability,
owner: self.owner.unwrap_or_else(|| "agentplane".to_owned()),
budget: self.budget,
cases: self.cases,
events: self.events,
tasks: self.tasks,
timers: self.timers,
batches: self.batches,
policy: self.policy,
identity: self.identity,
replanner: self.replanner,
calendar: self.calendar.unwrap_or_else(|| Arc::new(WallClock)),
}
}
}
impl Runtime {
pub async fn deliver(&self, event: &InboundEvent) -> Result<Delivery, RuntimeError> {
let events = self.events.as_ref().ok_or_else(|| {
RuntimeError::PlanContract(
"this runtime has no event store — build it with `.events(store)`".into(),
)
})?;
let now = now_for_admission();
if !events
.buffer(event, now)
.await
.map_err(RuntimeError::from_store)?
{
return Ok(Delivery::Duplicate);
}
let Some(sub) = events
.match_waiter(event, now)
.await
.map_err(RuntimeError::from_store)?
else {
return Ok(Delivery::Buffered);
};
let lease = self
.store
.acquire(sub.run, &self.owner, LEASE_TTL)
.await
.map_err(RuntimeError::from_store)?;
self.store
.append(
lease.epoch,
vec![{
let mut a = Append::new(
sub.run,
RecordKind::EffectDone {
output: event.payload.clone(),
spend: crate::core::Spend::default(),
},
)
.effect(sub.effect)
.step(sub.step)
.phase(sub.phase);
if let Some(c) = sub.case {
a = a.case(c);
}
a
}],
)
.await
.map_err(RuntimeError::from_store)?;
events
.unsubscribe(sub.run, sub.effect)
.await
.map_err(RuntimeError::from_store)?;
self.replay(sub.run, Mode::Resume).await?;
Ok(Delivery::Resumed { run: sub.run })
}
pub async fn sweep_events(&self, grace: time::Duration) -> Result<usize, RuntimeError> {
let events = self
.events
.as_ref()
.ok_or_else(|| RuntimeError::PlanContract("this runtime has no event store".into()))?;
let cutoff = now_for_admission() - grace;
let retired = events
.sweep_unclaimed(cutoff, "no run claimed this event within the grace window")
.await
.map_err(RuntimeError::from_store)?;
if retired > 0 {
tracing::error!(target: telemetry::DEAD_LETTERED, count = retired, %cutoff);
metrics::count_by(metrics::DEAD_LETTERS, "", retired as u64);
}
Ok(retired)
}
}