#[cfg(feature = "manifest")]
use std::collections::HashSet;
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, EffectDescriptor,
InboundEvent, Ledger, Outcome, Phase, PlanIR, PlanNode, PolicyBundleIdentity, RunId,
RuntimeError, Skill, SkillDescriptor, Spend, StepId, Tainted, WallClock,
};
use crate::journal::{Append, JournalStore, Record, RecordKind, ReplayCursor, StepCursor};
use crate::runtime::BuildError;
use super::ctx::{CaseContext, Mode, StepCtx};
use super::metrics;
use super::telemetry;
use tracing::Instrument;
#[derive(Debug, Clone)]
pub(crate) enum CaseBinding {
Correlate {
kind: String,
keys: Vec<CorrelationKey>,
},
Existing(crate::core::CaseId),
}
pub const LEASE_TTL: Duration = Duration::from_secs(30);
pub const MIN_LEASE_TTL: Duration = Duration::from_secs(2);
#[derive(Debug, Clone)]
pub struct RunOutcome {
pub run_id: RunId,
pub status: RunStatus,
pub spend: Spend,
pub chain_head: Digest,
pub output: Option<Tainted<Value>>,
}
impl RunOutcome {
#[must_use]
pub fn reason(&self) -> Option<std::borrow::Cow<'_, str>> {
self.status.reason()
}
pub fn success(self) -> Result<Tainted<Value>, RunFailure> {
match self.status {
RunStatus::Succeeded => {
Ok(self.output.unwrap_or_else(|| Tainted::trusted(Value::Null)))
}
status => Err(RunFailure {
run_id: self.run_id,
status,
}),
}
}
}
#[derive(Debug, Clone)]
pub struct RunFailure {
pub run_id: RunId,
pub status: RunStatus,
}
impl std::fmt::Display for RunFailure {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "run {} did not succeed: ", self.run_id)?;
match &self.status {
RunStatus::Succeeded => unreachable!("a success is not a failure"),
RunStatus::Failed(reason) => write!(f, "it failed — {reason}"),
RunStatus::Suspended(reason) => write!(f, "it is suspended — {reason:?}"),
RunStatus::Exhausted(limit) => write!(f, "a budget stopped it — {limit}"),
RunStatus::Quarantined(reason) => write!(f, "it is quarantined — {reason}"),
RunStatus::Replanning(reason) => write!(f, "it asked to replan — {reason}"),
RunStatus::Cancelled { actor, reason } => {
write!(f, "{actor} cancelled it — {reason}")
}
}
}
}
impl std::error::Error for RunFailure {}
#[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 reason(&self) -> Option<std::borrow::Cow<'_, str>> {
use std::borrow::Cow;
match self {
Self::Succeeded => None,
Self::Failed(reason) | Self::Quarantined(reason) | Self::Replanning(reason) => {
Some(Cow::Borrowed(reason.as_str()))
}
Self::Cancelled { reason, .. } => Some(Cow::Borrowed(reason.as_str())),
Self::Suspended(reason) => Some(Cow::Owned(reason.to_string())),
Self::Exhausted(exceeded) => Some(Cow::Owned(exceeded.to_string())),
}
}
#[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 { .. })
}
#[must_use]
pub fn seals(&self) -> bool {
matches!(
self,
Self::Succeeded | Self::Quarantined(_) | Self::Cancelled { .. }
)
}
}
pub const SEALED_OUTCOMES: &[&str] = &[
"succeeded",
"quarantined",
"cancelled",
super::sweeper::SWEEP_OUTCOME,
BREAK_GLASS_OUTCOME,
];
struct Admitted {
run: RunId,
epoch: crate::core::Epoch,
budget: Budget,
agent: String,
plan: PlanIR,
input: Tainted<Value>,
case: Option<CaseContext>,
}
struct Heartbeat(tokio::task::JoinHandle<()>);
impl Drop for Heartbeat {
fn drop(&mut self) {
self.0.abort();
}
}
#[derive(Debug, Clone)]
pub struct Runtime {
store: Arc<dyn JournalStore>,
skills: HashMap<String, Arc<dyn Skill>>,
by_capability: HashMap<Capability, String>,
self_ref: std::sync::Weak<Runtime>,
#[cfg(feature = "manifest")]
governed_by: HashMap<String, Arc<crate::manifest::Manifest>>,
tenant: crate::core::TenantId,
#[cfg(feature = "manifest")]
published_by: HashMap<String, crate::core::KeyId>,
#[cfg(feature = "manifest")]
providers: HashMap<String, Arc<dyn crate::model::ModelProvider>>,
owner: String,
lease_ttl: Duration,
memories: Option<Arc<dyn crate::memory::MemoryStore>>,
semantic: Option<Arc<super::SemanticMemory>>,
authorities: Option<Arc<dyn crate::authority::AuthorityStore>>,
#[cfg(feature = "manifest")]
tools: Option<(
Arc<crate::tools::ToolCatalog>,
Arc<dyn crate::tools::ToolClient>,
)>,
meter: super::metrics::Meter,
quotas: Option<Arc<dyn crate::quota::QuotaStore>>,
quota: crate::quota::TenantQuota,
budget: Budget,
cases: Option<Arc<dyn CaseStore>>,
events: Option<Arc<dyn EventStore>>,
tasks: Option<Arc<dyn TaskStore>>,
timers: Option<Arc<dyn TimerStore>>,
blobs: Option<Arc<dyn crate::blob::BlobStore>>,
#[cfg(feature = "keyring")]
keyring: Option<Arc<dyn crate::keyring::KeyRing>>,
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>>,
#[cfg(feature = "push")]
outbox: Option<Arc<crate::push::Outbox>>,
}
pub trait FullBackend:
JournalStore + CaseStore + TaskStore + EventStore + TimerStore + crate::memory::MemoryStore
{
}
impl<B> FullBackend for B where
B: JournalStore + CaseStore + TaskStore + EventStore + TimerStore + crate::memory::MemoryStore
{
}
impl Runtime {
#[must_use]
pub fn builder_on<B: FullBackend + 'static>(store: Arc<B>) -> RuntimeBuilder {
Self::builder(Arc::clone(&store) as Arc<dyn JournalStore>)
.cases(Arc::clone(&store) as Arc<dyn CaseStore>)
.tasks(Arc::clone(&store) as Arc<dyn TaskStore>)
.events(Arc::clone(&store) as Arc<dyn EventStore>)
.timers(Arc::clone(&store) as Arc<dyn TimerStore>)
.memory(store as Arc<dyn crate::memory::MemoryStore>)
}
#[must_use]
pub fn builder(store: Arc<dyn JournalStore>) -> RuntimeBuilder {
RuntimeBuilder {
store,
signer: None,
skills: Vec::new(),
owner: None,
lease_ttl: LEASE_TTL,
memories: None,
semantic: None,
authorities: None,
metric_tenant: super::metrics::TenantLabel::default(),
quotas: None,
quota: crate::quota::TenantQuota::default(),
budget: Budget::unlimited(),
cases: None,
events: None,
tasks: None,
timers: None,
blobs: None,
#[cfg(feature = "keyring")]
keyring: None,
#[cfg(feature = "manifest")]
tools: None,
tenant: crate::core::TenantId::default(),
batches: None,
policy: None,
identity: None,
replanner: None,
calendar: None,
#[cfg(feature = "manifest")]
toolbox: None,
#[cfg(feature = "manifest")]
tool_servers: Vec::new(),
#[cfg(feature = "manifest")]
agents: Vec::new(),
#[cfg(feature = "manifest")]
providers: HashMap::new(),
#[cfg(feature = "push")]
outbox: None,
}
}
#[cfg(feature = "manifest")]
pub(crate) fn model_provider(
&self,
name: &str,
) -> Option<Arc<dyn crate::model::ModelProvider>> {
self.providers.get(name).map(Arc::clone)
}
#[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 const fn budget(&self) -> &Budget {
&self.budget
}
#[must_use]
pub fn blobs(&self) -> Option<&Arc<dyn crate::blob::BlobStore>> {
self.blobs.as_ref()
}
pub async fn drill(&self) -> Result<crate::drill::DrillReport, RuntimeError> {
let cases = self.cases().ok_or_else(|| {
RuntimeError::PlanContract(
"this runtime has no case store — the drill walks cases, so there is \
nothing to drill; build it with `.cases(store)`"
.into(),
)
})?;
let stores = crate::drill::Stores {
cases,
blobs: self.blobs(),
#[cfg(feature = "keyring")]
keys: self.keyring.as_ref(),
};
crate::drill::drill(&stores)
.await
.map_err(RuntimeError::from_store)
}
#[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 {
match self.replay(run, Mode::Resume).await {
Ok(_) | Err(RuntimeError::LeaseHeld { .. }) => {}
Err(e) => return Err(e),
}
}
Ok(fresh)
}
pub async fn set_halt(&self, reason: Option<&str>) -> Result<(), RuntimeError> {
let quotas = self.quotas.as_ref().ok_or_else(|| {
RuntimeError::Store(crate::core::StoreError::Backend(
"an emergency stop needs a quota store to keep the flag in — an \
in-process one is forgotten by a restart and never seen by a \
second instance"
.to_owned(),
))
})?;
quotas.set_halt(reason).await.map_err(RuntimeError::Store)
}
pub async fn halted(&self) -> Result<Option<String>, RuntimeError> {
let quotas = self.quotas.as_ref().ok_or_else(|| {
RuntimeError::Store(crate::core::StoreError::Backend(
"no quota store is wired, so no emergency stop can be set or read".to_owned(),
))
})?;
quotas.halted().await.map_err(RuntimeError::Store)
}
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 tenant(&self) -> &crate::core::TenantId {
&self.tenant
}
#[must_use]
pub(crate) fn meter(&self) -> &super::metrics::Meter {
&self.meter
}
pub async fn record_break_glass(
&self,
actor: &str,
roles: &[String],
reason: &str,
) -> Result<RunId, RuntimeError> {
if reason.trim().is_empty() {
return Err(RuntimeError::PlanContract(
"break-glass needs a reason: an unexplained crossing of the tenant \
boundary is what this record exists to prevent"
.to_owned(),
));
}
let run = RunId::generate();
let epoch = BREAK_GLASS_EPOCH;
self.store
.append(
epoch,
vec![Append::new(
run,
RecordKind::BreakGlass {
actor: actor.to_owned(),
roles: roles.to_vec(),
reason: reason.to_owned(),
},
)],
)
.await
.map_err(RuntimeError::from_store)?;
let head = self
.store
.head(run)
.await
.map_err(RuntimeError::from_store)?;
self.store
.append(
epoch,
vec![Append::new(
run,
RecordKind::RunSealed {
outcome: BREAK_GLASS_OUTCOME.to_owned(),
chain_head: head.hash,
},
)],
)
.await
.map_err(RuntimeError::from_store)?;
self.store
.seal(run, epoch, BREAK_GLASS_OUTCOME)
.await
.map_err(RuntimeError::from_store)?;
tracing::warn!(
tenant = %self.tenant(),
%actor,
%run,
reason,
"break-glass: an operator crossed the tenant boundary"
);
Ok(run)
}
pub fn journal(&self) -> &Arc<dyn JournalStore> {
&self.store
}
pub async fn case_of(&self, run: RunId) -> Result<Option<crate::core::CaseId>, RuntimeError> {
Ok(self
.store
.read(run, 1)
.await
.map_err(RuntimeError::from_store)?
.first()
.and_then(|record| record.body.case))
}
#[must_use]
pub fn store(&self) -> &Arc<dyn JournalStore> {
&self.store
}
#[must_use]
pub fn owner_id(&self) -> &str {
&self.owner
}
#[cfg(feature = "manifest")]
fn governing(&self, skill: &dyn Skill) -> Option<Arc<crate::manifest::Manifest>> {
self.governed_by.get(&skill.descriptor().name).cloned()
}
fn heartbeat(&self, run: RunId, epoch: crate::core::Epoch) -> Heartbeat {
let store = Arc::clone(&self.store);
let owner = self.owner.clone();
let ttl = self.lease_ttl;
let period = ttl / 3;
Heartbeat(tokio::spawn(async move {
loop {
tokio::time::sleep(period).await;
if store.renew(run, &owner, epoch, ttl).await.is_err() {
return;
}
}
}))
}
async fn check_quota(&self, run: RunId) -> Result<(), RuntimeError> {
let Some(quotas) = self.quotas.as_ref() else {
return Ok(());
};
match quotas.halted().await {
Ok(Some(reason)) => {
return Err(RuntimeError::QuotaExceeded(
crate::quota::QuotaError::Halted {
tenant: self.tenant.as_str().to_owned(),
reason,
},
));
}
Ok(None) => {}
Err(e) => {
return Err(RuntimeError::QuotaExceeded(
crate::quota::QuotaError::Unavailable(e.to_string()),
));
}
}
if self.quota.is_unlimited() {
return Ok(());
}
if self.quota.bounds_spend() {
let period = self.quota.period.key_for(now_for_admission());
let spent = quotas.spent(&period).await.map_err(|e| {
RuntimeError::QuotaExceeded(crate::quota::QuotaError::Unavailable(e.to_string()))
})?;
crate::quota::check_spend(self.tenant.as_str(), &period, &self.quota, spent)
.map_err(RuntimeError::QuotaExceeded)?;
}
quotas
.reserve(run, self.quota.max_concurrent_runs, now_for_admission())
.await
.map_err(RuntimeError::QuotaExceeded)
}
async fn settle_quota(&self, run: RunId, spend: Spend) {
let Some(quotas) = self.quotas.as_ref() else {
return;
};
if let Err(e) = quotas.release(run).await {
tracing::debug!(%run, error = %e, "could not release the quota slot");
}
if self.quota.bounds_spend() {
let period = self.quota.period.key_for(now_for_admission());
if let Err(e) = quotas.accrue(&period, spend).await {
tracing::warn!(
%run, error = %e,
"could not record this run's spend against the tenant ceiling — \
the period will under-count"
);
}
}
}
fn budget_for(&self, target: &str) -> Budget {
#[cfg(feature = "manifest")]
if let Ok(skill) = self.resolve(target)
&& let Some(m) = self.governing(skill.as_ref())
{
return m.budget();
}
let _ = target;
self.budget
}
#[cfg(feature = "manifest")]
fn identity_for(&self, target: &str) -> Option<crate::journal::AgentIdentity> {
let skill = self.resolve(target).ok()?;
let m = self.governing(skill.as_ref())?;
Some(crate::journal::AgentIdentity {
name: m.metadata.name.clone(),
version: m.metadata.version.clone(),
digest: m.digest().ok()?,
publisher: self.published_by.get(&m.metadata.name).cloned(),
})
}
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));
}
let mut available: Vec<String> = self
.by_capability
.keys()
.map(std::string::ToString::to_string)
.collect();
available.sort();
Err(RuntimeError::NoProvider {
target: target.to_owned(),
available,
})
}
pub async fn run(
&self,
target: &str,
input: Tainted<Value>,
) -> Result<RunOutcome, RuntimeError> {
self.admit(target, input, None).await
}
pub async fn spawn(
self: &Arc<Self>,
target: &str,
input: Tainted<Value>,
) -> Result<RunId, RuntimeError> {
self.spawn_bound(target, input, None).await
}
async fn spawn_bound(
self: &Arc<Self>,
target: &str,
input: Tainted<Value>,
case: Option<CaseBinding>,
) -> Result<RunId, RuntimeError> {
let skill = self.resolve(target)?;
let capability = first_capability(&skill.descriptor());
let run = RunId::generate();
let admitted = self
.admit_only(run, PlanIR::single(capability), input, case)
.await?;
let plane = Arc::clone(self);
tokio::spawn(async move { plane.execute_admitted(admitted).await });
Ok(run)
}
pub async fn spawn_in_case(
self: &Arc<Self>,
target: &str,
input: Tainted<Value>,
case: crate::core::CaseId,
) -> Result<RunId, RuntimeError> {
self.spawn_bound(target, input, Some(CaseBinding::Existing(case)))
.await
}
pub async fn spawn_correlated(
self: &Arc<Self>,
target: &str,
input: Tainted<Value>,
case_kind: &str,
keys: &[CorrelationKey],
) -> Result<RunId, RuntimeError> {
self.spawn_bound(
target,
input,
Some(CaseBinding::Correlate {
kind: case_kind.to_owned(),
keys: keys.to_vec(),
}),
)
.await
}
pub async fn run_in_case(
&self,
target: &str,
input: Tainted<Value>,
case: crate::core::CaseId,
) -> Result<RunOutcome, RuntimeError> {
self.admit(target, input, Some(CaseBinding::Existing(case)))
.await
}
pub async fn run_correlated(
&self,
target: &str,
input: Tainted<Value>,
case_kind: &str,
keys: &[CorrelationKey],
) -> Result<RunOutcome, RuntimeError> {
self.admit(
target,
input,
Some(CaseBinding::Correlate {
kind: case_kind.to_owned(),
keys: keys.to_vec(),
}),
)
.await
}
pub async fn run_plan(
&self,
plan: PlanIR,
input: Tainted<Value>,
) -> Result<RunOutcome, RuntimeError> {
self.admit_plan(plan, input, None).await
}
pub async fn run_plan_correlated(
&self,
plan: PlanIR,
input: Tainted<Value>,
case_kind: &str,
keys: &[CorrelationKey],
) -> Result<RunOutcome, RuntimeError> {
self.admit_plan(
plan,
input,
Some(CaseBinding::Correlate {
kind: case_kind.to_owned(),
keys: keys.to_vec(),
}),
)
.await
}
fn recorded_case(&self, records: &[Record]) -> Option<CaseContext> {
let (case_id, correlation) = records.iter().find_map(|r| match r.kind() {
RecordKind::CaseBound { correlation, .. } => {
r.body.case.map(|case_id| (case_id, correlation.clone()))
}
_ => None,
})?;
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,
correlation,
})
}
pub(crate) fn contract(&self) -> crate::plan::Contract {
crate::plan::Contract::new(self.by_capability.keys().cloned())
}
async fn admit(
&self,
target: &str,
input: Tainted<Value>,
case: Option<CaseBinding>,
) -> Result<RunOutcome, RuntimeError> {
let skill = self.resolve(target)?;
let capability = first_capability(&skill.descriptor());
self.admit_plan(PlanIR::single(capability), input, case)
.await
}
async fn admit_plan(
&self,
plan: PlanIR,
input: Tainted<Value>,
case: Option<CaseBinding>,
) -> 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,
capability: &str,
governed_by: Option<&crate::journal::AgentIdentity>,
input: &Value,
) -> Result<(), RuntimeError> {
let Some(engine) = self.policy.as_ref() else {
return Ok(());
};
let mut context = serde_json::json!({ "input": input, "tenant": self.tenant.as_str() });
if let Some(id) = governed_by {
let mut agent = serde_json::json!({
"name": id.name,
"version": id.version,
"digest": id.digest.to_hex(),
});
if let Some(publisher) = id.publisher.as_ref() {
agent["publisher"] = serde_json::to_value(publisher)?;
}
context["agent"] = agent;
}
super::ctx::merge_identity(&mut context, self.identity.as_ref());
let principal = self
.identity
.as_ref()
.map_or(capability, |chain| chain.subject().id.as_str());
let request = crate::core::PolicyRequest {
principal,
action: crate::core::ACTION_ADMIT,
resource: capability,
context: &context,
};
let decision = engine.authorize(&request);
let malformed = decision.is_malformed();
let Some(reason) = decision.reason().map(ToOwned::to_owned) else {
return Ok(());
};
tracing::error!(
target: telemetry::POLICY_DENIED,
action = crate::core::ACTION_ADMIT,
resource = %capability,
policy_error = malformed,
%reason,
);
self.meter
.count(metrics::POLICY_DENIALS, crate::core::ACTION_ADMIT);
Err(RuntimeError::PolicyDenied(
crate::core::PolicyError::Denied {
principal: principal.to_owned(),
action: crate::core::ACTION_ADMIT.to_owned(),
resource: capability.to_owned(),
},
))
}
fn admission(
&self,
capability: &str,
governed_by: Option<crate::journal::AgentIdentity>,
input: &Tainted<Value>,
) -> RecordKind {
RecordKind::RunAdmitted {
capability: capability.to_owned(),
governed_by,
input: input.peek().clone(),
input_label: input.label().clone(),
policy_bundle: self.policy.as_ref().map(|p| p.bundle()),
canon: crate::core::canon::VERSION,
}
}
#[allow(clippy::too_many_lines)]
async fn admit_only(
&self,
run: RunId,
plan: PlanIR,
input: Tainted<Value>,
case: Option<CaseBinding>,
) -> Result<Admitted, 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());
#[cfg(feature = "manifest")]
let governed_by = self.identity_for(&agent);
#[cfg(not(feature = "manifest"))]
let governed_by: Option<crate::journal::AgentIdentity> = None;
self.authorize_scope(&plan)?;
self.authorize_admission(&agent, governed_by.as_ref(), input.peek())?;
self.check_quota(run).await?;
let out = self
.admit_reserved(run, plan, input, case, agent, governed_by)
.await;
if out.is_err()
&& let Some(quotas) = self.quotas.as_ref()
&& let Err(e) = quotas.release(run).await
{
tracing::debug!(%run, error = %e, "could not release the quota slot of a failed admission");
}
out
}
async fn admit_reserved(
&self,
run: RunId,
plan: PlanIR,
input: Tainted<Value>,
case: Option<CaseBinding>,
agent: String,
governed_by: Option<crate::journal::AgentIdentity>,
) -> Result<Admitted, RuntimeError> {
let lease = self
.store
.acquire(run, &self.owner, self.lease_ttl)
.await
.map_err(RuntimeError::from_store)?;
let out = self
.admit_under_lease(run, &lease, plan, input, case, &agent, governed_by)
.await;
if out.is_err()
&& let Err(e) = self.store.release_lease(run, lease.epoch).await
{
tracing::debug!(
%run,
error = %e,
"could not release the lease of a failed admission; it will expire on its own"
);
}
out
}
#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
async fn admit_under_lease(
&self,
run: RunId,
lease: &crate::journal::Lease,
plan: PlanIR,
input: Tainted<Value>,
case: Option<CaseBinding>,
agent: &str,
governed_by: Option<crate::journal::AgentIdentity>,
) -> Result<Admitted, RuntimeError> {
let mut records = vec![
Append::new(run, self.admission(agent, governed_by, &input)),
Append::new(
run,
RecordKind::PlanFrozen {
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(CaseBinding::Correlate { 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)?;
let bound_keys = case_correlation(cases.as_ref(), case_id, &keys).await?;
for r in &mut records {
r.case = Some(case_id);
}
records.push(
Append::new(
run,
RecordKind::CaseBound {
case_kind: kind,
opened: correlation.is_new(),
correlation: bound_keys.clone(),
},
)
.case(case_id),
);
Some(CaseContext {
cases: Arc::clone(cases),
tasks: self.tasks.clone(),
events: self.events.clone(),
calendar: Arc::clone(&self.calendar),
case_id,
correlation: bound_keys,
})
}
(Some(CaseBinding::Existing(case_id)), Some(cases)) => {
let existing = cases
.case(case_id)
.await
.map_err(RuntimeError::from_store)?
.ok_or_else(|| {
RuntimeError::PlanContract(format!("no such case: {case_id}"))
})?;
if existing.status.is_closed() {
return Err(RuntimeError::PlanContract(format!(
"case '{case_id}' is closed and cannot accept another run"
)));
}
cases
.attach_run(case_id, run)
.await
.map_err(RuntimeError::from_store)?;
for record in &mut records {
record.case = Some(case_id);
}
records.push(
Append::new(
run,
RecordKind::CaseBound {
case_kind: existing.kind,
opened: false,
correlation: existing.correlation.clone(),
},
)
.case(case_id),
);
Some(CaseContext {
cases: Arc::clone(cases),
tasks: self.tasks.clone(),
events: self.events.clone(),
calendar: Arc::clone(&self.calendar),
case_id,
correlation: existing.correlation,
})
}
(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)?;
#[cfg(feature = "push")]
if let Some(outbox) = &self.outbox
&& let Err(open_failed) = outbox.open(run).await
{
let head = self
.store
.head(run)
.await
.map_err(RuntimeError::from_store)?;
self.store
.append(
lease.epoch,
vec![
Append::new(
run,
RecordKind::Note {
text: format!(
"admission failed after its records were written: the \
outbox registration was refused ({open_failed}) — the \
run is concluded failed and was never executed"
),
},
),
Append::new(
run,
RecordKind::RunSealed {
outcome: RunStatus::Failed(String::new()).as_str().to_owned(),
chain_head: head.hash,
},
),
],
)
.await
.map_err(RuntimeError::from_store)?;
return Err(RuntimeError::from_store(open_failed));
}
Ok(Admitted {
run,
epoch: lease.epoch,
budget: self.budget_for(agent),
agent: agent.to_owned(),
plan,
input,
case: case_ctx,
})
}
async fn execute_admitted(&self, a: Admitted) -> Result<RunOutcome, RuntimeError> {
let mut cursor = ReplayCursor::default();
let _heartbeat = self.heartbeat(a.run, a.epoch);
self.execute(
Execution {
run: a.run,
epoch: a.epoch,
plan: &a.plan,
input: a.input,
mode: Mode::Live,
case: a.case,
budget: a.budget,
agent: a.agent,
refusal: None,
successors: Vec::new(),
started: BTreeSet::new(),
finished: BTreeSet::new(),
recorded_groups: BTreeMap::new(),
},
&mut cursor,
)
.await
}
pub(crate) async fn admit_plan_as(
&self,
run: RunId,
plan: PlanIR,
input: Tainted<Value>,
case: Option<CaseBinding>,
) -> Result<RunOutcome, RuntimeError> {
let admitted = self.admit_only(run, plan, input, case).await?;
self.execute_admitted(admitted).await
}
fn ensure_resume_policy_bundle(&self, records: &[Record]) -> Result<(), RuntimeError> {
let recorded = records
.iter()
.find_map(|record| match record.kind() {
RecordKind::RunAdmitted { policy_bundle, .. } => Some(policy_bundle.clone()),
_ => None,
})
.ok_or_else(|| {
RuntimeError::PlanContract("journal has no RunAdmitted record".into())
})?;
let configured = self.policy.as_ref().map(|policy| policy.bundle());
if recorded != configured {
return Err(RuntimeError::PolicyBundleChanged {
recorded: recorded.as_ref().map(PolicyBundleIdentity::digest),
configured: configured.as_ref().map(PolicyBundleIdentity::digest),
});
}
Ok(())
}
pub async fn replay(&self, run: RunId, mode: Mode) -> Result<RunOutcome, RuntimeError> {
let lease = if mode == Mode::Resume {
Some(
self.store
.acquire(run, &self.owner, self.lease_ttl)
.await
.map_err(RuntimeError::from_store)?,
)
} else {
None
};
self.replay_releasing(run, mode, lease).await
}
pub(crate) async fn resume_holding(
&self,
run: RunId,
lease: crate::journal::Lease,
) -> Result<RunOutcome, RuntimeError> {
self.replay_releasing(run, Mode::Resume, Some(lease)).await
}
async fn replay_releasing(
&self,
run: RunId,
mode: Mode,
lease: Option<crate::journal::Lease>,
) -> Result<RunOutcome, RuntimeError> {
let outcome = self.replay_under(run, mode, lease.as_ref()).await;
if let Some(lease) = &lease
&& let Err(e) = self.store.release_lease(run, lease.epoch).await
{
tracing::debug!(
%run,
error = %e,
"could not hand back the resume lease; it will expire on its own"
);
}
outcome
}
async fn replay_under(
&self,
run: RunId,
mode: Mode,
lease: Option<&crate::journal::Lease>,
) -> 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)?;
ensure_replayable_canon(&records)?;
let input = records.iter().find_map(recorded_input).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(),
});
}
if mode == Mode::Resume
&& let Some(outcome) = recorded_conclusion(&records)
&& records
.iter()
.any(|r| matches!(r.kind(), RecordKind::StepCompensated { .. }))
{
return Err(RuntimeError::PlanContract(format!(
"run {run} concluded '{outcome}' after compensating completed steps — \
the work was reversed, and resuming over undone work would report \
success about a world where it no longer stands. Start a fresh run"
)));
}
if mode == Mode::Resume {
self.ensure_resume_policy_bundle(&records)?;
}
let case_ctx = self.recorded_case(&records);
let mut cursor = ReplayCursor::from_records(&records);
let epoch = lease.map_or_else(|| records.last().map_or(1, |r| r.body.epoch), |l| l.epoch);
let _heartbeat = lease.map(|l| self.heartbeat(run, l.epoch));
self.execute(
Execution {
run,
epoch,
plan: &plan,
input,
mode,
case: case_ctx,
budget: {
let recorded = recorded_agent(&records);
self.budget_for(&recorded)
},
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(),
started: recorded_started_steps(&records),
finished: recorded_finished_steps(&records),
recorded_groups: recorded_groups(&records),
},
&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,
started,
finished,
recorded_groups,
} = 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(),
cursor,
Journalling {
run,
epoch,
writing,
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,
started: &started,
finished: &finished,
recorded_groups: &recorded_groups,
},
)
.await;
let outcomes = collect(dispatched, &ready, cursor)?;
if mode == Mode::Strict {
for &(step, ref status, _) in &outcomes {
if !matches!(status, RunStatus::Succeeded) {
continue;
}
if let Some(key) = cursor.unconsumed_in(step, Phase::Forward) {
self.meter.count(metrics::DIVERGENCES, "");
tracing::error!(
target: telemetry::NONDETERMINISM,
%step, %key, unconsumed = true,
);
return self
.conclude(
run,
epoch,
RunStatus::Quarantined(format!(
"strict replay verified less than the run recorded: \
step {step} finished with journaled effect {key} never \
requested — this build performs fewer effects than the \
recorded one"
)),
None,
writing,
case_id,
Spend::default(),
Spend::default(),
)
.await;
}
}
}
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;
}
}
if mode == Mode::Strict
&& let Some((step, phase, key)) = cursor.first_unconsumed()
{
self.meter.count(metrics::DIVERGENCES, "");
tracing::error!(
target: telemetry::NONDETERMINISM,
%step, %key, unconsumed = true,
);
return self
.conclude(
run,
epoch,
RunStatus::Quarantined(format!(
"strict replay verified less than the run recorded: journaled \
effect {key} (step {step}, {phase:?} phase) was never requested \
— this build performs fewer effects than the recorded one"
)),
None,
writing,
case_id,
Spend::default(),
Spend::default(),
)
.await;
}
let (spend, live_spend) = {
let l = ledger.lock().expect("budget mutex");
(l.consumed().spend, l.live_spend())
};
self.conclude(
run,
epoch,
completion(¤t, &done),
run_output(¤t, &outputs),
writing,
case_id,
spend,
live_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,
already_started: batch.started.contains(&step),
already_finished: batch.finished.contains(&step),
recorded_groups: batch
.recorded_groups
.iter()
.filter(|((s, p, _), _)| *s == step && *p == Phase::Forward)
.map(|((_, _, name), n)| (name.clone(), *n))
.collect(),
},
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(
&self.meter,
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 {
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)>,
cursor: &ReplayCursor,
journal: Journalling<'_>,
) -> Result<(Vec<StepId>, Option<crate::core::BudgetExceeded>), RuntimeError> {
let mut admitted = Vec::new();
for &step in ready {
if let Some((at, limit, used)) = recorded
&& *at == step
{
if mode == Mode::Strict {
return Ok((
admitted,
Some(crate::core::BudgetExceeded::Recorded {
limit: limit.clone(),
used: used.clone(),
}),
));
}
if let Err(exceeded) = ledger.lock().expect("budget mutex").admit_step() {
return Ok((admitted, Some(exceeded)));
}
if journal.writing {
self.store
.append(
journal.epoch,
vec![(journal.stamp)(
Append::new(
journal.run,
RecordKind::BudgetReadmitted {
limit: limit.clone(),
},
)
.step(step),
)],
)
.await
.map_err(RuntimeError::from_store)?;
}
admitted.push(step);
continue;
}
let verdict = match mode {
Mode::Strict => Ok(()),
Mode::Resume if !cursor.exhausted(step, Phase::Forward) => Ok(()),
Mode::Live | Mode::Resume => 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?;
if !writing
&& !matches!(unwound, RunStatus::Quarantined(_))
&& let Some((step, phase, key)) = cursor.first_unconsumed()
{
self.meter.count(metrics::DIVERGENCES, "");
tracing::error!(
target: telemetry::NONDETERMINISM,
%step, %key, unconsumed = true,
);
return self
.conclude(
run,
epoch,
RunStatus::Quarantined(format!(
"strict replay verified less than the run recorded: journaled \
effect {key} (step {step}, {phase:?} phase) was never requested \
— this build performs fewer effects than the recorded one"
)),
None,
writing,
case_id,
Spend::default(),
Spend::default(),
)
.await;
}
let (spend, live_spend) = {
let l = ledger.lock().expect("budget mutex");
(l.consumed().spend, l.live_spend())
};
self.conclude(
run, epoch, unwound, output, writing, case_id, spend, live_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::Cancelled { .. } => {}
other => return Ok(other),
}
let (list, evidence) = match self.gated_unwind_list(&cx, &status, completed).await? {
Ok(scope) => scope,
Err(quarantine) => return Ok(quarantine),
};
let UnwindEvidence {
mutated,
undone: already_undone,
recorded_groups,
} = evidence;
let completed = &list[..];
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,
recorded_groups
.iter()
.filter(|((s, p, _), _)| *s == step && *p == Phase::Compensating)
.map(|((_, _, name), n)| (name.clone(), *n))
.collect(),
)
.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);
self.meter.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,
);
self.meter.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<UnwindEvidence, RuntimeError> {
let records = self
.store
.read(run, 1)
.await
.map_err(RuntimeError::from_store)?;
let mut touching: BTreeMap<crate::core::EffectKey, (StepId, bool)> = BTreeMap::new();
let mut undone = BTreeSet::new();
for r in &records {
let Some(step) = r.body.step else { continue };
let key = r.effect_key();
match r.kind() {
RecordKind::EffectStarted { mutates: true, .. } if r.body.phase.is_forward() => {
if let Some(key) = key {
touching.insert(key, (step, true));
}
}
RecordKind::EffectFailed { disposition, .. }
| RecordKind::EffectReconciled { disposition, .. } => {
if let Some(entry) = key.and_then(|k| touching.get_mut(&k)) {
entry.1 = *disposition != crate::core::Disposition::DidNotHappen;
}
}
RecordKind::StepCompensated { .. } => {
undone.insert(step);
}
_ => {}
}
}
let mutated = touching
.into_values()
.filter_map(|(step, touched)| touched.then_some(step))
.collect();
Ok(UnwindEvidence {
mutated,
undone,
recorded_groups: recorded_groups(&records),
})
}
async fn run_compensation(
&self,
cx: &Unwind<'_>,
step: StepId,
skill: &dyn Skill,
outputs: &BTreeMap<StepId, Tainted<Value>>,
cursor: &mut ReplayCursor,
recorded_groups: BTreeMap<String, super::ctx::RecordedGroup>,
) -> 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(),
blobs: self.blobs.clone(),
memories: self.memories.clone(),
semantic: self.semantic.clone(),
authorities: self.authorities.clone(),
#[cfg(feature = "manifest")]
tools: self.tools.clone(),
meter: self.meter.clone(),
#[cfg(feature = "keyring")]
keyring: self.keyring.clone(),
tenant: self.tenant.clone(),
ledger: Arc::clone(cx.ledger),
policy: self.policy.clone(),
identity: self.identity.clone(),
agent: cx.agent.to_owned(),
plane: self.self_ref.clone(),
#[cfg(feature = "manifest")]
manifest: self.governing(skill),
signer: self.signer.clone(),
recorded_groups,
},
);
let result = skill.compensate(&mut ctx, &output).await;
cursor.restore(step, Phase::Compensating, ctx.into_cursor());
result
}
async fn gated_unwind_list(
&self,
cx: &Unwind<'_>,
status: &RunStatus,
completed: &[(StepId, Capability)],
) -> Result<Result<(Vec<(StepId, Capability)>, UnwindEvidence), RunStatus>, RuntimeError> {
let evidence = self.unwind_evidence(cx.run).await?;
let list = if status.is_cancelled() || self.will_compensate(completed, &evidence.mutated) {
match self
.stop_list(cx.run, completed, cx.ir, &evidence.mutated)
.await?
{
Ok(list) => list,
Err(quarantine) => return Ok(Err(quarantine)),
}
} else {
completed.to_vec()
};
Ok(Ok((list, evidence)))
}
async fn stop_list(
&self,
run: RunId,
completed: &[(StepId, Capability)],
ir: &PlanIR,
mutated: &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} holds a mutating effect whose outcome is unknown — it \
was announced and never concluded, or concluded in doubt with no \
reconciliation — so the run cannot be unwound: compensating around \
it would undo everything except the one thing nobody can account for"
))));
}
Ok(Ok(Self::with_interrupted_steps(completed, ir, mutated)))
}
fn will_compensate(
&self,
completed: &[(StepId, Capability)],
mutated: &BTreeSet<StepId>,
) -> bool {
completed.iter().any(|(step, capability)| {
self.resolve(&capability.0)
.map_or(true, |skill| match skill.compensation() {
crate::core::Compensation::Compensatable => true,
crate::core::Compensation::Undeclared => mutated.contains(step),
crate::core::Compensation::Pivot | crate::core::Compensation::Unnecessary => {
false
}
})
})
}
fn with_interrupted_steps(
completed: &[(StepId, Capability)],
ir: &PlanIR,
mutated: &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)) {
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, EffectDescriptor)> =
BTreeMap::new();
let mut doubts: BTreeMap<crate::core::EffectKey, (StepId, EffectDescriptor)> =
BTreeMap::new();
for r in &records {
let Some(key) = r.effect_key() else { continue };
match r.kind() {
RecordKind::EffectStarted {
mutates: true,
descriptor,
..
} => {
if let Some(step) = r.body.step {
open.insert(key, (step, descriptor.clone()));
}
}
RecordKind::EffectDone { .. } => {
if let Some((step, descriptor)) = open.remove(&key) {
doubts.retain(|_, (s, d)| !(*s == step && *d == descriptor));
}
}
RecordKind::EffectFailed { disposition, .. } => {
if let Some((step, descriptor)) = open.remove(&key)
&& *disposition == crate::core::Disposition::InDoubt
{
doubts.insert(key, (step, descriptor));
}
}
RecordKind::EffectReconciled { disposition, .. } => {
let settled = open.remove(&key);
match disposition {
crate::core::Disposition::Landed
| crate::core::Disposition::DidNotHappen => {
doubts.remove(&key);
}
crate::core::Disposition::InDoubt => {
if let Some(entry) = settled {
doubts.insert(key, entry);
}
}
}
}
_ => {}
}
}
Ok(open
.values()
.map(|(step, _)| *step)
.chain(doubts.values().map(|(step, _)| *step))
.min())
}
async fn run_step(
&self,
ctx: StepRun<'_>,
run_input: &Tainted<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: &Tainted<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,
already_started,
already_finished,
recorded_groups,
} = ctx;
let step = node.id;
let skill = self.resolve(&node.capability.0)?;
let step_input = assemble(node, run_input, outputs)?;
let announce = writes_step_record(mode, !already_started);
if announce {
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(),
blobs: self.blobs.clone(),
memories: self.memories.clone(),
semantic: self.semantic.clone(),
authorities: self.authorities.clone(),
#[cfg(feature = "manifest")]
tools: self.tools.clone(),
meter: self.meter.clone(),
#[cfg(feature = "keyring")]
keyring: self.keyring.clone(),
tenant: self.tenant.clone(),
ledger: Arc::clone(ledger),
policy: self.policy.clone(),
identity: self.identity.clone(),
agent: agent.to_owned(),
plane: self.self_ref.clone(),
#[cfg(feature = "manifest")]
manifest: self.governing(skill.as_ref()),
signer: self.signer.clone(),
recorded_groups,
},
);
let result = skill.invoke(&mut cx, step_input).await;
let result = settle_abandoned_group(&mut cx, result).await;
let wrote = cx.wrote_records();
let cursor = cx.into_cursor();
ledger.lock().expect("budget mutex").record_step();
let (status, output) = classify(&self.meter, result);
tracing::Span::current().record(telemetry::OUTCOME, status.as_str());
if let RunStatus::Quarantined(why) = &status {
tracing::error!(target: telemetry::QUARANTINED, %step, reason = %why);
}
let record_ending = writes_step_record(mode, announce || wrote || !already_finished);
if writing && record_ending {
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,
live_spend: Spend,
) -> Result<RunOutcome, RuntimeError> {
if let RunStatus::Failed(reason) = &status {
tracing::warn!(target: telemetry::RUN_FAILED, %run, reason = %reason);
}
let chain_head = if writing && !status.is_suspended() {
let before = self
.store
.head(run)
.await
.map_err(RuntimeError::from_store)?;
let repeated = !status.seals()
&& self
.store
.read(run, before.seq)
.await
.map_err(RuntimeError::from_store)?
.last()
.is_some_and(|r| {
matches!(
r.kind(),
RecordKind::RunSealed { outcome, .. }
if outcome == status.as_str()
)
});
if repeated {
before.hash
} else {
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);
}
let concluded = self
.store
.append(epoch, vec![sealed])
.await
.map_err(RuntimeError::from_store)?;
if status.seals() {
self.store
.seal(run, epoch, status.as_str())
.await
.map_err(RuntimeError::from_store)?
} else {
concluded.last().map_or(before.hash, |r| r.hash)
}
}
} else {
self.store
.head(run)
.await
.map_err(RuntimeError::from_store)?
.hash
};
if writing {
if let Err(e) = self.store.release_lease(run, epoch).await {
tracing::debug!(
%run,
error = %e,
"could not hand back the lease; it will expire on its own"
);
}
self.settle_quota(run, live_spend).await;
announce(&self.meter, run, &status);
}
Ok(RunOutcome {
run_id: run,
status,
chain_head,
spend,
output,
})
}
}
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(meter: &super::metrics::Meter, run: RunId, next: &PlanIR, reason: &str) {
tracing::info!(
target: telemetry::REPLANNED,
%run,
from = next.derived_from.map(Digest::to_hex),
version = next.version,
%reason,
);
meter.count(metrics::REPLANS, "");
}
fn announce(meter: &super::metrics::Meter, run: RunId, status: &RunStatus) {
meter.count(metrics::RUNS, status.as_str());
match status {
RunStatus::Quarantined(why) => {
tracing::error!(target: telemetry::QUARANTINED, %run, reason = %why);
meter.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 Tainted<Value>,
outputs: &'a BTreeMap<StepId, Tainted<Value>>,
started: &'a BTreeSet<StepId>,
finished: &'a BTreeSet<StepId>,
recorded_groups: &'a BTreeMap<(StepId, Phase, String), super::ctx::RecordedGroup>,
}
struct UnwindEvidence {
mutated: BTreeSet<StepId>,
undone: BTreeSet<StepId>,
recorded_groups: BTreeMap<(StepId, Phase, String), super::ctx::RecordedGroup>,
}
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,
already_started: bool,
already_finished: bool,
recorded_groups: BTreeMap<String, super::ctx::RecordedGroup>,
}
fn recorded_step_refusal(records: &[Record]) -> Option<(StepId, String, String)> {
let mut standing: Option<(StepId, String, String)> = None;
for r in records {
match r.kind() {
RecordKind::BudgetRefused { limit, used } if r.effect_key().is_none() => {
if let Some(step) = r.body.step {
standing = Some((step, limit.clone(), used.clone()));
}
}
RecordKind::BudgetReadmitted { .. }
if standing
.as_ref()
.is_some_and(|(s, _, _)| Some(*s) == r.body.step) =>
{
standing = None;
}
_ => {}
}
}
standing
}
fn recorded_conclusion(records: &[Record]) -> Option<String> {
records.iter().rev().find_map(|r| match r.kind() {
RecordKind::RunSealed { outcome, .. } => Some(outcome.clone()),
_ => None,
})
}
const fn writes_step_record(mode: Mode, new_fact: bool) -> bool {
match mode {
Mode::Live => true,
Mode::Resume => new_fact,
Mode::Strict => false,
}
}
fn recorded_started_steps(records: &[Record]) -> BTreeSet<StepId> {
records
.iter()
.filter_map(|r| match r.kind() {
RecordKind::StepStarted { .. } => r.body.step,
_ => None,
})
.collect()
}
fn recorded_finished_steps(records: &[Record]) -> BTreeSet<StepId> {
records
.iter()
.filter_map(|r| match r.kind() {
RecordKind::StepFinished { .. } => r.body.step,
_ => None,
})
.collect()
}
fn recorded_groups(
records: &[Record],
) -> BTreeMap<(StepId, Phase, String), super::ctx::RecordedGroup> {
let mut recorded: BTreeMap<(StepId, Phase, String), super::ctx::RecordedGroup> =
BTreeMap::new();
for r in records {
let (group, opened) = match r.kind() {
RecordKind::GroupOpened { group, .. } => (group, true),
RecordKind::GroupSettled { group, .. } => (group, false),
_ => continue,
};
let Some(step) = r.body.step else { continue };
let entry = recorded
.entry((step, r.body.phase, group.clone()))
.or_default();
if opened {
entry.opened += 1;
} else {
entry.settled += 1;
}
}
recorded
}
#[cfg(feature = "manifest")]
fn check_declaration_matches_skills(
m: &crate::manifest::Manifest,
mine: &HashSet<Capability>,
) -> Result<(), BuildError> {
let missing: Vec<String> = m
.spec
.capabilities
.provides
.iter()
.filter(|c| !mine.contains(&Capability::new(c.as_str())))
.cloned()
.collect();
if !missing.is_empty() {
return Err(BuildError::AdvertisesWhatItCannotProvide {
agent: m.metadata.name.clone(),
missing,
});
}
let undeclared: Vec<String> = {
let declared: HashSet<Capability> = m
.spec
.capabilities
.provides
.iter()
.map(|c| Capability::new(c.as_str()))
.collect();
let mut extra: Vec<String> = mine
.iter()
.filter(|c| !declared.contains(c))
.map(|c| c.0.clone())
.collect();
extra.sort();
extra
};
if !undeclared.is_empty() {
return Err(BuildError::ProvidesWhatItDoesNotAdvertise {
agent: m.metadata.name.clone(),
undeclared,
});
}
Ok(())
}
fn register_skill(
skill: Arc<dyn Skill>,
caps: &mut HashMap<Capability, String>,
skills: &mut HashMap<String, Arc<dyn Skill>>,
) -> Result<String, BuildError> {
let d = skill.descriptor();
if let Some(existing) = skills.get(&d.name)
&& !Arc::ptr_eq(existing, &skill)
{
return Err(BuildError::DuplicateSkillName { name: d.name });
}
for cap in d.capabilities() {
if let Some(first) = caps.get(&cap)
&& first != &d.name
{
return Err(BuildError::CapabilityClaimedTwice {
capability: cap.0,
first: first.clone(),
second: d.name,
});
}
caps.insert(cap, d.name.clone());
}
skills.insert(d.name.clone(), skill);
Ok(d.name)
}
fn first_capability(descriptor: &SkillDescriptor) -> Capability {
descriptor
.capabilities()
.into_iter()
.next()
.unwrap_or_else(|| Capability::new(descriptor.name.clone()))
}
async fn case_correlation(
cases: &dyn CaseStore,
case_id: crate::core::CaseId,
admitted: &[CorrelationKey],
) -> Result<Vec<CorrelationKey>, RuntimeError> {
let stored = cases
.case(case_id)
.await
.map_err(RuntimeError::from_store)?
.map(|case| case.correlation)
.unwrap_or_default();
let mut keys = if stored.is_empty() {
admitted.to_vec()
} else {
stored
};
keys.sort();
keys.dedup();
Ok(keys)
}
fn default_owner() -> String {
use std::collections::hash_map::RandomState;
use std::hash::{BuildHasher, Hasher};
use std::sync::atomic::{AtomicU64, Ordering};
static SEED: std::sync::OnceLock<u64> = std::sync::OnceLock::new();
static SEQ: AtomicU64 = AtomicU64::new(0);
let seed = *SEED.get_or_init(|| RandomState::new().build_hasher().finish());
let n = SEQ.fetch_add(1, Ordering::Relaxed);
format!("agentplane-{seed:016x}-{n}")
}
fn recorded_input(r: &Record) -> Option<Tainted<Value>> {
match r.kind() {
RecordKind::RunAdmitted {
input, input_label, ..
} => Some(Tainted::with_label(input.clone(), input_label.clone())),
_ => None,
}
}
fn ensure_replayable_canon(records: &[Record]) -> Result<(), RuntimeError> {
if let Some(recorded) = records.iter().find_map(recorded_canon)
&& recorded != crate::core::canon::VERSION
{
return Err(RuntimeError::CanonicalizationChanged {
recorded,
implemented: crate::core::canon::VERSION,
});
}
Ok(())
}
fn recorded_canon(r: &Record) -> Option<u16> {
match r.kind() {
RecordKind::RunAdmitted { canon, .. } => Some(*canon),
_ => None,
}
}
fn recorded_agent(records: &[Record]) -> String {
records
.iter()
.find_map(|r| match r.kind() {
RecordKind::RunAdmitted { capability, .. } => Some(capability.clone()),
_ => None,
})
.unwrap_or_default()
}
fn assemble(
node: &PlanNode,
run_input: &Tainted<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 fields = Vec::with_capacity(node.args.len());
for (name, source) in &node.args {
let v = resolve_arg(node, source, run_input, outputs)?;
fields.push((name.clone(), v));
}
Ok(Tainted::object(fields))
}
fn resolve_arg(
node: &PlanNode,
source: &ArgSource,
run_input: &Tainted<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::with_label(pick(run_input.peek(), field), run_input.label().clone())
}
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
))
})?;
match field {
Some(field) => upstream
.project_field(field)
.unwrap_or_else(|| Tainted::with_label(Value::Null, upstream.label().clone())),
None => upstream.clone(),
}
}
})
}
async fn settle_abandoned_group(
cx: &mut StepCtx<'_>,
result: Result<Outcome, crate::core::SkillError>,
) -> Result<Outcome, crate::core::SkillError> {
use crate::core::{SkillError, StepError};
let Some(name) = cx.open_group().map(|g| g.name.clone()) else {
return result;
};
if matches!(&result, Err(SkillError::Step(StepError::Suspended(_)))) {
return result;
}
let doubt = match &result {
Err(SkillError::Step(e)) => crate::runtime::group::may_have_externalised(e),
_ => false,
};
if doubt {
let detail = match &result {
Err(e) => e.to_string(),
Ok(_) => String::new(),
};
let settled = cx
.settle_open_group(crate::core::GroupOutcome::Quarantined, Some(&detail))
.await;
return match settled {
Ok(()) => result,
Err(e) => Err(SkillError::Step(e)),
};
}
match cx
.abort_open_group("the step ended without settling the group")
.await
{
Err(e) => Err(SkillError::Step(e)),
Ok(()) => match result {
Err(e) => Err(e),
failed @ Ok(Outcome::Fail { .. }) => failed,
Ok(_) => Err(SkillError::Step(StepError::GroupAborted {
what: format!(
"step made progress with group '{name}' still open — it was \
reversed, because a group that commits by being forgotten is worse \
than one that does not commit at all"
),
})),
},
}
}
fn classify(
meter: &super::metrics::Meter,
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),
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,
);
meter.count(metrics::DIVERGENCES, "");
}
SkillError::Step(StepError::ReplayOverrun { actual }) => {
tracing::error!(target: telemetry::NONDETERMINISM, %actual, overrun = true);
meter.count(metrics::DIVERGENCES, "");
}
SkillError::Step(StepError::Undecidable { key, detail, .. }) => {
tracing::error!(target: telemetry::UNDECIDABLE, %key, %detail);
meter.count(metrics::UNDECIDABLE, "");
}
_ => {}
}
let untrustworthy = matches!(
e,
SkillError::Step(
StepError::NonDeterminism { .. }
| StepError::ReplayOverrun { .. }
| StepError::Undecidable { .. }
| StepError::GroupUnsettled { .. }
)
);
if untrustworthy {
(RunStatus::Quarantined(msg), None)
} else {
(RunStatus::Failed(msg), None)
}
}
}
}
struct Execution<'a> {
refusal: Option<(StepId, String, String)>,
successors: Vec<PlanIR>,
started: BTreeSet<StepId>,
finished: BTreeSet<StepId>,
recorded_groups: BTreeMap<(StepId, Phase, String), super::ctx::RecordedGroup>,
run: RunId,
epoch: u64,
plan: &'a PlanIR,
input: Tainted<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(),
}),
"failed" | "exhausted" => None,
other => Some(RunStatus::Quarantined(format!(
"recorded as '{other}', which this build does not recognise as resumable"
))),
}
}
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()
}
const BREAK_GLASS_EPOCH: crate::core::Epoch = 1;
const BREAK_GLASS_OUTCOME: &str = "broke-glass";
#[cfg(feature = "manifest")]
#[derive(Debug, Default)]
pub struct Agent {
manifest: Option<Arc<crate::manifest::Manifest>>,
publisher: Option<crate::core::KeyId>,
skills: Vec<Arc<dyn Skill>>,
}
#[cfg(feature = "manifest")]
impl Agent {
#[must_use]
pub fn new(manifest: &crate::manifest::Manifest) -> Self {
Self {
manifest: Some(Arc::new(manifest.clone())),
publisher: None,
skills: Vec::new(),
}
}
#[must_use]
pub fn published_by(mut self, key_id: impl Into<crate::core::KeyId>) -> Self {
self.publisher = Some(key_id.into());
self
}
#[must_use]
pub fn skill(mut self, skill: impl Skill + 'static) -> Self {
self.skills.push(Arc::new(skill));
self
}
}
#[derive(Debug)]
pub struct RuntimeBuilder {
store: Arc<dyn JournalStore>,
signer: Option<Arc<dyn crate::core::Signer>>,
skills: Vec<Arc<dyn Skill>>,
#[cfg(feature = "manifest")]
tools: Option<(
Arc<crate::tools::ToolCatalog>,
Arc<dyn crate::tools::ToolClient>,
)>,
tenant: crate::core::TenantId,
owner: Option<String>,
lease_ttl: Duration,
memories: Option<Arc<dyn crate::memory::MemoryStore>>,
semantic: Option<Arc<super::SemanticMemory>>,
authorities: Option<Arc<dyn crate::authority::AuthorityStore>>,
metric_tenant: super::metrics::TenantLabel,
quotas: Option<Arc<dyn crate::quota::QuotaStore>>,
quota: crate::quota::TenantQuota,
budget: Budget,
#[cfg(feature = "manifest")]
toolbox: Option<crate::tools::ToolBox>,
#[cfg(feature = "manifest")]
tool_servers: Vec<(String, Arc<dyn crate::tools::ToolClient>)>,
#[cfg(feature = "manifest")]
agents: Vec<Agent>,
#[cfg(feature = "manifest")]
providers: HashMap<String, Arc<dyn crate::model::ModelProvider>>,
cases: Option<Arc<dyn CaseStore>>,
events: Option<Arc<dyn EventStore>>,
tasks: Option<Arc<dyn TaskStore>>,
timers: Option<Arc<dyn TimerStore>>,
blobs: Option<Arc<dyn crate::blob::BlobStore>>,
#[cfg(feature = "keyring")]
keyring: Option<Arc<dyn crate::keyring::KeyRing>>,
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>>,
#[cfg(feature = "push")]
outbox: Option<Arc<crate::push::Outbox>>,
}
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 lease_ttl(mut self, ttl: Duration) -> Self {
self.lease_ttl = ttl;
self
}
#[must_use]
pub const fn metric_tenant(mut self, label: super::metrics::TenantLabel) -> Self {
self.metric_tenant = label;
self
}
#[must_use]
pub fn memory(mut self, memories: Arc<dyn crate::memory::MemoryStore>) -> Self {
self.memories = Some(memories);
self
}
#[must_use]
pub fn semantic_memory(
mut self,
embedder: Arc<dyn crate::memory::Embedder>,
retriever: Arc<dyn crate::memory::SemanticRetriever>,
) -> Self {
self.semantic = Some(Arc::new(super::SemanticMemory {
embedder,
retriever,
}));
self
}
#[must_use]
pub fn authorities(mut self, authorities: Arc<dyn crate::authority::AuthorityStore>) -> Self {
self.authorities = Some(authorities);
self
}
#[must_use]
pub fn quota(
mut self,
quotas: Arc<dyn crate::quota::QuotaStore>,
quota: crate::quota::TenantQuota,
) -> Self {
self.quotas = Some(quotas);
self.quota = quota;
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
}
#[cfg(feature = "push")]
#[must_use]
pub fn outbox(mut self, outbox: Arc<crate::push::Outbox>) -> Self {
self.outbox = Some(outbox);
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
}
#[cfg(feature = "manifest")]
#[must_use]
pub fn agent(mut self, agent: Agent) -> Self {
self.agents.push(agent);
self
}
#[cfg(feature = "manifest")]
#[must_use]
pub fn provider(
mut self,
name: impl Into<String>,
provider: Arc<dyn crate::model::ModelProvider>,
) -> Self {
self.providers.insert(name.into(), provider);
self
}
#[must_use]
pub fn blobs(mut self, blobs: Arc<dyn crate::blob::BlobStore>) -> Self {
self.blobs = Some(blobs);
self
}
#[cfg(feature = "manifest")]
#[must_use]
pub fn tools(
mut self,
catalog: Arc<crate::tools::ToolCatalog>,
client: Arc<dyn crate::tools::ToolClient>,
) -> Self {
self.tools = Some((catalog, client));
self
}
#[cfg(feature = "manifest")]
#[must_use]
pub fn toolbox(mut self, tools: crate::tools::ToolBox) -> Self {
self.toolbox = Some(tools);
self
}
#[cfg(feature = "manifest")]
#[must_use]
pub fn tool_server(
mut self,
name: impl Into<String>,
client: Arc<dyn crate::tools::ToolClient>,
) -> Self {
self.tool_servers.push((name.into(), client));
self
}
#[must_use]
pub fn tenant(mut self, tenant: crate::core::TenantId) -> Self {
self.tenant = tenant;
self
}
#[cfg(feature = "keyring")]
#[must_use]
pub fn keyring(mut self, keyring: Arc<dyn crate::keyring::KeyRing>) -> Self {
self.keyring = Some(keyring);
self
}
#[cfg(feature = "keyring")]
fn seal_stores(&mut self) {
let Some(keys) = self.keyring.clone() else {
return;
};
let tenant = self.tenant.clone();
self.store = crate::keyring::SealedJournal::wrap(
Arc::clone(&self.store),
Arc::clone(&keys),
tenant.clone(),
);
if let Some(cases) = self.cases.take() {
self.cases = Some(crate::keyring::SealedCases::wrap(
cases,
Arc::clone(&keys),
tenant.clone(),
));
}
if let Some(events) = self.events.take() {
self.events = Some(crate::keyring::SealedEvents::wrap(
events,
Arc::clone(&keys),
tenant.clone(),
));
}
#[cfg(feature = "push")]
if let Some(outbox) = self.outbox.take() {
self.outbox = Some(Arc::new(
outbox
.as_ref()
.clone()
.sealed(Arc::clone(&keys), tenant.clone()),
));
}
if let Some(tasks) = self.tasks.take() {
self.tasks = Some(crate::keyring::SealedTasks::wrap(tasks, keys, tenant));
}
}
#[cfg(not(feature = "keyring"))]
#[allow(clippy::unused_self)]
fn seal_stores(&mut 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
}
#[cfg(feature = "manifest")]
fn settle_toolbox(&mut self) -> Result<(), BuildError> {
let servers = std::mem::take(&mut self.tool_servers);
let tools = self.toolbox.take();
if tools.is_none() && servers.is_empty() {
return Ok(());
}
let tools = tools.unwrap_or_default();
let remote_servers: std::collections::BTreeSet<String> =
servers.iter().map(|(name, _)| name.clone()).collect();
if remote_servers.contains(crate::tools::AGENT_SERVER)
|| tools
.servers()
.any(|server| server == crate::tools::AGENT_SERVER)
{
return Err(BuildError::ReservedToolServer);
}
if remote_servers.len() != servers.len() {
let mut seen = std::collections::BTreeSet::new();
let duplicate = servers
.iter()
.map(|(name, _)| name)
.find(|name| !seen.insert((*name).clone()))
.cloned()
.unwrap_or_default();
return Err(BuildError::DuplicateToolServer { server: duplicate });
}
if self.tools.is_some() {
return Err(BuildError::ToolsWiredTwice);
}
let mut catalog = crate::tools::ToolCatalog::new();
let mut declared = 0usize;
let mut source: BTreeMap<crate::tools::ToolId, (String, crate::tools::ToolSafety)> =
BTreeMap::new();
for agent in &self.agents {
let Some(manifest) = agent.manifest.as_ref() else {
continue;
};
declared += 1;
tools
.check_against(manifest, &remote_servers)
.map_err(|problems| BuildError::ToolDrift {
agent: manifest.metadata.name.clone(),
problems,
})?;
for (id, safety) in crate::tools::ToolCatalog::from_manifest(manifest).entries() {
if let Some((first, existing)) = source.get(&id) {
if existing != &safety {
return Err(BuildError::ToolDeclaredTwoWays {
tool: id.reference(),
first: first.clone(),
second: manifest.metadata.name.clone(),
});
}
continue;
}
source.insert(id.clone(), (manifest.metadata.name.clone(), safety.clone()));
catalog = catalog.allow(id, safety);
}
}
if declared == 0 {
return Err(BuildError::ToolsWithoutDeclaration);
}
for id in tools.ids() {
let (description, schema, _) = tools
.declared(id)
.expect("every registered typed tool has a declaration");
let reviewed_description = catalog
.declaration(id)
.map_or_else(|| description.to_owned(), |(text, _)| text.to_owned());
catalog = catalog.declare(id.clone(), reviewed_description, schema.clone());
}
let router = servers.into_iter().fold(
crate::tools::ToolRouter::new().toolbox(&Arc::new(tools)),
|router, (name, client)| router.server(name, client),
);
self.tools = Some((
Arc::new(catalog),
Arc::new(router) as Arc<dyn crate::tools::ToolClient>,
));
Ok(())
}
#[cfg(feature = "manifest")]
fn settle_tools(&mut self) -> Result<(), BuildError> {
self.settle_toolbox()?;
self.check_catalogue_not_laxer_than_grants()
}
#[cfg(feature = "manifest")]
fn check_catalogue_not_laxer_than_grants(&self) -> Result<(), BuildError> {
let Some((catalog, _)) = self.tools.as_ref() else {
return Ok(());
};
let mut problems = Vec::new();
for agent in &self.agents {
let Some(manifest) = agent.manifest.as_ref() else {
continue;
};
for grant in &manifest.spec.tools {
if !grant.mutates {
continue;
}
let Some(id) = crate::tools::ToolId::parse(&grant.reference) else {
continue;
};
if catalog.safety(&id).is_some_and(|s| !s.mutates) {
problems.push(format!(
"agent '{}' grants '{}' as mutating and the stated catalogue \
calls it read-only",
manifest.metadata.name, grant.reference
));
}
}
}
if problems.is_empty() {
Ok(())
} else {
Err(BuildError::CatalogueLaxerThanGrant { problems })
}
}
#[must_use]
pub fn build(self) -> Arc<Runtime> {
match self.try_build() {
Ok(runtime) => runtime,
Err(error) => panic!("{error}"),
}
}
#[cfg_attr(not(feature = "manifest"), allow(unused_mut))]
#[allow(clippy::too_many_lines)]
pub fn try_build(mut self) -> Result<Arc<Runtime>, BuildError> {
if self.lease_ttl < MIN_LEASE_TTL {
return Err(BuildError::LeaseUnrenewable {
ttl: self.lease_ttl,
minimum: MIN_LEASE_TTL,
});
}
if let Some(field) = self.budget.bricked_ceiling() {
return Err(BuildError::BudgetPermitsNothing { field });
}
if let Some(engine) = self.policy.as_ref() {
preflight_policy(engine.as_ref(), self.identity.as_ref())?;
}
#[cfg(feature = "push")]
let push = self.outbox.as_ref().map(|o| o.store_tenant());
#[cfg(not(feature = "push"))]
let push: Option<&str> = None;
check_same_tenant(
self.store.as_ref(),
self.blobs.as_ref(),
self.memories.as_ref(),
&[
("case", self.cases.as_ref().map(|s| s.tenant())),
("event", self.events.as_ref().map(|s| s.tenant())),
("task", self.tasks.as_ref().map(|s| s.tenant())),
("memory", self.memories.as_ref().map(|s| s.tenant())),
("push", push),
],
&self.tenant,
)?;
if let Some(semantic) = &self.semantic {
let embedder = semantic.embedder.revision();
let index = semantic.retriever.index().query_revision;
if embedder != index {
return Err(BuildError::EmbeddingSpaceMismatch { embedder, index });
}
if self.memories.is_none() {
return Err(BuildError::SemanticMemoryWithoutStore);
}
}
self.seal_stores();
#[cfg(feature = "manifest")]
self.settle_tools()?;
let mut skills = HashMap::new();
let mut by_capability = HashMap::new();
#[cfg(feature = "manifest")]
let mut governed_by: HashMap<String, Arc<crate::manifest::Manifest>> = HashMap::new();
#[cfg(feature = "manifest")]
let mut published_by: HashMap<String, crate::core::KeyId> = HashMap::new();
for s in self.skills {
register_skill(s, &mut by_capability, &mut skills)?;
}
#[cfg(feature = "manifest")]
for agent in self.agents {
if let (Some(m), Some(key)) = (agent.manifest.as_ref(), agent.publisher.clone()) {
published_by.insert(m.metadata.name.clone(), key);
}
let Some(m) = agent.manifest.clone() else {
for s in agent.skills {
register_skill(s, &mut by_capability, &mut skills)?;
}
continue;
};
let mut mine: HashSet<Capability> = HashSet::new();
for s in agent.skills {
mine.extend(s.descriptor().capabilities());
let name = register_skill(s, &mut by_capability, &mut skills)?;
governed_by.insert(name, Arc::clone(&m));
}
if let Some(execution) = &m.spec.execution {
let model = m
.spec
.models
.as_ref()
.and_then(|x| x.privileged.as_ref())
.ok_or_else(|| BuildError::DeclarativeWithoutModel {
agent: m.metadata.name.clone(),
})?;
let provider = self
.providers
.get(&model.provider)
.map(Arc::clone)
.ok_or_else(|| BuildError::UnknownProvider {
agent: m.metadata.name.clone(),
provider: model.provider.clone(),
})?;
if m.spec.capabilities.provides.is_empty() {
return Err(BuildError::DeclarativeProvidesNothing {
agent: m.metadata.name.clone(),
});
}
let tools = self.tools.clone().or_else(|| {
let all_agent = !m.spec.tools.is_empty()
&& m.spec.tools.iter().all(|g| {
crate::tools::ToolId::parse(&g.reference)
.is_some_and(|id| id.server == crate::tools::AGENT_SERVER)
});
all_agent.then(|| {
(
Arc::new(crate::tools::ToolCatalog::from_manifest(&m)),
Arc::new(crate::tools::ToolRouter::new())
as Arc<dyn crate::tools::ToolClient>,
)
})
});
if tools.is_none() {
let needs = match execution.kind {
crate::manifest::ExecutionKind::ToolCalling => {
Some(if m.spec.tools.is_empty() {
"no tool grants".to_owned()
} else {
format!("{} tool grant(s)", m.spec.tools.len())
})
}
crate::manifest::ExecutionKind::Planned if !m.spec.tools.is_empty() => {
Some(format!("{} tool grant(s)", m.spec.tools.len()))
}
_ => None,
};
if let Some(grants) = needs {
return Err(BuildError::DeclarativeToolsUnreachable {
agent: m.metadata.name.clone(),
kind: execution.kind.as_str(),
grants,
});
}
}
let declared = if m.spec.oversight.is_some() {
Some("`spec.oversight`".to_owned())
} else if m.spec.tools.iter().any(|g| g.requires_approval) {
Some("a grant with `requires_approval: true`".to_owned())
} else {
None
};
if let Some(declared) = declared {
let missing = if self.cases.is_none() {
Some(("case store", "cases"))
} else if self.tasks.is_none() {
Some(("worklist", "tasks"))
} else {
None
};
if let Some((missing, remedy)) = missing {
return Err(BuildError::OversightUnreachable {
agent: m.metadata.name.clone(),
declared,
missing,
remedy,
});
}
}
if let Some(memory) = &m.spec.memory {
let declared: [(&'static str, Option<&crate::manifest::MemorySubject>); 2] = [
(
"spec.memory.recall",
memory.recall.as_ref().map(|r| &r.subject),
),
(
"spec.memory.formation",
memory.formation.as_ref().map(|f| &f.subject),
),
];
for (field, subject) in declared {
let Some(subject) = subject else { continue };
if self.memories.is_none() {
return Err(BuildError::MemoryWithoutStore {
agent: m.metadata.name.clone(),
declared: field,
});
}
if subject.needs_case() && self.cases.is_none() {
return Err(BuildError::MemorySubjectUnbindable {
agent: m.metadata.name.clone(),
subject: subject.as_written(),
});
}
}
}
for cap in &m.spec.capabilities.provides {
let skill: Arc<dyn Skill> = Arc::new(super::declarative::Declarative::new(
execution.kind,
cap.clone(),
m.metadata.name.clone(),
Arc::clone(&provider),
tools.clone(),
execution.max_turns,
));
mine.insert(Capability::new(cap.as_str()));
let name = register_skill(skill, &mut by_capability, &mut skills)?;
governed_by.insert(name, Arc::clone(&m));
}
}
check_declaration_matches_skills(&m, &mine)?;
}
#[cfg(feature = "manifest")]
{
let mut checked = std::collections::BTreeSet::new();
for m in governed_by.values() {
if !checked.insert(m.metadata.name.clone()) {
continue;
}
for grant in &m.spec.tools {
let Some(id) = crate::tools::ToolId::parse(&grant.reference) else {
continue;
};
if id.server != crate::tools::AGENT_SERVER {
continue;
}
if m.spec.capabilities.provides.contains(&id.tool) {
return Err(BuildError::AgentToolSelfReference {
agent: m.metadata.name.clone(),
capability: id.tool,
});
}
if !by_capability.contains_key(&Capability::new(id.tool.as_str())) {
return Err(BuildError::AgentToolUnknownCapability {
agent: m.metadata.name.clone(),
capability: id.tool,
});
}
}
}
}
Ok(Arc::new_cyclic(|self_ref| Runtime {
self_ref: self_ref.clone(),
signer: self.signer,
store: self.store,
skills,
by_capability,
#[cfg(feature = "manifest")]
published_by,
#[cfg(feature = "manifest")]
providers: self.providers,
meter: super::metrics::Meter::new(self.metric_tenant, &self.tenant),
tenant: self.tenant,
owner: self.owner.unwrap_or_else(default_owner),
lease_ttl: self.lease_ttl,
memories: self.memories,
semantic: self.semantic,
authorities: self.authorities,
#[cfg(feature = "manifest")]
tools: self.tools,
quotas: self.quotas,
quota: self.quota,
budget: self.budget,
cases: self.cases,
events: self.events,
tasks: self.tasks,
timers: self.timers,
blobs: self.blobs,
#[cfg(feature = "keyring")]
keyring: self.keyring,
batches: self.batches,
policy: self.policy,
identity: self.identity,
replanner: self.replanner,
calendar: self.calendar.unwrap_or_else(|| Arc::new(WallClock)),
#[cfg(feature = "push")]
outbox: self.outbox,
#[cfg(feature = "manifest")]
governed_by,
}))
}
}
#[cfg_attr(not(feature = "manifest"), allow(dead_code))]
fn preflight_policy(
engine: &dyn crate::core::PolicyEngine,
identity: Option<&crate::core::Delegation>,
) -> Result<(), BuildError> {
use crate::core::{
ACTION_ADMIT, ACTION_DECLARED, ACTION_EGRESS, ACTION_PERFORM, ACTION_RELEASE,
};
let run = "run_00000000000000000000000000";
let effect = serde_json::json!({
"run": run,
"step": 0,
"tenant": "preflight",
"mutates": false,
"args": {},
});
let release = serde_json::json!({
"run": run,
"step": 0,
"release": {},
"label": {},
});
let admit = serde_json::json!({
"tenant": "preflight",
"input": {},
});
let (mut effect, mut release, mut admit) = (effect, release, admit);
for context in [&mut effect, &mut release, &mut admit] {
super::ctx::merge_identity(context, identity);
}
let probes = [
(ACTION_PERFORM, "preflight.effect", &effect),
(ACTION_DECLARED, "preflight.effect", &effect),
(ACTION_EGRESS, "preflight.effect", &effect),
(ACTION_RELEASE, "information_flow.label", &release),
(ACTION_ADMIT, "preflight.capability", &admit),
];
let requests: Vec<crate::core::PolicyRequest<'_>> = probes
.iter()
.map(|(action, resource, context)| crate::core::PolicyRequest {
principal: "preflight",
action,
resource,
context,
})
.collect();
let problems = engine.preflight(&requests);
if problems.is_empty() {
return Ok(());
}
Err(BuildError::PolicyUnevaluable {
problems: problems.join("; "),
})
}
fn check_same_tenant(
store: &dyn JournalStore,
blobs: Option<&Arc<dyn crate::blob::BlobStore>>,
memories: Option<&Arc<dyn crate::memory::MemoryStore>>,
state: &[(&'static str, Option<&str>)],
tenant: &crate::core::TenantId,
) -> Result<(), BuildError> {
if let Some(blobs) = blobs
&& blobs.tenant() != tenant.as_str()
{
return Err(BuildError::BlobStoreTenant {
plane: tenant.to_string(),
store: blobs.tenant().to_owned(),
});
}
if store.tenant() != tenant.as_str() {
return Err(BuildError::JournalStoreTenant {
plane: tenant.to_string(),
store: store.tenant().to_owned(),
});
}
for &(store, serves) in state {
if let Some(serves) = serves
&& serves != tenant.as_str()
{
return Err(BuildError::StateStoreTenant {
store,
plane: tenant.to_string(),
tenant: serves.to_owned(),
});
}
}
if store.is_shared() && memories.is_some_and(|m| m.erasure_is_distributed() == Some(false)) {
return Err(BuildError::ErasureCoordinatorNotShared);
}
Ok(())
}
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);
};
self.resume_subscription(events, sub, event).await
}
pub async fn deliver_to(
&self,
run: RunId,
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(),
)
})?;
match events
.deliver_to(run, event, now_for_admission())
.await
.map_err(RuntimeError::from_store)?
{
crate::case::TargetedDelivery::Duplicate => Ok(Delivery::Duplicate),
crate::case::TargetedDelivery::NotWaiting => Err(RuntimeError::PlanContract(format!(
"run {run} is not waiting for this input"
))),
crate::case::TargetedDelivery::Matched(sub) => {
self.resume_subscription(events, sub, event).await
}
}
}
async fn resume_subscription(
&self,
events: &Arc<dyn crate::case::EventStore>,
sub: crate::core::Subscription,
event: &InboundEvent,
) -> Result<Delivery, RuntimeError> {
let mut backoff = Duration::from_millis(25);
let lease = loop {
match self
.store
.acquire(sub.run, &self.owner, self.lease_ttl)
.await
{
Ok(lease) => break lease,
Err(crate::core::StoreError::LeaseHeld { .. })
if backoff < Duration::from_secs(1) =>
{
tokio::time::sleep(backoff).await;
backoff *= 2;
}
Err(crate::core::StoreError::LeaseHeld { .. }) => {
events
.subscribe(&sub, now_for_admission())
.await
.map_err(RuntimeError::from_store)?;
tracing::warn!(
run = %sub.run,
event = %event.id,
"an event was claimed for a run whose owner did not conclude \
within the retry window; the sweep will deliver it"
);
return Ok(Delivery::Buffered);
}
Err(e) => return Err(RuntimeError::from_store(e)),
}
};
let already_recorded = self
.store
.read(sub.run, 1)
.await
.map_err(RuntimeError::from_store)?
.iter()
.any(|record| {
record.effect_key() == Some(sub.effect)
&& matches!(record.kind(), RecordKind::EffectDone { .. })
});
if !already_recorded {
self.store
.append(
lease.epoch,
vec![{
let mut a = Append::new(
sub.run,
RecordKind::EffectDone {
output: event.payload.clone(),
source: Some(event.source.clone()),
spend: crate::core::Spend::default(),
declared: crate::core::DeclaredOutput::untrusted(),
},
)
.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)?;
match self.resume_holding(sub.run, lease).await {
Ok(_) | Err(RuntimeError::LeaseHeld { .. }) => {}
Err(e) => return Err(e),
}
Ok(Delivery::Resumed { run: sub.run })
}
pub(crate) async fn redeliver_claimed(&self, limit: usize) -> Result<usize, RuntimeError> {
let Some(events) = self.events.as_ref() else {
return Ok(0);
};
let waiting = events
.waiting(limit)
.await
.map_err(RuntimeError::from_store)?;
let mut delivered = 0usize;
for sub in waiting {
let Some(buffered) = events
.claim_for(&sub, now_for_admission())
.await
.map_err(RuntimeError::from_store)?
else {
continue;
};
match self
.resume_subscription(events, sub.clone(), &buffered.event)
.await
{
Ok(Delivery::Resumed { .. }) => delivered += 1,
Ok(_) => {}
Err(error) => {
tracing::error!(
run = %sub.run,
%error,
"a claimed event's redelivery failed; retried next tick",
);
}
}
}
Ok(delivered)
}
pub async fn sweep_events(&self, grace: std::time::Duration) -> Result<usize, RuntimeError> {
let events = self
.events
.as_ref()
.ok_or_else(|| RuntimeError::PlanContract("this runtime has no event store".into()))?;
let grace = time::Duration::try_from(grace).unwrap_or(time::Duration::MAX);
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);
self.meter
.count_by(metrics::DEAD_LETTERS, "", retired as u64);
}
Ok(retired)
}
}
#[cfg(test)]
pub(crate) fn every_status() -> Vec<RunStatus> {
use crate::core::{BudgetExceeded, CorrelationKey, SuspendReason, Timestamp};
vec![
RunStatus::Succeeded,
RunStatus::Failed("because".into()),
RunStatus::Suspended(SuspendReason::AwaitingEvent {
kind: "reply".into(),
correlation: vec![CorrelationKey::new("claim", "CLM-1")],
until: Timestamp::from_unix_timestamp(1_760_000_000).expect("time"),
}),
RunStatus::Exhausted(BudgetExceeded::Steps { allowed: 1 }),
RunStatus::Quarantined("unknown outcome".into()),
RunStatus::Replanning("try again".into()),
RunStatus::Cancelled {
actor: "ops".into(),
reason: "stop".into(),
},
]
}
#[cfg(test)]
mod resume_agreement_tests {
use super::{RunStatus, resume_is_closed};
use crate::journal::RecordKind;
use crate::runtime::every_status;
#[test]
fn a_sealing_conclusion_is_never_resumable() {
let statuses = every_status();
assert_eq!(
statuses.len(),
7,
"a RunStatus variant was added or removed — decide whether it seals \
and whether a resume may continue from it, then update this list"
);
for status in &statuses {
let records = sealed_as(status.as_str());
let verdict = resume_is_closed(&records);
if status.seals() {
assert!(
verdict.is_some(),
"'{}' seals the journal and enters the Merkle log, yet a resume \
is permitted from it — the resume would grow the history past \
the leaf every later checkpoint attests",
status.as_str()
);
}
}
}
#[test]
fn a_failed_or_exhausted_run_may_still_be_resumed() {
for outcome in ["failed", "exhausted"] {
assert!(
resume_is_closed(&sealed_as(outcome)).is_none(),
"'{outcome}' is a conclusion a resume must be able to continue \
from — its completed effects are read back from history, which \
is the point of having a journal"
);
}
}
#[test]
fn the_sealed_outcome_list_agrees_with_the_sealing_rule() {
for status in super::every_status() {
assert_eq!(
super::SEALED_OUTCOMES.contains(&status.as_str()),
status.seals(),
"'{}' disagrees between RunStatus::seals and SEALED_OUTCOMES — \
one rule, two spellings, and the export reads the list",
status.as_str()
);
}
for special in ["swept", "broke-glass"] {
assert!(
super::SEALED_OUTCOMES.contains(&special),
"'{special}' is sealed at birth by the sweeper or a break-glass \
crossing and must be exportable"
);
}
}
#[test]
fn an_unrecognised_conclusion_fails_closed() {
let verdict = resume_is_closed(&sealed_as("swept"));
assert!(
matches!(verdict, Some(RunStatus::Quarantined(_))),
"an unrecognised conclusion must quarantine rather than resume: {verdict:?}"
);
}
fn sealed_as(outcome: &str) -> Vec<crate::journal::Record> {
use crate::core::{Digest, Epoch};
use crate::journal::{Record, RecordBody};
let body = RecordBody {
seq: 1,
run: crate::core::RunId::generate(),
case: None,
step: None,
phase: super::Phase::Forward,
epoch: Epoch::default(),
v: 1,
effect_key: None,
kind: RecordKind::RunSealed {
outcome: outcome.to_owned(),
chain_head: Digest::of(b""),
},
};
vec![Record::seal(body, Digest::of(b"")).expect("a sealed record")]
}
}