#[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, Consumed, CorrelationKey, Delivery, Digest,
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),
}
#[derive(Debug, Clone, Default)]
pub struct RunTerms {
case: Option<CaseBinding>,
idempotency_key: Option<String>,
acting_as: Option<crate::core::Delegation>,
}
impl RunTerms {
fn bound(case: CaseBinding) -> Self {
Self {
case: Some(case),
idempotency_key: None,
acting_as: None,
}
}
#[must_use]
pub fn in_case(mut self, case: crate::core::CaseId) -> Self {
self.case = Some(CaseBinding::Existing(case));
self
}
#[must_use]
pub fn correlated(mut self, kind: &str, keys: &[CorrelationKey]) -> Self {
self.case = Some(CaseBinding::Correlate {
kind: kind.to_owned(),
keys: keys.to_vec(),
});
self
}
#[must_use]
pub fn once(mut self, idempotency_key: &str) -> Self {
self.idempotency_key = Some(idempotency_key.to_owned());
self
}
#[must_use]
pub fn acting_as(mut self, chain: crate::core::Delegation) -> Self {
self.acting_as = Some(chain);
self
}
fn keyed(self, key: &str) -> Self {
self.once(key)
}
}
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 consumed: Consumed,
pub chain_head: Digest,
pub output: Option<Tainted<Value>>,
}
impl RunOutcome {
#[must_use]
pub const fn spend(&self) -> Spend {
self.consumed.spend
}
#[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 enum Admission {
Fresh(RunOutcome),
Replayed(RunOutcome),
InFlight(RunId),
}
impl Admission {
#[must_use]
pub const fn run_id(&self) -> RunId {
match self {
Self::Fresh(outcome) | Self::Replayed(outcome) => outcome.run_id,
Self::InFlight(run) => *run,
}
}
#[must_use]
pub const fn is_fresh(&self) -> bool {
matches!(self, Self::Fresh(_))
}
#[must_use]
pub const fn outcome(&self) -> Option<&RunOutcome> {
match self {
Self::Fresh(outcome) | Self::Replayed(outcome) => Some(outcome),
Self::InFlight(_) => None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Spawned {
pub run: RunId,
pub fresh: bool,
}
#[derive(Clone)]
pub struct Stores {
pub journal: std::sync::Arc<dyn crate::journal::JournalStore>,
pub cases: std::sync::Arc<dyn crate::case::CaseStore>,
pub tasks: std::sync::Arc<dyn crate::case::TaskStore>,
pub events: std::sync::Arc<dyn crate::case::EventStore>,
pub timers: std::sync::Arc<dyn crate::case::TimerStore>,
pub memory: std::sync::Arc<dyn crate::memory::MemoryStore>,
}
impl std::fmt::Debug for Stores {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("Stores { journal, cases, tasks, events, timers, memory }")
}
}
impl Stores {
#[must_use]
pub fn on<B: FullBackend + 'static>(store: std::sync::Arc<B>) -> Self {
Self {
journal: std::sync::Arc::clone(&store) as _,
cases: std::sync::Arc::clone(&store) as _,
tasks: std::sync::Arc::clone(&store) as _,
events: std::sync::Arc::clone(&store) as _,
timers: std::sync::Arc::clone(&store) as _,
memory: store as _,
}
}
}
const ATTRIBUTION_RECORDS: usize = 8;
#[derive(Debug, Clone)]
pub struct LiveRun {
pub run: RunId,
pub governed_by: Option<crate::journal::AgentIdentity>,
pub subject: Option<String>,
pub stranded: bool,
}
#[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 {
if matches!(self.status, RunStatus::Succeeded) {
return write!(
f,
"run {} is recorded as a failure carrying a success — \
whatever built this value did so wrongly",
self.run_id
);
}
write!(f, "run {} did not succeed: ", self.run_id)?;
match &self.status {
RunStatus::Succeeded => write!(f, "it succeeded"),
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}")
}
RunStatus::Swept => write!(f, "the plane recorded a pass over its own state"),
RunStatus::BrokeGlass { actor, reason } => {
write!(f, "{actor} crossed into this tenant — {reason}")
}
RunStatus::Abandoned { actor, reason } => {
write!(
f,
"{actor} abandoned it, unresolved and not unwound — {reason}"
)
}
RunStatus::Withheld { subject, reason } => {
write!(
f,
"the authority it acts for — '{subject}' — was withdrawn: {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: crate::core::Operator,
reason: String,
},
Abandoned {
actor: crate::core::Operator,
reason: String,
},
Swept,
BrokeGlass {
actor: crate::core::Operator,
reason: String,
},
Withheld {
subject: String,
reason: String,
},
}
impl RunStatus {
#[must_use]
pub fn reason(&self) -> Option<std::borrow::Cow<'_, str>> {
use std::borrow::Cow;
match self {
Self::Succeeded | Self::Swept => None,
Self::Failed(reason)
| Self::Quarantined(reason)
| Self::Replanning(reason)
| Self::Cancelled { reason, .. }
| Self::Abandoned { reason, .. }
| Self::BrokeGlass { reason, .. }
| Self::Withheld { 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",
Self::Abandoned { .. } => "abandoned",
Self::Swept => super::sweeper::SWEEP_OUTCOME,
Self::BrokeGlass { .. } => BREAK_GLASS_OUTCOME,
Self::Withheld { .. } => WITHHELD_OUTCOME,
}
}
#[must_use]
pub fn actor(&self) -> Option<&crate::core::Operator> {
match self {
Self::Cancelled { actor, .. }
| Self::Abandoned { actor, .. }
| Self::BrokeGlass { actor, .. } => Some(actor),
Self::Succeeded
| Self::Failed(_)
| Self::Suspended(_)
| Self::Exhausted(_)
| Self::Quarantined(_)
| Self::Replanning(_)
| Self::Withheld { .. }
| Self::Swept => None,
}
}
#[must_use]
pub fn is_cancelled(&self) -> bool {
matches!(self, Self::Cancelled { .. })
}
#[must_use]
pub fn seals(&self) -> bool {
matches!(
self,
Self::Succeeded
| Self::Cancelled { .. }
| Self::Abandoned { .. }
| Self::Swept
| Self::BrokeGlass { .. }
)
}
}
pub const SEALED_OUTCOMES: &[&str] = &[
"succeeded",
"cancelled",
"abandoned",
super::sweeper::SWEEP_OUTCOME,
BREAK_GLASS_OUTCOME,
];
pub const OUTCOMES_OF_RECORD: &[&str] = &[
"succeeded",
"cancelled",
"abandoned",
"quarantined",
super::sweeper::SWEEP_OUTCOME,
BREAK_GLASS_OUTCOME,
];
struct Admitted {
run: RunId,
epoch: crate::core::Epoch,
quota: QuotaPass,
budget: Budget,
agent: String,
plan: PlanIR,
input: Tainted<Value>,
case: Option<CaseContext>,
identity: Option<crate::core::Delegation>,
}
#[derive(Debug, Clone, Default)]
struct QuotaPass {
period: Option<String>,
enabled: bool,
release_slot: bool,
}
impl QuotaPass {
fn at(
quota: &crate::quota::TenantQuota,
at: crate::core::Timestamp,
release_slot: bool,
) -> Self {
Self {
period: quota.bounds_spend().then(|| quota.period.key_for(at)),
enabled: true,
release_slot,
}
}
const fn disabled() -> Self {
Self {
period: None,
enabled: false,
release_slot: false,
}
}
fn period(&self) -> Option<&str> {
self.period.as_deref()
}
fn started(&self) -> Option<RecordKind> {
self.enabled.then(|| RecordKind::QuotaPassStarted {
period: self.period.clone(),
release_slot: self.release_slot,
})
}
fn settlement(
&self,
run: RunId,
epoch: crate::core::Epoch,
spend: Spend,
) -> Option<crate::quota::QuotaSettlement> {
self.enabled.then(|| crate::quota::QuotaSettlement {
run,
epoch,
period: self.period.clone(),
spend,
release_slot: self.release_slot,
})
}
}
#[derive(Debug)]
pub(crate) struct PeerWiring {
pub(crate) registry: crate::peers::PeerRegistry,
pub(crate) client: Arc<dyn crate::peers::PeerClient>,
}
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<Self>,
#[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>>,
peers: Option<Arc<PeerWiring>>,
#[cfg(feature = "manifest")]
tools: Option<(
Arc<crate::tools::ToolCatalog>,
Arc<dyn crate::tools::ToolClient>,
)>,
journal_ceiling: Option<crate::core::Sensitivity>,
#[cfg(feature = "manifest")]
egress: Option<Arc<crate::core::Egress>>,
meter: super::metrics::Meter,
quotas: Option<Arc<dyn crate::quota::QuotaStore>>,
quota: crate::quota::TenantQuota,
budget: Budget,
require_verifier: bool,
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 = "push")]
push: Option<Arc<dyn crate::push::PushStore>>,
#[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>>,
witnesses: Option<Arc<Witnessing>>,
inflight: Arc<super::drain::InFlight>,
}
#[derive(Debug)]
pub(crate) struct Witnessing {
pub(crate) witnesses: Vec<Arc<dyn crate::journal::Witness>>,
pub(crate) quorum: crate::journal::WitnessQuorum,
pub(crate) submitted: std::sync::atomic::AtomicU64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Entry {
Outside,
Nested,
}
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_with(Stores::on(store))
}
#[must_use]
pub fn builder_with(stores: Stores) -> RuntimeBuilder {
Self::builder(stores.journal)
.cases(stores.cases)
.tasks(stores.tasks)
.events(stores.events)
.timers(stores.timers)
.memory(stores.memory)
}
#[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,
peers: None,
tenant_label: super::telemetry::TenantLabel::default(),
quotas: None,
quota: crate::quota::TenantQuota::default(),
budget: Budget::unlimited(),
require_verifier: false,
cases: None,
events: None,
tasks: None,
timers: None,
blobs: None,
#[cfg(feature = "push")]
push: None,
#[cfg(feature = "keyring")]
keyring: None,
#[cfg(feature = "manifest")]
tools: None,
journal_ceiling: None,
#[cfg(feature = "manifest")]
egress: 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,
witnesses: Vec::new(),
quorum: 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(),
tenant: &self.tenant,
};
crate::drill::drill(&stores)
.await
.map_err(RuntimeError::from_store)
}
pub async fn retain(
&self,
older_than: crate::core::Timestamp,
at: crate::core::Timestamp,
reason: &str,
) -> Result<crate::retention::RetentionReport, RuntimeError> {
let cases = self.cases().ok_or_else(|| {
RuntimeError::PlanContract(
"this runtime has no case store — retention walks cases, so there is \
nothing to retain; build it with `.cases(store)`"
.into(),
)
})?;
let stores = crate::retention::Stores {
cases,
blobs: self.blobs(),
#[cfg(feature = "keyring")]
keys: self.keyring.as_ref(),
tenant: &self.tenant,
};
crate::retention::retain(&stores, older_than, at, reason)
.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()
}
#[cfg(feature = "push")]
#[must_use]
pub fn push(&self) -> Option<&Arc<dyn crate::push::PushStore>> {
self.push.as_ref()
}
#[cfg(feature = "push")]
pub async fn rearm_push(&self, task: RunId, id: &str) -> Result<bool, RuntimeError> {
let push = self
.push
.as_ref()
.ok_or_else(|| RuntimeError::PlanContract("this plane has no push store".into()))?;
let at = now_for_admission().unix_timestamp().max(0).unsigned_abs();
push.unpark(task, id, at)
.await
.map_err(RuntimeError::from_store)
}
pub async fn request_cancel(
&self,
run: RunId,
actor: &crate::core::Operator,
reason: &str,
) -> Result<bool, 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(),
)));
}
if recorded_conclusion(&records).as_deref() == Some("quarantined") {
return Err(RuntimeError::CannotUnwind {
run: 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 decide_quarantine(
&self,
run: RunId,
decider: &crate::core::Operator,
reason: &str,
decision: crate::core::QuarantineDecision,
) -> Result<RunOutcome, RuntimeError> {
if reason.trim().is_empty() {
return Err(RuntimeError::PlanContract(
"answering a quarantine needs a decider and a reason — who looked, and what they \
found"
.to_owned(),
));
}
let lease = self
.store
.acquire(run, &self.owner, self.lease_ttl)
.await
.map_err(RuntimeError::from_store)?;
let appended = self
.record_decision(run, &lease, decider, reason, decision)
.await;
if let Err(e) = appended {
let _ = self.store.release_lease(run, lease.epoch).await;
return Err(e);
}
self.resume_holding(run, lease).await
}
async fn record_decision(
&self,
run: RunId,
lease: &crate::journal::Lease,
decider: &crate::core::Operator,
reason: &str,
decision: crate::core::QuarantineDecision,
) -> Result<(), 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(),
)));
}
let status = recorded_conclusion(&records);
if status.as_deref() != Some("quarantined") {
return Err(RuntimeError::NotQuarantined {
run: run.to_string(),
status: status.unwrap_or_else(|| "still running".to_owned()),
});
}
let mut append = Append::new(
run,
RecordKind::QuarantineDecided {
decider: decider.clone(),
reason: reason.to_owned(),
decision,
},
);
if let Some(case) = records.iter().find_map(|r| r.body.case) {
append = append.case(case);
}
self.store
.append(lease.epoch, vec![append])
.await
.map_err(RuntimeError::from_store)?;
Ok(())
}
pub async fn reconcile_effect(
&self,
run: RunId,
effect: crate::core::EffectKey,
assertion: crate::core::Assertion,
asserted_by: &crate::core::Operator,
note: &str,
) -> Result<(), RuntimeError> {
if note.trim().is_empty() {
return Err(RuntimeError::PlanContract(
"an assertion about an effect needs a name and a note — who established this, and \
how"
.to_owned(),
));
}
let lease = self
.store
.acquire(run, &self.owner, self.lease_ttl)
.await
.map_err(RuntimeError::from_store)?;
let result = self
.record_assertion(run, &lease, effect, assertion, asserted_by, note)
.await;
let _ = self.store.release_lease(run, lease.epoch).await;
result
}
async fn record_assertion(
&self,
run: RunId,
lease: &crate::journal::Lease,
effect: crate::core::EffectKey,
assertion: crate::core::Assertion,
asserted_by: &crate::core::Operator,
note: &str,
) -> Result<(), 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(),
)));
}
let Some(undecided) = crate::journal::undecided_effects(&records)
.into_iter()
.find(|u| u.effect == effect)
else {
return Err(RuntimeError::NotUndecided {
run: run.to_string(),
effect: effect.to_string(),
});
};
let output = match &assertion {
crate::core::Assertion::Landed(value) => Some(value.clone()),
crate::core::Assertion::DidNotHappen => None,
};
let mut append = Append::new(
run,
RecordKind::EffectReconciled {
disposition: assertion.disposition(),
declared: output
.is_some()
.then(crate::core::DeclaredOutput::untrusted),
output,
spend: Spend::ZERO,
detail: Some(note.to_owned()),
asserted_by: Some(asserted_by.clone()),
},
)
.step(undecided.step)
.phase(undecided.phase)
.effect(effect);
if let Some(case) = records.iter().find_map(|r| r.body.case) {
append = append.case(case);
}
self.store
.append(lease.epoch, vec![append])
.await
.map_err(RuntimeError::from_store)?;
tracing::info!(
target: telemetry::RECONCILED,
%run,
step = %undecided.step,
verdict = ?assertion.disposition(),
actor = %asserted_by,
);
self.meter
.count(metrics::RECONCILIATIONS, assertion.disposition().as_str());
Ok(())
}
pub async fn undecided(&self, run: RunId) -> Result<Vec<crate::core::Undecided>, RuntimeError> {
let records = self
.store
.read(run, 1)
.await
.map_err(RuntimeError::from_store)?;
Ok(crate::journal::undecided_effects(&records))
}
pub async fn set_halt(
&self,
scope: &crate::quota::HaltScope,
by: &crate::core::Operator,
at: crate::core::Timestamp,
reason: &str,
) -> Result<(), RuntimeError> {
let quotas = self.quota_store()?;
quotas
.set_halt(scope, by, at, reason)
.await
.map_err(RuntimeError::Store)
}
pub async fn lift_halt(&self, scope: &crate::quota::HaltScope) -> Result<bool, RuntimeError> {
let quotas = self.quota_store()?;
quotas.lift_halt(scope).await.map_err(RuntimeError::Store)
}
fn quota_store(&self) -> Result<&std::sync::Arc<dyn crate::quota::QuotaStore>, RuntimeError> {
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(),
))
})
}
pub async fn halts(&self) -> Result<Vec<crate::quota::Halt>, 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.halts().await.map_err(RuntimeError::Store)
}
pub async fn running_runs(&self, limit: usize) -> Result<Vec<LiveRun>, RuntimeError> {
let quotas = self.quotas.as_ref().ok_or_else(|| {
RuntimeError::Store(crate::core::StoreError::Backend(
"no quota store is wired, so nothing tracks which runs hold this \
tenant's admission slots"
.to_owned(),
))
})?;
let live = quotas
.running_runs(limit)
.await
.map_err(RuntimeError::Store)?;
let stranded: std::collections::BTreeSet<RunId> = self
.store
.abandoned_runs(limit)
.await
.map_err(RuntimeError::Store)?
.into_iter()
.collect();
let mut out = Vec::with_capacity(live.len());
for run in live {
let head = self
.store
.read_page(run, 1, ATTRIBUTION_RECORDS)
.await
.map_err(RuntimeError::Store)?;
let governed_by = head.iter().find_map(|r| match r.kind() {
crate::journal::RecordKind::RunAdmitted { governed_by, .. } => {
governed_by.as_ref().map(|id| (**id).clone())
}
_ => None,
});
let subject = head.iter().find_map(|r| match r.kind() {
crate::journal::RecordKind::IdentityBound { chain } => {
chain.first().map(|p| p.id.clone())
}
_ => None,
});
out.push(LiveRun {
stranded: stranded.contains(&run),
run,
governed_by,
subject,
});
}
Ok(out)
}
pub async fn drain(&self, grace: std::time::Duration) -> super::DrainReport {
self.inflight.drain(grace).await
}
#[must_use]
pub fn is_draining(&self) -> bool {
self.inflight.is_draining()
}
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: &crate::core::Operator,
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.clone(),
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::RunConcluded {
outcome: BREAK_GLASS_OUTCOME.to_owned(),
reason: Some(reason.to_owned()),
exhaustion: None,
live_spend: Spend::default(),
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(crate) fn witnessing(&self) -> Option<&Arc<Witnessing>> {
self.witnesses.as_ref()
}
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 withdrawn_authority(
&self,
identity: Option<&crate::core::Delegation>,
) -> Result<Option<Withdrawal>, RuntimeError> {
let (Some(quotas), Some(chain)) = (self.quotas.as_ref(), identity) else {
return Ok(None);
};
let subject = chain.subject().id.as_str();
let halts = quotas.halts().await.map_err(RuntimeError::Store)?;
Ok(halts.into_iter().find_map(|halt| {
(halt.scope.withdrawn_subject() == Some(subject)).then(|| Withdrawal {
subject: subject.to_owned(),
reason: halt.reason,
by: halt.by,
})
}))
}
#[allow(clippy::too_many_arguments)]
async fn withhold(
&self,
cx: Unwind<'_>,
withdrawal: Withdrawal,
record: bool,
completed: &[(StepId, Capability)],
outputs: &BTreeMap<StepId, Tainted<Value>>,
cursor: &mut ReplayCursor,
case_id: Option<crate::core::CaseId>,
quota: &QuotaPass,
) -> Result<RunOutcome, RuntimeError> {
let Withdrawal {
subject,
reason,
by,
} = withdrawal;
if record {
let stamped = (cx.stamp)(Append::new(
cx.run,
RecordKind::AuthorityWithheld {
subject: subject.clone(),
reason: reason.clone(),
by: by.clone(),
},
));
self.store
.append(cx.epoch, vec![stamped])
.await
.map_err(RuntimeError::from_store)?;
}
self.stop(
cx,
RunStatus::Withheld { subject, reason },
completed,
outputs,
cursor,
case_id,
quota,
)
.await
}
async fn check_quota(
&self,
run: RunId,
governed_by: Option<&crate::journal::AgentIdentity>,
subject: Option<&str>,
at: crate::core::Timestamp,
) -> Result<QuotaPass, RuntimeError> {
let Some(quotas) = self.quotas.as_ref() else {
return Ok(QuotaPass::disabled());
};
let pass = QuotaPass::at(&self.quota, at, true);
match quotas.halts().await {
Ok(halts) => {
if let Some(halt) = halts
.into_iter()
.filter(|halt| halt.scope.covers(governed_by, subject))
.max_by(|a, b| a.scope.cmp(&b.scope))
{
return Err(RuntimeError::QuotaExceeded(
crate::quota::QuotaError::Halted {
tenant: self.tenant.as_str().to_owned(),
scope: halt.scope,
reason: halt.reason,
},
));
}
}
Err(e) => {
return Err(RuntimeError::QuotaExceeded(
crate::quota::QuotaError::Unavailable(e.to_string()),
));
}
}
if let Some(period) = pass.period() {
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, at)
.await
.map_err(RuntimeError::QuotaExceeded)?;
Ok(pass)
}
async fn settle_quota(
&self,
run: RunId,
epoch: crate::core::Epoch,
spend: Spend,
pass: &QuotaPass,
) -> Result<(), RuntimeError> {
if !pass.enabled {
return Ok(());
}
let Some(quotas) = self.quotas.as_ref() else {
return Ok(());
};
let settlement = pass
.settlement(run, epoch, spend)
.expect("an enabled quota pass has a settlement");
quotas
.settle(&settlement)
.await
.map_err(|error| RuntimeError::QuotaSettlementPending {
run: run.to_string(),
epoch,
detail: error.to_string(),
})
}
pub(crate) async fn release_empty_quota_reservation(
&self,
run: RunId,
) -> Result<(), crate::core::StoreError> {
match self.quotas.as_ref() {
Some(quotas) => quotas.release(run).await,
None => Ok(()),
}
}
async fn settle_recorded_quota_passes(
&self,
run: RunId,
records: &[Record],
) -> Result<(), RuntimeError> {
let Some(quotas) = self.quotas.as_ref() else {
if records
.iter()
.any(|record| matches!(record.kind(), RecordKind::QuotaPassStarted { .. }))
{
return Err(RuntimeError::PlanContract(
"this run has unsettled quota passes but the resuming plane has no quota store"
.to_owned(),
));
}
return Ok(());
};
for record in records {
let RecordKind::QuotaPassStarted {
period,
release_slot,
} = record.kind()
else {
continue;
};
let epoch = record.body.epoch;
let spend = records
.iter()
.rev()
.find_map(|candidate| (candidate.body.epoch == epoch).then_some(candidate.kind()))
.and_then(|kind| match kind {
RecordKind::RunConcluded { live_spend, .. } => Some(*live_spend),
_ => None,
})
.unwrap_or_else(|| spend_recorded_in_epoch(records, epoch));
quotas
.settle(&crate::quota::QuotaSettlement {
run,
epoch,
period: period.clone(),
spend,
release_slot: *release_slot,
})
.await
.map_err(|error| RuntimeError::QuotaSettlementPending {
run: run.to_string(),
epoch,
detail: error.to_string(),
})?;
}
Ok(())
}
async fn recover_quota_before_resume(
&self,
run: RunId,
records: &[Record],
lease: &crate::journal::Lease,
recovering: bool,
) -> Result<Option<RunOutcome>, RuntimeError> {
self.settle_recorded_quota_passes(run, records).await?;
if !recovering {
return Ok(None);
}
let Some(last) = records.last() else {
return Ok(None);
};
let same_quota_pass = records.iter().any(|record| {
record.body.epoch == last.body.epoch
&& matches!(record.kind(), RecordKind::QuotaPassStarted { .. })
});
let RecordKind::RunConcluded {
outcome,
reason,
exhaustion,
..
} = last.kind()
else {
return Ok(None);
};
if !same_quota_pass {
return Ok(None);
}
let status = recorded_status(outcome, reason.as_deref(), exhaustion.as_ref(), records);
let chain_head = if SEALED_OUTCOMES.contains(&outcome.as_str()) {
self.store
.seal(run, lease.epoch, outcome)
.await
.map_err(RuntimeError::from_store)?
} else {
self.store
.head(run)
.await
.map_err(RuntimeError::from_store)?
.hash
};
Ok(Some(RunOutcome {
run_id: run,
status,
chain_head,
consumed: Consumed::default(),
output: None,
}))
}
async fn closed_resume_outcome(
&self,
run: RunId,
records: &[Record],
lease: &crate::journal::Lease,
) -> Result<Option<RunOutcome>, RuntimeError> {
let Some(status) = resume_is_closed(records) else {
return Ok(None);
};
let outcome = recorded_conclusion(records).unwrap_or_else(|| status.as_str().to_owned());
let chain_head = if SEALED_OUTCOMES.contains(&outcome.as_str()) {
self.store
.seal(run, lease.epoch, &outcome)
.await
.map_err(RuntimeError::from_store)?
} else {
self.store
.head(run)
.await
.map_err(RuntimeError::from_store)?
.hash
};
Ok(Some(RunOutcome {
run_id: run,
status,
chain_head,
output: None,
consumed: Consumed::default(),
}))
}
async fn abandon_if_decided(
&self,
run: RunId,
records: &[Record],
lease: &crate::journal::Lease,
) -> Result<Option<RunOutcome>, RuntimeError> {
let Some(status) = pending_abandonment(records) else {
return Ok(None);
};
self.conclude(
run,
lease.epoch,
status,
None,
true,
self.recorded_case(records).map(|c| c.id()),
Consumed::default(),
Spend::ZERO,
&QuotaPass::disabled(),
)
.await
.map(Some)
}
async fn start_replay_quota_pass(
&self,
run: RunId,
mode: Mode,
lease: Option<&crate::journal::Lease>,
) -> Result<QuotaPass, RuntimeError> {
if mode != Mode::Resume || self.quotas.is_none() {
return Ok(QuotaPass::disabled());
}
let pass = QuotaPass::at(&self.quota, now_for_admission(), false);
self.store
.append(
lease.expect("resume holds a lease").epoch,
vec![Append::new(
run,
pass.started().expect("an enabled pass has a marker"),
)],
)
.await
.map_err(RuntimeError::from_store)?;
Ok(pass)
}
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(),
})
}
#[cfg(feature = "mcp-server")]
pub(crate) fn provides(&self, target: &str) -> bool {
self.resolve(target).is_ok()
}
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, RunTerms::default()).await
}
pub async fn spawn(
self: &Arc<Self>,
target: &str,
input: Tainted<Value>,
) -> Result<RunId, RuntimeError> {
self.spawn_bound(target, input, RunTerms::default()).await
}
async fn spawn_bound(
self: &Arc<Self>,
target: &str,
input: Tainted<Value>,
terms: RunTerms,
) -> Result<RunId, RuntimeError> {
let skill = self.resolve(target)?;
let capability = first_capability(&skill.descriptor());
let run = RunId::generate();
let ticket = self.inflight.enter(run);
let admitted = self
.admit_only(
run,
PlanIR::single(capability),
input,
terms,
Entry::Outside,
)
.await?;
let plane = Arc::clone(self);
tokio::spawn(async move {
let _ticket = ticket;
plane.execute_admitted(admitted).await
});
Ok(run)
}
pub async fn run_in_case(
&self,
target: &str,
input: Tainted<Value>,
case: crate::core::CaseId,
) -> Result<RunOutcome, RuntimeError> {
self.admit(target, input, RunTerms::bound(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,
RunTerms::default().correlated(case_kind, keys),
)
.await
}
pub async fn run_once(
&self,
target: &str,
input: Tainted<Value>,
idempotency_key: &str,
) -> Result<Admission, RuntimeError> {
self.admit_once(target, input, RunTerms::default().keyed(idempotency_key))
.await
}
pub async fn run_correlated_once(
&self,
target: &str,
input: Tainted<Value>,
case_kind: &str,
keys: &[CorrelationKey],
idempotency_key: &str,
) -> Result<Admission, RuntimeError> {
self.admit_once(
target,
input,
RunTerms::default()
.correlated(case_kind, keys)
.keyed(idempotency_key),
)
.await
}
pub async fn spawn_correlated_once(
self: &Arc<Self>,
target: &str,
input: Tainted<Value>,
case_kind: &str,
keys: &[CorrelationKey],
idempotency_key: &str,
) -> Result<Spawned, RuntimeError> {
self.spawn_bound_once(
target,
input,
RunTerms::default()
.correlated(case_kind, keys)
.keyed(idempotency_key),
)
.await
}
pub async fn run_plan(
&self,
plan: PlanIR,
input: Tainted<Value>,
) -> Result<RunOutcome, RuntimeError> {
self.admit_plan(plan, input, RunTerms::default()).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, RunTerms::default().correlated(case_kind, keys))
.await
}
pub async fn run_under(
&self,
target: &str,
input: Tainted<Value>,
terms: RunTerms,
) -> Result<Admission, RuntimeError> {
if terms.idempotency_key.is_some() {
self.admit_once(target, input, terms).await
} else {
self.admit(target, input, terms).await.map(Admission::Fresh)
}
}
pub async fn spawn_under(
self: &Arc<Self>,
target: &str,
input: Tainted<Value>,
terms: RunTerms,
) -> Result<Spawned, RuntimeError> {
if terms.idempotency_key.is_some() {
self.spawn_bound_once(target, input, terms).await
} else {
let run = self.spawn_bound(target, input, terms).await?;
Ok(Spawned { run, fresh: true })
}
}
pub async fn run_plan_under(
&self,
plan: PlanIR,
input: Tainted<Value>,
terms: RunTerms,
) -> Result<RunOutcome, RuntimeError> {
self.admit_plan(plan, input, terms).await
}
async fn admit_once(
&self,
target: &str,
input: Tainted<Value>,
terms: RunTerms,
) -> Result<Admission, RuntimeError> {
let key = terms
.idempotency_key
.clone()
.expect("admit_once is only reached with a key");
validate_admission_key(&key)?;
if let Some(held) = self.holder_of(&key).await? {
return self.answer_with(held).await;
}
match self.admit(target, input, terms).await {
Ok(outcome) => Ok(Admission::Fresh(outcome)),
Err(RuntimeError::Store(crate::core::StoreError::DuplicateAdmission {
run, ..
})) => self.answer_with(parse_holder(&run)?).await,
Err(e) => Err(e),
}
}
async fn spawn_bound_once(
self: &Arc<Self>,
target: &str,
input: Tainted<Value>,
terms: RunTerms,
) -> Result<Spawned, RuntimeError> {
let key = terms
.idempotency_key
.clone()
.expect("spawn_bound_once is only reached with a key");
validate_admission_key(&key)?;
if let Some(run) = self.holder_of(&key).await? {
return Ok(Spawned { run, fresh: false });
}
match self.spawn_bound(target, input, terms).await {
Ok(run) => Ok(Spawned { run, fresh: true }),
Err(RuntimeError::Store(crate::core::StoreError::DuplicateAdmission {
run, ..
})) => Ok(Spawned {
run: parse_holder(&run)?,
fresh: false,
}),
Err(e) => Err(e),
}
}
async fn holder_of(&self, key: &str) -> Result<Option<RunId>, RuntimeError> {
self.store
.admitted_as(key)
.await
.map_err(RuntimeError::from_store)
}
async fn answer_with(&self, held: RunId) -> Result<Admission, RuntimeError> {
Ok(match self.recorded_outcome(held).await? {
Some(outcome) => Admission::Replayed(outcome),
None => Admission::InFlight(held),
})
}
pub async fn recorded_outcome(&self, run: RunId) -> Result<Option<RunOutcome>, RuntimeError> {
let records = self
.store
.read(run, 1)
.await
.map_err(RuntimeError::from_store)?;
let Some(last) = records.last() else {
return Ok(None);
};
let Some(status) = observed_status(&records) else {
return Ok(None);
};
Ok(Some(RunOutcome {
run_id: run,
status,
chain_head: last.hash,
consumed: Consumed::default(),
output: None,
}))
}
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 {
let contract = crate::plan::Contract::new(self.by_capability.keys().cloned());
if self.require_verifier {
contract.require_verifier()
} else {
contract
}
}
async fn admit(
&self,
target: &str,
input: Tainted<Value>,
terms: RunTerms,
) -> Result<RunOutcome, RuntimeError> {
let skill = self.resolve(target)?;
let capability = first_capability(&skill.descriptor());
self.admit_plan(PlanIR::single(capability), input, terms)
.await
}
async fn admit_plan(
&self,
plan: PlanIR,
input: Tainted<Value>,
terms: RunTerms,
) -> Result<RunOutcome, RuntimeError> {
self.admit_plan_as(RunId::generate(), plan, input, terms, Entry::Outside)
.await
}
pub(crate) async fn commission_run(
&self,
target: &str,
input: Tainted<Value>,
terms: RunTerms,
) -> Result<RunOutcome, RuntimeError> {
let skill = self.resolve(target)?;
let capability = first_capability(&skill.descriptor());
self.admit_plan_as(
RunId::generate(),
PlanIR::single(capability),
input,
terms,
Entry::Nested,
)
.await
}
fn bind_identity(
run: RunId,
chain: Option<&crate::core::Delegation>,
records: &mut Vec<Append>,
) {
if let Some(chain) = chain {
records.push(Append::new(
run,
RecordKind::IdentityBound {
chain: chain.links().cloned().collect(),
},
));
}
}
fn authorize_scope(
&self,
plan: &PlanIR,
chain: Option<&crate::core::Delegation>,
) -> Result<(), RuntimeError> {
let Some(chain) = chain else {
return Ok(());
};
chain.admissible(self.tenant.as_str(), now_for_admission())?;
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>,
chain: Option<&crate::core::Delegation>,
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, chain);
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>,
idempotency_key: Option<String>,
) -> RecordKind {
RecordKind::RunAdmitted {
capability: capability.to_owned(),
governed_by: governed_by.map(Box::new),
input: input.peek().clone(),
input_label: input.label().clone(),
policy_bundle: self.policy.as_ref().map(|p| Box::new(p.bundle())),
canon: crate::core::canon::VERSION,
idempotency_key,
}
}
#[allow(clippy::too_many_lines)]
async fn admit_only(
&self,
run: RunId,
plan: PlanIR,
input: Tainted<Value>,
terms: RunTerms,
entry: Entry,
) -> Result<Admitted, RuntimeError> {
if entry == Entry::Outside && self.inflight.is_draining() {
return Err(RuntimeError::Draining);
}
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;
let chain = terms.acting_as.as_ref().or(self.identity.as_ref());
self.authorize_scope(&plan, chain)?;
self.authorize_admission(&agent, governed_by.as_ref(), chain, input.peek())?;
self.admit_reserved(run, plan, input, terms, agent, governed_by)
.await
}
async fn admit_reserved(
&self,
run: RunId,
plan: PlanIR,
input: Tainted<Value>,
terms: RunTerms,
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 quota = match self
.check_quota(
run,
governed_by.as_ref(),
terms
.acting_as
.as_ref()
.or(self.identity.as_ref())
.map(|c| c.subject().id.as_str()),
now_for_admission(),
)
.await
{
Ok(pass) => pass,
Err(error) => {
let _ = self.store.release_lease(run, lease.epoch).await;
return Err(error);
}
};
let out = self
.admit_under_lease(run, &lease, plan, input, terms, &agent, governed_by, quota)
.await;
if out.is_err() {
if let Some(quotas) = self.quotas.as_ref()
&& let Err(error) = quotas.release(run).await
{
tracing::debug!(%run, %error, "could not release the quota slot of a failed admission");
}
if let Err(error) = self.store.release_lease(run, lease.epoch).await {
tracing::debug!(
%run,
%error,
"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>,
terms: RunTerms,
agent: &str,
governed_by: Option<crate::journal::AgentIdentity>,
quota: QuotaPass,
) -> Result<Admitted, RuntimeError> {
let RunTerms {
case,
idempotency_key,
acting_as,
} = terms;
let chain = acting_as.as_ref().or(self.identity.as_ref());
let mut records = vec![
Append::new(
run,
self.admission(agent, governed_by, &input, idempotency_key),
),
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, chain, &mut records);
if let Some(started) = quota.started() {
records.push(Append::new(run, started));
}
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,
};
if let Err(e) = self.store.append(lease.epoch, records).await {
if let Some(ctx) = case_ctx.as_ref()
&& let Err(detach) = ctx.cases.detach_run(ctx.case_id, run).await
{
tracing::warn!(
%run,
case = %ctx.case_id,
error = %detach,
"a failed admission stayed attached to its case; the case lists a \
run with no journal"
);
}
return Err(RuntimeError::from_store(e));
}
#[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::RunConcluded {
outcome: RunStatus::Failed(String::new()).as_str().to_owned(),
reason: Some(format!(
"the outbox registration was refused ({open_failed}), so \
this run was concluded without ever executing"
)),
exhaustion: None,
live_spend: Spend::default(),
chain_head: head.hash,
},
),
],
)
.await
.map_err(RuntimeError::from_store)?;
return Err(RuntimeError::from_store(open_failed));
}
Ok(Admitted {
run,
epoch: lease.epoch,
quota,
budget: self.budget_for(agent),
agent: agent.to_owned(),
plan,
input,
case: case_ctx,
identity: chain.cloned(),
})
}
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,
quota: a.quota,
plan: &a.plan,
input: a.input,
mode: Mode::Live,
case: a.case,
budget: a.budget,
agent: a.agent,
identity: a.identity,
refusal: None,
withheld_subject: 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>,
terms: RunTerms,
entry: Entry,
) -> Result<RunOutcome, RuntimeError> {
let admitted = self.admit_only(run, plan, input, terms, entry).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.as_deref().cloned())
}
_ => 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, false).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), false)
.await
}
pub(crate) async fn recover_abandoned_run(
&self,
run: RunId,
) -> Result<RunOutcome, RuntimeError> {
let lease = self
.store
.acquire(run, &self.owner, self.lease_ttl)
.await
.map_err(RuntimeError::from_store)?;
self.replay_releasing(run, Mode::Resume, Some(lease), true)
.await
}
async fn replay_releasing(
&self,
run: RunId,
mode: Mode,
lease: Option<crate::journal::Lease>,
recovering: bool,
) -> Result<RunOutcome, RuntimeError> {
let outcome = self
.replay_under(run, mode, lease.as_ref(), recovering)
.await;
let settlement_pending =
matches!(outcome, Err(RuntimeError::QuotaSettlementPending { .. }));
if !settlement_pending
&& 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>,
recovering: bool,
) -> 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)?;
if mode == Mode::Resume
&& let Some(outcome) = self
.recover_quota_before_resume(
run,
&records,
lease.expect("resume holds a lease"),
recovering,
)
.await?
{
return Ok(outcome);
}
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(outcome) = self
.closed_resume_outcome(run, &records, lease.expect("resume holds a lease"))
.await?
{
return Ok(outcome);
}
if mode == Mode::Resume
&& let Some(outcome) = self
.abandon_if_decided(run, &records, lease.expect("resume holds a lease"))
.await?
{
return Ok(outcome);
}
if mode == Mode::Resume {
refuse_resume_over_undone_work(run, &records)?;
}
if mode == Mode::Resume {
self.ensure_resume_policy_bundle(&records)?;
}
let quota = self.start_replay_quota_pass(run, mode, lease).await?;
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,
quota,
plan: &plan,
input,
mode,
case: case_ctx,
budget: {
let recorded = recorded_agent(&records);
self.budget_for(&recorded)
},
agent: recorded_agent(&records),
identity: recorded_chain(&records)?,
refusal: recorded_step_refusal(&records),
withheld_subject: standing_withholding(&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::GEN_AI_AGENT_NAME } = plan.agent.as_str(),
{ telemetry::MODE } = telemetry::mode_str(plan.mode),
{ telemetry::GEN_AI_CONVERSATION_ID } = plan
.case
.as_ref()
.map(|c| super::ctx::CaseContext::id(c).to_string()),
{ telemetry::TENANT } = tracing::field::Empty,
{ telemetry::OUTCOME } = tracing::field::Empty,
{ telemetry::SEMCONV } = telemetry::SEMCONV_VERSION,
);
if !self.meter.tenant().is_empty() {
span.record(telemetry::TENANT, self.meter.tenant());
}
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,
quota,
plan: ir,
input,
mode,
case,
budget,
agent,
identity,
refusal: recorded_refusal,
withheld_subject,
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 mut withheld: Option<String> = withheld_subject;
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,
identity: identity.as_ref(),
run,
epoch,
ir: ¤t,
mode,
case: case.clone(),
ledger: &ledger,
writing,
stamp: &stamp,
},
status,
&completed,
&outputs,
cursor,
case_id,
"a,
)
.await;
}
if writing {
let standing = self.withdrawn_authority(identity.as_ref()).await?;
match (standing, withheld.take()) {
(None, Some(prior)) => {
self.store
.append(
epoch,
vec![stamp(Append::new(
run,
RecordKind::AuthorityRestored { subject: prior },
))],
)
.await
.map_err(RuntimeError::from_store)?;
}
(Some(withdrawal), Some(_prior)) => {
return self
.withhold(
Unwind {
agent: &agent,
identity: identity.as_ref(),
run,
epoch,
ir: ¤t,
mode,
case: case.clone(),
ledger: &ledger,
writing,
stamp: &stamp,
},
withdrawal,
false,
&completed,
&outputs,
cursor,
case_id,
"a,
)
.await;
}
(Some(withdrawal), None) => {
return self
.withhold(
Unwind {
agent: &agent,
identity: identity.as_ref(),
run,
epoch,
ir: ¤t,
mode,
case: case.clone(),
ledger: &ledger,
writing,
stamp: &stamp,
},
withdrawal,
true,
&completed,
&outputs,
cursor,
case_id,
"a,
)
.await;
}
(None, None) => {}
}
}
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,
identity: identity.as_ref(),
run,
epoch,
ir: ¤t,
mode,
case: &case,
ledger: &ledger,
writing,
stamp: &stamp,
input: &input,
outputs: &outputs,
started: &started,
finished: &finished,
recorded_groups: &recorded_groups,
parallelism: budget.parallelism(),
},
)
.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,
Consumed::default(),
Spend::default(),
"a,
)
.await;
}
}
}
let mut stopped = apply(¤t, outcomes, &mut done, &mut completed, &mut outputs)
.or_else(|| 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,
identity: identity.as_ref(),
run,
epoch,
ir: ¤t,
mode,
case: case.clone(),
ledger: &ledger,
writing,
stamp: &stamp,
},
status,
&completed,
&outputs,
cursor,
case_id,
"a,
)
.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,
Consumed::default(),
Spend::default(),
"a,
)
.await;
}
let (consumed, live_spend) = {
let l = ledger.lock().expect("budget mutex");
(l.consumed(), l.live_spend())
};
self.conclude(
run,
epoch,
completion(¤t, &done),
run_output(¤t, &outputs),
writing,
case_id,
consumed,
live_spend,
"a,
)
.await
}
async fn observe_clock(
cx: &mut StepCtx<'_>,
ledger: &Arc<std::sync::Mutex<Ledger>>,
) -> Result<(), crate::core::StepError> {
if !ledger.lock().expect("budget mutex").tracks_wallclock() {
return Ok(());
}
let at = cx.now().await?;
ledger.lock().expect("budget mutex").observe_clock(at);
Ok(())
}
async fn dispatch(
&self,
admitted: &[StepId],
cursor: &mut ReplayCursor,
batch: Batch<'_>,
) -> Vec<Dispatched> {
use futures_util::StreamExt;
let slices: Vec<(StepId, StepCursor)> = admitted
.iter()
.map(|&s| (s, cursor.take(s, Phase::Forward)))
.collect();
let width = batch.parallelism.max(1);
futures_util::stream::iter(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,
identity: batch.identity,
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))
}))
.buffered(width)
.collect()
.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(),
));
}
if next.reason.as_deref().is_none_or(str::is_empty) {
return Err(RunStatus::Failed(
"the successor plan carries no reason — use `PlanIR::succeed_with`, or the \
journal freezes a plan that replaced another with nothing saying why"
.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(admitted.len())
{
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(admitted.len()),
};
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>,
quota: &QuotaPass,
) -> 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,
Consumed::default(),
Spend::default(),
quota,
)
.await;
}
let (consumed, live_spend) = {
let l = ledger.lock().expect("budget mutex");
(l.consumed(), l.live_spend())
};
self.conclude(
run, epoch, unwound, output, writing, case_id, consumed, live_spend, quota,
)
.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),
})
}
fn plane_frame(&self, step: StepFrame) -> super::ctx::Frame {
let StepFrame {
run,
epoch,
step,
phase,
mode,
case,
ledger,
identity,
agent,
#[cfg(feature = "manifest")]
manifest,
recorded_groups,
} = step;
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(),
peers: self.peers.clone(),
#[cfg(feature = "manifest")]
tools: self.tools.clone(),
journal_ceiling: self.journal_ceiling,
#[cfg(feature = "manifest")]
egress: self.egress.clone(),
meter: self.meter.clone(),
#[cfg(feature = "keyring")]
keyring: self.keyring.clone(),
tenant: self.tenant.clone(),
ledger,
policy: self.policy.clone(),
identity,
agent,
plane: self.self_ref.clone(),
#[cfg(feature = "manifest")]
manifest,
signer: self.signer.clone(),
recorded_groups,
}
}
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),
self.plane_frame(StepFrame {
run: cx.run,
epoch: cx.epoch,
step,
phase: Phase::Compensating,
mode: cx.mode,
case: cx.case.clone(),
ledger: Arc::clone(cx.ledger),
identity: cx.identity.cloned(),
agent: cx.agent.to_owned(),
#[cfg(feature = "manifest")]
manifest: self.governing(skill),
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)?;
Ok(crate::journal::undecided_effects(&records)
.into_iter()
.map(|u| u.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 } = ctx.phase.as_str(),
{ 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 {
identity,
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,
self.plane_frame(StepFrame {
run,
epoch,
step,
phase,
mode,
case,
ledger: Arc::clone(ledger),
identity: identity.cloned(),
agent: agent.to_owned(),
#[cfg(feature = "manifest")]
manifest: self.governing(skill.as_ref()),
recorded_groups,
}),
);
let result = match Self::observe_clock(&mut cx, ledger).await {
Ok(()) => skill.invoke(&mut cx, step_input).await,
Err(e) => Err(crate::core::SkillError::Step(e)),
};
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>,
consumed: Consumed,
live_spend: Spend,
quota: &QuotaPass,
) -> 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::RunConcluded { outcome, .. }
if outcome == status.as_str()
)
});
if repeated {
before.hash
} else {
let mut sealed = Append::new(
run,
RecordKind::RunConcluded {
outcome: status.as_str().to_owned(),
reason: status.reason().map(std::borrow::Cow::into_owned),
exhaustion: match &status {
RunStatus::Exhausted(limit) => Some(limit.clone()),
_ => None,
},
live_spend,
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)?;
concluded.last().map_or(before.hash, |r| r.hash)
}
} else {
self.store
.head(run)
.await
.map_err(RuntimeError::from_store)?
.hash
};
if writing {
self.settle_quota(run, epoch, live_spend, quota).await?;
let chain_head = if status.seals() {
self.store
.seal(run, epoch, status.as_str())
.await
.map_err(RuntimeError::from_store)?
} else {
chain_head
};
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"
);
}
announce(&self.meter, run, &status);
return Ok(RunOutcome {
run_id: run,
status,
chain_head,
consumed,
output,
});
}
Ok(RunOutcome {
run_id: run,
status,
chain_head,
consumed,
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 spend_recorded_in_epoch(records: &[Record], epoch: crate::core::Epoch) -> Spend {
records
.iter()
.filter(|record| record.body.epoch == epoch)
.fold(Spend::default(), |total, record| {
let spend = match record.kind() {
RecordKind::EffectDone { spend, .. }
| RecordKind::EffectFailed { spend, .. }
| RecordKind::EffectReconciled { spend, .. } => *spend,
_ => Spend::default(),
};
total.plus(spend)
})
}
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::Abandoned { actor, reason } => {
tracing::error!(target: telemetry::ABANDONED, %run, %actor, %reason);
}
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::Abandoned { .. } | RunStatus::Failed(_) => 2,
RunStatus::Exhausted(_) | RunStatus::Withheld { .. } => 1,
RunStatus::Replanning(_)
| RunStatus::Suspended(_)
| RunStatus::Succeeded
| RunStatus::Swept
| RunStatus::BrokeGlass { .. } => 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,
identity: Option<&'a crate::core::Delegation>,
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>,
parallelism: usize,
}
struct Withdrawal {
subject: String,
reason: String,
by: crate::core::Operator,
}
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> {
identity: Option<&'a crate::core::Delegation>,
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> {
identity: Option<&'a crate::core::Delegation>,
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>,
}
struct StepFrame {
run: RunId,
epoch: crate::core::Epoch,
step: StepId,
phase: Phase,
mode: Mode,
case: Option<CaseContext>,
ledger: Arc<std::sync::Mutex<Ledger>>,
identity: Option<crate::core::Delegation>,
agent: String,
#[cfg(feature = "manifest")]
manifest: Option<Arc<crate::manifest::Manifest>>,
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 standing_withholding(records: &[Record]) -> Option<String> {
let mut standing: Option<String> = None;
for r in records {
match r.kind() {
RecordKind::AuthorityWithheld { subject, .. } => standing = Some(subject.clone()),
RecordKind::AuthorityRestored { .. } => standing = None,
_ => {}
}
}
standing
}
fn recorded_conclusion(records: &[Record]) -> Option<String> {
records.iter().rev().find_map(|r| match r.kind() {
RecordKind::RunConcluded { 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_chain(records: &[Record]) -> Result<Option<crate::core::Delegation>, RuntimeError> {
records
.iter()
.find_map(|r| match r.kind() {
RecordKind::IdentityBound { chain } => Some(chain.clone()),
_ => None,
})
.map(crate::core::Delegation::rehydrate)
.transpose()
.map_err(RuntimeError::Delegation)
}
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, "");
}
SkillError::Step(StepError::Unreproducible { what, detail }) => {
tracing::error!(target: telemetry::UNREPRODUCIBLE, %what, %detail);
meter.count(metrics::UNREPRODUCIBLE, "");
}
_ => {}
}
let untrustworthy = matches!(
e,
SkillError::Step(
StepError::NonDeterminism { .. }
| StepError::ReplayOverrun { .. }
| StepError::Undecidable { .. }
| StepError::GroupUnsettled { .. }
| StepError::Unreproducible { .. }
)
);
if untrustworthy {
(RunStatus::Quarantined(msg), None)
} else {
(RunStatus::Failed(msg), None)
}
}
}
}
struct Execution<'a> {
refusal: Option<(StepId, String, String)>,
withheld_subject: Option<String>,
successors: Vec<PlanIR>,
started: BTreeSet<StepId>,
finished: BTreeSet<StepId>,
recorded_groups: BTreeMap<(StepId, Phase, String), super::ctx::RecordedGroup>,
run: RunId,
epoch: u64,
quota: QuotaPass,
plan: &'a PlanIR,
input: Tainted<Value>,
mode: Mode,
case: Option<CaseContext>,
budget: Budget,
agent: String,
identity: Option<crate::core::Delegation>,
}
const UNATTRIBUTED: &str = "the chain records this ending and not who asked for it — an \
operator act with nobody on it cannot be answered by reading further";
fn resume_is_closed(records: &[Record]) -> Option<RunStatus> {
let outcome = records.iter().rev().find_map(|r| match r.kind() {
RecordKind::RunConcluded { outcome, .. } => Some(outcome.as_str()),
_ => None,
})?;
match outcome {
"succeeded" => Some(RunStatus::Succeeded),
"quarantined" => quarantine_decision(records).is_none().then(|| {
RunStatus::Quarantined(
"recorded as quarantined; a named person has to answer it before it can run \
again — reopen it once the doubt is resolved, or abandon it"
.into(),
)
}),
"abandoned" => Some(recorded_decider(records).map_or_else(
|| RunStatus::Quarantined(UNATTRIBUTED.to_owned()),
|actor| {
RunStatus::Abandoned {
actor,
reason: "recorded as abandoned; its outcome was never established and nothing \
was unwound"
.into(),
}
},
)),
"cancelled" => Some(recorded_canceller(records).map_or_else(
|| RunStatus::Quarantined(UNATTRIBUTED.to_owned()),
|actor| RunStatus::Cancelled {
actor,
reason: "recorded as cancelled; an operator stopped this run".into(),
},
)),
"failed" | "exhausted" | WITHHELD_OUTCOME => None,
other => Some(RunStatus::Quarantined(format!(
"recorded as '{other}', which this build does not recognise as resumable"
))),
}
}
fn refuse_resume_over_undone_work(run: RunId, records: &[Record]) -> Result<(), RuntimeError> {
let Some(outcome) = recorded_conclusion(records) else {
return Ok(());
};
if !records
.iter()
.any(|r| matches!(r.kind(), RecordKind::StepCompensated { .. }))
{
return Ok(());
}
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"
)))
}
fn quarantine_decision(records: &[Record]) -> Option<(crate::core::QuarantineDecision, &Record)> {
records
.iter()
.rev()
.take_while(|r| !matches!(r.kind(), RecordKind::RunConcluded { .. }))
.find_map(|r| match r.kind() {
RecordKind::QuarantineDecided { decision, .. } => Some((*decision, r)),
_ => None,
})
}
fn pending_abandonment(records: &[Record]) -> Option<RunStatus> {
let (decision, record) = quarantine_decision(records)?;
if decision != crate::core::QuarantineDecision::Abandon {
return None;
}
match record.kind() {
RecordKind::QuarantineDecided {
decider, reason, ..
} => Some(RunStatus::Abandoned {
actor: decider.clone(),
reason: reason.clone(),
}),
_ => None,
}
}
pub const MAX_ADMISSION_KEY_BYTES: usize = 512;
fn validate_admission_key(key: &str) -> Result<(), RuntimeError> {
if key.trim().is_empty() {
return Err(RuntimeError::PlanContract(
"an admission key is empty — an unset header or variable arrives here as \
an empty string, and accepting it would answer every later message with \
the first one's run. Use the message's own identity, such as \
`InboundEvent::dedup_key()`"
.into(),
));
}
if key.len() > MAX_ADMISSION_KEY_BYTES {
return Err(RuntimeError::PlanContract(format!(
"an admission key of {} bytes exceeds the {MAX_ADMISSION_KEY_BYTES}-byte \
limit — the key is chosen by the emitter and becomes part of a storage \
key, so its size is not theirs to decide",
key.len()
)));
}
Ok(())
}
fn parse_holder(run: &str) -> Result<RunId, RuntimeError> {
RunId::parse(run).map_err(|e| {
RuntimeError::PlanContract(format!(
"an admission key is held by run '{run}', which does not parse ({e}) — refusing rather than admitting a second run under a key whose holder this build cannot read"
))
})
}
pub(crate) fn observed_status(records: &[Record]) -> Option<RunStatus> {
Some(match records.last()?.kind() {
RecordKind::RunSuspended { reason } => RunStatus::Suspended(reason.clone()),
RecordKind::RunConcluded {
outcome,
reason,
exhaustion,
..
} => recorded_status(outcome, reason.as_deref(), exhaustion.as_ref(), records),
_ => return None,
})
}
fn recorded_status(
outcome: &str,
reason: Option<&str>,
exhaustion: Option<&crate::core::BudgetExceeded>,
records: &[Record],
) -> RunStatus {
let said = || {
reason.map_or_else(
|| "concluded without a recorded reason".to_owned(),
ToOwned::to_owned,
)
};
match outcome {
"succeeded" => RunStatus::Succeeded,
"failed" => RunStatus::Failed(said()),
"exhausted" => exhaustion.cloned().map_or_else(
|| {
RunStatus::Quarantined(
"recorded as exhausted without the typed ceiling verdict".to_owned(),
)
},
RunStatus::Exhausted,
),
"quarantined" => RunStatus::Quarantined(said()),
"cancelled" => recorded_canceller(records).map_or_else(
|| RunStatus::Quarantined(UNATTRIBUTED.to_owned()),
|actor| RunStatus::Cancelled {
actor,
reason: said(),
},
),
"abandoned" => recorded_decider(records).map_or_else(
|| RunStatus::Quarantined(UNATTRIBUTED.to_owned()),
|actor| RunStatus::Abandoned {
actor,
reason: said(),
},
),
super::sweeper::SWEEP_OUTCOME => RunStatus::Swept,
BREAK_GLASS_OUTCOME => recorded_crosser(records).map_or_else(
|| RunStatus::Quarantined(UNATTRIBUTED.to_owned()),
|actor| RunStatus::BrokeGlass {
actor,
reason: said(),
},
),
other => RunStatus::Quarantined(format!(
"recorded as '{other}', which this build does not recognise"
)),
}
}
fn recorded_canceller(records: &[Record]) -> Option<crate::core::Operator> {
records.iter().rev().find_map(|r| match r.kind() {
RecordKind::RunCancelled { actor, .. } => Some(actor.clone()),
_ => None,
})
}
fn recorded_decider(records: &[Record]) -> Option<crate::core::Operator> {
records.iter().rev().find_map(|r| match r.kind() {
RecordKind::QuarantineDecided { decider, .. } => Some(decider.clone()),
_ => None,
})
}
fn recorded_crosser(records: &[Record]) -> Option<crate::core::Operator> {
records.iter().rev().find_map(|r| match r.kind() {
RecordKind::BreakGlass { 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";
pub const WITHHELD_OUTCOME: &str = "withheld";
#[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>,
witnesses: Vec<Arc<dyn crate::journal::Witness>>,
quorum: Option<crate::journal::WitnessQuorum>,
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>,
)>,
#[cfg(feature = "manifest")]
egress: Option<Arc<crate::core::Egress>>,
journal_ceiling: Option<crate::core::Sensitivity>,
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>>,
peers: Option<Arc<PeerWiring>>,
tenant_label: super::telemetry::TenantLabel,
quotas: Option<Arc<dyn crate::quota::QuotaStore>>,
quota: crate::quota::TenantQuota,
budget: Budget,
require_verifier: bool,
#[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 = "push")]
push: Option<Arc<dyn crate::push::PushStore>>,
#[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 tenant_label(mut self, label: super::telemetry::TenantLabel) -> Self {
self.tenant_label = 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 require_verifier(mut self) -> Self {
self.require_verifier = true;
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
}
#[cfg(feature = "push")]
#[must_use]
pub fn push(mut self, push: Arc<dyn crate::push::PushStore>) -> Self {
self.push = Some(push);
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
}
#[cfg(feature = "manifest")]
#[must_use]
pub fn egress(mut self, egress: crate::core::Egress) -> Self {
self.egress = Some(Arc::new(egress));
self
}
#[must_use]
pub const fn max_sensitivity_journaled(mut self, ceiling: crate::core::Sensitivity) -> Self {
self.journal_ceiling = Some(ceiling);
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, clippy::needless_pass_by_ref_mut)]
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 peers(
mut self,
registry: crate::peers::PeerRegistry,
client: Arc<dyn crate::peers::PeerClient>,
) -> Self {
self.peers = Some(Arc::new(PeerWiring { registry, client }));
self
}
#[must_use]
pub fn witnesses(
mut self,
witnesses: Vec<Arc<dyn crate::journal::Witness>>,
quorum: crate::journal::WitnessQuorum,
) -> Self {
self.witnesses = witnesses;
self.quorum = Some(quorum);
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();
let peer_names: std::collections::BTreeSet<String> = self
.peers
.as_ref()
.map(|wiring| wiring.registry.peers().map(ToString::to_string).collect())
.unwrap_or_default();
if peer_names.contains(crate::tools::AGENT_SERVER) {
return Err(BuildError::ReservedToolServer);
}
for name in &peer_names {
if servers.iter().any(|(server, _)| server == name)
|| tools
.as_ref()
.is_some_and(|t| t.servers().any(|s| s == name))
{
return Err(BuildError::PeerIsAlsoAToolServer {
server: name.clone(),
});
}
}
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();
let reachable_elsewhere: std::collections::BTreeSet<String> =
remote_servers.union(&peer_names).cloned().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, &reachable_elsewhere)
.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
.push
.as_ref()
.map(|s| s.tenant())
.or_else(|| 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())),
("quota", self.quotas.as_ref().map(|s| s.tenant())),
("timer", self.timers.as_ref().map(|s| s.tenant())),
("batch", self.batches.as_ref().map(|s| s.tenant())),
("authority", self.authorities.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
|| self.peers.as_ref().is_some_and(|wiring| {
wiring
.registry
.grant(&crate::peers::PeerId::new(&id.server))
.is_some()
})
})
});
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 let Some(wiring) = self.peers.as_ref()
&& let Some(peer_grant) = wiring
.registry
.grant(&crate::peers::PeerId::new(&id.server))
{
if !peer_grant.scope.permits(&Capability::new(id.tool.as_str())) {
return Err(BuildError::PeerGrantOutsideScope {
agent: m.metadata.name.clone(),
peer: id.server,
capability: id.tool,
});
}
let cannot_delegate = m
.spec
.topology
.as_ref()
.is_some_and(|t| t.role == crate::manifest::Role::Specialist)
|| m.spec.security.max_delegation_depth == Some(0);
if cannot_delegate {
return Err(BuildError::PeerGrantOnASpecialist {
agent: m.metadata.name.clone(),
peer: id.server,
});
}
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,
});
}
}
}
}
let witnesses = match (self.witnesses.is_empty(), self.quorum) {
(true, None) => None,
(false, Some(quorum)) if quorum.required() <= self.witnesses.len() => {
Some(Arc::new(Witnessing {
witnesses: self.witnesses,
quorum,
submitted: std::sync::atomic::AtomicU64::new(0),
}))
}
(empty, quorum) => {
return Err(BuildError::WitnessQuorumUnreachable {
declared: quorum.map_or(0, |q| q.required()),
configured: if empty { 0 } else { self.witnesses.len() },
});
}
};
Ok(Arc::new_cyclic(|self_ref| Runtime {
self_ref: self_ref.clone(),
witnesses,
inflight: Arc::default(),
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.tenant_label, &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,
peers: self.peers,
#[cfg(feature = "manifest")]
tools: self.tools,
journal_ceiling: self.journal_ceiling,
#[cfg(feature = "manifest")]
egress: self.egress,
quotas: self.quotas,
quota: self.quota,
budget: self.budget,
require_verifier: self.require_verifier,
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")]
push: self.push,
#[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<crate::runtime::Redelivered, RuntimeError> {
let Some(events) = self.events.as_ref() else {
return Ok(crate::runtime::Redelivered::default());
};
let waiting = events
.waiting(limit)
.await
.map_err(RuntimeError::from_store)?;
let examined = waiting.len();
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(crate::runtime::Redelivered {
finished: delivered,
examined,
})
}
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()
.checked_sub(grace)
.unwrap_or_else(crate::core::first_instant);
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: crate::core::Operator::asserted("ops").expect("a name"),
reason: "stop".into(),
},
RunStatus::Abandoned {
actor: crate::core::Operator::asserted("ops").expect("a name"),
reason: "nobody could establish what happened".into(),
},
RunStatus::Swept,
RunStatus::BrokeGlass {
actor: crate::core::Operator::asserted("ops").expect("a name"),
reason: "INC-42".into(),
},
RunStatus::Withheld {
subject: "alice".into(),
reason: "credential withdrawn".into(),
},
]
}
#[cfg(test)]
mod resume_agreement_tests {
use super::{QuotaPass, 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(),
11,
"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", super::WITHHELD_OUTCOME] {
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()
);
}
}
#[test]
fn every_outcome_this_build_writes_is_one_it_can_read_back() {
for outcome in super::OUTCOMES_OF_RECORD {
let status = super::recorded_status(outcome, Some("why"), None, &sealed_as(outcome));
assert_eq!(
&status.as_str(),
outcome,
"this build seals runs as '{outcome}' and reads that back as \
'{}' — a catch-all answering `Quarantined` about a record this \
plane wrote itself",
status.as_str()
);
}
}
#[test]
fn an_ending_with_no_attribution_record_is_not_given_a_name() {
for outcome in ["cancelled", "abandoned", super::BREAK_GLASS_OUTCOME] {
let status =
super::recorded_status(outcome, Some("why"), None, &concluded_only(outcome));
assert!(
matches!(status, RunStatus::Quarantined(_)),
"'{outcome}' with no record of who asked read back as {status:?} — \
an operator act with nobody on it was given a name"
);
}
}
#[test]
fn a_break_glass_crossing_reports_the_reason_it_was_given() {
let status = super::recorded_status(
super::BREAK_GLASS_OUTCOME,
Some("INC-42: stuck settlement"),
None,
&sealed_as(super::BREAK_GLASS_OUTCOME),
);
assert_eq!(
status.reason().as_deref(),
Some("INC-42: stuck settlement"),
"a crossing must surface the reason it refused to be recorded without"
);
}
#[test]
fn the_offline_sweep_covers_every_ending_and_the_quarantine_backlog() {
for outcome in super::SEALED_OUTCOMES {
assert!(
super::OUTCOMES_OF_RECORD.contains(outcome),
"'{outcome}' is an ending this plane can reach and the export's \
default sweep would silently omit it"
);
}
assert!(
super::OUTCOMES_OF_RECORD.contains(&"quarantined"),
"a quarantine is the plane's outstanding unresolved risk and the one \
thing an auditor is most likely to have come for"
);
for open in ["failed", "exhausted", "suspended"] {
assert!(
!super::OUTCOMES_OF_RECORD.contains(&open),
"'{open}' is work in progress that a resume moves on its own — \
exporting it is a snapshot of a moving target, not a record"
);
}
}
#[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 concluded_only(outcome: &str) -> Vec<crate::journal::Record> {
let full = sealed_as(outcome);
full.into_iter()
.filter(|r| matches!(r.kind(), RecordKind::RunConcluded { .. }))
.collect()
}
fn sealed_as(outcome: &str) -> Vec<crate::journal::Record> {
use crate::core::{Digest, Epoch};
use crate::journal::{Record, RecordBody};
let run = crate::core::RunId::generate();
let who = || crate::core::Operator::asserted("ops").expect("a name");
let mut kinds: Vec<RecordKind> = match outcome {
"cancelled" => vec![RecordKind::RunCancelled {
actor: who(),
reason: "stop".into(),
}],
"abandoned" => vec![RecordKind::QuarantineDecided {
decider: who(),
reason: "nobody could establish what happened".into(),
decision: crate::core::QuarantineDecision::Abandon,
}],
super::BREAK_GLASS_OUTCOME => vec![RecordKind::BreakGlass {
actor: who(),
roles: vec!["oncall".into()],
reason: "INC-42".into(),
}],
_ => Vec::new(),
};
kinds.push(RecordKind::RunConcluded {
outcome: outcome.to_owned(),
chain_head: Digest::of(b""),
reason: None,
exhaustion: None,
live_spend: crate::core::Spend::default(),
});
let mut prev = Digest::of(b"");
let mut out = Vec::new();
for (i, kind) in kinds.into_iter().enumerate() {
let body = RecordBody {
seq: u64::try_from(i + 1).expect("a small fixture"),
run,
case: None,
step: None,
phase: super::Phase::Forward,
epoch: Epoch::default(),
v: 1,
effect_key: None,
kind,
};
let record = Record::seal(body, prev).expect("a sealed record");
prev = record.hash;
out.push(record);
}
out
}
#[test]
fn a_quota_pass_keeps_the_period_it_started_in() {
let quota = crate::quota::TenantQuota {
max_tokens_per_period: Some(1_000),
period: crate::quota::Period::Daily,
..crate::quota::TenantQuota::default()
};
let before_midnight = crate::core::Timestamp::from_unix_timestamp(86_399).unwrap();
let after_midnight = crate::core::Timestamp::from_unix_timestamp(86_400).unwrap();
let pass = QuotaPass::at("a, before_midnight, true);
assert_eq!(
pass.period(),
Some(quota.period.key_for(before_midnight).as_str())
);
assert_ne!(
pass.period(),
Some(quota.period.key_for(after_midnight).as_str()),
"settlement recomputed the period at completion, so this pass was \
authorized against one ledger and charged to another"
);
}
}