use std::error::Error;
use std::fmt;
use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use std::time::Duration;
use serde::Serialize;
use serde::de::DeserializeOwned;
use serde_json::Value;
use crate::error::RuntimeError;
use crate::failure::FailureClass;
use crate::id::{EffectId, EffectKey, IdempotencyKey};
use crate::kind::EffectKind;
use crate::policy::{Capabilities, RiskLevel};
use crate::retry::RetryPolicy;
use crate::runtime::Runtime;
use crate::store::{EffectStore, ErrorRecord};
use crate::verification::{NoVerification, Verification, VerificationMode, Verifier, VerifyWith};
#[derive(Clone, Debug)]
pub struct EffectContext {
pub(crate) id: EffectId,
pub(crate) key: EffectKey,
pub(crate) attempt: u32,
}
impl EffectContext {
pub fn effect_id(&self) -> EffectId {
self.id
}
pub fn key(&self) -> &EffectKey {
&self.key
}
pub fn idempotency_key(&self) -> IdempotencyKey {
self.key.idempotency_key()
}
pub fn attempt(&self) -> u32 {
self.attempt
}
}
#[derive(Debug)]
pub struct EffectFailure {
class: FailureClass,
request_sent: Option<bool>,
source: Box<dyn Error + Send + Sync>,
}
impl EffectFailure {
pub fn new(class: FailureClass, source: impl Into<Box<dyn Error + Send + Sync>>) -> Self {
Self {
class,
request_sent: None,
source: source.into(),
}
}
pub fn transient(source: impl Into<Box<dyn Error + Send + Sync>>) -> Self {
Self::new(FailureClass::Transient, source)
}
pub fn permanent(source: impl Into<Box<dyn Error + Send + Sync>>) -> Self {
Self::new(FailureClass::Permanent, source)
}
pub fn ambiguous(source: impl Into<Box<dyn Error + Send + Sync>>) -> Self {
Self::new(FailureClass::Ambiguous, source)
}
pub fn rate_limited(
retry_after: Option<Duration>,
source: impl Into<Box<dyn Error + Send + Sync>>,
) -> Self {
Self::new(FailureClass::RateLimited { retry_after }, source)
}
pub fn authentication(source: impl Into<Box<dyn Error + Send + Sync>>) -> Self {
Self::new(FailureClass::Authentication, source)
}
pub fn authorization(source: impl Into<Box<dyn Error + Send + Sync>>) -> Self {
Self::new(FailureClass::Authorization, source)
}
pub fn validation(source: impl Into<Box<dyn Error + Send + Sync>>) -> Self {
Self::new(FailureClass::Validation, source)
}
#[must_use]
pub fn request_sent(mut self, sent: bool) -> Self {
self.request_sent = Some(sent);
self
}
pub fn class(&self) -> FailureClass {
self.class.with_request_sent(self.request_sent)
}
pub(crate) fn to_record(&self) -> ErrorRecord {
ErrorRecord {
class: Some(self.class()),
message: self.source.to_string(),
}
}
}
impl fmt::Display for EffectFailure {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "{:?} failure: {}", self.class(), self.source)
}
}
impl Error for EffectFailure {
fn source(&self) -> Option<&(dyn Error + 'static)> {
Some(self.source.as_ref())
}
}
#[derive(Debug, Clone, PartialEq)]
#[non_exhaustive]
pub enum EffectOutcome<T> {
Committed(T),
Failed(ErrorRecord),
Rejected(ErrorRecord),
Unknown {
id: EffectId,
},
NeedsIntervention {
id: EffectId,
},
InProgress {
id: EffectId,
},
AwaitingApproval {
id: EffectId,
},
Compensated {
id: EffectId,
},
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum Precondition {
Satisfied,
Rejected {
reason: String,
},
RetryLater {
after: Duration,
reason: String,
},
}
impl Precondition {
pub fn reject(reason: impl Into<String>) -> Self {
Self::Rejected {
reason: reason.into(),
}
}
pub fn retry_later(after: Duration, reason: impl Into<String>) -> Self {
Self::RetryLater {
after,
reason: reason.into(),
}
}
}
pub(crate) type PreconditionFn =
Arc<dyn Fn(EffectContext) -> Pin<Box<dyn Future<Output = Precondition> + Send>> + Send + Sync>;
#[must_use = "an effect does nothing until `run` is awaited"]
pub struct EffectBuilder<S, V = NoVerification> {
runtime: Runtime<S>,
name: String,
key: String,
kind: EffectKind,
remote_idempotency: bool,
input: Option<Result<Value, serde_json::Error>>,
actor: Option<String>,
retry: RetryPolicy,
attempt_timeout: Option<Duration>,
precondition: Option<PreconditionFn>,
require_approval: bool,
risk: RiskLevel,
verification: VerificationMode,
verifier: V,
}
impl<S: EffectStore> EffectBuilder<S> {
pub(crate) fn new(runtime: Runtime<S>, name: String, key: String) -> Self {
let retry = runtime.default_retry();
Self {
runtime,
name,
key,
kind: EffectKind::IrreversibleWrite,
remote_idempotency: false,
input: None,
actor: None,
retry,
attempt_timeout: None,
precondition: None,
require_approval: false,
risk: RiskLevel::Low,
verification: VerificationMode::None,
verifier: NoVerification,
}
}
pub fn verify<T, F, Fut>(self, check: F) -> EffectBuilder<S, VerifyWith<F>>
where
F: Fn(EffectContext) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<Verification<T>, EffectFailure>> + Send + 'static,
{
self.with_verifier(VerificationMode::Authoritative, VerifyWith(check))
}
pub fn verify_eventually<T, F, Fut>(
self,
settle: Duration,
check: F,
) -> EffectBuilder<S, VerifyWith<F>>
where
F: Fn(EffectContext) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<Verification<T>, EffectFailure>> + Send + 'static,
{
self.with_verifier(
VerificationMode::EventuallyConsistent { settle },
VerifyWith(check),
)
}
fn with_verifier<V>(self, mode: VerificationMode, verifier: V) -> EffectBuilder<S, V> {
EffectBuilder {
runtime: self.runtime,
name: self.name,
key: self.key,
kind: self.kind,
remote_idempotency: self.remote_idempotency,
input: self.input,
actor: self.actor,
retry: self.retry,
attempt_timeout: self.attempt_timeout,
precondition: self.precondition,
require_approval: self.require_approval,
risk: self.risk,
verification: mode,
verifier,
}
}
}
impl<S: EffectStore, V> EffectBuilder<S, V> {
pub fn kind(mut self, kind: EffectKind) -> Self {
self.kind = kind;
self
}
pub fn remote_idempotency(mut self, supported: bool) -> Self {
self.remote_idempotency = supported;
self
}
pub fn input<I: Serialize + ?Sized>(mut self, input: &I) -> Self {
self.input = Some(serde_json::to_value(input));
self
}
pub fn actor(mut self, actor: impl Into<String>) -> Self {
self.actor = Some(actor.into());
self
}
pub fn retry(mut self, policy: RetryPolicy) -> Self {
self.retry = policy;
self
}
pub fn attempt_timeout(mut self, timeout: Duration) -> Self {
self.attempt_timeout = Some(timeout);
self
}
pub fn precondition<F, Fut>(mut self, check: F) -> Self
where
F: Fn(EffectContext) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Precondition> + Send + 'static,
{
self.precondition = Some(Arc::new(move |ctx| Box::pin(check(ctx))));
self
}
pub fn risk(mut self, risk: RiskLevel) -> Self {
self.risk = risk;
self
}
pub fn require_approval(mut self) -> Self {
self.require_approval = true;
self
}
pub async fn run<T, F, Fut>(self, action: F) -> Result<EffectOutcome<T>, RuntimeError>
where
T: Serialize + DeserializeOwned + Send + 'static,
F: Fn(EffectContext) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<T, EffectFailure>> + Send + 'static,
V: Verifier<T>,
{
let key = EffectKey::new(
crate::id::EffectName::new(self.name)?,
crate::id::LogicalKey::new(self.key)?,
);
let input = self.input.transpose().map_err(RuntimeError::Input)?;
let spec = EffectSpec {
fingerprint: None,
input,
key,
capabilities: Capabilities {
kind: self.kind,
remote_idempotency: self.remote_idempotency,
verification: self.verification,
},
actor: self.actor,
retry: self.retry,
attempt_timeout: self.attempt_timeout,
precondition: self.precondition,
require_approval: self.require_approval,
risk: self.risk,
automatic_retry: true,
input_stored: false,
};
self.runtime.execute(spec, action, self.verifier).await
}
}
pub(crate) struct EffectSpec {
pub(crate) key: EffectKey,
pub(crate) capabilities: Capabilities,
pub(crate) input: Option<Value>,
pub(crate) fingerprint: Option<String>,
pub(crate) actor: Option<String>,
pub(crate) retry: RetryPolicy,
pub(crate) attempt_timeout: Option<Duration>,
pub(crate) precondition: Option<PreconditionFn>,
pub(crate) require_approval: bool,
pub(crate) risk: RiskLevel,
pub(crate) automatic_retry: bool,
pub(crate) input_stored: bool,
}