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