runifold_workflow/worker/
execution.rs1use 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 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 #[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 #[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 #[must_use]
60 pub fn with_sleeper(mut self, sleeper: Arc<dyn WorkflowWorkerSleeper>) -> Self {
61 self.sleeper = sleeper;
62 self
63 }
64
65 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}