1use std::collections::hash_map::RandomState;
8use std::fmt::Display;
9use std::future::Future;
10use std::hash::BuildHasher;
11use std::sync::Arc;
12use std::sync::atomic::{AtomicU32, Ordering};
13use std::time::{Duration, SystemTime};
14
15use serde::Serialize;
16use serde::de::DeserializeOwned;
17use serde_json::{Value, json};
18use tracing::{Instrument, Span, debug, field, info_span, warn};
19
20use crate::approval::{ApprovalDecision, ApprovalProvider, ApprovalRequest, ErasedApproval};
21use crate::clock::{Clock, SystemClock};
22use crate::effect::{
23 EffectBuilder, EffectContext, EffectFailure, EffectOutcome, EffectSpec, Precondition,
24};
25use crate::error::RuntimeError;
26use crate::failure::{Disposition, FailureClass};
27#[cfg(feature = "fault-injection")]
28use crate::fault::FaultInjector;
29use crate::fault::FaultPoint;
30use crate::fingerprint::fingerprint;
31use crate::handler::{
32 CompensationSubmission, EffectHandler, Handler, Registered, Registry, Resume, Submission,
33};
34use crate::id::{EffectId, EffectName, WorkerId};
35use crate::kind::EffectKind;
36use crate::observer::{EffectObserver, Observation};
37use crate::policy::{RiskPolicy, UnknownPlan};
38use crate::redaction::{Field, Redactor};
39use crate::retention::RetentionPolicy;
40use crate::retry::RetryPolicy;
41use crate::state::{EffectStatus, Transition};
42use crate::store::{
43 EffectRecord, EffectStore, ErrorRecord, Lease, NewEffect, StoreError, TransitionRequest,
44};
45use crate::verification::{NotFoundReading, Verification, VerificationMode, Verifier};
46
47const MAX_ROUNDS: usize = 4;
50
51pub struct Runtime<S> {
53 inner: Arc<Inner<S>>,
54}
55
56struct Inner<S> {
57 store: S,
58 clock: Arc<dyn Clock>,
59 worker: WorkerId,
60 lease_ttl: Duration,
61 retry: RetryPolicy,
62 handlers: Registry<S>,
63 approval: Option<Arc<dyn ErasedApproval>>,
64 policy: RiskPolicy,
65 redactor: Option<Arc<dyn Redactor>>,
66 observers: Vec<Arc<dyn EffectObserver>>,
67 retention: RetentionPolicy,
68 #[cfg(feature = "fault-injection")]
69 faults: Option<Arc<FaultInjector>>,
70}
71
72impl<S> Clone for Runtime<S> {
73 fn clone(&self) -> Self {
74 Self {
75 inner: Arc::clone(&self.inner),
76 }
77 }
78}
79
80#[must_use]
82pub struct RuntimeBuilder<S> {
83 store: S,
84 clock: Arc<dyn Clock>,
85 worker: Option<WorkerId>,
86 lease_ttl: Duration,
87 retry: RetryPolicy,
88 handlers: Registry<S>,
89 approval: Option<Arc<dyn ErasedApproval>>,
90 policy: RiskPolicy,
91 redactor: Option<Arc<dyn Redactor>>,
92 observers: Vec<Arc<dyn EffectObserver>>,
93 retention: RetentionPolicy,
94 #[cfg(feature = "fault-injection")]
95 faults: Option<Arc<FaultInjector>>,
96}
97
98impl<S: EffectStore> RuntimeBuilder<S> {
99 pub fn retention(mut self, policy: RetentionPolicy) -> Self {
103 self.retention = policy;
104 self
105 }
106
107 pub fn observer(mut self, observer: impl EffectObserver) -> Self {
111 self.observers.push(Arc::new(observer));
112 self
113 }
114
115 pub fn redactor(mut self, redactor: impl Redactor) -> Self {
118 self.redactor = Some(Arc::new(redactor));
119 self
120 }
121
122 pub fn risk_policy(mut self, policy: RiskPolicy) -> Self {
126 self.policy = policy;
127 self
128 }
129
130 pub fn approval_provider(mut self, provider: impl ApprovalProvider) -> Self {
134 self.approval = Some(Arc::new(provider));
135 self
136 }
137
138 pub fn register<H: EffectHandler>(mut self, handler: Handler<H>) -> Self {
146 assert!(
147 !self.handlers.contains_key(H::NAME),
148 "a handler is already registered for effect `{}`",
149 H::NAME
150 );
151 self.handlers.insert(H::NAME, Registered::new(handler));
152 self
153 }
154
155 pub fn clock(mut self, clock: impl Clock) -> Self {
159 self.clock = Arc::new(clock);
160 self
161 }
162
163 pub fn worker_id(mut self, worker: WorkerId) -> Self {
166 self.worker = Some(worker);
167 self
168 }
169
170 pub fn lease_ttl(mut self, ttl: Duration) -> Self {
176 self.lease_ttl = ttl.max(Duration::from_millis(3));
177 self
178 }
179
180 pub fn retry_policy(mut self, policy: RetryPolicy) -> Self {
184 self.retry = policy;
185 self
186 }
187
188 #[cfg(feature = "fault-injection")]
191 pub fn fault_injector(mut self, injector: Arc<FaultInjector>) -> Self {
192 self.faults = Some(injector);
193 self
194 }
195
196 pub fn build(self) -> Runtime<S> {
198 Runtime {
199 inner: Arc::new(Inner {
200 store: self.store,
201 clock: self.clock,
202 worker: self.worker.unwrap_or_else(WorkerId::random),
203 lease_ttl: self.lease_ttl,
204 retry: self.retry,
205 handlers: self.handlers,
206 approval: self.approval,
207 policy: self.policy,
208 redactor: self.redactor,
209 observers: self.observers,
210 retention: self.retention,
211 #[cfg(feature = "fault-injection")]
212 faults: self.faults,
213 }),
214 }
215 }
216}
217
218pub(crate) enum Interrupt {
220 LeaseLost,
222 Error(RuntimeError),
223}
224
225impl From<StoreError> for Interrupt {
226 fn from(error: StoreError) -> Self {
227 match error {
228 StoreError::LeaseLost => Self::LeaseLost,
229 other => Self::Error(other.into()),
230 }
231 }
232}
233
234impl<S: EffectStore> Runtime<S> {
235 pub fn new(store: S) -> Self {
237 Self::builder(store).build()
238 }
239
240 pub fn builder(store: S) -> RuntimeBuilder<S> {
242 RuntimeBuilder {
243 store,
244 clock: Arc::new(SystemClock),
245 worker: None,
246 lease_ttl: Duration::from_secs(30),
247 retry: RetryPolicy::default(),
248 handlers: Registry::new(),
249 approval: None,
250 policy: RiskPolicy::default(),
251 redactor: None,
252 observers: Vec::new(),
253 retention: RetentionPolicy::KEEP_ALL,
254 #[cfg(feature = "fault-injection")]
255 faults: None,
256 }
257 }
258
259 pub fn effect(&self, name: impl Into<String>, key: impl Display) -> EffectBuilder<S> {
263 EffectBuilder::new(self.clone(), name.into(), key.to_string())
264 }
265
266 pub fn submit<H: EffectHandler>(
273 &self,
274 key: impl Display,
275 input: H::Input,
276 ) -> Submission<'_, S, H> {
277 Submission::new(self, key, input)
278 }
279
280 pub fn compensate<H: EffectHandler>(
284 &self,
285 key: impl Display,
286 ) -> CompensationSubmission<'_, S, H> {
287 CompensationSubmission::new(self, key)
288 }
289
290 pub(crate) fn handler<H: EffectHandler>(&self) -> Option<Handler<H>> {
291 self.inner.handlers.get(H::NAME)?.typed::<H>()
292 }
293
294 pub(crate) fn resumer(&self, name: &str) -> Option<Resume<S>> {
295 self.inner.handlers.get(name).map(|r| Arc::clone(&r.resume))
296 }
297
298 pub fn store(&self) -> &S {
300 &self.inner.store
301 }
302
303 pub fn worker_id(&self) -> &WorkerId {
305 &self.inner.worker
306 }
307
308 pub async fn wait<T: DeserializeOwned>(
323 &self,
324 id: EffectId,
325 timeout: Duration,
326 ) -> Result<EffectOutcome<T>, RuntimeError> {
327 let deadline = tokio::time::Instant::now() + timeout;
328 let mut pause = Duration::from_millis(10);
329 loop {
330 let record = self
331 .store()
332 .get(id)
333 .await?
334 .ok_or(StoreError::NotFound(id))?;
335 let busy = !settled(record.status) && record.live_lease_owner(self.now()).is_some();
336 let now = tokio::time::Instant::now();
337 if !busy {
338 return report(&record, None);
339 }
340 if now >= deadline {
341 return Ok(EffectOutcome::InProgress { id });
342 }
343 tokio::time::sleep(pause.min(deadline - now)).await;
344 pause = (pause * 2).min(Duration::from_millis(250));
345 }
346 }
347
348 pub(crate) fn default_retry(&self) -> RetryPolicy {
349 self.inner.retry
350 }
351
352 pub(crate) fn now(&self) -> SystemTime {
353 self.inner.clock.now()
354 }
355
356 pub(crate) fn retention(&self) -> RetentionPolicy {
357 self.inner.retention
358 }
359
360 pub(crate) fn lease_ttl(&self) -> Duration {
361 self.inner.lease_ttl
362 }
363
364 pub(crate) async fn with_lease<Fut: Future>(
367 &self,
368 lease: &Lease,
369 future: Fut,
370 ) -> Result<Fut::Output, Interrupt> {
371 let ttl = self.inner.lease_ttl;
372 tokio::pin!(future);
373 loop {
374 tokio::select! {
375 output = &mut future => return Ok(output),
376 () = tokio::time::sleep(ttl / 3) => {
377 match self.store().renew_lease(lease, self.now(), ttl).await {
378 Ok(_) => {}
379 Err(StoreError::LeaseLost) => return Err(Interrupt::LeaseLost),
380 Err(e) => warn!(error = %e, "lease renewal failed; will retry"),
381 }
382 }
383 }
384 }
385 }
386
387 pub(crate) async fn transition_leased(
389 &self,
390 record: &EffectRecord,
391 lease: &Lease,
392 actor: Option<&str>,
393 transition: Transition,
394 customize: impl FnOnce(&mut TransitionRequest),
395 ) -> Result<EffectRecord, Interrupt> {
396 let mut request = TransitionRequest::new(record, Some(lease), transition, self.now());
397 request.actor = actor.map(str::to_owned);
398 customize(&mut request);
399 let record = self.commit_transition(record, request).await?;
400 debug!(%transition, status = %record.status, "effect transition");
401 Ok(record)
402 }
403
404 pub(crate) async fn commit_transition(
407 &self,
408 before: &EffectRecord,
409 mut request: TransitionRequest,
410 ) -> Result<EffectRecord, StoreError> {
411 self.redact_request(&mut request, &before.key.name);
412 let transition = request.transition;
413 let after = self.store().transition(request).await?;
414 if !self.inner.observers.is_empty() {
415 let observation = Observation {
416 record: &after,
417 transition,
418 from: before.status,
419 to: after.status,
420 in_previous_status: after
421 .updated_at
422 .duration_since(before.updated_at)
423 .unwrap_or_default(),
424 since_created: after
425 .updated_at
426 .duration_since(after.created_at)
427 .unwrap_or_default(),
428 };
429 self.notify(|observer| observer.on_transition(&observation));
430 }
431 Ok(after)
432 }
433
434 fn notify(&self, call: impl Fn(&dyn EffectObserver)) {
436 for observer in &self.inner.observers {
437 let observer: &dyn EffectObserver = observer.as_ref();
438 if std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| call(observer))).is_err() {
439 warn!("an effect observer panicked; ignoring it");
440 }
441 }
442 }
443
444 pub(crate) fn redact(&self, field: Field, effect: &EffectName, value: &mut Value) {
446 if let Some(redactor) = &self.inner.redactor {
447 redactor.redact(field, effect, value);
448 }
449 }
450
451 pub(crate) fn redact_request(&self, request: &mut TransitionRequest, effect: &EffectName) {
453 if self.inner.redactor.is_none() {
454 return;
455 }
456 if let Some(output) = request.output.as_mut() {
457 self.redact(Field::Output, effect, output);
458 }
459 if let Some(payload) = request.payload.as_mut() {
460 self.redact(Field::AuditPayload, effect, payload);
461 }
462 if let Some(error) = request.error.as_mut() {
463 let mut message = Value::String(std::mem::take(&mut error.message));
464 self.redact(Field::ErrorMessage, effect, &mut message);
465 error.message = match message {
466 Value::String(text) => text,
467 other => other.to_string(),
468 };
469 }
470 }
471
472 #[cfg_attr(not(feature = "fault-injection"), allow(clippy::unused_self))]
475 pub(crate) fn checkpoint(&self, point: FaultPoint) {
476 #[cfg(feature = "fault-injection")]
477 if let Some(faults) = &self.inner.faults {
478 faults.reach(point);
479 }
480 #[cfg(not(feature = "fault-injection"))]
481 let _ = point;
482 }
483
484 pub(crate) async fn execute<T, F, Fut, V>(
485 &self,
486 mut spec: EffectSpec,
487 action: F,
488 verifier: V,
489 ) -> Result<EffectOutcome<T>, RuntimeError>
490 where
491 T: Serialize + DeserializeOwned + Send + 'static,
492 F: Fn(EffectContext) -> Fut + Send + Sync + 'static,
493 Fut: Future<Output = Result<T, EffectFailure>> + Send + 'static,
494 V: Verifier<T>,
495 {
496 let required = self
498 .inner
499 .policy
500 .requirements(spec.risk, spec.capabilities.kind);
501 if required.verification && spec.capabilities.verification == VerificationMode::None {
502 return Err(RuntimeError::PolicyViolation {
503 key: spec.key.to_string(),
504 requirement: "verification",
505 });
506 }
507 spec.require_approval |= required.approval;
508 spec.automatic_retry &= !required.no_automatic_retry;
509
510 if !spec.input_stored
514 && let Some(input) = spec.input.as_mut()
515 {
516 self.redact(Field::Input, &spec.key.name, input);
517 spec.fingerprint = Some(fingerprint(input));
518 }
519
520 let span = info_span!(
521 "agent_effect.execute",
522 effect.name = %spec.key.name,
523 effect.logical_key = %spec.key.key,
524 effect.kind = spec.capabilities.kind.as_str(),
525 effect.risk_level = %spec.risk,
526 effect.id = field::Empty,
527 effect.status = field::Empty,
528 effect.attempt = field::Empty,
529 );
530 if spec.capabilities.unknown_always_escalates() {
531 span.in_scope(|| {
532 warn!(
533 "effect is neither idempotent nor verifiable: \
534 any unknown outcome will need an operator"
535 );
536 });
537 }
538 let runtime = self.clone();
541 let task = tokio::spawn(
542 async move { runtime.drive(spec, action, verifier).await }.instrument(span),
543 );
544 task.await
545 .unwrap_or_else(|e| Err(RuntimeError::Internal(e.to_string())))
546 }
547
548 async fn drive<T, F, Fut, V>(
549 &self,
550 spec: EffectSpec,
551 action: F,
552 verifier: V,
553 ) -> Result<EffectOutcome<T>, RuntimeError>
554 where
555 T: Serialize + DeserializeOwned + Send + 'static,
556 F: Fn(EffectContext) -> Fut + Send + Sync + 'static,
557 Fut: Future<Output = Result<T, EffectFailure>> + Send + 'static,
558 V: Verifier<T>,
559 {
560 self.checkpoint(FaultPoint::BeforeInsert);
561 let store = self.store();
562 let inserted = store
563 .insert_or_get(NewEffect {
564 id: EffectId::new(),
565 key: spec.key.clone(),
566 kind: spec.capabilities.kind,
567 input: spec.input.clone(),
568 input_fingerprint: spec.fingerprint.clone(),
569 created_by: spec.actor.clone(),
570 now: self.now(),
571 })
572 .await?;
573 let mut record = inserted.record;
574 if inserted.inserted {
575 self.notify(|observer| observer.on_created(&record));
576 }
577 Span::current().record("effect.id", field::display(record.id));
578 self.checkpoint(FaultPoint::AfterInsert);
579 if !inserted.inserted {
580 check_matches(&record, &spec)?;
581 }
582
583 let checks_made = AtomicU32::new(0);
586 for _ in 0..MAX_ROUNDS {
587 if let Some(outcome) = observe(&record)? {
588 return Ok(outcome);
589 }
590 let lease = match store
591 .acquire_lease(
592 record.id,
593 self.worker_id(),
594 self.now(),
595 self.inner.lease_ttl,
596 )
597 .await
598 {
599 Ok(lease) => lease,
600 Err(StoreError::LeaseHeld { .. }) => {
601 return Ok(EffectOutcome::InProgress { id: record.id });
602 }
603 Err(e) => return Err(e.into()),
604 };
605 let current = store
607 .get(record.id)
608 .await?
609 .ok_or(StoreError::NotFound(record.id))?;
610 let driver = Driver {
611 rt: self,
612 spec: &spec,
613 action: &action,
614 verifier: &verifier,
615 lease: &lease,
616 checks_made: &checks_made,
617 };
618 let advanced = driver.advance(current).await;
619 if let Err(e) = store.release_lease(&lease).await {
620 warn!(error = %e, "could not release lease; it will expire");
621 }
622 match advanced {
623 Ok((settled, output)) => {
624 Span::current().record("effect.status", settled.status.as_str());
625 return report(&settled, output);
626 }
627 Err(Interrupt::LeaseLost) => {
628 warn!("lease lost mid-effect; re-reading the record");
629 record = store
630 .get(record.id)
631 .await?
632 .ok_or(StoreError::NotFound(record.id))?;
633 }
634 Err(Interrupt::Error(e)) => return Err(e),
635 }
636 }
637 Ok(EffectOutcome::InProgress { id: record.id })
638 }
639}
640
641fn observe<T: DeserializeOwned>(
648 record: &EffectRecord,
649) -> Result<Option<EffectOutcome<T>>, RuntimeError> {
650 match record.status {
651 status if settled(status) => report(record, None).map(Some),
652 EffectStatus::Pending
653 | EffectStatus::AwaitingApproval
654 | EffectStatus::Executing
655 | EffectStatus::Verifying
656 | EffectStatus::Unknown => Ok(None),
657 _ => Ok(Some(EffectOutcome::InProgress { id: record.id })),
658 }
659}
660
661struct Driver<'a, S, F, V> {
663 rt: &'a Runtime<S>,
664 spec: &'a EffectSpec,
665 action: &'a F,
666 verifier: &'a V,
667 lease: &'a Lease,
668 checks_made: &'a AtomicU32,
670}
671
672impl<S: EffectStore, F, V> Driver<'_, S, F, V> {
673 async fn advance<T, Fut>(
677 &self,
678 mut record: EffectRecord,
679 ) -> Result<(EffectRecord, Option<T>), Interrupt>
680 where
681 T: Serialize + Send + 'static,
682 F: Fn(EffectContext) -> Fut,
683 Fut: Future<Output = Result<T, EffectFailure>> + Send + 'static,
684 V: Verifier<T>,
685 {
686 let mut output = None;
687 let mut verification_exhausted = false;
688 loop {
689 match record.status {
690 EffectStatus::Pending => {
691 if let Some(at) = record.next_attempt_at {
692 self.sleep_until(at).await?;
693 }
694 if record.attempt_count == 0
695 && let Some(reason) = self.check_precondition(&record).await?
696 {
697 record = self
698 .transition(&record, Transition::PreconditionRejected, |r| {
699 r.error = Some(reason);
700 })
701 .await?;
702 continue;
703 }
704 if record.attempt_count == 0 && self.spec.require_approval && !record.approved {
705 record = self
706 .transition(&record, Transition::RequestApproval, |_| {})
707 .await?;
708 self.rt.checkpoint(FaultPoint::AfterApprovalRequested);
709 continue;
710 }
711 record = self
712 .transition(&record, Transition::StartAttempt, |r| {
713 r.payload = Some(json!({ "worker": self.rt.worker_id() }));
714 })
715 .await?;
716 Span::current().record("effect.attempt", record.attempt_count);
717 self.rt.checkpoint(FaultPoint::AfterAttemptPersisted);
718 let (next, produced, exhausted) = self.attempt(record).await?;
719 record = next;
720 if produced.is_some() {
721 output = produced;
722 }
723 verification_exhausted = exhausted;
724 }
725 EffectStatus::Executing | EffectStatus::Verifying => {
728 record = self
729 .transition(&record, Transition::LeaseExpired, |_| {})
730 .await?;
731 }
732 EffectStatus::AwaitingApproval => match self.ask_approval(&record).await? {
735 ApprovalDecision::Approved { by } => {
736 record = self
737 .transition(&record, Transition::Approve, |r| r.actor = Some(by))
738 .await?;
739 }
740 ApprovalDecision::Denied { by, reason } => {
741 record = self
742 .transition(&record, Transition::Deny, |r| {
743 r.actor = Some(by);
744 r.error = Some(ErrorRecord {
745 class: None,
746 message: reason,
747 });
748 })
749 .await?;
750 }
751 ApprovalDecision::Deferred => break,
752 },
753 EffectStatus::Unknown => match self.spec.capabilities.unknown_plan() {
754 UnknownPlan::Verify if !verification_exhausted => {
755 let (next, verified, exhausted) = self.verify(record, None).await?;
756 record = next;
757 output = verified.or(output);
758 verification_exhausted = exhausted;
759 }
760 UnknownPlan::Verify => break,
763 UnknownPlan::Reexecute if self.may_retry(record.attempt_count) => {
764 record = self
765 .schedule_retry(&record, FailureClass::Ambiguous, None)
766 .await?;
767 }
768 UnknownPlan::Reexecute | UnknownPlan::Escalate => {
769 record = self
770 .transition(&record, Transition::Escalate, |_| {})
771 .await?;
772 }
773 },
774 _ => break,
775 }
776 }
777 Ok((record, output))
778 }
779
780 async fn attempt<T, Fut>(
783 &self,
784 record: EffectRecord,
785 ) -> Result<(EffectRecord, Option<T>, bool), Interrupt>
786 where
787 T: Serialize + Send + 'static,
788 F: Fn(EffectContext) -> Fut,
789 Fut: Future<Output = Result<T, EffectFailure>> + Send + 'static,
790 V: Verifier<T>,
791 {
792 let mut task = tokio::spawn((self.action)(context(&record)));
796 self.rt.checkpoint(FaultPoint::AfterActionStarted);
797 let joined = match self.spec.attempt_timeout {
798 None => self.leased(&mut task).await?,
799 Some(limit) => match self.leased(tokio::time::timeout(limit, &mut task)).await? {
800 Ok(joined) => joined,
801 Err(_elapsed) => {
802 task.abort();
803 Ok(Err(EffectFailure::ambiguous(format!(
804 "attempt timed out after {limit:?}"
805 ))))
806 }
807 },
808 };
809 self.rt.checkpoint(FaultPoint::AfterActionReturned);
810
811 match joined {
812 Ok(Ok(value)) if self.spec.capabilities.verification != VerificationMode::None => {
813 let (record, verified, exhausted) =
814 self.verify(record, Some(output_json(&value))).await?;
815 let output = match record.status {
816 EffectStatus::Committed => verified.or(Some(value)),
817 _ => None,
818 };
819 Ok((record, output, exhausted))
820 }
821 Ok(Ok(value)) => {
822 let (output, payload) = output_json(&value);
823 let record = self
824 .transition(&record, Transition::Succeeded, |r| {
825 r.output = output;
826 r.payload = payload;
827 })
828 .await?;
829 Ok((record, Some(value), false))
830 }
831 Ok(Err(failure)) => {
832 debug!(%failure, "action failed");
833 let class = failure.class();
834 let error = failure.to_record();
835 let record = match class.disposition() {
836 Disposition::Retry if self.may_retry(record.attempt_count) => {
837 self.schedule_retry(&record, class, Some(error)).await?
838 }
839 Disposition::Retry | Disposition::Fail
843 if record.may_have_applied && record.kind != EffectKind::Read =>
844 {
845 let unknown = self
846 .transition(&record, Transition::OutcomeUnknown, |r| {
847 r.error = Some(error);
848 r.payload =
849 Some(json!({ "earlier_attempt_may_have_applied": true }));
850 })
851 .await?;
852 if self.spec.capabilities.unknown_plan() == UnknownPlan::Verify {
853 unknown
854 } else {
855 self.transition(&unknown, Transition::Escalate, |_| {})
856 .await?
857 }
858 }
859 Disposition::Retry | Disposition::Fail => {
860 self.transition(&record, Transition::FailedDefinitively, |r| {
861 r.error = Some(error);
862 })
863 .await?
864 }
865 Disposition::Unknown => {
866 self.transition(&record, Transition::OutcomeUnknown, |r| {
867 r.error = Some(error);
868 })
869 .await?
870 }
871 };
872 Ok((record, None, false))
873 }
874 Err(join_error) => {
875 let record = self
877 .transition(&record, Transition::OutcomeUnknown, |r| {
878 r.error = Some(ErrorRecord {
879 class: Some(FailureClass::Ambiguous),
880 message: format!("action did not complete: {join_error}"),
881 });
882 })
883 .await?;
884 Ok((record, None, false))
885 }
886 }
887 }
888
889 async fn verify<T>(
897 &self,
898 record: EffectRecord,
899 succeeded: Option<(Option<Value>, Option<Value>)>,
900 ) -> Result<(EffectRecord, Option<T>, bool), Interrupt>
901 where
902 T: Serialize + Send + 'static,
903 V: Verifier<T>,
904 {
905 let mut record = self
906 .transition(&record, Transition::StartVerification, |r| {
907 if let Some((output, payload)) = succeeded {
908 (r.output, r.payload) = (output, payload);
909 }
910 })
911 .await?;
912 self.rt.checkpoint(FaultPoint::AfterVerificationStarted);
913 let mode = self.spec.capabilities.verification;
914 let mut checks = 0;
915 let mut last_problem = String::from("no check completed");
916 let verifying_since = record.updated_at;
922 let ended_before = record.attempt_ended_at.map_or(Duration::ZERO, |ended| {
923 verifying_since.duration_since(ended).unwrap_or_default()
924 });
925 let started = tokio::time::Instant::now();
926
927 while self.take_check() {
928 let check = checks;
929 checks += 1;
930 let Some(future) = self.verifier.check(context(&record)) else {
931 break;
932 };
933 let found = match self.leased(tokio::spawn(future)).await? {
934 Ok(Ok(found)) => found,
935 Ok(Err(failure)) => {
936 last_problem = format!("check failed: {failure}");
937 Verification::Inconclusive
938 }
939 Err(join_error) => {
940 last_problem = format!("check did not complete: {join_error}");
941 Verification::Inconclusive
942 }
943 };
944 match found {
945 Verification::Confirmed(value) => {
946 let (output, payload) = output_json(&value);
947 record = self
948 .transition(&record, Transition::VerificationConfirmed, |r| {
949 r.output = output;
950 r.payload = payload;
951 })
952 .await?;
953 return Ok((record, Some(value), false));
954 }
955 Verification::Conflict { details } => {
956 record = self
957 .transition(&record, Transition::VerificationConflict, |r| {
958 r.error = Some(ErrorRecord {
959 class: None,
960 message: details,
961 });
962 })
963 .await?;
964 return Ok((record, None, false));
965 }
966 Verification::NotApplied => {
967 match mode.read_not_found(ended_before + started.elapsed()) {
968 NotFoundReading::NotApplied => {
969 return Ok((self.not_applied(&record).await?, None, false));
970 }
971 NotFoundReading::TooEarly { wait } => {
972 last_problem = "not visible yet within the settle delay".into();
973 if self.has_checks_left() {
974 self.sleep(wait).await?;
975 }
976 }
977 }
978 }
979 Verification::Inconclusive => {
980 if last_problem == "no check completed" {
981 last_problem = "remote system could not tell".into();
982 }
983 if self.has_checks_left() {
984 let delay =
985 self.retry()
986 .delay(check, FailureClass::Transient, jitter_sample());
987 self.sleep(delay).await?;
988 }
989 }
990 }
991 }
992
993 if checks == 0 {
994 last_problem = "this call's checks were already spent".into();
995 }
996 record = self
997 .transition(&record, Transition::OutcomeUnknown, |r| {
998 r.error = Some(ErrorRecord {
999 class: Some(FailureClass::Ambiguous),
1000 message: format!(
1001 "verification inconclusive after {checks} checks: {last_problem}"
1002 ),
1003 });
1004 })
1005 .await?;
1006 Ok((record, None, true))
1007 }
1008
1009 fn take_check(&self) -> bool {
1012 self.checks_made.fetch_add(1, Ordering::Relaxed) < self.retry().max_attempts.max(1)
1013 }
1014
1015 fn has_checks_left(&self) -> bool {
1016 self.checks_made.load(Ordering::Relaxed) < self.retry().max_attempts.max(1)
1017 }
1018
1019 async fn not_applied(&self, record: &EffectRecord) -> Result<EffectRecord, Interrupt> {
1022 let error = ErrorRecord {
1023 class: None,
1024 message: "verification found that the effect did not apply".into(),
1025 };
1026 if self.may_retry(record.attempt_count) {
1027 self.schedule_retry(record, FailureClass::Transient, Some(error))
1028 .await
1029 } else {
1030 self.transition(record, Transition::VerificationNotApplied, |r| {
1031 r.error = Some(error);
1032 })
1033 .await
1034 }
1035 }
1036
1037 async fn check_precondition(
1039 &self,
1040 record: &EffectRecord,
1041 ) -> Result<Option<ErrorRecord>, Interrupt> {
1042 let Some(precondition) = &self.spec.precondition else {
1043 return Ok(None);
1044 };
1045 let max_checks = self.retry().max_attempts.max(1);
1046 let mut last_reason = String::new();
1047 for check in 1..=max_checks {
1048 let rejection = |message: String| {
1049 Some(ErrorRecord {
1050 class: None,
1051 message,
1052 })
1053 };
1054 match self
1055 .leased(tokio::spawn(precondition(context(record))))
1056 .await?
1057 {
1058 Ok(Precondition::Satisfied) => return Ok(None),
1059 Ok(Precondition::Rejected { reason }) => return Ok(rejection(reason)),
1060 Ok(Precondition::RetryLater { after, reason }) => {
1061 last_reason = reason;
1062 if check < max_checks {
1063 self.sleep(after).await?;
1064 }
1065 }
1066 Err(join_error) => {
1068 return Ok(rejection(format!(
1069 "precondition check did not complete: {join_error}"
1070 )));
1071 }
1072 }
1073 }
1074 Ok(Some(ErrorRecord {
1075 class: None,
1076 message: format!("precondition not satisfied after {max_checks} checks: {last_reason}"),
1077 }))
1078 }
1079
1080 async fn schedule_retry(
1084 &self,
1085 record: &EffectRecord,
1086 class: FailureClass,
1087 error: Option<ErrorRecord>,
1088 ) -> Result<EffectRecord, Interrupt> {
1089 let retry = record.attempt_count.saturating_sub(1);
1090 let delay = self.retry().delay(retry, class, jitter_sample());
1091 let at = self.rt.now() + delay;
1092 debug!(?delay, ?class, "retry scheduled");
1093 self.transition(record, Transition::ScheduleRetry, |r| {
1094 r.next_attempt_at = Some(at);
1095 r.error = error;
1096 r.payload =
1097 Some(json!({ "delay_ms": u64::try_from(delay.as_millis()).unwrap_or(u64::MAX) }));
1098 })
1099 .await
1100 }
1101
1102 async fn ask_approval(&self, record: &EffectRecord) -> Result<ApprovalDecision, Interrupt> {
1105 let Some(provider) = self.rt.inner.approval.clone() else {
1106 return Ok(ApprovalDecision::Deferred);
1107 };
1108 let request = ApprovalRequest {
1109 effect_id: record.id,
1110 key: record.key.clone(),
1111 kind: record.kind,
1112 risk: self.spec.risk,
1113 input: record.input.clone(),
1114 requested_by: record.created_by.clone(),
1115 };
1116 let task = tokio::spawn(async move { provider.request_boxed(request).await });
1117 Ok(self.leased(task).await?.unwrap_or_else(|e| {
1118 warn!(error = %e, "approval provider failed; deferring");
1119 ApprovalDecision::Deferred
1120 }))
1121 }
1122
1123 async fn sleep_until(&self, at: SystemTime) -> Result<(), Interrupt> {
1124 let wait = at.duration_since(self.rt.now()).unwrap_or_default();
1125 self.sleep(wait).await
1126 }
1127
1128 async fn sleep(&self, duration: Duration) -> Result<(), Interrupt> {
1129 if duration.is_zero() {
1130 return Ok(());
1131 }
1132 self.leased(tokio::time::sleep(duration)).await
1133 }
1134
1135 async fn leased<Fut: Future>(&self, future: Fut) -> Result<Fut::Output, Interrupt> {
1136 self.rt.with_lease(self.lease, future).await
1137 }
1138
1139 async fn transition(
1140 &self,
1141 record: &EffectRecord,
1142 transition: Transition,
1143 customize: impl FnOnce(&mut TransitionRequest),
1144 ) -> Result<EffectRecord, Interrupt> {
1145 self.rt
1146 .transition_leased(
1147 record,
1148 self.lease,
1149 self.spec.actor.as_deref(),
1150 transition,
1151 customize,
1152 )
1153 .await
1154 }
1155
1156 fn retry(&self) -> &RetryPolicy {
1157 &self.spec.retry
1158 }
1159
1160 fn may_retry(&self, attempts: u32) -> bool {
1163 self.spec.automatic_retry && self.retry().allows_another(attempts)
1164 }
1165}
1166
1167fn context(record: &EffectRecord) -> EffectContext {
1168 EffectContext {
1169 id: record.id,
1170 key: record.key.clone(),
1171 attempt: record.attempt_count,
1172 }
1173}
1174
1175fn output_json<T: Serialize>(value: &T) -> (Option<Value>, Option<Value>) {
1179 match serde_json::to_value(value) {
1180 Ok(output) => (Some(output), None),
1181 Err(e) => (None, Some(json!({ "output_not_stored": e.to_string() }))),
1182 }
1183}
1184
1185pub(crate) fn jitter_sample() -> f64 {
1188 let bits = RandomState::new().hash_one(0_u8) >> 11;
1189 #[allow(clippy::cast_precision_loss)]
1190 let sample = bits as f64 / (1_u64 << 53) as f64;
1191 sample
1192}
1193
1194fn settled(status: EffectStatus) -> bool {
1196 matches!(
1197 status,
1198 EffectStatus::Committed
1199 | EffectStatus::Failed
1200 | EffectStatus::Rejected
1201 | EffectStatus::NeedsIntervention
1202 | EffectStatus::Compensated
1203 | EffectStatus::CompensationFailed
1204 )
1205}
1206
1207fn check_matches(record: &EffectRecord, spec: &EffectSpec) -> Result<(), RuntimeError> {
1208 if record.kind != spec.capabilities.kind {
1209 return Err(RuntimeError::KindMismatch {
1210 id: record.id,
1211 stored: record.kind,
1212 requested: spec.capabilities.kind,
1213 });
1214 }
1215 if record.input_fingerprint != spec.fingerprint {
1216 return Err(RuntimeError::InputMismatch { id: record.id });
1217 }
1218 Ok(())
1219}
1220
1221fn report<T: DeserializeOwned>(
1224 record: &EffectRecord,
1225 fresh: Option<T>,
1226) -> Result<EffectOutcome<T>, RuntimeError> {
1227 let id = record.id;
1228 Ok(match record.status {
1229 EffectStatus::Committed => match fresh {
1230 Some(value) => EffectOutcome::Committed(value),
1231 None => EffectOutcome::Committed(
1232 serde_json::from_value(record.output.clone().unwrap_or(Value::Null))
1233 .map_err(|source| RuntimeError::Output { id, source })?,
1234 ),
1235 },
1236 EffectStatus::Failed => EffectOutcome::Failed(last_error(record)),
1237 EffectStatus::Rejected => EffectOutcome::Rejected(last_error(record)),
1238 EffectStatus::NeedsIntervention | EffectStatus::CompensationFailed => {
1240 EffectOutcome::NeedsIntervention { id }
1241 }
1242 EffectStatus::Unknown => EffectOutcome::Unknown { id },
1243 EffectStatus::Compensated => EffectOutcome::Compensated { id },
1244 EffectStatus::AwaitingApproval => EffectOutcome::AwaitingApproval { id },
1245 _ => EffectOutcome::InProgress { id },
1246 })
1247}
1248
1249pub(crate) fn last_error(record: &EffectRecord) -> ErrorRecord {
1250 record.last_error.clone().unwrap_or_else(|| ErrorRecord {
1251 class: None,
1252 message: "no error was recorded".into(),
1253 })
1254}