use async_trait::async_trait;
use serde::{Serialize, de::DeserializeOwned};
use serde_json::Value;
use crate::core::{
EffectError, EffectKey, Phase, Provenance, RetryPolicy, Sensitivity, Spend, StepId, Trust,
canon,
};
impl EffectKey {
#[doc(hidden)]
#[must_use]
pub fn for_effect(
step: StepId,
phase: Phase,
ordinal: u32,
attempt: u32,
d: &EffectDescriptor,
) -> Self {
Self::derive(
step,
phase,
ordinal,
attempt,
&d.kind,
&canon::value_bytes(&d.args),
)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, serde::Deserialize)]
pub struct EffectDescriptor {
pub kind: String,
pub args: Value,
}
impl EffectDescriptor {
pub fn new(kind: impl Into<String>, args: Value) -> Self {
Self {
kind: kind.into(),
args,
}
}
pub fn nullary(kind: impl Into<String>) -> Self {
Self {
kind: kind.into(),
args: Value::Null,
}
}
}
#[derive(Debug)]
pub enum Reconciliation<T> {
Landed(T),
DidNotHappen,
Inconclusive,
}
impl<T> Reconciliation<T> {
#[must_use]
pub fn disposition(&self) -> crate::core::Disposition {
use crate::core::Disposition;
match self {
Self::Landed(_) => Disposition::Landed,
Self::DidNotHappen => Disposition::DidNotHappen,
Self::Inconclusive => Disposition::InDoubt,
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, serde::Deserialize)]
#[serde(rename_all = "snake_case", tag = "mode")]
pub enum Recovery {
Retry,
Idempotent { key: String },
Reconcile,
#[default]
RequiresOperator,
}
#[async_trait]
pub trait Effect: Send + Sync {
type Output: Serialize + DeserializeOwned + Send;
fn descriptor(&self) -> EffectDescriptor;
fn attach(&mut self, _provenance: &Provenance) {}
fn mutates(&self) -> bool {
true
}
fn recovery(&self) -> Recovery {
if self.mutates() {
Recovery::RequiresOperator
} else {
Recovery::Retry
}
}
fn retry(&self) -> RetryPolicy {
RetryPolicy::default()
}
fn max_sensitivity(&self) -> Sensitivity {
Sensitivity::Public
}
fn trust(&self) -> Trust {
Trust::Untrusted
}
fn output_sensitivity(&self) -> Sensitivity {
Sensitivity::Public
}
fn spend(&self, _output: &Self::Output) -> Spend {
Spend::default()
}
async fn perform(&self) -> Result<Self::Output, EffectError>;
async fn reconcile(&self) -> Result<Reconciliation<Self::Output>, EffectError> {
Ok(Reconciliation::Inconclusive)
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn descriptors_with_reordered_args_are_equal_after_canonicalization() {
let a = EffectDescriptor::new("mcp.tools/call", json!({"b": 2, "a": 1}));
let b = EffectDescriptor::new("mcp.tools/call", json!({"a": 1, "b": 2}));
assert_eq!(
crate::core::canon::value_bytes(&a.args),
crate::core::canon::value_bytes(&b.args),
"argument order must not change an effect's identity"
);
}
#[test]
fn mutating_effects_default_to_operator_recovery() {
struct Mutating;
#[async_trait]
impl Effect for Mutating {
type Output = ();
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::nullary("test.mutate")
}
async fn perform(&self) -> Result<(), EffectError> {
Ok(())
}
}
assert!(matches!(Mutating.recovery(), Recovery::RequiresOperator));
}
#[test]
fn read_only_effects_default_to_retry() {
struct ReadOnly;
#[async_trait]
impl Effect for ReadOnly {
type Output = ();
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::nullary("test.read")
}
fn mutates(&self) -> bool {
false
}
async fn perform(&self) -> Result<(), EffectError> {
Ok(())
}
}
assert!(matches!(ReadOnly.recovery(), Recovery::Retry));
}
}