#[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: Acting,
admitted_by: Option<String>,
expected_declaration: Option<Digest>,
subjects: Vec<String>,
}
#[derive(Debug, Clone, Default)]
enum Acting {
#[default]
Plane,
Caller(crate::core::Delegation),
Nobody,
PlaneDelegate(crate::core::Delegation),
}
impl Acting {
const fn resolve<'a>(
&'a self,
plane: Option<&'a crate::core::Delegation>,
) -> Option<&'a crate::core::Delegation> {
match self {
Self::Plane => plane,
Self::Caller(chain) | Self::PlaneDelegate(chain) => Some(chain),
Self::Nobody => None,
}
}
const fn ceiling<'a>(
&'a self,
plane: Option<&'a crate::core::Delegation>,
) -> Option<&'a crate::core::Delegation> {
match self {
Self::Plane | Self::Nobody => plane,
Self::Caller(chain) | Self::PlaneDelegate(chain) => Some(chain),
}
}
}
impl RunTerms {
fn bound(case: CaseBinding) -> Self {
Self {
case: Some(case),
idempotency_key: None,
acting_as: Acting::Plane,
admitted_by: None,
expected_declaration: None,
subjects: Vec::new(),
}
}
#[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 = Acting::Caller(chain);
self
}
pub(crate) fn acting_as_plane(mut self, chain: crate::core::Delegation) -> Self {
self.acting_as = Acting::PlaneDelegate(chain);
self
}
#[must_use]
pub fn served(mut self, chain: Option<crate::core::Delegation>) -> Self {
self.acting_as = chain.map_or(Acting::Nobody, Acting::Caller);
self
}
#[must_use]
pub fn admitted_by(mut self, actor: impl Into<String>) -> Self {
self.admitted_by = Some(actor.into());
self
}
#[must_use]
pub fn subject(mut self, subject: impl Into<String>) -> Self {
self.subjects.push(subject.into());
self
}
#[must_use]
pub const fn expect_declaration(mut self, digest: Digest) -> Self {
self.expected_declaration = Some(digest);
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,
}
}
#[must_use]
pub fn into_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>,
pub quotas: std::sync::Arc<dyn crate::quota::QuotaStore>,
pub disclosures: std::sync::Arc<dyn crate::disclosure::DisclosureRegister>,
}
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, quotas, disclosures }")
}
}
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: std::sync::Arc::clone(&store) as _,
quotas: std::sync::Arc::clone(&store) as _,
disclosures: 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::Observed => write!(
f,
"it is a session this plane observed rather than executed, so it \
neither succeeded nor failed"
),
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::HaltLifted { actor } => write!(f, "{actor} lifted an emergency stop"),
RunStatus::HoldReleased { actor } => write!(f, "{actor} released a legal hold"),
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,
},
HaltLifted {
actor: crate::core::Operator,
},
HoldReleased {
actor: crate::core::Operator,
},
Withheld {
subject: String,
reason: String,
},
Observed,
}
impl RunStatus {
#[must_use]
pub fn reason(&self) -> Option<std::borrow::Cow<'_, str>> {
use std::borrow::Cow;
match self {
Self::Succeeded
| Self::Swept
| Self::Observed
| Self::HaltLifted { .. }
| Self::HoldReleased { .. } => 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 own_account(&self) -> Option<std::borrow::Cow<'_, str>> {
use std::borrow::Cow;
match self {
Self::Failed(reason) | Self::Quarantined(reason) | Self::Replanning(reason) => {
Some(Cow::Borrowed(reason.as_str()))
}
Self::Cancelled { .. }
| Self::Abandoned { .. }
| Self::BrokeGlass { .. }
| Self::HaltLifted { .. }
| Self::HoldReleased { .. }
| Self::Withheld { .. }
| Self::Succeeded
| Self::Swept
| Self::Observed
| Self::Suspended(_)
| Self::Exhausted(_) => None,
}
}
#[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::HaltLifted { .. } => HALT_LIFTED_OUTCOME,
Self::HoldReleased { .. } => HOLD_RELEASED_OUTCOME,
Self::Withheld { .. } => WITHHELD_OUTCOME,
Self::Observed => OBSERVED_OUTCOME,
}
}
#[must_use]
pub fn actor(&self) -> Option<&crate::core::Operator> {
match self {
Self::Cancelled { actor, .. }
| Self::Abandoned { actor, .. }
| Self::BrokeGlass { actor, .. }
| Self::HaltLifted { actor }
| Self::HoldReleased { actor } => Some(actor),
Self::Succeeded
| Self::Failed(_)
| Self::Suspended(_)
| Self::Exhausted(_)
| Self::Quarantined(_)
| Self::Replanning(_)
| Self::Withheld { .. }
| Self::Swept
| Self::Observed => 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 { .. }
| Self::HaltLifted { .. }
| Self::HoldReleased { .. }
| Self::Observed
)
}
}
pub const SEALED_OUTCOMES: &[&str] = &[
"succeeded",
"cancelled",
"abandoned",
super::sweeper::SWEEP_OUTCOME,
BREAK_GLASS_OUTCOME,
HALT_LIFTED_OUTCOME,
HOLD_RELEASED_OUTCOME,
OBSERVED_OUTCOME,
];
pub const OUTCOMES_OF_RECORD: &[&str] = &[
"succeeded",
"cancelled",
"abandoned",
"quarantined",
super::sweeper::SWEEP_OUTCOME,
BREAK_GLASS_OUTCOME,
HALT_LIFTED_OUTCOME,
HOLD_RELEASED_OUTCOME,
OBSERVED_OUTCOME,
];
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Claim {
Standing,
Unlisted,
Recovering,
}
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>,
subjects: BTreeSet<crate::core::SubjectRef>,
}
#[derive(Debug, Clone, Default)]
struct QuotaPass {
period: Option<String>,
enabled: bool,
release_slot: bool,
width: Option<usize>,
}
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,
width: None,
}
}
const fn disabled() -> Self {
Self {
period: None,
enabled: false,
release_slot: false,
width: None,
}
}
fn sized_for(mut self, width: usize) -> Self {
if self.period.is_some() {
self.width = Some(width);
}
self
}
fn parallelism(&self, budget: &Budget) -> usize {
self.width
.map_or_else(|| budget.parallelism(), |w| budget.parallelism().min(w))
}
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,
concludes: bool,
) -> Option<crate::quota::QuotaSettlement> {
self.enabled.then(|| crate::quota::QuotaSettlement {
run,
epoch,
period: self.period.clone(),
spend,
release_slot: self.release_slot,
concludes,
})
}
}
fn reserved_width(budget: &Budget, plan: &PlanIR) -> usize {
budget.parallelism().min(plan.nodes.len()).max(1)
}
#[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>>,
#[cfg(feature = "manifest")]
checkers: crate::content::Checkers,
streams: Option<Arc<dyn RunStreamObserver>>,
meter: super::metrics::Meter,
quotas: Option<Arc<dyn crate::quota::QuotaStore>>,
#[cfg(feature = "manifest")]
rates: Option<Arc<super::ctx::Rates>>,
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>>,
disclosures: Option<Arc<dyn crate::disclosure::DisclosureRegister>>,
#[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,
pub(crate) interval: Option<std::time::Duration>,
pub(crate) last_met: std::sync::Mutex<Option<crate::core::Timestamp>>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Entry {
Outside,
Nested,
}
pub trait FullBackend:
JournalStore
+ CaseStore
+ TaskStore
+ EventStore
+ TimerStore
+ crate::memory::MemoryStore
+ crate::quota::QuotaStore
+ crate::disclosure::DisclosureRegister
{
}
impl<B> FullBackend for B where
B: JournalStore
+ CaseStore
+ TaskStore
+ EventStore
+ TimerStore
+ crate::memory::MemoryStore
+ crate::quota::QuotaStore
+ crate::disclosure::DisclosureRegister
{
}
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)
.quota(stores.quotas, crate::quota::TenantQuota::default())
.disclosures(stores.disclosures)
}
#[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,
disclosures: None,
#[cfg(feature = "push")]
push: None,
#[cfg(feature = "keyring")]
keyring: None,
#[cfg(feature = "manifest")]
tools: None,
journal_ceiling: None,
#[cfg(feature = "manifest")]
egress: None,
#[cfg(feature = "manifest")]
checkers: crate::content::Checkers::default(),
streams: 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,
witness_interval: 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,
at: crate::core::Timestamp,
) -> 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,
};
let report = crate::drill::drill(&stores)
.await
.map_err(RuntimeError::from_store)?;
let checkpoint = self
.store
.checkpoint()
.await
.map_err(RuntimeError::from_store)?;
cases
.record_drill(&crate::case::DrillRecord {
at,
sound: report.is_sound(),
cases: report.cases as u64,
findings: report.findings.len() as u64,
not_checked: report.not_checked.len() as u64,
origin: checkpoint.origin.clone(),
size: checkpoint.size,
})
.await
.map_err(RuntimeError::from_store)?;
Ok(report)
}
pub async fn last_drill(&self) -> Result<Option<crate::case::DrillRecord>, RuntimeError> {
self.cases()
.ok_or_else(|| {
RuntimeError::PlanContract(
"this runtime has no case store, so it has never drilled one".into(),
)
})?
.last_drill()
.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,
disclosures: self.disclosures.as_ref(),
};
crate::retention::retain(&stores, self.store.as_ref(), 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> {
if reason.trim().is_empty() {
return Err(RuntimeError::PlanContract(
"stopping a run needs a reason — what the person stopping it saw".to_owned(),
));
}
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(),
});
}
if let Some(status) = resume_is_closed(&records) {
if self.cancellation(run).await?.is_some() {
return Ok(false);
}
return Err(RuntimeError::AlreadyConcluded {
run: run.to_string(),
outcome: recorded_conclusion(&records)
.unwrap_or_else(|| status.as_str().to_owned()),
});
}
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 { .. }
| RuntimeError::NoProvider { .. }
| RuntimeError::NoCaseStore { .. }
| RuntimeError::PayloadsSealed { .. },
) => {}
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.replay_releasing(run, Mode::Resume, Some(lease), Claim::Standing)
.await
}
pub async fn record_quarantine_decision(
&self,
run: RunId,
decider: &crate::core::Operator,
reason: &str,
decision: crate::core::QuarantineDecision,
) -> Result<(), 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;
let _ = self.store.release_lease(run, lease.epoch).await;
appended
}
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") && !payloads_erased(&records) {
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: None,
note: 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)?;
self.forget_credentials(scope);
Ok(())
}
fn forget_credentials(&self, scope: &crate::quota::HaltScope) {
if let (crate::quota::HaltScope::Subject { id }, Some(wiring)) =
(scope, self.peers.as_ref())
{
wiring.registry.forget(id);
}
}
pub async fn lift_halt(
&self,
scope: &crate::quota::HaltScope,
by: &crate::core::Operator,
at: crate::core::Timestamp,
) -> Result<Option<ControlLifted>, RuntimeError> {
let quotas = self.quota_store()?;
let standing = quotas.halts().await.map_err(RuntimeError::Store)?;
let Some(halt) = standing.into_iter().find(|h| &h.scope == scope) else {
return Ok(None);
};
let run = self
.seal_operator_record(
RecordKind::HaltLifted {
scope: scope.key(),
by: by.clone(),
at,
reason: halt.reason.clone(),
thrown_by: halt.by.clone(),
thrown_at: halt.at,
},
None,
HALT_LIFTED_OUTCOME,
None,
)
.await?;
let removed = quotas
.lift_halt_if(&halt)
.await
.map_err(|e| removal_failed(run, &e))?;
self.forget_credentials(scope);
if !removed {
let now = quotas.halts().await.map_err(|e| removal_failed(run, &e))?;
if now.iter().any(|h| &h.scope == scope) {
return Err(RuntimeError::ControlStands {
run: run.to_string(),
detail: format!(
"the halt on '{}' was re-thrown after this lift read it, and the new \
halt still stands",
scope.key()
),
});
}
}
Ok(Some(ControlLifted {
record: run,
removed,
}))
}
pub async fn release_hold(
&self,
case: crate::core::CaseId,
by: &crate::core::Operator,
at: crate::core::Timestamp,
) -> Result<Option<ControlLifted>, RuntimeError> {
let cases = self.cases.as_ref().ok_or_else(|| {
RuntimeError::PlanContract("releasing a hold needs a case store".to_owned())
})?;
let Some(standing) = cases.hold(case).await.map_err(RuntimeError::Store)? else {
return match cases.case(case).await.map_err(RuntimeError::Store)? {
Some(_) => Ok(None),
None => Err(RuntimeError::Store(crate::core::StoreError::NotFound(
format!("case {case}"),
))),
};
};
let run = self
.seal_operator_record(
RecordKind::HoldReleased {
by: by.clone(),
at,
placed_by: standing.by.clone(),
placed_at: standing.placed_at,
},
Some(case),
HOLD_RELEASED_OUTCOME,
None,
)
.await?;
let removed = cases
.release_hold_if(case, &standing)
.await
.map_err(|e| removal_failed(run, &e))?;
if !removed
&& cases
.hold(case)
.await
.map_err(|e| removal_failed(run, &e))?
.is_some()
{
return Err(RuntimeError::ControlStands {
run: run.to_string(),
detail: format!(
"the hold on case {case} was placed again after this release read it, \
and the new hold still stands"
),
});
}
Ok(Some(ControlLifted {
record: run,
removed,
}))
}
pub async fn lifted_halts(&self, limit: usize) -> Result<Vec<Record>, RuntimeError> {
self.operator_records(HALT_LIFTED_OUTCOME, limit).await
}
pub async fn released_holds(&self, limit: usize) -> Result<Vec<Record>, RuntimeError> {
self.operator_records(HOLD_RELEASED_OUTCOME, limit).await
}
async fn operator_records(
&self,
outcome: &str,
limit: usize,
) -> Result<Vec<Record>, RuntimeError> {
let runs = self
.store
.runs_by_outcome(outcome, limit)
.await
.map_err(RuntimeError::from_store)?;
let mut found = Vec::with_capacity(runs.len());
for run in runs {
let first = self
.store
.read_page(run, 1, 1)
.await
.map_err(RuntimeError::from_store)?
.into_iter()
.next();
match first {
Some(record)
if matches!(
(outcome, record.kind()),
(HALT_LIFTED_OUTCOME, RecordKind::HaltLifted { .. })
| (HOLD_RELEASED_OUTCOME, RecordKind::HoldReleased { .. })
) =>
{
found.push(record);
}
_ => {
return Err(RuntimeError::Store(crate::core::StoreError::Corrupt {
seq: 1,
detail: format!(
"run {run} is sealed as '{outcome}' but holds no such record"
),
}));
}
}
}
Ok(found)
}
pub(crate) fn quota_store_if_wired(&self) -> Option<&Arc<dyn crate::quota::QuotaStore>> {
self.quotas.as_ref()
}
pub async fn reservations(
&self,
limit: usize,
) -> Result<Vec<crate::quota::Held>, RuntimeError> {
let quotas = self.quotas.as_ref().ok_or_else(|| {
RuntimeError::Store(crate::core::StoreError::Backend(
"no quota store is wired, so no run holds spend against a period".to_owned(),
))
})?;
quotas
.reservations(limit)
.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 waiting_runs(
&self,
limit: usize,
) -> Result<Vec<crate::journal::WaitingRun>, RuntimeError> {
self.store
.waiting_runs(limit)
.await
.map_err(RuntimeError::from_store)
}
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
}
#[cfg(any(feature = "a2a-server", feature = "mcp-server"))]
pub(crate) fn served_policy_problems(&self) -> Vec<String> {
match (&self.policy, &self.identity) {
(Some(engine), Some(_)) => runtime_policy_problems(engine.as_ref(), None),
_ => Vec::new(),
}
}
pub(crate) fn own_depth(&self) -> usize {
self.identity
.as_ref()
.map_or(0, crate::core::Delegation::depth)
}
#[must_use]
pub(crate) fn meter(&self) -> &super::metrics::Meter {
&self.meter
}
async fn seal_operator_record(
&self,
kind: RecordKind,
case: Option<crate::core::CaseId>,
outcome: &str,
reason: Option<&str>,
) -> Result<RunId, RuntimeError> {
let run = RunId::generate();
let epoch = BREAK_GLASS_EPOCH;
let mut act = Append::new(run, kind);
if let Some(case) = case {
act = act.case(case);
}
self.store
.append(epoch, vec![act])
.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: outcome.to_owned(),
reason: reason.map(str::to_owned),
exhaustion: None,
live_spend: Spend::default(),
chain_head: head.hash,
},
)],
)
.await
.map_err(RuntimeError::from_store)?;
self.store
.seal(run, epoch, outcome)
.await
.map_err(RuntimeError::from_store)?;
Ok(run)
}
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 = self
.seal_operator_record(
RecordKind::BreakGlass {
actor: actor.clone(),
roles: roles.to_vec(),
reason: reason.to_owned(),
},
None,
BREAK_GLASS_OUTCOME,
Some(reason),
)
.await?;
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 admit_content(
&self,
agent: &str,
input: Tainted<Value>,
) -> Result<Tainted<Value>, RuntimeError> {
let Some(manifest) = self
.resolve(agent)
.ok()
.and_then(|skill| self.governing(skill.as_ref()))
else {
return Ok(input);
};
let outcome = super::ctx::judged(
manifest.spec.security.content.as_ref(),
crate::content::At::Admission,
input.peek(),
)
.map_err(RuntimeError::PolicyDenied)?;
if let Some(hit) = outcome.refused.first() {
return Err(RuntimeError::PolicyDenied(
crate::core::PolicyError::Content {
rule: hit.rule.clone(),
pointer: hit.pointer.clone(),
},
));
}
let raised = outcome.raise(input.label().sensitivity);
if raised == input.label().sensitivity {
return Ok(input);
}
Ok(input.raised_to(raised))
}
#[cfg(not(feature = "manifest"))]
#[allow(clippy::unused_self, clippy::unnecessary_wraps)]
fn admit_content(
&self,
_agent: &str,
input: Tainted<Value>,
) -> Result<Tainted<Value>, RuntimeError> {
Ok(input)
}
#[cfg(feature = "manifest")]
fn delivered_verdict(
&self,
history: &[crate::journal::Record],
payload: &serde_json::Value,
) -> Option<crate::core::ContentVerdict> {
let capability = history.iter().find_map(|r| match r.kind() {
RecordKind::RunAdmitted { capability, .. } => Some(capability.clone()),
_ => None,
})?;
let manifest = self.governing(self.resolve(&capability).ok()?.as_ref())?;
super::ctx::verdict(super::ctx::judged(
manifest.spec.security.content.as_ref(),
crate::content::At::Source(super::ctx::AWAIT_KIND),
payload,
))
}
#[cfg(not(feature = "manifest"))]
#[allow(clippy::unused_self)]
fn delivered_verdict(
&self,
_history: &[crate::journal::Record],
_payload: &serde_json::Value,
) -> Option<crate::core::ContentVerdict> {
None
}
#[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;
match store.renew(run, &owner, epoch, ttl).await {
Ok(_) => {}
Err(crate::core::StoreError::LeaseNotHeld { .. }) => return,
Err(error) => tracing::warn!(
target: "agentplane::lease",
run = %run,
error = %error,
"a lease renewal failed; the next one retries"
),
}
}
}))
}
pub(crate) async fn withdrawal_against(
&self,
identity: Option<&crate::core::Delegation>,
) -> Result<Option<Withdrawal>, crate::core::StoreError> {
let (Some(quotas), Some(chain)) = (self.quotas.as_ref(), identity) else {
return Ok(None);
};
let halts = quotas.halts().await?;
Ok(halts.into_iter().find_map(|halt| {
let subject = halt.scope.withdraws_from(chain)?.to_owned();
Some(Withdrawal {
subject,
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>,
chain: Option<&crate::core::Delegation>,
at: crate::core::Timestamp,
(budget, width): (&Budget, usize),
) -> Result<QuotaPass, RuntimeError> {
let Some(quotas) = self.quotas.as_ref() else {
return Ok(QuotaPass::disabled());
};
let pass = QuotaPass::at(&self.quota, at, true).sized_for(width);
match quotas.halts().await {
Ok(halts) => {
if let Some(halt) = halts
.into_iter()
.filter(|halt| halt.scope.covers(governed_by, chain))
.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()),
));
}
}
let hold = match crate::quota::reservation(&self.quota, budget, width as u64) {
Ok(amount) => {
amount
.zip(pass.period())
.map(|(amount, period)| crate::quota::SpendHold {
period: period.to_owned(),
amount,
})
}
Err(field) => {
return Err(RuntimeError::QuotaExceeded(
crate::quota::QuotaError::Unbounded {
tenant: self.tenant.as_str().to_owned(),
unit: if field.contains("minor") {
"money"
} else {
"tokens"
},
field,
},
));
}
};
quotas
.reserve(run, &self.quota, hold.as_ref(), at)
.await
.map_err(RuntimeError::QuotaExceeded)?;
Ok(pass)
}
async fn settle_quota(
&self,
run: RunId,
epoch: crate::core::Epoch,
spend: Spend,
pass: &QuotaPass,
concludes: bool,
) -> Result<(), RuntimeError> {
let Some(quotas) = self.quotas.as_ref() else {
return Ok(());
};
let pending = |error: crate::core::StoreError| RuntimeError::QuotaSettlementPending {
run: run.to_string(),
epoch,
detail: error.to_string(),
};
if !pass.enabled {
if concludes {
quotas.release(run).await.map_err(pending)?;
}
return Ok(());
}
let settlement = pass
.settlement(run, epoch, spend, concludes)
.expect("an enabled quota pass has a settlement");
quotas.settle(&settlement).await.map_err(pending)
}
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(()),
}
}
pub(crate) async fn release_sealed_slots(&self, limit: usize) -> Result<usize, RuntimeError> {
let Some(quotas) = self.quotas.as_ref() else {
return Ok(0);
};
let held = quotas
.running_runs(limit)
.await
.map_err(RuntimeError::Store)?;
let positions = self
.store
.log_positions(&held)
.await
.map_err(RuntimeError::from_store)?;
let mut released = 0;
for (run, sealed) in held.into_iter().zip(positions) {
if sealed.is_none() {
continue;
}
let records = self
.store
.read(run, 1)
.await
.map_err(RuntimeError::from_store)?;
match self.settle_recorded_quota_passes(run, &records).await {
Ok(()) => released += 1,
Err(error) => tracing::warn!(
%run,
%error,
"a sealed run's admission slot could not be released; the next sweep retries"
),
}
}
Ok(released)
}
pub(crate) async fn seal_unsealed_conclusions(
&self,
limit: usize,
) -> Result<usize, RuntimeError> {
let mut sealed = 0;
for &outcome in LEASE_FREE_OUTCOMES {
let runs = self
.store
.runs_by_outcome(outcome, limit)
.await
.map_err(RuntimeError::from_store)?;
let positions = self
.store
.log_positions(&runs)
.await
.map_err(RuntimeError::from_store)?;
for (run, position) in runs.into_iter().zip(positions) {
if position.is_some() {
continue;
}
match self.store.seal(run, LEASE_FREE_EPOCH, outcome).await {
Ok(_) => sealed += 1,
Err(error) => tracing::warn!(
%run,
%error,
"a concluded run's seal could not be finished; the next sweep retries"
),
}
}
}
Ok(sealed)
}
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 {
period: Some(_),
..
}
)
}) {
return Err(RuntimeError::PlanContract(
"this run has passes that accrued spend to a quota period, and the resuming \
plane has no quota store to settle them"
.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,
concludes: concluded_in(records, epoch),
})
.await
.map_err(|error| RuntimeError::QuotaSettlementPending {
run: run.to_string(),
epoch,
detail: error.to_string(),
})?;
}
if records.iter().any(is_sealing_conclusion) {
quotas
.release(run)
.await
.map_err(|error| RuntimeError::QuotaSettlementPending {
run: run.to_string(),
epoch: records.last().map_or(0, |r| r.body.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 RecordKind::RunConcluded {
outcome,
reason,
exhaustion,
..
} = last.kind()
else {
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()) {
let head = self
.store
.seal(run, lease.epoch, outcome)
.await
.map_err(RuntimeError::from_store)?;
self.retire_waits(run).await;
head
} 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()) {
let head = self
.store
.seal(run, lease.epoch, &outcome)
.await
.map_err(RuntimeError::from_store)?;
self.retire_waits(run).await;
head
} 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,
recorded_case_id(records),
Consumed::default(),
Spend::ZERO,
&QuotaPass::disabled(),
)
.await
.map(Some)
}
fn marking_pass(&self, run: RunId, quota: &QuotaPass) -> Option<Self> {
let marker = quota.started()?;
Some(Self {
store: Arc::new(super::pass_journal::PassJournal::new(
Arc::clone(&self.store),
run,
marker,
)),
..self.clone()
})
}
async fn start_replay_quota_pass(
&self,
run: RunId,
mode: Mode,
width: usize,
) -> Result<QuotaPass, RuntimeError> {
let (Mode::Resume, Some(quotas)) = (mode, self.quotas.as_ref()) else {
return Ok(QuotaPass::disabled());
};
let pass = QuotaPass::at(&self.quota, now_for_admission(), false).sized_for(width);
if let Some(period) = pass.period() {
quotas
.carry(run, period)
.await
.map_err(RuntimeError::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()?;
self.identity_of(self.governing(skill.as_ref())?.as_ref())
}
#[cfg(feature = "manifest")]
pub(crate) fn reach_of(&self, capability: &str) -> Option<crate::core::Reach> {
let skill = self.resolve(capability).ok()?;
let declaration = self.governing(skill.as_ref()).and_then(|manifest| {
let digest = manifest.digest().ok()?;
Some(super::reach::declared(&manifest, digest, |server| {
self.peers.as_ref().is_some_and(|wiring| {
wiring
.registry
.grant(&crate::peers::PeerId::new(server))
.is_some()
})
}))
});
Some(crate::core::Reach {
capability: capability.to_owned(),
declaration,
})
}
#[cfg(feature = "manifest")]
fn identity_of(&self, m: &crate::manifest::Manifest) -> Option<crate::journal::AgentIdentity> {
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()
}
#[cfg(feature = "mcp-server")]
pub(crate) fn wires_peer(&self, server: &str) -> bool {
self.peers.as_ref().is_some_and(|wiring| {
wiring
.registry
.grant(&crate::peers::PeerId::new(server))
.is_some()
})
}
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<Admission, RuntimeError> {
let Some(key) = terms.idempotency_key.clone() else {
return Ok(Admission::Fresh(self.admit_plan(plan, input, terms).await?));
};
validate_admission_key(&key)?;
if let Some(held) = self.holder_of(&key).await? {
return self.answer_with(held).await;
}
match self.admit_plan(plan, 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 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,
ceiling: Option<&crate::core::Delegation>,
acts_under: bool,
) -> Result<(), RuntimeError> {
let Some(chain) = ceiling 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: if acts_under {
chain.subject().id.clone()
} else {
node.capability.to_string()
},
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 request = crate::policy::requests::admission(
&crate::policy::requests::Acting {
tenant: self.tenant.as_str(),
capability,
agent: governed_by,
chain,
},
input,
);
let decision = engine.authorize(&request.as_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: request.principal,
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>,
admitted_by: Option<String>,
acting_as: &Acting,
) -> 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,
admitted_by,
served_unchained: matches!(acting_as, Acting::Nobody),
plane_chain: match acting_as {
Acting::Plane => self.identity.is_some(),
Acting::PlaneDelegate(_) => true,
Acting::Caller(_) | Acting::Nobody => false,
},
}
}
#[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;
if let Some(expected) = terms.expected_declaration {
let found = governed_by.as_ref().map(|id| id.digest);
if found != Some(expected) {
return Err(RuntimeError::DeclarationPinMismatch {
agent,
expected,
found,
});
}
}
let chain = terms.acting_as.resolve(self.identity.as_ref());
self.authorize_scope(
&plan,
terms.acting_as.ceiling(self.identity.as_ref()),
chain.is_some(),
)?;
self.authorize_admission(&agent, governed_by.as_ref(), chain, input.peek())?;
let input = self.admit_content(&agent, input)?;
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 budget = self.budget_for(&agent);
let quota = match Box::pin(self.check_quota(
run,
governed_by.as_ref(),
terms.acting_as.resolve(self.identity.as_ref()),
now_for_admission(),
(&budget, reserved_width(&budget, &plan)),
))
.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,
admitted_by,
expected_declaration: _,
subjects,
} = terms;
let chain = acting_as.resolve(self.identity.as_ref());
let mut records = vec![
Append::new(
run,
self.admission(
agent,
governed_by,
&input,
idempotency_key,
admitted_by,
&acting_as,
),
),
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));
}
#[cfg(feature = "manifest")]
let declared: Vec<crate::journal::SubjectBinding> = self
.resolve(agent)
.ok()
.and_then(|skill| self.governing(skill.as_ref()))
.map(|m| {
m.spec
.data_subjects
.iter()
.map(crate::manifest::DataSubject::binding)
.collect()
})
.unwrap_or_default();
#[cfg(not(feature = "manifest"))]
let declared: Vec<crate::journal::SubjectBinding> = Vec::new();
let bound = bind_subjects(&declared, subjects, &input, case.is_some())?;
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();
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"
)));
}
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,
};
let bound = fill_case_subjects(bound, case_ctx.as_ref());
let subject_refs = subject_refs(run, &bound);
if !bound.is_empty() {
let mut append = Append::new(run, RecordKind::DataSubjectBound { bindings: bound });
append.case = case_ctx.as_ref().map(CaseContext::id);
records.push(append);
}
self.store
.append(lease.epoch, records)
.await
.map_err(RuntimeError::from_store)?;
if let Err(open_failed) = self
.register(run, case_ctx.as_ref().map(CaseContext::id))
.await
{
let head = self
.store
.head(run)
.await
.map_err(RuntimeError::from_store)?;
self.store
.append(
lease.epoch,
vec![Append::new(
run,
RecordKind::RunConcluded {
outcome: RunStatus::Failed(String::new()).as_str().to_owned(),
reason: Some(format!(
"the run's 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(),
subjects: subject_refs,
})
}
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,
subjects: a.subjects,
refusal: None,
withheld_subject: None,
successors: Vec::new(),
started: BTreeSet::new(),
finished: BTreeSet::new(),
succeeded: 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 resume_refusal(&self, records: &[Record]) -> Result<Option<RuntimeError>, RuntimeError> {
if let Some(refusal) = self.policy_bundle_refusal(records)? {
return Ok(Some(refusal));
}
#[cfg(feature = "manifest")]
if let Some(refusal) = self.declaration_refusal(records)? {
return Ok(Some(refusal));
}
Ok(None)
}
fn refuse_undrivable<'p>(
&self,
plans: impl IntoIterator<Item = &'p PlanIR>,
) -> Result<(), RuntimeError> {
for node in plans.into_iter().flat_map(|p| &p.nodes) {
self.resolve(&node.capability.0)?;
}
Ok(())
}
async fn register(
&self,
run: RunId,
case: Option<crate::core::CaseId>,
) -> Result<(), crate::core::StoreError> {
if let (Some(case), Some(cases)) = (case, self.cases.as_ref()) {
cases.attach_run(case, run).await?;
}
#[cfg(feature = "push")]
if let Some(outbox) = &self.outbox {
outbox.open(run).await?;
}
Ok(())
}
fn refuse_caseless(&self, run: RunId, records: &[Record]) -> Result<(), RuntimeError> {
match recorded_case_id(records) {
Some(case) if self.cases.is_none() => Err(RuntimeError::NoCaseStore {
run: run.to_string(),
case: case.to_string(),
}),
_ => Ok(()),
}
}
async fn quarantine_refused(
&self,
run: RunId,
records: &[Record],
lease: &crate::journal::Lease,
refusal: &RuntimeError,
) -> Result<RunOutcome, RuntimeError> {
self.conclude(
run,
lease.epoch,
RunStatus::Quarantined(refusal.to_string()),
None,
true,
recorded_case_id(records),
Consumed::default(),
Spend::ZERO,
&QuotaPass::disabled(),
)
.await
}
fn policy_bundle_refusal(
&self,
records: &[Record],
) -> Result<Option<RuntimeError>, 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 Ok(Some(RuntimeError::PolicyBundleChanged {
recorded: recorded.as_ref().map(PolicyBundleIdentity::digest),
configured: configured.as_ref().map(PolicyBundleIdentity::digest),
}));
}
Ok(None)
}
#[cfg(feature = "manifest")]
fn declaration_refusal(
&self,
records: &[Record],
) -> Result<Option<RuntimeError>, RuntimeError> {
let Some(recorded) = records.iter().find_map(|record| match record.kind() {
RecordKind::RunAdmitted { governed_by, .. } => governed_by.as_deref().cloned(),
_ => None,
}) else {
return Ok(None);
};
let Some(current) = self
.governed_by
.values()
.find(|m| m.metadata.name == recorded.name)
else {
return Ok(None);
};
let configured = current
.digest()
.map_err(|e| RuntimeError::PlanContract(e.to_string()))?;
if configured != recorded.digest {
return Ok(Some(RuntimeError::DeclarationChanged {
agent: recorded.name,
recorded: recorded.digest,
configured,
}));
}
Ok(None)
}
fn sealed_or_erased(&self, error: RuntimeError) -> RuntimeError {
match error {
RuntimeError::PayloadsErased { run } if !self.holds_key_ring() => {
RuntimeError::PayloadsSealed { run }
}
other => other,
}
}
#[cfg_attr(not(feature = "keyring"), allow(clippy::unused_self))]
const fn holds_key_ring(&self) -> bool {
#[cfg(feature = "keyring")]
{
self.keyring.is_some()
}
#[cfg(not(feature = "keyring"))]
{
false
}
}
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, Claim::Standing)
.await
}
pub async fn verify(&self, run: RunId) -> Result<super::Verdict, RuntimeError> {
use super::{CannotReplay, Finding};
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 recorded = records.iter().find_map(|r| match r.kind() {
RecordKind::RunAdmitted { governed_by, .. } => governed_by.as_deref().cloned(),
_ => None,
});
let candidate = self.candidate_revision(recorded.as_ref());
let verdict = |finding| super::Verdict {
run,
recorded: recorded.clone(),
candidate: candidate.clone(),
finding,
};
if let Err(RuntimeError::CanonicalizationChanged {
recorded,
implemented,
}) = ensure_replayable_canon(&records)
{
return Ok(verdict(Finding::CannotReplay(
CannotReplay::CanonicalizationChanged {
recorded,
implemented,
},
)));
}
let Some(ending) = recorded_ending(&records) else {
return Ok(verdict(Finding::CannotReplay(CannotReplay::NotConcluded)));
};
if ended_by_an_act(&ending) {
return Ok(verdict(Finding::CannotReplay(CannotReplay::EndedByAnAct {
outcome: ending,
})));
}
let mut dispatched = false;
let mut divergence = None;
let replayed = self
.replay_under(
run,
Mode::Strict,
None,
Claim::Standing,
&mut dispatched,
&mut divergence,
)
.await;
let finding = match (replayed, divergence) {
(Ok(_), Some(d)) => Finding::Diverged(d),
(Ok(outcome), None) if outcome.status.as_str() == ending => {
Finding::Verified { outcome: ending }
}
(Ok(outcome), None) => Finding::OutcomeDiffers {
recorded: ending,
replayed: outcome.status.as_str().to_owned(),
},
(Err(RuntimeError::PayloadsErased { .. }), _) => {
Finding::CannotReplay(CannotReplay::Erased)
}
(Err(RuntimeError::PayloadsSealed { .. }), _) => {
Finding::CannotReplay(CannotReplay::KeyAbsent)
}
(Err(RuntimeError::NoProvider { target, .. }), _) => {
Finding::CannotReplay(CannotReplay::EntryPointRemoved { capability: target })
}
(Err(other), _) => return Err(other),
};
Ok(verdict(finding))
}
#[cfg(feature = "manifest")]
fn candidate_revision(
&self,
recorded: Option<&crate::journal::AgentIdentity>,
) -> Option<crate::journal::AgentIdentity> {
let name = &recorded?.name;
let held = self
.governed_by
.values()
.find(|m| &m.metadata.name == name)?;
self.identity_of(held)
}
#[cfg(not(feature = "manifest"))]
#[allow(clippy::unused_self)]
const fn candidate_revision(
&self,
_recorded: Option<&crate::journal::AgentIdentity>,
) -> Option<crate::journal::AgentIdentity> {
None
}
pub(crate) async fn resume_holding(
&self,
run: RunId,
lease: crate::journal::Lease,
) -> Result<RunOutcome, RuntimeError> {
self.replay_releasing(run, Mode::Resume, Some(lease), Claim::Unlisted)
.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)?;
match self
.replay_releasing(run, Mode::Resume, Some(lease.clone()), Claim::Recovering)
.await
{
Err(error) if recurs_on_every_resume(&error) => {
let records = self
.store
.read(run, 1)
.await
.map_err(RuntimeError::from_store)?;
let refusal = RuntimeError::PlanContract(format!(
"recovery refused, and would be refused the same way on every attempt: \
{error}"
));
self.quarantine_refused(run, &records, &lease, &refusal)
.await
}
other => other,
}
}
async fn replay_releasing(
&self,
run: RunId,
mode: Mode,
lease: Option<crate::journal::Lease>,
claim: Claim,
) -> Result<RunOutcome, RuntimeError> {
let mut dispatched = false;
let outcome = self
.replay_under(run, mode, lease.as_ref(), claim, &mut dispatched, &mut None)
.await;
let release = match &outcome {
Ok(_) => true,
Err(RuntimeError::QuotaSettlementPending { .. }) => false,
Err(_) => claim == Claim::Standing && !dispatched,
};
if release
&& 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
}
#[allow(clippy::too_many_lines)]
async fn replay_under(
&self,
run: RunId,
mode: Mode,
lease: Option<&crate::journal::Lease>,
claim: Claim,
dispatched: &mut bool,
divergence: &mut Option<crate::journal::Divergence>,
) -> 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"),
claim == Claim::Recovering,
)
.await?
{
return Ok(outcome);
}
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)?;
}
let input = records.iter().find_map(recorded_input).ok_or_else(|| {
RuntimeError::PlanContract("journal has no RunAdmitted record".into())
})?;
let (plan, successors) =
frozen_plans(run, &records).map_err(|e| self.sealed_or_erased(e))?;
if mode == Mode::Resume {
self.refuse_undrivable(std::iter::once(&plan).chain(&successors))?;
self.refuse_caseless(run, &records)?;
}
if mode == Mode::Resume
&& let Some(refusal) = self.resume_refusal(&records)?
{
let lease = lease.expect("resume holds a lease");
return self
.quarantine_refused(run, &records, lease, &refusal)
.await;
}
if mode == Mode::Resume {
self.register(run, recorded_case_id(&records))
.await
.map_err(RuntimeError::from_store)?;
}
*dispatched = true;
let budget = self.budget_for(&recorded_agent(&records));
let quota = self
.start_replay_quota_pass(run, mode, reserved_width(&budget, &plan))
.await?;
let through = self.marking_pass(run, "a);
let this = through.as_ref().unwrap_or(self);
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));
this.execute(
Execution {
run,
epoch,
quota,
plan: &plan,
input,
mode,
case: case_ctx,
budget,
agent: recorded_agent(&records),
identity: recorded_chain(&records)?,
subjects: recorded_subjects(run, &records),
refusal: recorded_step_refusal(&records),
withheld_subject: standing_withholding(&records),
successors,
started: recorded_started_steps(&records),
finished: recorded_finished_steps(&records),
succeeded: recorded_succeeded_steps(&records),
recorded_groups: recorded_groups(&records),
},
&mut cursor,
)
.await
.inspect(|_| *divergence = cursor.divergence().cloned())
}
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,
subjects,
refusal: recorded_refusal,
withheld_subject,
successors,
started,
finished,
succeeded,
recorded_groups,
} = plan;
let input = input.attributed(&subjects);
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 ran: BTreeMap<StepId, PlanNode> = BTreeMap::new();
let mut outputs: BTreeMap<StepId, Tainted<Value>> = BTreeMap::new();
let mut deferred: Option<Vec<StepId>> = None;
loop {
let mut ready = current.ready(&done);
if let Some(rest) = deferred.take() {
let held: Vec<StepId> =
ready.iter().copied().filter(|s| rest.contains(s)).collect();
if !held.is_empty() {
ready = held;
}
}
if writing
&& ready.iter().any(|s| succeeded.contains(s))
&& ready.iter().any(|s| !succeeded.contains(s))
{
let (recorded, rest) = ready.into_iter().partition(|s| succeeded.contains(s));
ready = recorded;
deferred = Some(rest);
}
let at_frontier = ready.is_empty() || ready.iter().any(|s| !succeeded.contains(s));
if writing
&& at_frontier
&& 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)?;
let abandon: Vec<StepId> = ready
.iter()
.copied()
.filter(|&s| {
recorded_groups.iter().any(|((step, phase, _), g)| {
*step == s && *phase == Phase::Forward && g.opened > g.settled
})
})
.collect();
if !abandon.is_empty() {
let dispatched = self
.dispatch(
&abandon,
cursor,
Batch {
agent: &agent,
identity: identity.as_ref(),
subjects: &subjects,
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: 1,
abandoning: true,
},
)
.await;
collect(dispatched, &abandon, cursor)?;
}
return self
.stop(
Unwind {
agent: &agent,
identity: identity.as_ref(),
subjects: &subjects,
run,
epoch,
ir: ¤t,
mode,
case: case.clone(),
ledger: &ledger,
writing,
stamp: &stamp,
},
status,
&completed,
&outputs,
cursor,
case_id,
"a,
)
.await;
}
if writing && at_frontier {
let standing = self
.withdrawal_against(identity.as_ref())
.await
.map_err(RuntimeError::Store)?;
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(),
subjects: &subjects,
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(),
subjects: &subjects,
run,
epoch,
ir: ¤t,
mode,
case: case.clone(),
ledger: &ledger,
writing,
stamp: &stamp,
},
withdrawal,
true,
&completed,
&outputs,
cursor,
case_id,
"a,
)
.await;
}
(None, None) => {}
}
}
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(),
subjects: &subjects,
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: quota.parallelism(&budget),
abandoning: false,
},
)
.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(missing) = cursor.unconsumed_in(step, Phase::Forward) {
cursor.note_divergence(crate::journal::Divergence::unconsumed(
step,
Phase::Forward,
&missing,
));
self.meter.count(metrics::DIVERGENCES, "");
tracing::error!(
target: telemetry::NONDETERMINISM,
%step, key = %missing.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 {missing} \
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 ran,
&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 {
run,
current: ¤t,
ran: &ran,
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(),
subjects: &subjects,
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, missing)) = cursor.first_unconsumed()
{
cursor.note_divergence(crate::journal::Divergence::unconsumed(
step, phase, &missing,
));
self.meter.count(metrics::DIVERGENCES, "");
tracing::error!(
target: telemetry::NONDETERMINISM,
%step, key = %missing.key, unconsumed = true,
);
return self
.conclude(
run,
epoch,
RunStatus::Quarantined(format!(
"strict replay verified less than the run recorded: journaled \
effect {missing} at 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,
subjects: batch.subjects,
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(),
abandoning: batch.abandoning,
},
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 started = if cx.recorded.is_none() {
let records = self
.store
.read(cx.run, 1)
.await
.map_err(RuntimeError::from_store)?;
let effectful: BTreeSet<StepId> = records
.iter()
.filter(|r| r.body.phase.is_forward())
.filter_map(|r| match r.kind() {
RecordKind::EffectStarted { .. } => r.body.step,
_ => None,
})
.collect();
let mut started = started_skills(&records);
started.retain(|step, _| effectful.contains(step));
started
} else {
BTreeMap::new()
};
let next = match self.successor(cx, outputs, completed, &started).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)],
started: &BTreeMap<StepId, String>,
) -> 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, node_ran) in cx.ran {
let Some(node) = next.node(*step) else {
continue;
};
if node.capability != node_ran.capability {
return Err(RunStatus::Failed(format!(
"the successor plan reuses step {step} — which already ran \
as '{}' — for '{}'. Keep a completed step's node or leave \
the step out; effect keys are derived from the step id, so \
new work at a used id cannot be replayed",
node_ran.capability.0, node.capability.0
)));
}
if node != node_ran {
return Err(RunStatus::Failed(format!(
"the successor plan redeclares completed step {step} ('{}') \
with different arguments, dependencies or flags. A \
completed step counts as done by its id, so the new node \
would be satisfied by work done under the old one — carry \
it over unchanged, or give the new work a fresh id",
node.capability.0
)));
}
}
for (step, skill) in started {
if cx.ran.contains_key(step) {
continue;
}
let Some(node) = next.node(*step) else {
continue;
};
let resolved = self
.resolve(&node.capability.0)
.map(|s| s.descriptor().name)
.ok();
if resolved.as_deref() != Some(skill.as_str()) {
return Err(RunStatus::Failed(format!(
"the successor plan reuses step {step} — which already \
started as '{skill}' and journaled effects — for '{}'. Keep the step's \
capability or leave the step out; effect keys are derived \
from the step id, and the unwind undoes a started step as \
the skill that ran",
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, missing)) = cursor.first_unconsumed()
{
cursor.note_divergence(crate::journal::Divergence::unconsumed(
step, phase, &missing,
));
self.meter.count(metrics::DIVERGENCES, "");
tracing::error!(
target: telemetry::NONDETERMINISM,
%step, key = %missing.key, unconsumed = true,
);
return self
.conclude(
run,
epoch,
RunStatus::Quarantined(format!(
"strict replay verified less than the run recorded: journaled \
effect {missing} at 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,
started: _,
undone: already_undone,
recorded_groups,
failed: _,
} = 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 let Err(fault) = &result {
tracing::error!(
target: telemetry::COMPENSATION_FAILED,
run = %cx.run,
%step,
error_type = fault.class(),
reason_digest = %telemetry::reason_digest(&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)
}
pub(super) async fn holds_landed_work(&self, run: RunId) -> Result<bool, RuntimeError> {
let evidence = self.unwind_evidence(run).await?;
Ok(!evidence.mutated.is_subset(&evidence.undone))
}
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 grouped: BTreeMap<StepId, Vec<crate::core::EffectKey>> = BTreeMap::new();
let mut undone = BTreeSet::new();
let mut failed = BTreeSet::new();
for r in &records {
let Some(step) = r.body.step else { continue };
let key = r.effect_key();
match r.kind() {
RecordKind::StepFinished { outcome } if r.body.phase.is_forward() => {
if outcome == RunStatus::Succeeded.as_str() {
failed.remove(&step);
} else {
failed.insert(step);
}
}
RecordKind::EffectStarted { mutates: true, .. } if r.body.phase.is_forward() => {
if let Some(key) = key {
touching.insert(key, (step, true));
if let Some(members) = grouped.get_mut(&step) {
members.push(key);
}
}
}
RecordKind::GroupOpened { .. } if r.body.phase.is_forward() => {
grouped.insert(step, Vec::new());
}
RecordKind::GroupSettled { outcome, .. } if r.body.phase.is_forward() => {
let members = grouped.remove(&step).unwrap_or_default();
if *outcome == crate::core::GroupOutcome::Aborted {
for member in members {
touching.remove(&member);
}
}
}
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,
started: started_skills(&records),
undone,
recorded_groups: recorded_groups(&records),
failed,
})
}
fn plane_frame(&self, step: StepFrame) -> super::ctx::Frame {
let StepFrame {
run,
epoch,
step,
phase,
mode,
case,
ledger,
identity,
subjects,
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(),
#[cfg(feature = "manifest")]
checkers: self.checkers.clone(),
streams: self.streams.clone(),
#[cfg(feature = "manifest")]
rates: self.rates.clone(),
meter: self.meter.clone(),
#[cfg(feature = "keyring")]
keyring: self.keyring.clone(),
tenant: self.tenant.clone(),
ledger,
policy: self.policy.clone(),
identity,
subjects,
agent,
plane: self.self_ref.clone(),
#[cfg(feature = "manifest")]
declaration: manifest.as_deref().and_then(|m| self.identity_of(m)),
#[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(),
subjects: cx.subjects.clone(),
agent: cx.agent.to_owned(),
#[cfg(feature = "manifest")]
manifest: self.governing(skill),
recorded_groups,
}),
);
let result = skill.compensate(&mut ctx, &output).await;
let mut slice = ctx.into_cursor();
if let Err(crate::core::SkillError::Step(error)) = &result
&& let Some(divergence) =
crate::journal::Divergence::of(step, Phase::Compensating, error, &slice)
{
slice.note_divergence(divergence);
}
cursor.restore(step, Phase::Compensating, slice);
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 candidates = Self::with_interrupted_siblings(completed, &evidence);
let list = if status.is_cancelled() || self.will_compensate(&candidates, &evidence.mutated)
{
match self
.stop_list(cx.run, completed, &evidence.started, &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)],
started: &BTreeMap<StepId, String>,
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, started, 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_siblings(
completed: &[(StepId, Capability)],
evidence: &UnwindEvidence,
) -> Vec<(StepId, Capability)> {
let mut out = completed.to_vec();
let done: BTreeSet<StepId> = out.iter().map(|(s, _)| *s).collect();
for step in evidence
.mutated
.iter()
.filter(|s| !done.contains(s) && !evidence.failed.contains(s))
{
if let Some(skill) = evidence.started.get(step) {
out.push((*step, Capability::new(skill.as_str())));
}
}
out
}
fn with_interrupted_steps(
completed: &[(StepId, Capability)],
started: &BTreeMap<StepId, String>,
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(skill) = started.get(step) {
out.push((*step, Capability::new(skill.as_str())));
}
}
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
}
#[allow(clippy::too_many_lines)]
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,
subjects,
run,
epoch,
node,
phase,
mode,
case,
ledger,
writing,
stamp,
agent,
already_started,
already_finished,
recorded_groups,
abandoning,
} = 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(),
subjects: subjects.clone(),
agent: agent.to_owned(),
#[cfg(feature = "manifest")]
manifest: self.governing(skill.as_ref()),
recorded_groups,
}),
);
cx.abandoning = abandoning;
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 mut cursor = cx.into_cursor();
if let Err(crate::core::SkillError::Step(error)) = &result
&& let Some(divergence) = crate::journal::Divergence::of(step, phase, error, &cursor)
{
cursor.note_divergence(divergence);
}
ledger.lock().expect("budget mutex").record_step();
let fault = result.as_ref().err().map(crate::core::SkillError::class);
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,
%run,
%step,
error_type = fault.unwrap_or("quarantined"),
reason_digest = %telemetry::reason_digest(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))
}
async fn last_is_conclusion(
&self,
run: RunId,
head: crate::core::Seq,
status: &RunStatus,
) -> Result<bool, RuntimeError> {
if head < 1 {
return Ok(false);
}
let page = self
.store
.read_page(run, head, 1)
.await
.map_err(RuntimeError::from_store)?;
Ok(matches!(
page.first().map(Record::kind),
Some(RecordKind::RunConcluded { outcome, .. }) if outcome == status.as_str()
))
}
#[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,
error_type = "failed",
reason_digest = %telemetry::reason_digest(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.last_is_conclusion(run, before.seq, &status).await?;
if repeated {
before.hash
} else {
let mut sealed = Append::new(
run,
RecordKind::RunConcluded {
outcome: status.as_str().to_owned(),
reason: status.own_account().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, status.seals())
.await?;
let chain_head = if status.seals() {
let head = self
.store
.seal(run, epoch, status.as_str())
.await
.map_err(RuntimeError::from_store)?;
self.retire_waits(run).await;
head
} 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 unanswered_waits(records: &[Record]) -> Vec<crate::core::EffectKey> {
let settled: BTreeSet<_> = records
.iter()
.filter(|r| {
matches!(
r.kind(),
RecordKind::EffectDone { .. }
| RecordKind::EffectFailed { .. }
| RecordKind::EffectReconciled { .. }
)
})
.filter_map(Record::effect_key)
.collect();
records
.iter()
.filter(|r| matches!(r.kind(), RecordKind::EffectStarted { .. }))
.filter_map(Record::effect_key)
.filter(|k| !settled.contains(k))
.collect()
}
fn concluded_in(records: &[Record], epoch: crate::core::Epoch) -> bool {
records
.iter()
.any(|r| r.body.epoch == epoch && is_sealing_conclusion(r))
}
fn is_sealing_conclusion(record: &Record) -> bool {
matches!(
record.kind(),
RecordKind::RunConcluded { outcome, .. } if SEALED_OUTCOMES.contains(&outcome.as_str())
)
}
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_digest = %telemetry::reason_digest(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,
error_type = "quarantined",
reason_digest = %telemetry::reason_digest(why),
);
meter.count(metrics::QUARANTINES, "");
}
RunStatus::Abandoned { actor, reason } => {
tracing::error!(
target: telemetry::ABANDONED,
%run,
%actor,
reason_digest = %telemetry::reason_digest(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)>,
ran: &mut BTreeMap<StepId, PlanNode>,
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()));
ran.insert(step, node.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::Observed
| RunStatus::BrokeGlass { .. }
| RunStatus::HaltLifted { .. }
| RunStatus::HoldReleased { .. } => 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>,
subjects: &'a BTreeSet<crate::core::SubjectRef>,
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,
abandoning: bool,
}
pub(crate) struct Withdrawal {
pub(crate) subject: String,
pub(crate) reason: String,
pub(crate) by: crate::core::Operator,
}
struct UnwindEvidence {
mutated: BTreeSet<StepId>,
started: BTreeMap<StepId, String>,
undone: BTreeSet<StepId>,
recorded_groups: BTreeMap<(StepId, Phase, String), super::ctx::RecordedGroup>,
failed: BTreeSet<StepId>,
}
struct Replan<'a> {
run: RunId,
current: &'a PlanIR,
ran: &'a BTreeMap<StepId, PlanNode>,
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>,
subjects: &'a BTreeSet<crate::core::SubjectRef>,
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,
}
#[allow(clippy::struct_excessive_bools)]
struct StepRun<'a> {
identity: Option<&'a crate::core::Delegation>,
subjects: &'a BTreeSet<crate::core::SubjectRef>,
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>,
abandoning: bool,
}
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>,
subjects: BTreeSet<crate::core::SubjectRef>,
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 started_skills(records: &[Record]) -> BTreeMap<StepId, String> {
records
.iter()
.filter(|r| r.body.phase.is_forward())
.filter_map(|r| match (r.kind(), r.body.step) {
(RecordKind::StepStarted { skill }, Some(step)) => Some((step, skill.clone())),
_ => 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_case_id(records: &[Record]) -> Option<crate::core::CaseId> {
records.iter().find_map(|r| match r.kind() {
RecordKind::CaseBound { .. } => r.body.case,
_ => None,
})
}
fn recorded_succeeded_steps(records: &[Record]) -> BTreeSet<StepId> {
let mut succeeded = BTreeSet::new();
for r in records.iter().filter(|r| r.body.phase.is_forward()) {
if let (RecordKind::StepFinished { outcome }, Some(step)) = (r.kind(), r.body.step) {
if outcome == RunStatus::Succeeded.as_str() {
succeeded.insert(step);
} else {
succeeded.remove(&step);
}
}
}
succeeded
}
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 bind_subjects(
declared: &[crate::journal::SubjectBinding],
named: Vec<String>,
input: &Tainted<Value>,
in_case: bool,
) -> Result<Vec<crate::journal::BoundSubject>, RuntimeError> {
use crate::journal::SubjectBinding;
let refuse = |binding: &SubjectBinding, reason: &str| RuntimeError::SubjectUnbound {
binding: binding.to_string(),
reason: reason.to_owned(),
};
let mut resolved = Vec::new();
let named = named
.into_iter()
.map(|subject| (SubjectBinding::Named, Some(subject)));
for (position, (binding, given)) in declared
.iter()
.cloned()
.map(|binding| (binding, None))
.chain(named)
.enumerate()
{
let (subject, trusted) = match (&binding, given) {
(_, Some(subject)) => (subject, true),
(SubjectBinding::Input { pointer }, None) => {
let Some(selected) = input.project_pointer(pointer) else {
return Err(refuse(&binding, "it selects nothing in the run's input"));
};
let subject = match selected.peek() {
Value::String(text) => text.clone(),
Value::Number(number) => number.to_string(),
_ => {
return Err(refuse(&binding, "it selects neither a string nor a number"));
}
};
(
subject,
selected.label().trust == crate::core::Trust::Trusted,
)
}
(SubjectBinding::Case, None) => {
if !in_case {
return Err(refuse(&binding, "the run belongs to no case"));
}
(CASE_PENDING.to_owned(), true)
}
(SubjectBinding::Named, None) => {
return Err(refuse(&binding, "no subject was named"));
}
};
if subject.trim().is_empty() {
return Err(refuse(&binding, "it resolves to an empty subject"));
}
let index = u16::try_from(position)
.map_err(|_| refuse(&binding, "a run binds at most 65536 subjects"))?;
resolved.push(crate::journal::BoundSubject {
index,
binding,
trusted,
subject,
});
}
Ok(resolved)
}
const CASE_PENDING: &str = "$case";
fn fill_case_subjects(
mut bound: Vec<crate::journal::BoundSubject>,
case: Option<&CaseContext>,
) -> Vec<crate::journal::BoundSubject> {
if let Some(case) = case {
for entry in &mut bound {
if entry.binding == crate::journal::SubjectBinding::Case {
entry.subject = case.case_id.to_string();
}
}
}
bound
}
fn subject_refs(
run: RunId,
bound: &[crate::journal::BoundSubject],
) -> BTreeSet<crate::core::SubjectRef> {
bound
.iter()
.map(|b| crate::core::SubjectRef {
run,
index: b.index,
})
.collect()
}
fn recorded_subjects(run: RunId, records: &[Record]) -> BTreeSet<crate::core::SubjectRef> {
records
.iter()
.find_map(|r| match r.kind() {
RecordKind::DataSubjectBound { bindings } => Some(subject_refs(run, bindings)),
_ => None,
})
.unwrap_or_default()
}
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)
}
const fn recurs_on_every_resume(error: &RuntimeError) -> bool {
matches!(
error,
RuntimeError::Encoding(_)
| RuntimeError::ChainBroken { .. }
| RuntimeError::CanonicalizationChanged { .. }
| 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 !cx.abandoning
&& matches!(&result, Err(SkillError::Step(e)) if crate::runtime::group::is_pause(e))
{
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(SkillError::Step(StepError::Withheld { subject, reason })) => {
(RunStatus::Withheld { subject, reason }, None)
}
Err(e) => {
let msg = e.to_string();
match &e {
SkillError::Step(StepError::NonDeterminism {
seq,
expected,
actual,
detail,
}) => {
tracing::error!(
target: telemetry::NONDETERMINISM,
%seq, %expected, %actual,
error_type = "nondeterminism",
reason_digest = %telemetry::reason_digest(detail),
);
meter.count(metrics::DIVERGENCES, "");
}
SkillError::Step(StepError::ReplayOverrun { actual, kind }) => {
tracing::error!(
target: telemetry::NONDETERMINISM,
%actual, %kind, overrun = true,
);
meter.count(metrics::DIVERGENCES, "");
}
SkillError::Step(StepError::Undecidable { key, detail, .. }) => {
tracing::error!(
target: telemetry::UNDECIDABLE,
%key,
error_type = "undecidable",
reason_digest = %telemetry::reason_digest(detail),
);
meter.count(metrics::UNDECIDABLE, "");
}
SkillError::Step(StepError::Unreproducible { what, detail }) => {
tracing::error!(
target: telemetry::UNREPRODUCIBLE,
error_type = "unreproducible",
reason_digest = %telemetry::reason_digest(&format!("{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>,
succeeded: 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,
subjects: BTreeSet<crate::core::SubjectRef>,
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";
pub(crate) 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_abandonment(records).map_or_else(
|| RunStatus::Quarantined(UNATTRIBUTED.to_owned()),
|(actor, reason)| RunStatus::Abandoned { actor, reason },
)),
"cancelled" => Some(recorded_cancellation(records).map_or_else(
|| RunStatus::Quarantined(UNATTRIBUTED.to_owned()),
|(actor, reason)| RunStatus::Cancelled { actor, reason },
)),
"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()),
"replanning" => RunStatus::Replanning(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_cancellation(records).map_or_else(
|| RunStatus::Quarantined(UNATTRIBUTED.to_owned()),
|(actor, reason)| RunStatus::Cancelled { actor, reason },
),
"abandoned" => recorded_abandonment(records).map_or_else(
|| RunStatus::Quarantined(UNATTRIBUTED.to_owned()),
|(actor, reason)| RunStatus::Abandoned { actor, reason },
),
WITHHELD_OUTCOME => recorded_withholding(records).map_or_else(
|| RunStatus::Quarantined(UNATTRIBUTED.to_owned()),
|(subject, reason)| RunStatus::Withheld { subject, reason },
),
super::sweeper::SWEEP_OUTCOME => RunStatus::Swept,
OBSERVED_OUTCOME => RunStatus::Observed,
BREAK_GLASS_OUTCOME => recorded_crossing(records).map_or_else(
|| RunStatus::Quarantined(UNATTRIBUTED.to_owned()),
|(actor, reason)| RunStatus::BrokeGlass { actor, reason },
),
HALT_LIFTED_OUTCOME => recorded_halt_lift(records).map_or_else(
|| RunStatus::Quarantined(UNATTRIBUTED.to_owned()),
|actor| RunStatus::HaltLifted { actor },
),
HOLD_RELEASED_OUTCOME => recorded_hold_release(records).map_or_else(
|| RunStatus::Quarantined(UNATTRIBUTED.to_owned()),
|actor| RunStatus::HoldReleased { actor },
),
other => RunStatus::Quarantined(format!(
"recorded as '{other}', which this build does not recognise"
)),
}
}
fn sealed_payload(value: &serde_json::Value) -> bool {
#[cfg(feature = "keyring")]
{
crate::journal::payload::is_sealed(value)
}
#[cfg(not(feature = "keyring"))]
{
let _ = value;
false
}
}
fn recorded_ending(records: &[Record]) -> Option<String> {
records.iter().rev().find_map(|r| match r.kind() {
RecordKind::RunConcluded { outcome, .. } => Some(outcome.clone()),
RecordKind::RunSuspended { .. } => Some("suspended".to_owned()),
_ => None,
})
}
fn ended_by_an_act(outcome: &str) -> bool {
matches!(outcome, "cancelled" | "abandoned")
|| [
super::sweeper::SWEEP_OUTCOME,
BREAK_GLASS_OUTCOME,
HALT_LIFTED_OUTCOME,
HOLD_RELEASED_OUTCOME,
WITHHELD_OUTCOME,
OBSERVED_OUTCOME,
]
.contains(&outcome)
}
fn frozen_plans(run: RunId, records: &[Record]) -> Result<(PlanIR, Vec<PlanIR>), RuntimeError> {
let mut plans = records
.iter()
.filter_map(|r| match r.kind() {
RecordKind::PlanFrozen { plan, .. } => Some(plan),
_ => None,
})
.map(|plan| {
if sealed_payload(plan) {
return Err(RuntimeError::PayloadsErased {
run: run.to_string(),
});
}
serde_json::from_value(plan.clone()).map_err(RuntimeError::Encoding)
})
.collect::<Result<Vec<PlanIR>, _>>()?
.into_iter();
let admitted = plans
.next()
.ok_or_else(|| RuntimeError::PlanContract("journal has no PlanFrozen record".into()))?;
Ok((admitted, plans.collect()))
}
fn payloads_erased(records: &[Record]) -> bool {
records.iter().any(|r| match r.kind() {
RecordKind::PlanFrozen { plan, .. } => sealed_payload(plan),
_ => false,
})
}
fn recorded_cancellation(records: &[Record]) -> Option<(crate::core::Operator, String)> {
records.iter().rev().find_map(|r| match r.kind() {
RecordKind::RunCancelled { actor, reason } => Some((actor.clone(), reason.clone())),
_ => None,
})
}
fn recorded_abandonment(records: &[Record]) -> Option<(crate::core::Operator, String)> {
records.iter().rev().find_map(|r| match r.kind() {
RecordKind::QuarantineDecided {
decider, reason, ..
} => Some((decider.clone(), reason.clone())),
_ => None,
})
}
fn recorded_crossing(records: &[Record]) -> Option<(crate::core::Operator, String)> {
records.iter().rev().find_map(|r| match r.kind() {
RecordKind::BreakGlass { actor, reason, .. } => Some((actor.clone(), reason.clone())),
_ => None,
})
}
fn recorded_halt_lift(records: &[Record]) -> Option<crate::core::Operator> {
records.iter().rev().find_map(|r| match r.kind() {
RecordKind::HaltLifted { by, .. } => Some(by.clone()),
_ => None,
})
}
fn recorded_hold_release(records: &[Record]) -> Option<crate::core::Operator> {
records.iter().rev().find_map(|r| match r.kind() {
RecordKind::HoldReleased { by, .. } => Some(by.clone()),
_ => None,
})
}
fn recorded_withholding(records: &[Record]) -> Option<(String, String)> {
records.iter().rev().find_map(|r| match r.kind() {
RecordKind::AuthorityWithheld {
subject, reason, ..
} => Some((subject.clone(), reason.clone())),
_ => None,
})
}
#[allow(clippy::disallowed_methods)]
pub(super) fn now_for_admission() -> crate::core::Timestamp {
crate::core::Timestamp::now_utc()
}
fn removal_failed(run: RunId, e: &crate::core::StoreError) -> RuntimeError {
RuntimeError::ControlStands {
run: run.to_string(),
detail: e.to_string(),
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ControlLifted {
pub record: RunId,
pub removed: bool,
}
const BREAK_GLASS_EPOCH: crate::core::Epoch = LEASE_FREE_EPOCH;
const BREAK_GLASS_OUTCOME: &str = "broke-glass";
pub const HALT_LIFTED_OUTCOME: &str = "halt-lifted";
pub const HOLD_RELEASED_OUTCOME: &str = "hold-released";
pub const WITHHELD_OUTCOME: &str = "withheld";
pub const OBSERVED_OUTCOME: &str = "observed";
pub(crate) const LEASE_FREE_EPOCH: crate::core::Epoch = 1;
pub(crate) const LEASE_FREE_OUTCOMES: &[&str] = &[
super::sweeper::SWEEP_OUTCOME,
BREAK_GLASS_OUTCOME,
HALT_LIFTED_OUTCOME,
HOLD_RELEASED_OUTCOME,
OBSERVED_OUTCOME,
];
#[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
}
}
pub trait RunStreamObserver: Send + Sync + std::fmt::Debug {
fn event(&self, run: RunId, event: Tainted<crate::model::ModelStreamEvent>);
}
#[derive(Debug)]
pub struct RuntimeBuilder {
store: Arc<dyn JournalStore>,
witnesses: Vec<Arc<dyn crate::journal::Witness>>,
quorum: Option<crate::journal::WitnessQuorum>,
witness_interval: Option<std::time::Duration>,
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>>,
#[cfg(feature = "manifest")]
checkers: crate::content::Checkers,
streams: Option<Arc<dyn RunStreamObserver>>,
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>>,
disclosures: Option<Arc<dyn crate::disclosure::DisclosureRegister>>,
#[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
}
fn states_a_money_ceiling(&self) -> bool {
#[cfg(feature = "manifest")]
if self
.agents
.iter()
.filter_map(|agent| agent.manifest.as_ref())
.any(|manifest| manifest.budget().max_minor_units.is_some())
{
return true;
}
self.budget.max_minor_units.is_some()
}
#[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
}
#[must_use]
pub fn disclosures(mut self, register: Arc<dyn crate::disclosure::DisclosureRegister>) -> Self {
self.disclosures = Some(register);
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 fn observe_model_streams(mut self, observer: Arc<dyn RunStreamObserver>) -> Self {
self.streams = Some(observer);
self
}
#[cfg(feature = "manifest")]
#[must_use]
pub fn content_checker(mut self, checker: Arc<dyn crate::content::ContentChecker>) -> Self {
Arc::make_mut(&mut self.checkers).insert(checker.name().to_owned(), checker);
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 witness_interval(mut self, interval: std::time::Duration) -> Self {
self.witness_interval = Some(interval);
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(feature = "manifest")]
fn check_content_checkers(
checkers: &crate::content::Checkers,
m: &crate::manifest::Manifest,
) -> Result<(), BuildError> {
let Some(content) = m.spec.security.content.as_ref() else {
return Ok(());
};
for check in &content.checks {
let refused = |detail: String| BuildError::ContentCheck {
agent: m.metadata.name.clone(),
check: check.id.clone(),
detail,
};
let Some(checker) = checkers.get(&check.checker) else {
return Err(refused(format!(
"no checker named '{}' is registered on this plane",
check.checker
)));
};
if let Some(unknown) = check
.on
.keys()
.find(|category| !checker.categories().contains(*category))
{
return Err(refused(format!(
"'{unknown}' is not a category checker '{}' declares",
check.checker
)));
}
}
Ok(())
}
#[cfg(feature = "manifest")]
fn declared_provider(
providers: &HashMap<String, Arc<dyn crate::model::ModelProvider>>,
m: &crate::manifest::Manifest,
) -> Result<Arc<dyn crate::model::ModelProvider>, BuildError> {
let model = m
.spec
.models
.as_ref()
.and_then(|x| x.privileged.as_ref())
.ok_or_else(|| BuildError::DeclarativeWithoutModel {
agent: m.metadata.name.clone(),
})?;
providers
.get(&model.provider)
.map(Arc::clone)
.ok_or_else(|| BuildError::UnknownProvider {
agent: m.metadata.name.clone(),
provider: model.provider.clone(),
})
}
#[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 semantic.embedder.pricing().is_none() && self.states_a_money_ceiling() {
return Err(BuildError::UnpricedEmbedder { embedder });
}
let Some(memories) = self.memories.take() else {
return Err(BuildError::SemanticMemoryWithoutStore);
};
let reached = memories.erasure_index();
let indexed: Arc<dyn crate::memory::MemoryStore> = match reached {
Some(index) if Arc::ptr_eq(&index, &semantic.retriever) => memories,
Some(_) => return Err(BuildError::MemoryIndexedElsewhere),
None if memories.seals() => {
return Err(BuildError::SealedMemoryMissesIndex);
}
None => Arc::new(crate::memory::IndexedMemoryStore::new(
Arc::clone(&memories),
Arc::clone(&semantic.retriever),
)),
};
self.memories = Some(indexed);
}
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 provider = if execution.kind == crate::manifest::ExecutionKind::Call {
None
} else {
Some(Self::declared_provider(&self.providers, &m)?)
};
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
| crate::manifest::ExecutionKind::Call
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(),
provider.clone(),
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;
}
Self::check_content_checkers(&self.checkers, m)?;
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,
});
}
}
}
}
#[cfg(feature = "manifest")]
let rates = {
let mut declarations: Vec<&Arc<crate::manifest::Manifest>> =
governed_by.values().collect();
declarations.sort_by(|a, b| a.metadata.name.cmp(&b.metadata.name));
let mut ceilings: BTreeMap<String, Vec<crate::quota::RateCeiling>> = BTreeMap::new();
let mut first = None;
for m in declarations {
for grant in &m.spec.tools {
let Some(rate) = grant.rate_limit else {
continue;
};
first.get_or_insert_with(|| (m.metadata.name.clone(), grant.reference.clone()));
let stated = ceilings.entry(grant.reference.clone()).or_default();
let ceiling = crate::quota::RateCeiling::from(rate);
if !stated.contains(&ceiling) {
stated.push(ceiling);
stated.sort();
}
}
}
if ceilings.is_empty() {
None
} else if let Some(store) = self.quotas.clone() {
Some(Arc::new(super::ctx::Rates { store, ceilings }))
} else {
let (agent, grant) = first.unwrap_or_default();
return Err(BuildError::RateLimitWithoutQuotaStore { agent, grant });
}
};
if self.witness_interval.is_some() && self.witnesses.is_empty() {
return Err(BuildError::IntervalWithoutWitnesses);
}
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),
interval: self.witness_interval,
last_met: std::sync::Mutex::new(None),
}))
}
(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,
#[cfg(feature = "manifest")]
checkers: self.checkers,
streams: self.streams,
quotas: self.quotas,
#[cfg(feature = "manifest")]
rates,
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,
disclosures: self.disclosures,
#[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> {
let problems = runtime_policy_problems(engine, identity);
if problems.is_empty() {
return Ok(());
}
Err(BuildError::PolicyUnevaluable {
problems: problems.join("; "),
})
}
fn runtime_policy_problems(
engine: &dyn crate::core::PolicyEngine,
identity: Option<&crate::core::Delegation>,
) -> Vec<String> {
use crate::policy::requests;
let run = RunId(ulid::Ulid::nil());
let step = StepId(0);
let label = crate::core::Label::untrusted(crate::core::SourceId::new("preflight"));
let release = crate::core::Release::whole(
crate::core::ReleaseScope::trust(),
"preflight",
"tool://preflight/probe",
["preflight"],
);
let tool_args = serde_json::json!({
"server": "preflight",
"tool": "preflight",
"arguments": {},
"protected_fields": [],
});
let caller = crate::core::Delegation::root(crate::core::Principal::new(
"preflight",
crate::core::Scope::root(),
));
let chains: Vec<Option<&crate::core::Delegation>> = match identity {
Some(chain) => vec![Some(chain)],
None => vec![None, Some(&caller)],
};
let mut built = Vec::new();
for chain in chains {
let acting = requests::Acting {
tenant: "preflight",
capability: "preflight.capability",
agent: None,
chain,
};
for (kind, args) in [
("preflight.effect", serde_json::json!({})),
("tool.call", tool_args.clone()),
] {
for (mutates, labelled) in [(false, false), (true, false), (false, true), (true, true)]
{
built.push(requests::effect(
&acting,
run,
step,
kind,
&args,
mutates,
labelled.then_some(&label),
));
}
}
built.push(requests::release(&acting, run, step, &release, &label));
built.push(requests::admission(&acting, &serde_json::json!({})));
}
let requests: Vec<crate::core::PolicyRequest<'_>> = built
.iter()
.map(requests::GatedRequest::as_request)
.collect();
let mut problems = engine.preflight(&requests);
problems.sort();
problems.dedup();
problems
}
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(())
}
fn refuse_reserved_kind(event: &InboundEvent) -> Result<(), RuntimeError> {
if crate::core::is_reserved_kind(&event.kind) {
return Err(RuntimeError::ReservedEventKind {
kind: event.kind.clone(),
});
}
Ok(())
}
impl Runtime {
pub async fn deliver(&self, event: &InboundEvent) -> Result<Delivery, RuntimeError> {
refuse_reserved_kind(event)?;
self.deliver_buffered(event, true).await
}
pub(crate) async fn deliver_minted(
&self,
event: &InboundEvent,
) -> Result<Delivery, RuntimeError> {
self.deliver_buffered(event, false).await
}
async fn deliver_buffered(
&self,
event: &InboundEvent,
rematch: bool,
) -> 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();
let fresh = events
.buffer(event, now)
.await
.map_err(RuntimeError::from_store)?;
if !fresh && !rematch {
return Ok(Delivery::Duplicate);
}
let Some(sub) = events
.match_waiter(event, now)
.await
.map_err(RuntimeError::from_store)?
else {
return Ok(if fresh {
Delivery::Buffered
} else {
Delivery::Duplicate
});
};
self.resume_subscription(events, sub, event).await
}
pub async fn deliver_to(
&self,
run: RunId,
event: &InboundEvent,
) -> Result<Delivery, RuntimeError> {
refuse_reserved_kind(event)?;
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
.park_wait(&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 history = self
.store
.read(sub.run, 1)
.await
.map_err(RuntimeError::from_store)?;
let already_recorded = history.iter().any(|record| {
record.effect_key() == Some(sub.effect)
&& matches!(record.kind(), RecordKind::EffectDone { .. })
});
let closed = resume_is_closed(&history).is_some();
if !already_recorded && !closed {
let content = self.delivered_verdict(&history, &event.payload);
self.store
.append(
lease.epoch,
vec![{
let mut a = Append::new(
sub.run,
RecordKind::EffectDone {
output: event.payload.clone(),
source: Some(event.source.clone()),
by: event.by.clone(),
spend: crate::core::Spend::default(),
declared: crate::core::DeclaredOutput::untrusted(),
content,
elapsed_ms: None,
},
)
.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)?;
}
if !closed {
events
.unsubscribe(sub.run, sub.effect)
.await
.map_err(RuntimeError::from_store)?;
}
match self.resume_holding(sub.run, lease).await {
Ok(_) | Err(RuntimeError::Fenced { .. }) => {}
Err(
RuntimeError::NoProvider { .. }
| RuntimeError::NoCaseStore { .. }
| RuntimeError::PayloadsSealed { .. },
) => {
return Ok(Delivery::Buffered);
}
Err(e) => return Err(e),
}
if closed {
return Ok(Delivery::Buffered);
}
Ok(Delivery::Resumed { run: sub.run })
}
async fn retire_waits(&self, run: RunId) {
if let Some(timers) = &self.timers
&& let Err(error) = timers.disarm_run(run).await
{
tracing::error!(%run, %error, "a closed run's timers could not be retired");
}
let unanswered = match self.store.read(run, 1).await {
Ok(records) => unanswered_waits(&records),
Err(error) => {
tracing::error!(%run, %error, "a closed run's waits could not be read to retire them");
return;
}
};
if let Some(events) = &self.events {
match events.unsubscribe_run(run, &unanswered).await {
Ok(retired) => {
for event in retired.released {
self.offer_released(events, &event).await;
}
}
Err(error) => {
tracing::error!(%run, %error, "a closed run's subscriptions could not be retired");
}
}
}
if let Some(tasks) = &self.tasks {
let awaited: Vec<_> = unanswered
.iter()
.map(|effect| crate::core::TaskId::derive(run, *effect))
.collect();
if let Err(error) = tasks.withdraw_run(run, &awaited).await {
tracing::error!(%run, %error, "a closed run's open tasks could not be withdrawn");
}
}
}
async fn offer_released(
&self,
events: &Arc<dyn crate::case::EventStore>,
event: &InboundEvent,
) {
let sub = match events.match_waiter(event, now_for_admission()).await {
Ok(Some(sub)) => sub,
Ok(None) => return,
Err(error) => {
tracing::error!(event = %event.id, %error, "a released message could not be offered");
return;
}
};
if let Err(error) = Box::pin(self.resume_subscription(events, sub, event)).await {
tracing::error!(event = %event.id, %error, "a released message could not be delivered");
}
}
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 parked = events
.parked_waits(limit)
.await
.map_err(RuntimeError::from_store)?;
let examined = parked.len();
let mut delivered = 0usize;
for sub in parked {
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::HaltLifted {
actor: crate::core::Operator::asserted("ops").expect("a name"),
},
RunStatus::HoldReleased {
actor: crate::core::Operator::asserted("ops").expect("a name"),
},
RunStatus::Withheld {
subject: "alice".into(),
reason: "credential withdrawn".into(),
},
RunStatus::Observed,
]
}
#[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(),
14,
"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 written in super::every_status() {
if written.is_suspended() {
continue;
}
let outcome = written.as_str();
let status = super::recorded_status(
outcome,
written.own_account().as_deref(),
match &written {
RunStatus::Exhausted(limit) => Some(limit),
_ => None,
},
&sealed_as(outcome),
);
assert_eq!(
status.as_str(),
outcome,
"this build concludes 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_a_person_caused_reports_their_own_words() {
for (outcome, expected) in [
("cancelled", "stop"),
("abandoned", "nobody could establish what happened"),
(super::BREAK_GLASS_OUTCOME, "INC-42"),
(super::WITHHELD_OUTCOME, "credential withdrawn"),
] {
let status = super::recorded_status(outcome, None, None, &sealed_as(outcome));
assert_eq!(
status.reason().as_deref(),
Some(expected),
"'{outcome}' lost the reason the person gave: {status:?}"
);
}
}
#[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 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(),
}],
super::HALT_LIFTED_OUTCOME => vec![RecordKind::HaltLifted {
scope: "tenant".into(),
by: who(),
at: crate::core::Timestamp::UNIX_EPOCH,
reason: "INC-42".into(),
thrown_by: who(),
thrown_at: crate::core::Timestamp::UNIX_EPOCH,
}],
super::HOLD_RELEASED_OUTCOME => vec![RecordKind::HoldReleased {
by: who(),
at: crate::core::Timestamp::UNIX_EPOCH,
placed_by: who(),
placed_at: crate::core::Timestamp::UNIX_EPOCH,
}],
super::WITHHELD_OUTCOME => vec![RecordKind::AuthorityWithheld {
subject: "alice".into(),
reason: "credential withdrawn".into(),
by: who(),
}],
_ => 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"
);
}
}