use serde::{Deserialize, Serialize};
use crate::core::Spend;
use crate::core::{Digest, EffectKey, Sensitivity, Seq};
macro_rules! debug_is_display {
($($t:ty),+ $(,)?) => {$(
impl ::core::fmt::Debug for $t {
fn fmt(&self, f: &mut ::core::fmt::Formatter<'_>) -> ::core::fmt::Result {
::core::fmt::Display::fmt(self, f)
}
}
)+};
}
pub(crate) use debug_is_display;
fn rate_limited_message(detail: &str, retry_after: Option<std::time::Duration>) -> String {
match retry_after {
Some(window) => {
format!("effect rate limited: {detail} (the peer asked for {window:?})")
}
None => format!("effect rate limited: {detail}"),
}
}
const LISTED_CAPABILITIES: usize = 10;
fn no_provider_message(target: &str, available: &[String]) -> String {
if available.is_empty() {
return format!(
"no skill provides capability '{target}', and this plane has none at all — \
register one with `RuntimeBuilder::skill(..)`, or an agent's with \
`.agent(Agent::new(&manifest))`"
);
}
let shown = available
.iter()
.take(LISTED_CAPABILITIES)
.cloned()
.collect::<Vec<_>>()
.join(", ");
let rest = available.len().saturating_sub(LISTED_CAPABILITIES);
let and_more = if rest == 0 {
String::new()
} else {
format!(", and {rest} more")
};
format!(
"no skill provides capability '{target}' — this plane provides: {shown}{and_more}. \
`run` takes a capability, not a skill name; a skill declares its own with \
`SkillDescriptor::new(..).provides(..)`"
)
}
#[derive(thiserror::Error)]
#[non_exhaustive]
pub enum RuntimeError {
#[error("policy denied: {0}")]
PolicyDenied(#[from] PolicyError),
#[error(transparent)]
Delegation(#[from] crate::core::DelegationError),
#[error(transparent)]
TaskClaim(#[from] crate::core::ClaimError),
#[error(
"task {task} holds a proposal this plane cannot show — {reason} — so an approval \
would be of arguments nobody was shown; decide it where the key ring that sealed \
it is wired, or reject it"
)]
ProposalWithheld {
task: String,
reason: crate::core::Withheld,
},
#[error("plan contract violation: {0}")]
PlanContract(String),
#[error(
"this process serves no plane for tenant '{0}' — refused rather than \
defaulted, because a fallback would serve another tenant's data"
)]
UnknownTenant(String),
#[error(
"event kind '{kind}' is in the `agentplane.` namespace, which only this plane mints — \
a task is decided on the worklist, never by posting its answer as an event"
)]
ReservedEventKind { kind: String },
#[error(
"the policy bundle changed under an open run: admitted under {}, and this plane holds {} \
— the run is quarantined; reopen it with `quarantine` and `replay` it on a plane \
holding the recorded bundle, or abandon it with `quarantine`",
bundle_named(.recorded.as_ref()),
bundle_named(.configured.as_ref())
)]
PolicyBundleChanged {
recorded: Option<crate::core::Digest>,
configured: Option<crate::core::Digest>,
},
#[error(
"the declaration for `{agent}` changed under an open run: admitted under {recorded}, \
and this plane holds {configured} — the run is quarantined; reopen it with \
`quarantine` and `replay` it under the revision that wrote the journal, or abandon \
it with `quarantine`"
)]
DeclarationChanged {
agent: String,
recorded: crate::core::Digest,
configured: crate::core::Digest,
},
#[error(
"this run's derived digests were produced by canonicalization rule \
{recorded} and this build implements {implemented}, so its effect keys \
cannot be recomputed here. The journal is intact — the chain hashes \
stored bytes, not re-canonicalized ones — and this is not a divergence"
)]
CanonicalizationChanged { recorded: u16, implemented: u16 },
#[error(
"run {run}'s recorded plan is sealed to a destroyed key: its payloads were \
erased, so it cannot be replayed or resumed. The journal is intact and \
still verifies, and nothing is wrong with this build"
)]
PayloadsErased { run: String },
#[error(
"run {run}'s recorded plan is sealed and this plane holds no key ring to open \
it: nothing is known to be erased — replay it where the key ring it was sealed \
under is wired"
)]
PayloadsSealed { run: String },
#[error("{}", no_provider_message(target, available))]
NoProvider {
target: String,
available: Vec<String>,
},
#[error(transparent)]
QuotaExceeded(#[from] crate::quota::QuotaError),
#[error(
"this instance is draining and did not admit the run — retry; \
nothing was written and another instance can take it"
)]
Draining,
#[error(
"quota settlement for run {run} at epoch {epoch} is pending: {detail} — \
the run remains leased so recovery can retry without charging twice"
)]
QuotaSettlementPending {
run: String,
epoch: u64,
detail: String,
},
#[error("journal integrity broken at seq {seq}: {detail}")]
ChainBroken { seq: Seq, detail: String },
#[error("fenced at run {run}: held epoch {held}, store is at {current}")]
Fenced {
run: String,
held: u64,
current: u64,
},
#[error("run {run} is leased by '{owner}' for another {remaining_secs}s")]
LeaseHeld {
run: String,
owner: String,
remaining_secs: u64,
},
#[error(
"run {run} is '{status}', not quarantined — only a run the runtime could not decide can \
be reopened or abandoned"
)]
NotQuarantined { run: String, status: String },
#[error(
"effect {effect} in run {run} is not in doubt — its outcome is on the record, and an \
assertion may supply a missing fact but never replace a recorded one"
)]
NotUndecided { run: String, effect: String },
#[error(
"run {run} is quarantined, and cancelling unwinds — which is exactly what a run holding \
an unknown outcome must not do. Answer the doubt and reopen it, or abandon it, which \
closes it without unwinding"
)]
CannotUnwind { run: String },
#[error("run {run} already concluded as '{outcome}'; there is nothing left to stop")]
AlreadyConcluded { run: String, outcome: String },
#[error(transparent)]
Store(#[from] StoreError),
#[error(transparent)]
Encoding(#[from] serde_json::Error),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Disposition {
DidNotHappen,
InDoubt,
Landed,
}
impl Disposition {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::DidNotHappen => "did_not_happen",
Self::InDoubt => "in_doubt",
Self::Landed => "landed",
}
}
#[must_use]
pub fn is_definitely_safe_to_repeat(self) -> bool {
matches!(self, Self::DidNotHappen)
}
}
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum EffectError {
#[error("driver '{driver}' unavailable: {detail}")]
Unavailable { driver: String, detail: String },
#[error("effect rejected: {0}")]
Rejected(String),
#[error("{}", rate_limited_message(detail, *retry_after))]
RateLimited {
detail: String,
retry_after: Option<std::time::Duration>,
},
#[error("effect refused: {0}")]
Refused(String),
#[error("driver '{driver}' did not answer within {waited_ms}ms")]
Timeout { driver: String, waited_ms: u64 },
#[error("driver '{driver}' interrupted: {detail}")]
Interrupted { driver: String, detail: String },
#[error("effect consumed resources and failed: {detail}")]
Metered {
detail: String,
spend: Spend,
disposition: Disposition,
},
#[error("effect performed and failed: {0}")]
Performed(String),
#[error("effect output did not match its declared type: {0}")]
OutputShape(#[from] serde_json::Error),
#[error("{detail}")]
Final {
detail: String,
disposition: Disposition,
},
#[error("{0}")]
Other(String),
}
impl EffectError {
#[must_use]
pub fn spend(&self) -> Spend {
match self {
Self::Metered { spend, .. } => *spend,
Self::Unavailable { .. }
| Self::Rejected(_)
| Self::RateLimited { .. }
| Self::Refused(_)
| Self::Timeout { .. }
| Self::Interrupted { .. }
| Self::Performed(_)
| Self::OutputShape(_)
| Self::Final { .. }
| Self::Other(_) => Spend::default(),
}
}
#[must_use]
pub const fn retry_after(&self) -> Option<std::time::Duration> {
match self {
Self::RateLimited { retry_after, .. } => *retry_after,
_ => None,
}
}
#[must_use]
pub const fn class(&self) -> &'static str {
match self {
Self::Unavailable { .. } => "unavailable",
Self::Rejected(_) => "rejected",
Self::RateLimited { .. } => "rate_limited",
Self::Refused(_) => "refused",
Self::Timeout { .. } => "timeout",
Self::Interrupted { .. } => "interrupted",
Self::Metered { .. } => "metered",
Self::Performed(_) => "performed",
Self::OutputShape(_) => "output_shape",
Self::Final { .. } => "final",
Self::Other(_) => "other",
}
}
#[must_use]
pub fn disposition(&self) -> Disposition {
match self {
Self::Metered { disposition, .. } | Self::Final { disposition, .. } => *disposition,
Self::Unavailable { .. }
| Self::Rejected(_)
| Self::RateLimited { .. }
| Self::Refused(_) => Disposition::DidNotHappen,
Self::OutputShape(_) | Self::Performed(_) => Disposition::Landed,
Self::Timeout { .. } | Self::Interrupted { .. } | Self::Other(_) => {
Disposition::InDoubt
}
}
}
}
#[derive(thiserror::Error)]
#[non_exhaustive]
pub enum SkillError {
#[error("input did not match the declared schema: {0}")]
Input(String),
#[error(transparent)]
Step(#[from] StepError),
#[error(transparent)]
Tool(#[from] crate::tools::ToolError),
#[error("{0}")]
Other(String),
}
#[derive(thiserror::Error)]
#[non_exhaustive]
pub enum StepError {
#[error(transparent)]
Effect(#[from] EffectError),
#[error(transparent)]
Policy(#[from] PolicyError),
#[error(transparent)]
Store(#[from] StoreError),
#[error("{0}")]
Encoding(#[from] serde_json::Error),
#[error(transparent)]
Tool(#[from] crate::tools::ToolError),
#[error(
"effect {key} is undecidable ({detail}); recovery mode {recovery:?} forbids \
guessing — run quarantined"
)]
Undecidable {
key: EffectKey,
recovery: crate::core::Recovery,
detail: String,
},
#[error(
"effect {key} completed ({disposition:?}) but its outcome could not be recorded: {detail}"
)]
Unrecorded {
key: EffectKey,
disposition: crate::core::Disposition,
detail: String,
},
#[error(
"{what} is no longer what this run pinned it to ({detail}) — the history \
cannot be reproduced, so the run is quarantined rather than failed"
)]
Unreproducible { what: String, detail: String },
#[error("non-determinism at seq {seq}: {detail} (history {expected}, this build {actual})")]
NonDeterminism {
seq: Seq,
expected: EffectKey,
actual: EffectKey,
detail: String,
},
#[error(transparent)]
Budget(#[from] crate::core::BudgetExceeded),
#[error("suspended: {0}")]
Suspended(crate::core::SuspendReason),
#[error("policy denied '{action}' on '{resource}': {reason}")]
Denied {
action: String,
resource: String,
reason: String,
},
#[error("effect group '{group}': {detail}")]
GroupFootprint { group: String, detail: String },
#[error("effect group aborted and fully reversed: {what}")]
GroupAborted { what: String },
#[error("effect group '{group}' could not be settled: {detail} — run quarantined")]
GroupUnsettled { group: String, detail: String },
#[error(
"replay overrun: journal is exhausted but the run requested `{kind}` ({actual}) — \
this build performs more effects than the recorded one"
)]
ReplayOverrun { actual: EffectKey, kind: String },
}
impl SkillError {
#[must_use]
pub const fn class(&self) -> &'static str {
match self {
Self::Input(_) => "input",
Self::Step(step) => step.class(),
Self::Tool(_) => "tool",
Self::Other(_) => "other",
}
}
}
impl StepError {
#[must_use]
pub const fn class(&self) -> &'static str {
match self {
Self::Effect(e) => e.class(),
Self::Policy(_) => "policy",
Self::Store(_) => "store",
Self::Encoding(_) => "encoding",
Self::Tool(_) => "tool",
Self::Undecidable { .. } => "undecidable",
Self::Unrecorded { .. } => "unrecorded",
Self::Unreproducible { .. } => "unreproducible",
Self::NonDeterminism { .. } => "nondeterminism",
Self::Budget(_) => "budget",
Self::Suspended(_) => "suspended",
Self::Denied { .. } => "denied",
Self::GroupFootprint { .. } => "group_footprint",
Self::GroupAborted { .. } => "group_aborted",
Self::GroupUnsettled { .. } => "group_unsettled",
Self::ReplayOverrun { .. } => "replay_overrun",
}
}
}
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum PolicyError {
#[error("principal '{principal}' may not '{action}' on '{resource}'")]
Denied {
principal: String,
action: String,
resource: String,
},
#[error("{reason}")]
Recorded { reason: String },
#[error("untrusted data may not reach mutating sink '{sink}' without an authorized release")]
TaintGate { sink: String },
#[error("sink '{sink}' does not bind the arguments it sends to the value checked by policy")]
UnboundSinkArguments { sink: String },
#[error("sink '{sink}' must be dispatched with StepCtx::sink so its outbound value is checked")]
SinkGateRequired { sink: String },
#[error(
"sink '{sink}' attempted to send arguments other than the labeled value policy \
checked — they first differ at '{at}' (bound {bound}, sent {sent})"
)]
SinkArgumentsMismatch {
sink: String,
at: String,
bound: Digest,
sent: Digest,
},
#[error("sink '{sink}' requires protected field '{path}', but the argument is absent")]
ProtectedFieldMissing { sink: String, path: String },
#[error("untrusted data may not select protected field '{path}' of sink '{sink}'")]
ProtectedFieldTaint { sink: String, path: String },
#[error(
"untrusted data may not reach mutating sink '{sink}': the value's release \
names destination '{granted}', and this sink is '{actual}'"
)]
ReleaseDestination {
sink: String,
granted: String,
actual: String,
},
#[error(
"untrusted data may not select protected field '{path}' of sink '{sink}': \
the field's release names destination '{granted}', and this sink is '{actual}'"
)]
ProtectedFieldReleaseDestination {
sink: String,
path: String,
granted: String,
actual: String,
},
#[error(
"protected field '{path}' of sink '{sink}' derives from undeclared source '{actual_source}'"
)]
ProtectedFieldSource {
sink: String,
path: String,
actual_source: String,
},
#[error(
"protected field '{path}' of sink '{sink}' carries a value outside the \
declared set — the manifest enumerates what may stand in this field, \
and this value is not one of them"
)]
ProtectedFieldValue { sink: String, path: String },
#[error(
"protected field '{path}' sensitivity {actual} exceeds sink '{sink}' field ceiling {ceiling}"
)]
ProtectedFieldSensitivity {
sink: String,
path: String,
actual: Sensitivity,
ceiling: Sensitivity,
},
#[error(
"release scope contains a missing or untracked field; use Tainted::object/array before releasing selected fields"
)]
UntrackedReleaseField,
#[error("invalid release: {detail}")]
InvalidRelease { detail: String },
#[error(
"sensitivity {actual} exceeds the journal ceiling {ceiling} for sink \
'{sink}' — the journal is append-only, so this argument could not be \
removed afterwards. Put the bytes in a blob and pass the digest, or \
configure a key ring so payloads are sealed under a key erasure destroys"
)]
JournalCeiling {
sink: String,
actual: crate::core::Sensitivity,
ceiling: crate::core::Sensitivity,
},
#[error("sensitivity {actual} exceeds sink '{sink}' ceiling {ceiling}")]
EgressCeiling {
sink: String,
actual: Sensitivity,
ceiling: Sensitivity,
},
#[error("delegation depth {actual} exceeds sink '{sink}' ceiling {ceiling}")]
DelegationDepth {
sink: String,
actual: usize,
ceiling: usize,
},
}
pub const REFUSED: &str = "this action was not permitted";
impl PolicyError {
#[must_use]
pub const fn for_model(&self) -> &'static str {
REFUSED
}
}
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum StoreError {
#[error("backend: {0}")]
Backend(String),
#[error("not found: {0}")]
NotFound(String),
#[error(
"record of {bytes} bytes exceeds the {limit}-byte journal limit — \
journal a digest and keep the bytes outside the chain"
)]
RecordTooLarge { bytes: usize, limit: usize },
#[error("effect {0} already started in this run")]
DuplicateEffect(EffectKey),
#[error("admission key '{key}' is already held by run {run}")]
DuplicateAdmission { key: String, run: String },
#[error(
"batch {batch} is bound to plan {stored}, and this resume offers {offered} — \
one batch runs one frozen plan; start a new batch for a new plan"
)]
BatchPlanChanged {
batch: String,
stored: String,
offered: String,
},
#[error("case {case} has moved to {current}; the write was made against {expected}")]
CaseConflict {
case: String,
expected: u64,
current: u64,
},
#[error(
"case {case} still has {outstanding} open obligation(s); closure is \
when people stop looking, so meet or cancel them before closing"
)]
ObligationsOutstanding { case: String, outstanding: usize },
#[error(
"case {case} is closed; reopen it before registering an obligation, because \
closure is when people stop looking"
)]
CaseClosed { case: String },
#[error(
"blob {digest} was erased at {at} ({reason}); storing these bytes again would \
put back what somebody asked to have removed"
)]
BlobErased {
digest: String,
at: i64,
reason: String,
},
#[error(
"obligation '{obligation}' on case {case} is {from} and may not become \
{to}; how an obligation ended is not editable, so record a late answer \
as an account of the breach rather than as a state"
)]
DeadlineFinal {
case: String,
obligation: String,
from: String,
to: String,
},
#[error("obligation '{obligation}' on case {case} is {state}, not breached")]
NotBreached {
case: String,
obligation: String,
state: String,
},
#[error("fenced: run {run} is owned at epoch {current}, writer held {held}")]
Fenced {
run: String,
held: u64,
current: u64,
},
#[error(
"the transaction's outcome is unknown — COMMIT may or may not have \
been applied: {detail}"
)]
CommitUnknown { detail: String },
#[error("run {run} is sealed as '{outcome}'; a sealed journal accepts no appends")]
RunSealed { run: String, outcome: String },
#[error("memory '{id}' is under legal hold, so nothing was erased")]
UnderLegalHold { id: String },
#[error("run {run} is leased by '{owner}' at epoch {epoch} for another {remaining_secs}s")]
LeaseHeld {
run: String,
owner: String,
epoch: u64,
remaining_secs: u64,
},
#[error(
"run {run}: the lease is not held at epoch {epoch} by this owner — it was \
released, lapsed, or taken over, and a renewal never claims"
)]
LeaseNotHeld { run: String, epoch: u64 },
#[error("corrupt record at seq {seq}: {detail}")]
Corrupt { seq: Seq, detail: String },
#[error(
"record {kind} is v{version} and this build reads v{reads} — the bytes are intact and \
hash as written; what is missing is a reader that knows the shape. If v{version} is the \
newer one, deploy readers before writers"
)]
UnknownRecordVersion {
kind: String,
version: u16,
reads: u16,
},
#[error(
"record {kind} is v{version}, the version this build writes, and its shape does not \
parse here: {detail}. The bytes are intact and hash as written, so another build \
wrote this journal — run the build that wrote it, or read this history from its \
export"
)]
UnreadableRecordShape {
kind: String,
version: u16,
detail: String,
},
#[error(transparent)]
Encoding(#[from] serde_json::Error),
}
fn bundle_named(digest: Option<&crate::core::Digest>) -> String {
digest.map_or_else(|| "no policy engine".to_owned(), ToString::to_string)
}
impl RuntimeError {
#[must_use]
pub fn from_store(e: StoreError) -> Self {
match e {
StoreError::Fenced { run, held, current } => Self::Fenced { run, held, current },
StoreError::LeaseHeld {
run,
owner,
remaining_secs,
..
} => Self::LeaseHeld {
run,
owner,
remaining_secs,
},
StoreError::Corrupt { seq, detail } => Self::ChainBroken { seq, detail },
other => Self::Store(other),
}
}
}
#[cfg(test)]
mod tests {
use super::{Disposition, LISTED_CAPABILITIES, RuntimeError, SkillError, StepError};
#[test]
fn only_a_call_that_never_left_is_safe_to_repeat_on_its_own_terms() {
assert!(Disposition::DidNotHappen.is_definitely_safe_to_repeat());
assert!(!Disposition::InDoubt.is_definitely_safe_to_repeat());
assert!(!Disposition::Landed.is_definitely_safe_to_repeat());
}
#[test]
fn a_failure_debugs_as_the_message_it_carries() {
let e = RuntimeError::NoProvider {
target: "demo.greet".to_owned(),
available: vec!["demo.other".to_owned()],
};
assert_eq!(format!("{e:?}"), e.to_string());
assert!(
!format!("{e:?}").starts_with("NoProvider"),
"the derived Debug is back: {e:?}"
);
let skill = SkillError::Other("boom".to_owned());
assert_eq!(format!("{skill:?}"), skill.to_string());
let step = StepError::Encoding(serde_json::from_str::<i32>("x").unwrap_err());
assert_eq!(format!("{step:?}"), step.to_string());
}
#[test]
fn an_unknown_capability_is_told_what_exists() {
let e = RuntimeError::NoProvider {
target: "demo.greeet".to_owned(),
available: vec!["demo.greet".to_owned(), "demo.sum".to_owned()],
};
let msg = e.to_string();
assert!(msg.contains("demo.greeet"), "{msg}");
assert!(msg.contains("demo.greet, demo.sum"), "{msg}");
}
#[test]
fn an_empty_plane_says_it_has_no_skills() {
let e = RuntimeError::NoProvider {
target: "demo.greet".to_owned(),
available: Vec::new(),
};
let msg = e.to_string();
assert!(msg.contains("has none at all"), "{msg}");
assert!(msg.contains("RuntimeBuilder::skill"), "{msg}");
}
#[test]
fn a_capped_capability_list_admits_the_cap() {
let available: Vec<String> = (0..LISTED_CAPABILITIES + 3)
.map(|i| format!("cap.{i}"))
.collect();
let msg = RuntimeError::NoProvider {
target: "nope".to_owned(),
available,
}
.to_string();
assert!(msg.contains("and 3 more"), "{msg}");
}
}
debug_is_display!(RuntimeError, SkillError, StepError);
#[cfg(any(feature = "a2a-server", feature = "mcp-server"))]
pub(crate) fn withheld_fault(
surface: &'static str,
doing: &str,
error: &dyn std::fmt::Display,
) -> &'static str {
tracing::error!(
target: "agentplane::served",
surface,
doing,
%error,
"a served request failed internally"
);
"internal error"
}