1use std::any::Any;
61use std::collections::HashMap;
62use std::fmt::Display;
63use std::future::{Future, IntoFuture};
64use std::pin::Pin;
65use std::sync::Arc;
66use std::time::Duration;
67
68use serde::Serialize;
69use serde::de::DeserializeOwned;
70use serde_json::Value;
71
72use crate::compensation::{
73 CompensationContext, CompensationOutcome, CompensationSpec, Compensator,
74};
75use crate::effect::{EffectContext, EffectFailure, EffectOutcome, EffectSpec, Precondition};
76use crate::error::RuntimeError;
77use crate::id::{EffectKey, EffectName, LogicalKey};
78use crate::kind::EffectKind;
79use crate::policy::{Capabilities, RiskLevel};
80use crate::retry::RetryPolicy;
81use crate::runtime::Runtime;
82use crate::state::EffectStatus;
83use crate::store::{EffectRecord, EffectStore, StoreError};
84use crate::verification::{NoVerification, Verification, VerificationMode, VerifyWith};
85
86pub trait EffectHandler: Send + Sync + 'static {
92 const NAME: &'static str;
94
95 type Input: Serialize + DeserializeOwned + Send + Sync + 'static;
98
99 type Output: Serialize + DeserializeOwned + Send + 'static;
101
102 type Error: Into<EffectFailure> + Send + 'static;
104
105 fn kind(&self) -> EffectKind {
108 EffectKind::IrreversibleWrite
109 }
110
111 fn remote_idempotency(&self) -> bool {
114 false
115 }
116
117 fn retry_policy(&self) -> Option<RetryPolicy> {
119 None
120 }
121
122 fn attempt_timeout(&self) -> Option<Duration> {
124 None
125 }
126
127 fn risk(&self) -> RiskLevel {
130 RiskLevel::Low
131 }
132
133 fn requires_approval(&self) -> bool {
136 false
137 }
138
139 fn precondition(
142 &self,
143 ctx: &EffectContext,
144 input: &Self::Input,
145 ) -> impl Future<Output = Precondition> + Send {
146 let _ = (ctx, input);
147 async { Precondition::Satisfied }
148 }
149
150 fn execute(
152 &self,
153 ctx: &EffectContext,
154 input: &Self::Input,
155 ) -> impl Future<Output = Result<Self::Output, Self::Error>> + Send;
156}
157
158pub trait VerifiableEffect: EffectHandler {
164 fn verification_mode(&self) -> VerificationMode {
168 VerificationMode::Authoritative
169 }
170
171 fn verify(
173 &self,
174 ctx: &EffectContext,
175 input: &Self::Input,
176 ) -> impl Future<Output = Result<Verification<Self::Output>, Self::Error>> + Send;
177}
178
179pub trait CompensableEffect: EffectHandler {
186 fn compensate(
189 &self,
190 ctx: &CompensationContext,
191 input: &Self::Input,
192 output: Option<&Self::Output>,
193 ) -> impl Future<Output = Result<(), Self::Error>> + Send;
194}
195
196type BoxFuture<T> = Pin<Box<dyn Future<Output = T> + Send>>;
197
198type VerifyFn<H> = Arc<
199 dyn Fn(
200 Arc<H>,
201 EffectContext,
202 Arc<<H as EffectHandler>::Input>,
203 ) -> BoxFuture<Result<Verification<<H as EffectHandler>::Output>, EffectFailure>>
204 + Send
205 + Sync,
206>;
207
208pub struct Handler<H: EffectHandler> {
211 effect: Arc<H>,
212 verify: Option<(VerificationMode, VerifyFn<H>)>,
213 compensate: Option<Compensator>,
214}
215
216impl<H: EffectHandler> Clone for Handler<H> {
217 fn clone(&self) -> Self {
218 Self {
219 effect: Arc::clone(&self.effect),
220 verify: self.verify.clone(),
221 compensate: self.compensate.clone(),
222 }
223 }
224}
225
226impl<H: EffectHandler> Handler<H> {
227 pub fn new(handler: H) -> Self {
229 Self {
230 effect: Arc::new(handler),
231 verify: None,
232 compensate: None,
233 }
234 }
235}
236
237impl<H: CompensableEffect> Handler<H> {
238 #[must_use]
243 pub fn compensable(mut self) -> Self {
244 let handler = Arc::clone(&self.effect);
245 self.compensate = Some(Arc::new(move |ctx, input, output| {
246 let handler = Arc::clone(&handler);
247 Box::pin(async move {
248 let input: H::Input = serde_json::from_value(input.unwrap_or(Value::Null))
249 .map_err(|e| {
250 EffectFailure::permanent(format!("stored input does not match: {e}"))
251 })?;
252 let output: Option<H::Output> = output
253 .map(serde_json::from_value)
254 .transpose()
255 .map_err(|e| {
256 EffectFailure::permanent(format!("stored output does not match: {e}"))
257 })?;
258 handler
259 .compensate(&ctx, &input, output.as_ref())
260 .await
261 .map_err(Into::into)
262 })
263 }));
264 self
265 }
266}
267
268impl<H: VerifiableEffect> Handler<H> {
269 #[must_use]
272 pub fn verifiable(mut self) -> Self {
273 let mode = self.effect.verification_mode();
274 let verify: VerifyFn<H> = Arc::new(|handler, ctx, input| {
275 Box::pin(async move { handler.verify(&ctx, &input).await.map_err(Into::into) })
276 });
277 self.verify = Some((mode, verify));
278 self
279 }
280}
281
282pub(crate) struct Registered<S> {
285 typed: Arc<dyn Any + Send + Sync>,
287 pub(crate) resume: Resume<S>,
289}
290
291pub(crate) type Resume<S> = Arc<
292 dyn Fn(Runtime<S>, EffectRecord) -> BoxFuture<Result<EffectStatus, RuntimeError>> + Send + Sync,
293>;
294
295pub(crate) type Registry<S> = HashMap<&'static str, Registered<S>>;
296
297impl<S: EffectStore> Registered<S> {
298 pub(crate) fn new<H: EffectHandler>(handler: Handler<H>) -> Self {
299 let typed: Arc<dyn Any + Send + Sync> = Arc::new(handler.clone());
300 let resume: Resume<S> = Arc::new(move |runtime: Runtime<S>, record: EffectRecord| {
301 let handler = handler.clone();
302 Box::pin(async move { resume(&runtime, &handler, record).await })
303 });
304 Self { typed, resume }
305 }
306
307 pub(crate) fn typed<H: EffectHandler>(&self) -> Option<Handler<H>> {
308 self.typed.downcast_ref::<Handler<H>>().cloned()
309 }
310}
311
312#[must_use = "a submission does nothing until it is awaited"]
315pub struct Submission<'a, S, H: EffectHandler> {
316 runtime: &'a Runtime<S>,
317 key: String,
318 input: H::Input,
319 actor: Option<String>,
320}
321
322impl<'a, S, H: EffectHandler> Submission<'a, S, H> {
323 pub(crate) fn new(runtime: &'a Runtime<S>, key: impl Display, input: H::Input) -> Self {
324 Self {
325 runtime,
326 key: key.to_string(),
327 input,
328 actor: None,
329 }
330 }
331
332 pub fn actor(mut self, actor: impl Into<String>) -> Self {
334 self.actor = Some(actor.into());
335 self
336 }
337}
338
339impl<'a, S: EffectStore, H: EffectHandler> IntoFuture for Submission<'a, S, H> {
340 type Output = Result<EffectOutcome<H::Output>, RuntimeError>;
341 type IntoFuture = Pin<Box<dyn Future<Output = Self::Output> + Send + 'a>>;
342
343 fn into_future(self) -> Self::IntoFuture {
344 Box::pin(async move {
345 let handler = self
346 .runtime
347 .handler::<H>()
348 .ok_or(RuntimeError::NotRegistered { name: H::NAME })?;
349 let key = EffectKey::new(EffectName::new(H::NAME)?, LogicalKey::new(self.key)?);
350 let json = serde_json::to_value(&self.input).map_err(RuntimeError::Input)?;
351 let stored = Stored {
352 fingerprint: None,
354 json,
355 actor: self.actor,
356 from_record: false,
357 };
358 run(self.runtime, &handler, key, Arc::new(self.input), stored).await
359 })
360 }
361}
362
363async fn resume<S: EffectStore, H: EffectHandler>(
366 runtime: &Runtime<S>,
367 handler: &Handler<H>,
368 record: EffectRecord,
369) -> Result<EffectStatus, RuntimeError> {
370 let id = record.id;
371 if record.status == EffectStatus::Compensating {
372 let compensate = handler
373 .compensate
374 .clone()
375 .ok_or(RuntimeError::NotCompensable { name: H::NAME })?;
376 let spec = CompensationSpec {
377 key: record.key.clone(),
378 reason: None,
379 actor: Some(format!("recovery:{}", runtime.worker_id())),
380 retry: handler.effect.retry_policy(),
381 attempt_timeout: handler.effect.attempt_timeout(),
382 };
383 runtime.compensate_effect(spec, compensate).await?;
384 let settled = runtime
385 .store()
386 .get(id)
387 .await?
388 .ok_or(StoreError::NotFound(id))?;
389 return Ok(settled.status);
390 }
391 let json = record.input.clone().unwrap_or(Value::Null);
392 let input: H::Input = serde_json::from_value(json.clone())
393 .map_err(|source| RuntimeError::StoredInput { id, source })?;
394 let stored = Stored {
398 json,
399 fingerprint: record.input_fingerprint.clone(),
400 actor: record.created_by.clone(),
401 from_record: true,
402 };
403 run(
404 runtime,
405 handler,
406 record.key.clone(),
407 Arc::new(input),
408 stored,
409 )
410 .await?;
411 let settled = runtime
412 .store()
413 .get(id)
414 .await?
415 .ok_or(StoreError::NotFound(id))?;
416 Ok(settled.status)
417}
418
419struct Stored {
421 json: Value,
422 fingerprint: Option<String>,
423 actor: Option<String>,
424 from_record: bool,
426}
427
428async fn run<S: EffectStore, H: EffectHandler>(
430 runtime: &Runtime<S>,
431 handler: &Handler<H>,
432 key: EffectKey,
433 input: Arc<H::Input>,
434 stored: Stored,
435) -> Result<EffectOutcome<H::Output>, RuntimeError> {
436 let effect = Arc::clone(&handler.effect);
437 let precondition = {
438 let (effect, input) = (Arc::clone(&effect), Arc::clone(&input));
439 Arc::new(move |ctx: EffectContext| -> BoxFuture<Precondition> {
440 let (effect, input) = (Arc::clone(&effect), Arc::clone(&input));
441 Box::pin(async move { effect.precondition(&ctx, &input).await })
442 })
443 };
444 let spec = EffectSpec {
445 fingerprint: stored.fingerprint,
446 input: Some(stored.json),
447 key,
448 capabilities: Capabilities {
449 kind: effect.kind(),
450 remote_idempotency: effect.remote_idempotency(),
451 verification: handler
452 .verify
453 .as_ref()
454 .map_or(VerificationMode::None, |(mode, _)| *mode),
455 },
456 actor: stored.actor,
457 retry: effect
458 .retry_policy()
459 .unwrap_or_else(|| runtime.default_retry()),
460 attempt_timeout: effect.attempt_timeout(),
461 precondition: Some(precondition),
462 require_approval: effect.requires_approval(),
463 risk: effect.risk(),
464 automatic_retry: true,
465 input_stored: stored.from_record,
466 };
467 let action = {
468 let (effect, input) = (Arc::clone(&effect), Arc::clone(&input));
469 move |ctx: EffectContext| {
470 let (effect, input) = (Arc::clone(&effect), Arc::clone(&input));
471 async move { effect.execute(&ctx, &input).await.map_err(Into::into) }
472 }
473 };
474 match &handler.verify {
475 Some((_, verify)) => {
476 let verify = Arc::clone(verify);
477 let checker =
478 VerifyWith(move |ctx| verify(Arc::clone(&effect), ctx, Arc::clone(&input)));
479 runtime.execute(spec, action, checker).await
480 }
481 None => runtime.execute(spec, action, NoVerification).await,
482 }
483}
484
485#[must_use = "a compensation does nothing until it is awaited"]
488pub struct CompensationSubmission<'a, S, H: EffectHandler> {
489 runtime: &'a Runtime<S>,
490 key: String,
491 reason: Option<String>,
492 actor: Option<String>,
493 _handler: std::marker::PhantomData<fn() -> H>,
494}
495
496impl<'a, S, H: EffectHandler> CompensationSubmission<'a, S, H> {
497 pub(crate) fn new(runtime: &'a Runtime<S>, key: impl Display) -> Self {
498 Self {
499 runtime,
500 key: key.to_string(),
501 reason: None,
502 actor: None,
503 _handler: std::marker::PhantomData,
504 }
505 }
506
507 pub fn reason(mut self, reason: impl Into<String>) -> Self {
509 self.reason = Some(reason.into());
510 self
511 }
512
513 pub fn actor(mut self, actor: impl Into<String>) -> Self {
515 self.actor = Some(actor.into());
516 self
517 }
518}
519
520impl<'a, S: EffectStore, H: EffectHandler> IntoFuture for CompensationSubmission<'a, S, H> {
521 type Output = Result<CompensationOutcome, RuntimeError>;
522 type IntoFuture = Pin<Box<dyn Future<Output = Self::Output> + Send + 'a>>;
523
524 fn into_future(self) -> Self::IntoFuture {
525 Box::pin(async move {
526 let handler = self
527 .runtime
528 .handler::<H>()
529 .ok_or(RuntimeError::NotRegistered { name: H::NAME })?;
530 let compensate = handler
531 .compensate
532 .clone()
533 .ok_or(RuntimeError::NotCompensable { name: H::NAME })?;
534 let spec = CompensationSpec {
535 key: EffectKey::new(EffectName::new(H::NAME)?, LogicalKey::new(self.key)?),
536 reason: self.reason,
537 actor: self.actor,
538 retry: handler.effect.retry_policy(),
539 attempt_timeout: handler.effect.attempt_timeout(),
540 };
541 self.runtime.compensate_effect(spec, compensate).await
542 })
543 }
544}