Skip to main content

runifold_workflow/worker/
execution.rs

1use super::{
2    Arc, CheckpointErrorKind, Duration, Either, LeaseDuration, SystemWorkflowWorkerSleeper, Usage,
3    WorkerId, WorkflowCheckpoint, WorkflowDefinition, WorkflowDisposition, WorkflowError,
4    WorkflowExecution, WorkflowFailurePolicy, WorkflowLease, WorkflowRegistry, WorkflowStore,
5    WorkflowStoreError, WorkflowStoreErrorKind, WorkflowTask, WorkflowWake, WorkflowWorker,
6    WorkflowWorkerError, WorkflowWorkerOutcome, WorkflowWorkerSleeper, select,
7};
8
9impl WorkflowWorker {
10    /// Creates a worker with validated heartbeat timing.
11    ///
12    /// # Errors
13    ///
14    /// Rejects a zero heartbeat interval or one not shorter than the lease.
15    pub fn new(
16        store: Arc<dyn WorkflowStore>,
17        registry: WorkflowRegistry,
18        worker: WorkerId,
19        lease_duration: LeaseDuration,
20        heartbeat_interval: Duration,
21    ) -> Result<Self, WorkflowWorkerError> {
22        let heartbeat_ms = u64::try_from(heartbeat_interval.as_millis()).map_err(|_| {
23            WorkflowWorkerError::InvalidConfig(
24                "heartbeat interval exceeds the supported millisecond range".into(),
25            )
26        })?;
27        if heartbeat_ms == 0 || heartbeat_ms >= lease_duration.as_millis() {
28            return Err(WorkflowWorkerError::InvalidConfig(
29                "heartbeat interval must be positive and shorter than the lease".into(),
30            ));
31        }
32        Ok(Self {
33            store,
34            registry,
35            worker,
36            lease_duration,
37            heartbeat_interval,
38            missing_definition_retry: Duration::from_secs(5),
39            budget_denied_retry: Duration::from_secs(5),
40            sleeper: Arc::new(SystemWorkflowWorkerSleeper),
41        })
42    }
43
44    /// Overrides the retry delay for definitions absent from this worker.
45    #[must_use]
46    pub const fn with_missing_definition_retry(mut self, delay: Duration) -> Self {
47        self.missing_definition_retry = delay;
48        self
49    }
50
51    /// Overrides the retry delay when a tenant aggregate budget is exhausted.
52    #[must_use]
53    pub const fn with_budget_denied_retry(mut self, delay: Duration) -> Self {
54        self.budget_denied_retry = delay;
55        self
56    }
57
58    /// Overrides heartbeat sleeping for deterministic runtimes and tests.
59    #[must_use]
60    pub fn with_sleeper(mut self, sleeper: Arc<dyn WorkflowWorkerSleeper>) -> Self {
61        self.sleeper = sleeper;
62        self
63    }
64
65    /// Claims and processes at most one task.
66    ///
67    /// # Errors
68    ///
69    /// Returns infrastructure, checkpoint, or restored-budget failures.
70    pub async fn run_once(&self) -> Result<WorkflowWorkerOutcome, WorkflowWorkerError> {
71        let Some(claimed) = self
72            .store
73            .claim(self.worker.clone(), self.lease_duration)
74            .await?
75        else {
76            return Ok(WorkflowWorkerOutcome::Idle);
77        };
78        let id = claimed.task.checkpoint_id;
79        let Some(definition) = self.registry.get(&claimed.task).cloned() else {
80            return self
81                .finish_or_lose(
82                    claimed.lease,
83                    WorkflowDisposition::RetryAfter(self.missing_definition_retry),
84                    WorkflowWorkerOutcome::DefinitionUnavailable { checkpoint_id: id },
85                )
86                .await;
87        };
88        self.execute_claim(claimed.task, claimed.lease, claimed.wake, definition)
89            .await
90    }
91
92    async fn execute_claim(
93        &self,
94        task: WorkflowTask,
95        lease: WorkflowLease,
96        wake: Option<WorkflowWake>,
97        definition: Arc<WorkflowDefinition>,
98    ) -> Result<WorkflowWorkerOutcome, WorkflowWorkerError> {
99        let checkpoint = WorkflowCheckpoint::distributed(self.store.clone(), lease.clone());
100        let persisted = match checkpoint.load_async().await {
101            Ok((_, state)) => Some(state),
102            Err(error) if error.kind == CheckpointErrorKind::NotFound => None,
103            Err(error) => return Err(error.into()),
104        };
105        let usage = persisted
106            .as_ref()
107            .map_or_else(Usage::default, |state| state.usage);
108        let run = definition.run_context(usage)?;
109        match self
110            .store
111            .reserve_budget(lease.clone(), definition.budget, usage)
112            .await
113        {
114            Ok(_) => {}
115            Err(error) if error.kind == WorkflowStoreErrorKind::AdmissionDenied => {
116                return self
117                    .finish_or_lose(
118                        lease,
119                        WorkflowDisposition::RetryAfter(self.budget_denied_retry),
120                        WorkflowWorkerOutcome::Retried {
121                            checkpoint_id: task.checkpoint_id,
122                        },
123                    )
124                    .await;
125            }
126            Err(error) if error.kind == WorkflowStoreErrorKind::InvalidInput => {
127                return self
128                    .finish_or_lose(
129                        lease,
130                        WorkflowDisposition::Failed(
131                            "workflow definition is incompatible with tenant budget policy".into(),
132                        ),
133                        WorkflowWorkerOutcome::Failed {
134                            checkpoint_id: task.checkpoint_id,
135                        },
136                    )
137                    .await;
138            }
139            Err(error) => return Err(error.into()),
140        }
141        let execution = async {
142            if persisted.is_some() {
143                definition
144                    .workflow
145                    .resume_controlled(&checkpoint, &run, definition.resume_policy, wake)
146                    .await
147            } else {
148                definition
149                    .workflow
150                    .run_checkpointed_controlled(task.input, &run, &checkpoint)
151                    .await
152            }
153        };
154        let heartbeat = heartbeat_until_failure(
155            self.store.clone(),
156            lease.clone(),
157            self.lease_duration,
158            self.heartbeat_interval,
159            self.sleeper.clone(),
160        );
161        match select(Box::pin(execution), Box::pin(heartbeat)).await {
162            Either::Left((result, pending_heartbeat)) => {
163                drop(pending_heartbeat);
164                let usage = run.budget().usage();
165                self.finish_execution(lease, definition.failure_policy, result, usage)
166                    .await
167            }
168            Either::Right((_heartbeat_error, pending_execution)) => {
169                run.cancellation().cancel();
170                let _ = pending_execution.await;
171                Ok(WorkflowWorkerOutcome::LeaseLost {
172                    checkpoint_id: lease.checkpoint_id,
173                })
174            }
175        }
176    }
177
178    async fn finish_execution(
179        &self,
180        lease: WorkflowLease,
181        failure_policy: WorkflowFailurePolicy,
182        result: Result<WorkflowExecution, WorkflowError>,
183        usage: Usage,
184    ) -> Result<WorkflowWorkerOutcome, WorkflowWorkerError> {
185        let id = lease.checkpoint_id;
186        match self.store.settle_budget(lease.clone(), usage).await {
187            Ok(()) => {}
188            Err(error) if error.kind == WorkflowStoreErrorKind::LeaseLost => {
189                return Ok(WorkflowWorkerOutcome::LeaseLost { checkpoint_id: id });
190            }
191            Err(error) => return Err(error.into()),
192        }
193        match result {
194            Ok(WorkflowExecution::Completed(outcome)) => {
195                self.finish_or_lose(
196                    lease,
197                    WorkflowDisposition::Completed,
198                    WorkflowWorkerOutcome::Completed {
199                        checkpoint_id: id,
200                        outcome,
201                    },
202                )
203                .await
204            }
205            Ok(WorkflowExecution::Suspended(wait)) => {
206                self.finish_or_lose(
207                    lease,
208                    WorkflowDisposition::Suspend(wait),
209                    WorkflowWorkerOutcome::Suspended { checkpoint_id: id },
210                )
211                .await
212            }
213            Err(error) => match failure_policy {
214                WorkflowFailurePolicy::Fail => {
215                    self.finish_or_lose(
216                        lease,
217                        WorkflowDisposition::Failed(safe_failure_reason(&error)),
218                        WorkflowWorkerOutcome::Failed { checkpoint_id: id },
219                    )
220                    .await
221                }
222                WorkflowFailurePolicy::RetryAfter(delay) => {
223                    self.finish_or_lose(
224                        lease,
225                        WorkflowDisposition::RetryAfter(delay),
226                        WorkflowWorkerOutcome::Retried { checkpoint_id: id },
227                    )
228                    .await
229                }
230            },
231        }
232    }
233
234    async fn finish_or_lose(
235        &self,
236        lease: WorkflowLease,
237        disposition: WorkflowDisposition,
238        outcome: WorkflowWorkerOutcome,
239    ) -> Result<WorkflowWorkerOutcome, WorkflowWorkerError> {
240        let id = lease.checkpoint_id;
241        match self.store.finish(lease, disposition).await {
242            Ok(()) => Ok(outcome),
243            Err(error) if error.kind == WorkflowStoreErrorKind::LeaseLost => {
244                Ok(WorkflowWorkerOutcome::LeaseLost { checkpoint_id: id })
245            }
246            Err(error) => Err(error.into()),
247        }
248    }
249}
250
251impl std::fmt::Debug for WorkflowWorker {
252    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
253        formatter
254            .debug_struct("WorkflowWorker")
255            .field("registry", &self.registry)
256            .field("worker", &self.worker)
257            .field("lease_duration", &self.lease_duration)
258            .field("heartbeat_interval", &self.heartbeat_interval)
259            .finish_non_exhaustive()
260    }
261}
262
263async fn heartbeat_until_failure(
264    store: Arc<dyn WorkflowStore>,
265    mut lease: WorkflowLease,
266    extension: LeaseDuration,
267    interval: Duration,
268    sleeper: Arc<dyn WorkflowWorkerSleeper>,
269) -> WorkflowStoreError {
270    loop {
271        sleeper.sleep(interval).await;
272        match store.heartbeat(lease, extension).await {
273            Ok(renewed) => lease = renewed,
274            Err(error) => return error,
275        }
276    }
277}
278
279fn safe_failure_reason(error: &WorkflowError) -> String {
280    let message = error.to_string();
281    if message.len() <= 1_024 {
282        return message;
283    }
284    let mut boundary = 1_024;
285    while !message.is_char_boundary(boundary) {
286        boundary -= 1;
287    }
288    message[..boundary].to_owned()
289}