1use crate::worker::{
2 FatalError, RunFinished, Worker, WorkerContext, WorkerError, WorkerResult, WorkerResultOk,
3};
4use assert_matches::assert_matches;
5use chrono::{DateTime, Utc};
6use concepts::prefixed_ulid::{DeploymentId, RunId};
7use concepts::storage::{
8 AppendEventsToExecution, AppendRequest, AppendResponseToExecution, DbErrorGeneric,
9 DbErrorWrite, DbErrorWriteNonRetriable, DbExecutor, DbPool, ExecutionLog, LockedExecution,
10 Unlocked,
11};
12use concepts::time::{ClockFn, Sleep};
13use concepts::{
14 ComponentId, ComponentRetryConfig, ComponentType, FunctionMetadata, StrVariant,
15 SupportedFunctionReturnValue,
16};
17use concepts::{ExecutionFailureKind, JoinSetId};
18use concepts::{ExecutionId, FunctionFqn, prefixed_ulid::ExecutorId};
19use concepts::{
20 FinishedExecutionFailure,
21 storage::{ExecutionRequest, Version},
22};
23use std::{
24 sync::{
25 Arc,
26 atomic::{AtomicBool, Ordering},
27 },
28 time::Duration,
29};
30use tokio::task::JoinHandle;
31use tracing::{Instrument, Level, Span, debug, error, info, info_span, instrument, trace, warn};
32
33#[derive(Debug, Clone)]
34pub struct ExecConfig {
35 pub lock_expiry: Duration,
36 pub tick_sleep: Duration,
37 pub batch_size: u32,
38 pub component_id: ComponentId,
39 pub task_limiter_global: Option<Arc<tokio::sync::Semaphore>>,
40 pub task_limiter_local: Option<Arc<tokio::sync::Semaphore>>,
41 pub executor_id: ExecutorId,
42 pub retry_config: ComponentRetryConfig,
43 pub locking_strategy: LockingStrategy,
44}
45
46pub struct ExecTask {
47 worker: Arc<dyn Worker>,
48 pub config: ExecConfig,
49 clock_fn: Box<dyn ClockFn>, db_pool: Arc<dyn DbPool>,
51 locking_strategy_holder: LockingStrategyHolder,
52 worker_count_tx: tokio::sync::watch::Sender<usize>,
53 executor_close_watcher: tokio::sync::watch::Receiver<bool>,
54}
55
56#[derive(derive_more::Debug, Default)]
57pub struct ExecutionProgress {
58 #[debug(skip)]
59 #[allow(dead_code)]
60 executions: Vec<(ExecutionId, JoinHandle<()>)>,
61}
62
63impl ExecutionProgress {
64 #[cfg(feature = "test")]
65 pub async fn wait_for_tasks(self) -> Vec<ExecutionId> {
66 let mut vec = Vec::new();
67 for (exe, join_handle) in self.executions {
68 vec.push(exe);
69 join_handle.await.unwrap();
70 }
71 vec
72 }
73}
74
75#[derive(Clone, Copy, Debug, PartialEq, Eq)]
76pub enum WorkerType {
77 Activity,
78 Workflow,
79}
80
81#[derive(derive_more::Debug)]
82pub struct ExecutorTaskHandle {
83 #[debug(skip)]
85 is_closing: Arc<AtomicBool>,
86 #[debug(skip)]
87 join_handle: JoinHandle<()>,
88 component_id: ComponentId,
89 executor_id: ExecutorId,
90 deployment_id: DeploymentId,
91 executor_closing_signal_sender: tokio::sync::watch::Sender<bool>,
93 worker_count_rx: tokio::sync::watch::Receiver<usize>,
95}
96
97pub struct WorkerTasksHandle {
98 component_id: ComponentId,
99 executor_id: ExecutorId,
100 deployment_id: DeploymentId,
101 executor_closing_signal_sender: tokio::sync::watch::Sender<bool>,
103 worker_count_rx: tokio::sync::watch::Receiver<usize>,
105}
106
107impl ExecutorTaskHandle {
108 #[instrument(name = "executor.close", skip_all, fields(executor_id = %self.executor_id, component_id = %self.component_id,
110 deployment_id = %self.deployment_id))]
111 pub async fn close_outer_task(mut self) -> WorkerTasksHandle {
112 debug!("Gracefully closing executor task");
113 self.is_closing.store(true, Ordering::Relaxed);
114
115 loop {
116 tokio::select! {
117 () = tokio::time::sleep(Duration::from_secs(5)) => {
118 info!("Waiting for executor to shut down");
119 }
120 _ = &mut self.join_handle => {
121 break;
122 }
123 }
124 }
125 let worker_tasks = *self.worker_count_rx.borrow();
126 debug!("Gracefully closed executor task, {worker_tasks} workers are in progress");
127 WorkerTasksHandle {
128 component_id: self.component_id,
129 executor_id: self.executor_id,
130 deployment_id: self.deployment_id,
131 executor_closing_signal_sender: self.executor_closing_signal_sender,
132 worker_count_rx: self.worker_count_rx,
133 }
134 }
135
136 #[must_use]
137 pub fn component_id(&self) -> &ComponentId {
138 &self.component_id
139 }
140}
141
142impl WorkerTasksHandle {
143 #[instrument(name = "WorkerTasksHandle.close", skip_all, fields(executor_id = %self.executor_id, component_id = %self.component_id,
145 deployment_id = %self.deployment_id))]
146 pub async fn close(self) {
147 debug!("Signaling worker tasks to unlock");
148 let _ = self.executor_closing_signal_sender.send(true);
149 let mut worker_count_rx = self.worker_count_rx.clone();
150 loop {
151 tokio::select! {
152 () = tokio::time::sleep(Duration::from_secs(1)) => {
153 info!("Waiting for {} workers to shut down", *self.worker_count_rx.borrow());
154 }
155 _ = worker_count_rx.wait_for(|&count| count == 0) => {
156 break;
157 }
158 }
159 }
160 debug!("All worker tasks are finished");
161 }
162}
163
164#[cfg(feature = "test")]
165pub fn extract_exported_ffqns_noext_test(worker: &dyn Worker) -> Arc<[FunctionFqn]> {
166 extract_exported_ffqns_noext(worker)
167}
168
169fn extract_exported_ffqns_noext(worker: &dyn Worker) -> Arc<[FunctionFqn]> {
170 worker
171 .exported_functions_noext()
172 .iter()
173 .map(|FunctionMetadata { ffqn, .. }| ffqn.clone())
174 .collect::<Arc<_>>()
175}
176
177#[derive(Debug, Clone, Copy, PartialEq, Eq)]
178pub enum LockingStrategy {
179 ByFfqns,
180 ByComponentDigest,
181 Auto,
182}
183impl LockingStrategy {
184 fn holder(&self, ffqns: Arc<[FunctionFqn]>) -> LockingStrategyHolder {
185 match self {
186 LockingStrategy::ByFfqns => LockingStrategyHolder::ByFfqns(ffqns),
187 LockingStrategy::ByComponentDigest => LockingStrategyHolder::ByComponentDigest,
188 LockingStrategy::Auto => LockingStrategyHolder::Auto { ffqns },
189 }
190 }
191}
192
193enum LockingStrategyHolder {
194 ByFfqns(Arc<[FunctionFqn]>),
195 ByComponentDigest,
196 Auto { ffqns: Arc<[FunctionFqn]> },
197}
198
199#[derive(Default)]
200#[expect(dead_code)] struct TaskLimiterPermit {
202 global: Option<tokio::sync::OwnedSemaphorePermit>,
203 local: Option<tokio::sync::OwnedSemaphorePermit>,
204}
205
206impl ExecTask {
207 #[cfg(feature = "test")]
208 pub fn new_test(
209 config: ExecConfig,
210 worker: Arc<dyn Worker>,
211 clock_fn: Box<dyn ClockFn>,
212 db_pool: Arc<dyn DbPool>,
213 ffqns: Arc<[FunctionFqn]>,
214 ) -> (Self, tokio::sync::watch::Sender<bool>) {
215 let (worker_count_tx, _) = tokio::sync::watch::channel(0usize);
216 let (close_tx, executor_close_watcher) = tokio::sync::watch::channel(false);
217 (
218 ExecTask {
219 worker,
220 locking_strategy_holder: config.locking_strategy.holder(ffqns),
221 config,
222 clock_fn,
223 db_pool,
224 worker_count_tx,
225 executor_close_watcher,
226 },
227 close_tx,
228 )
229 }
230
231 #[cfg(feature = "test")]
232 pub fn new_all_ffqns_test(
233 worker: Arc<dyn Worker>,
234 config: ExecConfig,
235 clock_fn: Box<dyn ClockFn>,
236 db_pool: Arc<dyn DbPool>,
237 ) -> (Self, tokio::sync::watch::Sender<bool>) {
238 let ffqns = extract_exported_ffqns_noext(worker.as_ref());
239 let (worker_count_tx, _) = tokio::sync::watch::channel(0usize);
240 let (close_tx, executor_close_watcher) = tokio::sync::watch::channel(false);
241 (
242 Self {
243 worker,
244 locking_strategy_holder: config.locking_strategy.holder(ffqns),
245 config,
246 clock_fn,
247 db_pool,
248 worker_count_tx,
249 executor_close_watcher,
250 },
251 close_tx,
252 )
253 }
254
255 #[cfg(feature = "test")]
256 pub fn new_all_ffqns_test_with_close_watcher(
257 worker: Arc<dyn Worker>,
258 config: ExecConfig,
259 clock_fn: Box<dyn ClockFn>,
260 db_pool: Arc<dyn DbPool>,
261 executor_close_watcher: tokio::sync::watch::Receiver<bool>,
262 ) -> Self {
263 let ffqns = extract_exported_ffqns_noext(worker.as_ref());
264 let (worker_count_tx, _) = tokio::sync::watch::channel(0usize);
265 Self {
266 worker,
267 locking_strategy_holder: config.locking_strategy.holder(ffqns),
268 config,
269 clock_fn,
270 db_pool,
271 worker_count_tx,
272 executor_close_watcher,
273 }
274 }
275
276 pub fn spawn_new(
278 deployment_id: DeploymentId,
279 worker: Arc<dyn Worker>,
280 config: ExecConfig,
281 clock_fn: Box<dyn ClockFn>,
282 db_pool: Arc<dyn DbPool>,
283 sleep: impl Sleep + Clone + 'static,
284 ) -> ExecutorTaskHandle {
285 let is_closing = Arc::new(AtomicBool::default());
286 let is_closing_inner = is_closing.clone();
287 let ffqns = extract_exported_ffqns_noext(worker.as_ref());
288 let component_id = config.component_id.clone();
289 let executor_id = config.executor_id;
290 let (worker_count_tx, worker_count_rx) = tokio::sync::watch::channel(0);
291 let (executor_closing_signal_sender, executor_close_watcher) =
292 tokio::sync::watch::channel(false);
293 let join_handle = tokio::spawn(async move {
294 debug!(executor_id = %config.executor_id, component_id = %config.component_id, "Spawned executor");
295 let lock_strategy_holder = config.locking_strategy.holder(ffqns);
296 let task = ExecTask {
297 worker,
298 config,
299 db_pool,
300 locking_strategy_holder: lock_strategy_holder,
301 clock_fn: clock_fn.clone_box(),
302 worker_count_tx,
303 executor_close_watcher,
304 };
305 let mut old_err = None;
306 while !is_closing_inner.load(Ordering::Relaxed) {
307 let res = task.db_pool.db_exec_conn().await;
308 let res = log_err_if_new(res, &mut old_err);
309 if let Ok(db_exec) = res {
310 let _ = task
311 .tick(
312 db_exec.as_ref(),
313 clock_fn.now(),
314 RunId::generate(),
315 deployment_id,
316 )
317 .await;
318 let timeout_fut = {
319 let sleep = sleep.clone();
320 Box::pin(async move { sleep.sleep(task.config.tick_sleep).await })
321 };
322 match &task.locking_strategy_holder {
323 LockingStrategyHolder::ByFfqns(ffqns) => {
324 db_exec
325 .wait_for_pending_by_ffqn(
326 clock_fn.now(),
327 ffqns.clone(),
328 None,
329 timeout_fut,
330 )
331 .await;
332 }
333 LockingStrategyHolder::ByComponentDigest => {
334 db_exec
335 .wait_for_pending_by_component_digest(
336 clock_fn.now(),
337 &task.config.component_id.component_digest,
338 timeout_fut,
339 )
340 .await;
341 }
342 LockingStrategyHolder::Auto { ffqns } => {
343 db_exec
344 .wait_for_pending_by_ffqn(
345 clock_fn.now(),
346 ffqns.clone(),
347 Some(task.config.component_id.component_digest.clone()),
348 timeout_fut,
349 )
350 .await;
351 }
352 }
353 } else {
354 sleep.sleep(task.config.tick_sleep).await;
355 }
356 }
357 });
358 ExecutorTaskHandle {
359 is_closing,
360 join_handle,
361 component_id,
362 executor_id,
363 deployment_id,
364 executor_closing_signal_sender,
365 worker_count_rx,
366 }
367 }
368
369 fn acquire_task_permits(&self) -> Vec<TaskLimiterPermit> {
370 let mut locks = Vec::with_capacity(
371 usize::try_from(self.config.batch_size).expect("16 bit systems are unsupported"),
372 );
373 for _ in 0..self.config.batch_size {
374 match (
375 &self.config.task_limiter_global,
376 &self.config.task_limiter_local,
377 ) {
378 (Some(global), Some(local)) => {
379 if let Ok(global) = global.clone().try_acquire_owned()
380 && let Ok(local) = local.clone().try_acquire_owned()
381 {
382 locks.push(TaskLimiterPermit {
383 global: Some(global),
384 local: Some(local),
385 });
386 } else {
387 break;
388 }
389 }
390 (Some(global), None) => {
391 if let Ok(global) = global.clone().try_acquire_owned() {
392 locks.push(TaskLimiterPermit {
393 global: Some(global),
394 local: None,
395 });
396 } else {
397 break;
398 }
399 }
400 (None, Some(local)) => {
401 if let Ok(local) = local.clone().try_acquire_owned() {
402 locks.push(TaskLimiterPermit {
403 global: None,
404 local: Some(local),
405 });
406 } else {
407 break;
408 }
409 }
410 (None, None) => {
411 locks.push(TaskLimiterPermit::default());
412 }
413 }
414 }
415 locks
416 }
417
418 #[cfg(feature = "test")]
419 pub async fn tick_test(&self, executed_at: DateTime<Utc>, run_id: RunId) -> ExecutionProgress {
420 use concepts::prefixed_ulid::DEPLOYMENT_ID_DUMMY;
421
422 let db_exec = self.db_pool.db_exec_conn().await.unwrap();
423 self.tick(db_exec.as_ref(), executed_at, run_id, DEPLOYMENT_ID_DUMMY)
424 .await
425 .unwrap()
426 }
427
428 #[cfg(feature = "test")]
429 pub async fn tick_test_await(
430 &self,
431 executed_at: DateTime<Utc>,
432 run_id: RunId,
433 ) -> Vec<ExecutionId> {
434 use concepts::prefixed_ulid::DEPLOYMENT_ID_DUMMY;
435
436 let db_exec = self.db_pool.db_exec_conn().await.unwrap();
437 self.tick(db_exec.as_ref(), executed_at, run_id, DEPLOYMENT_ID_DUMMY)
438 .await
439 .unwrap()
440 .wait_for_tasks()
441 .await
442 }
443
444 #[instrument(level = Level::TRACE, name = "executor.tick" skip_all, fields(executor_id = %self.config.executor_id, component_id = %self.config.component_id))]
445 async fn tick(
446 &self,
447 db_exec: &dyn DbExecutor,
448 executed_at: DateTime<Utc>,
449 run_id: RunId,
450 deployment_id: DeploymentId,
451 ) -> Result<ExecutionProgress, DbErrorWrite> {
452 let locked_executions = {
453 let mut permits = self.acquire_task_permits();
454 if permits.is_empty() {
455 return Ok(ExecutionProgress::default());
456 }
457 let lock_expires_at = executed_at + self.config.lock_expiry;
458 let batch_size = u32::try_from(permits.len()).expect("ExecConfig.batch_size is u32");
459 let locked_executions = match &self.locking_strategy_holder {
460 LockingStrategyHolder::Auto { ffqns } => {
461 db_exec
462 .lock_pending_by_ffqns_auto(
463 batch_size,
464 executed_at, ffqns.clone(),
466 executed_at, self.config.component_id.clone(),
468 deployment_id,
469 self.config.executor_id,
470 lock_expires_at,
471 run_id,
472 self.config.retry_config,
473 )
474 .await?
475 }
476 LockingStrategyHolder::ByFfqns(ffqns) => {
477 db_exec
478 .lock_pending_by_ffqns(
479 batch_size,
480 executed_at, ffqns.clone(),
482 executed_at, self.config.component_id.clone(),
484 deployment_id,
485 self.config.executor_id,
486 lock_expires_at,
487 run_id,
488 self.config.retry_config,
489 )
490 .await?
491 }
492 LockingStrategyHolder::ByComponentDigest => {
493 db_exec
494 .lock_pending_by_component_digest(
495 batch_size,
496 executed_at, &self.config.component_id,
498 deployment_id,
499 executed_at, self.config.executor_id,
501 lock_expires_at,
502 run_id,
503 self.config.retry_config,
504 )
505 .await?
506 }
507 };
508 while permits.len() > locked_executions.len() {
510 permits.pop();
511 }
512 assert_eq!(permits.len(), locked_executions.len());
513 locked_executions.into_iter().zip(permits)
514 };
515
516 let mut executions = Vec::with_capacity(locked_executions.len());
517 for (locked_execution, permit) in locked_executions {
518 let execution_id = locked_execution.execution_id.clone();
519 let join_handle = {
520 let worker = self.worker.clone();
521 let db_pool = self.db_pool.clone();
522 let clock_fn = self.clock_fn.clone_box();
523 let worker_span = info_span!(parent: None, "worker",
524 "otel.name" = format!("worker {}", locked_execution.ffqn),
525 %execution_id, %run_id,
526 ffqn = %locked_execution.ffqn,
527 executor_id = %self.config.executor_id,
528 component_id = %self.config.component_id,
529 %deployment_id,
530 );
531 locked_execution.metadata.enrich(&worker_span);
532 let component_type = self.config.component_id.component_type;
533 let worker_count_tx = self.worker_count_tx.clone();
534 worker_count_tx.send_modify(|n| *n += 1);
535 let executor_close_watcher = self.executor_close_watcher.clone();
536 tokio::spawn({
537 let worker_span2 = worker_span.clone();
538 let retry_config = self.config.retry_config;
539 async move {
540 let _permit = permit;
541 let res = Self::run_worker(
542 component_type,
543 worker,
544 db_pool,
545 clock_fn,
546 locked_execution,
547 retry_config,
548 worker_span2,
549 executor_close_watcher
550 )
551 .await;
552 if let Err(db_error) = res {
553 error!("Got db error `{db_error:?}`, expecting watcher to mark execution as timed out");
554 }
555 worker_count_tx.send_modify(|n| *n -= 1);
556 }
557 .instrument(worker_span)
558 })
559 };
560 executions.push((execution_id, join_handle));
561 }
562 Ok(ExecutionProgress { executions })
563 }
564
565 #[expect(clippy::too_many_arguments)]
566 async fn run_worker(
567 component_type: ComponentType,
568 worker: Arc<dyn Worker>,
569 db_pool: Arc<dyn DbPool>,
570 clock_fn: Box<dyn ClockFn>,
571 locked_execution: LockedExecution,
572 retry_config: ComponentRetryConfig,
573 worker_span: Span,
574 executor_close_watcher: tokio::sync::watch::Receiver<bool>,
575 ) -> Result<(), DbErrorWrite> {
576 debug!("Worker::run starting");
577 trace!(
578 version = %locked_execution.next_version,
579 params = ?locked_execution.params,
580 event_history = ?locked_execution.event_history,
581 "Worker::run starting"
582 );
583 let can_be_retried = ExecutionLog::can_be_retried_after(
584 locked_execution.intermittent_event_count + 1,
585 retry_config.max_retries,
586 retry_config.retry_exp_backoff,
587 );
588 let unlock_expiry_on_limit_reached =
589 ExecutionLog::compute_retry_duration_when_retrying_forever(
590 locked_execution.intermittent_event_count + 1,
591 retry_config.retry_exp_backoff,
592 );
593 let parent = locked_execution.parent.clone();
594 let execution_id = locked_execution.execution_id.clone();
595 let ctx = WorkerContext {
596 execution_id: locked_execution.execution_id.clone(),
597 metadata: locked_execution.metadata,
598 component_digest: locked_execution.component_digest,
599 ffqn: locked_execution.ffqn,
600 params: locked_execution.params,
601 event_history: locked_execution.event_history,
602 responses: locked_execution.responses,
603 parent: parent.clone(),
604 version: locked_execution.next_version,
605 can_be_retried: can_be_retried.is_some(),
606 locked_event: locked_execution.locked_event,
607 worker_span,
608 executor_close_watcher,
609 };
610 let worker_result = worker.run(ctx).await;
611 debug!("Worker::run finished {worker_result:?}");
612 let result_obtained_at = clock_fn.now();
613 match Self::worker_result_to_execution_event(
614 component_type,
615 execution_id,
616 worker_result,
617 result_obtained_at,
618 parent,
619 can_be_retried,
620 unlock_expiry_on_limit_reached,
621 )? {
622 Some(append) => {
623 trace!("Appending {append:?}");
624 let db_exec = db_pool.db_exec_conn().await?;
625 append.append(db_exec.as_ref()).await
626 }
627 None => Ok(()),
628 }
629 }
630
631 fn worker_result_to_execution_event(
633 component_type: ComponentType,
634 execution_id: ExecutionId,
635 worker_result: WorkerResult,
636 result_obtained_at: DateTime<Utc>,
637 parent: Option<(ExecutionId, JoinSetId)>,
638 can_be_retried: Option<Duration>,
639 unlock_expiry_on_limit_reached: Duration,
640 ) -> Result<Option<Append>, DbErrorWrite> {
641 Ok(match worker_result {
642 WorkerResult::Ok(WorkerResultOk::RunFinished(RunFinished {
643 retval: ref retval @ SupportedFunctionReturnValue::Err(ref result_err),
644 version,
645 http_client_traces,
646 })) if component_type == ComponentType::Activity
647 && can_be_retried.is_some()
648 && !retval.is_permanent_variant() =>
649 {
650 let detail = serde_json::to_string(result_err)
652 .expect("SupportedFunctionReturnValue should be serializable to JSON");
653 let duration = can_be_retried.expect(
654 "ActivityReturnedError must not be returned when retries are exhausted",
655 );
656 let expires_at = result_obtained_at + duration;
657 debug!("Retrying ActivityReturnedError after {duration:?} at {expires_at}");
658 let primary_event = ExecutionRequest::TemporarilyFailed {
659 backoff_expires_at: expires_at,
660 reason: StrVariant::Static("activity finished with error"),
661 detail: Some(detail),
662 http_client_traces,
663 };
664 Some(Append {
665 created_at: result_obtained_at,
666 primary_event: AppendRequest {
667 created_at: result_obtained_at,
668 event: primary_event,
669 },
670 execution_id,
671 version,
672 child_finished: None,
673 })
674 }
675
676 WorkerResult::Ok(WorkerResultOk::RunFinished(RunFinished {
677 retval: result,
678 version,
679 http_client_traces,
680 })) => {
681 info!("Execution finished: {result}");
682 let child_finished =
683 parent.map(
684 |(parent_execution_id, parent_join_set)| ChildFinishedResponse {
685 parent_execution_id,
686 parent_join_set,
687 result: result.clone(),
688 },
689 );
690 let primary_event = AppendRequest {
691 created_at: result_obtained_at,
692 event: ExecutionRequest::Finished {
693 retval: result,
694 http_client_traces,
695 },
696 };
697
698 Some(Append {
699 created_at: result_obtained_at,
700 primary_event,
701 execution_id,
702 version,
703 child_finished,
704 })
705 }
706
707 WorkerResult::Ok(WorkerResultOk::DbUpdatedByWorkerOrWatcher) => None,
708
709 WorkerResult::Err(err) => {
710 let reason_generic = err.to_string(); let (primary_event, child_finished, version) = match err {
713 WorkerError::ExecutorClosing(version) => {
714 let primary_event = ExecutionRequest::Unlocked(Unlocked {
715 unlocked_at: result_obtained_at,
716 reason: "executor closing".into(),
717 });
718 (primary_event, None, version)
719 }
720 WorkerError::TemporaryTimeout {
721 http_client_traces,
722 version,
723 } => {
724 if let Some(duration) = can_be_retried {
725 let backoff_expires_at = result_obtained_at + duration;
726 info!(
727 "Temporary timeout, retrying after {duration:?} at {backoff_expires_at}"
728 );
729 (
730 ExecutionRequest::TemporarilyTimedOut {
731 backoff_expires_at,
732 http_client_traces,
733 },
734 None,
735 version,
736 )
737 } else {
738 info!("Execution timed out");
739 let result = SupportedFunctionReturnValue::ExecutionFailure(
740 FinishedExecutionFailure {
741 kind: ExecutionFailureKind::TimedOut,
742 reason: None,
743 detail: None,
744 },
745 );
746 let child_finished =
747 parent.map(|(parent_execution_id, parent_join_set)| {
748 ChildFinishedResponse {
749 parent_execution_id,
750 parent_join_set,
751 result: result.clone(),
752 }
753 });
754 (
755 ExecutionRequest::Finished {
756 retval: result,
757 http_client_traces,
758 },
759 child_finished,
760 version,
761 )
762 }
763 }
764 WorkerError::DbError(db_error) => {
765 return Err(db_error);
766 }
767 WorkerError::ActivityTrap {
768 reason: _, trap_kind,
770 detail,
771 version,
772 http_client_traces,
773 } => {
774 if let Some(duration) = can_be_retried {
775 let expires_at = result_obtained_at + duration;
776 debug!(
777 "Retrying activity with `{trap_kind}` execution after {duration:?} at {expires_at}"
778 );
779 (
780 ExecutionRequest::TemporarilyFailed {
781 reason: StrVariant::from(reason_generic),
782 backoff_expires_at: expires_at,
783 detail,
784 http_client_traces,
785 },
786 None,
787 version,
788 )
789 } else {
790 info!(
791 "Activity with `{trap_kind}` marked as permanent failure - {reason_generic}"
792 );
793 let result = SupportedFunctionReturnValue::ExecutionFailure(
794 FinishedExecutionFailure {
795 reason: Some(reason_generic),
796 kind: ExecutionFailureKind::Uncategorized,
797 detail,
798 },
799 );
800 let child_finished =
801 parent.map(|(parent_execution_id, parent_join_set)| {
802 ChildFinishedResponse {
803 parent_execution_id,
804 parent_join_set,
805 result: result.clone(),
806 }
807 });
808 (
809 ExecutionRequest::Finished {
810 retval: result,
811 http_client_traces,
812 },
813 child_finished,
814 version,
815 )
816 }
817 }
818 WorkerError::LimitReached {
819 reason,
820 version: new_version,
821 } => {
822 let expires_at = result_obtained_at + unlock_expiry_on_limit_reached;
823 warn!(
824 "Limit reached: {reason}, unlocking after {unlock_expiry_on_limit_reached:?} at {expires_at}"
825 );
826 (
827 ExecutionRequest::Unlocked(Unlocked {
828 unlocked_at: expires_at,
829 reason: StrVariant::from(reason),
830 }),
831 None,
832 new_version,
833 )
834 }
835 WorkerError::FatalError(FatalError::Cancelled, _version) => {
836 unreachable!(
837 "activity workers must return DbUpdatedByWorkerOrWatcher, cancellation append happens in CancelRegistry::cancel_activity"
838 )
839 }
840 WorkerError::FatalError(fatal_error, version) => {
841 warn!("Fatal worker error - {fatal_error:?}");
842 let result = SupportedFunctionReturnValue::ExecutionFailure(
843 FinishedExecutionFailure::from(fatal_error),
844 );
845 let child_finished =
846 parent.map(|(parent_execution_id, parent_join_set)| {
847 ChildFinishedResponse {
848 parent_execution_id,
849 parent_join_set,
850 result: result.clone(),
851 }
852 });
853 (
854 ExecutionRequest::Finished {
855 retval: result,
856 http_client_traces: None,
857 },
858 child_finished,
859 version,
860 )
861 }
862 };
863 Some(Append {
864 created_at: result_obtained_at,
865 primary_event: AppendRequest {
866 created_at: result_obtained_at,
867 event: primary_event,
868 },
869 execution_id,
870 version,
871 child_finished,
872 })
873 }
874 })
875 }
876}
877
878#[derive(Debug, Clone)]
879pub(crate) struct ChildFinishedResponse {
880 pub(crate) parent_execution_id: ExecutionId,
881 pub(crate) parent_join_set: JoinSetId,
882 pub(crate) result: SupportedFunctionReturnValue,
883}
884
885#[derive(Debug, Clone)]
886pub(crate) struct Append {
887 pub(crate) created_at: DateTime<Utc>,
888 pub(crate) primary_event: AppendRequest,
889 pub(crate) execution_id: ExecutionId,
890 pub(crate) version: Version,
891 pub(crate) child_finished: Option<ChildFinishedResponse>,
892}
893
894impl Append {
895 pub(crate) async fn append(self, db_exec: &dyn DbExecutor) -> Result<(), DbErrorWrite> {
896 if let Some(child_finished) = self.child_finished {
897 assert_matches!(
898 &self.primary_event,
899 AppendRequest {
900 event: ExecutionRequest::Finished { .. },
901 ..
902 }
903 );
904 let child_execution_id = assert_matches!(self.execution_id.clone(), ExecutionId::Derived(derived) => derived);
905 let events = AppendEventsToExecution {
906 execution_id: self.execution_id,
907 version: self.version.clone(),
908 batch: vec![self.primary_event],
909 };
910 let response = AppendResponseToExecution {
911 parent_execution_id: child_finished.parent_execution_id,
912 created_at: self.created_at,
913 join_set_id: child_finished.parent_join_set,
914 child_execution_id,
915 finished_version: self.version, result: child_finished.result,
917 };
918
919 db_exec
920 .append_batch_respond_to_parent(events, response, self.created_at)
921 .await?;
922 } else {
923 let res = db_exec
924 .append(self.execution_id, self.version, self.primary_event)
925 .await;
926 match res {
927 Ok(_)
928 | Err(DbErrorWrite::NonRetriable(
929 DbErrorWriteNonRetriable::UnlockedCannotBeAppended(_),
930 )) => {
931 }
933 Err(err) => return Err(err),
934 }
935 }
936 Ok(())
937 }
938}
939
940fn log_err_if_new<T>(
941 res: Result<T, DbErrorGeneric>,
942 old_err: &mut Option<DbErrorGeneric>,
943) -> Result<T, ()> {
944 match (res, &old_err) {
945 (Ok(ok), _) => {
946 *old_err = None;
947 Ok(ok)
948 }
949 (Err(err), Some(old)) if err == *old => Err(()),
950 (Err(err), _) => {
951 warn!("Tick failed: {err:?}");
952 *old_err = Some(err);
953 Err(())
954 }
955 }
956}
957
958#[cfg(any(test, feature = "test"))]
959pub mod simple_worker {
960 use crate::worker::{Worker, WorkerContext, WorkerResult};
961 use async_trait::async_trait;
962 use concepts::{
963 FunctionFqn, FunctionMetadata, ParameterTypes, RETURN_TYPE_DUMMY,
964 storage::{HistoryEvent, Version},
965 };
966 use indexmap::IndexMap;
967 use std::sync::Arc;
968 use tracing::trace;
969
970 pub(crate) const FFQN_SOME: FunctionFqn = FunctionFqn::new_static("ns:pkg/ifc", "fn");
971 pub type SimpleWorkerResultMap =
972 Arc<std::sync::Mutex<IndexMap<Version, (Vec<HistoryEvent>, WorkerResult)>>>;
973
974 #[derive(Clone, Debug)]
975 pub struct SimpleWorker {
976 pub worker_results_rev: SimpleWorkerResultMap,
977 pub ffqn: FunctionFqn,
978 exported: [FunctionMetadata; 1],
979 }
980
981 impl SimpleWorker {
982 #[must_use]
983 pub fn with_single_result(res: WorkerResult) -> Self {
984 Self::with_worker_results_rev(Arc::new(std::sync::Mutex::new(IndexMap::from([(
985 Version::new(2),
986 (vec![], res),
987 )]))))
988 }
989
990 #[must_use]
991 pub fn with_ffqn(self, ffqn: FunctionFqn) -> Self {
992 Self {
993 worker_results_rev: self.worker_results_rev,
994 exported: [FunctionMetadata {
995 ffqn: ffqn.clone(),
996 parameter_types: ParameterTypes::default(),
997 return_type: RETURN_TYPE_DUMMY,
998 extension: None,
999 submittable: true,
1000 }],
1001 ffqn,
1002 }
1003 }
1004
1005 #[must_use]
1006 pub fn with_worker_results_rev(worker_results_rev: SimpleWorkerResultMap) -> Self {
1007 Self {
1008 worker_results_rev,
1009 ffqn: FFQN_SOME,
1010 exported: [FunctionMetadata {
1011 ffqn: FFQN_SOME,
1012 parameter_types: ParameterTypes::default(),
1013 return_type: RETURN_TYPE_DUMMY,
1014 extension: None,
1015 submittable: true,
1016 }],
1017 }
1018 }
1019 }
1020
1021 #[async_trait]
1022 impl Worker for SimpleWorker {
1023 async fn run(&self, ctx: WorkerContext) -> WorkerResult {
1024 let (expected_version, (expected_eh, worker_result)) =
1025 self.worker_results_rev.lock().unwrap().pop().unwrap();
1026 trace!(%expected_version, version = %ctx.version, ?expected_eh, eh = ?ctx.event_history, "Running SimpleWorker");
1027 assert_eq!(expected_version, ctx.version);
1028 assert_eq!(
1029 expected_eh,
1030 ctx.event_history
1031 .iter()
1032 .map(|(event, _version)| event.clone())
1033 .collect::<Vec<_>>()
1034 );
1035 worker_result
1036 }
1037
1038 fn exported_functions_noext(&self) -> &[FunctionMetadata] {
1039 &self.exported
1040 }
1041 }
1042}
1043
1044#[cfg(test)]
1045mod tests {
1046 use self::simple_worker::SimpleWorker;
1047 use super::*;
1048 use crate::{expired_timers_watcher, worker::WorkerResult};
1049 use assert_matches::assert_matches;
1050 use async_trait::async_trait;
1051 use concepts::prefixed_ulid::DEPLOYMENT_ID_DUMMY;
1052 use concepts::storage::{
1053 CreateRequest, DbConnectionTest, JoinSetRequest, JoinSetResponse, JoinSetResponseEvent,
1054 };
1055 use concepts::storage::{DbPoolCloseable, LockedBy};
1056 use concepts::storage::{
1057 ExecutionEvent, ExecutionRequest, HistoryEvent, PendingState, PendingStatePendingAt,
1058 };
1059 use concepts::time::{ConstClock, Now};
1060 use concepts::{
1061 FunctionMetadata, JoinSetKind, ParameterTypes, Params, RETURN_TYPE_DUMMY,
1062 SUPPORTED_RETURN_VALUE_OK_EMPTY, StrVariant, SupportedFunctionReturnValue, TrapKind,
1063 };
1064 use db_tests::Database;
1065 use indexmap::IndexMap;
1066 use rstest::rstest;
1067 use simple_worker::FFQN_SOME;
1068 use std::{fmt::Debug, future::Future, ops::Deref, sync::Arc};
1069 use test_db_macro::expand_enum_database;
1070 use test_utils::set_up;
1071 use test_utils::sim_clock::SimClock;
1072
1073 pub(crate) const FFQN_CHILD: FunctionFqn = FunctionFqn::new_static("ns:pkg/ifc", "fn-child");
1074
1075 async fn tick_fn<W: Worker + Debug>(
1076 config: ExecConfig,
1077 clock_fn: Box<dyn ClockFn>,
1078 db_pool: Arc<dyn DbPool>,
1079 worker: Arc<W>,
1080 executed_at: DateTime<Utc>,
1081 ) -> Vec<ExecutionId> {
1082 trace!("Ticking with {worker:?}");
1083 let ffqns = super::extract_exported_ffqns_noext(worker.as_ref());
1084 let (executor, _close_tx) = ExecTask::new_test(config, worker, clock_fn, db_pool, ffqns);
1085 executor
1086 .tick_test_await(executed_at, RunId::generate())
1087 .await
1088 }
1089
1090 #[expand_enum_database]
1091 #[rstest]
1092 #[tokio::test]
1093 async fn execute_simple_lifecycle_tick_based(
1094 database: Database,
1095 #[values(LockingStrategy::ByFfqns, LockingStrategy::ByComponentDigest)]
1096 locking_strategy: LockingStrategy,
1097 ) {
1098 set_up();
1099 let created_at = Now.now();
1100 let (_guard, db_pool, db_close) = database.set_up().await;
1101 let db_connection = db_pool.connection_test().await.unwrap();
1102 execute_simple_lifecycle_tick_based_inner(
1103 db_connection.as_ref(),
1104 db_pool.clone(),
1105 Box::new(ConstClock(created_at)),
1106 locking_strategy,
1107 )
1108 .await;
1109 drop(db_connection);
1110 db_close.close().await;
1111 }
1112
1113 async fn execute_simple_lifecycle_tick_based_inner(
1114 db_connection: &dyn DbConnectionTest,
1115 db_pool: Arc<dyn DbPool>,
1116 clock_fn: Box<dyn ClockFn>,
1117 locking_strategy: LockingStrategy,
1118 ) {
1119 let created_at = clock_fn.now();
1120 let exec_config = ExecConfig {
1121 batch_size: 1,
1122 lock_expiry: Duration::from_secs(1),
1123 tick_sleep: Duration::from_millis(100),
1124 component_id: ComponentId::dummy_activity(),
1125 task_limiter_global: None,
1126 task_limiter_local: None,
1127 executor_id: ExecutorId::generate(),
1128 retry_config: ComponentRetryConfig::ZERO,
1129 locking_strategy,
1130 };
1131
1132 let execution_log = create_and_tick(
1133 CreateAndTickConfig {
1134 execution_id: ExecutionId::generate(),
1135 created_at,
1136 executed_at: created_at,
1137 },
1138 clock_fn,
1139 db_connection,
1140 db_pool,
1141 exec_config,
1142 Arc::new(SimpleWorker::with_single_result(WorkerResult::Ok(
1143 WorkerResultOk::RunFinished(RunFinished {
1144 retval: SUPPORTED_RETURN_VALUE_OK_EMPTY,
1145 version: Version::new(2),
1146 http_client_traces: None,
1147 }),
1148 ))),
1149 tick_fn,
1150 )
1151 .await;
1152 assert_matches!(
1153 execution_log.events.get(2).unwrap(),
1154 ExecutionEvent {
1155 event: ExecutionRequest::Finished {
1156 retval: SupportedFunctionReturnValue::Ok(None),
1157 http_client_traces: None
1158 },
1159 created_at: _,
1160 backtrace_id: None,
1161 version: Version(2),
1162 }
1163 );
1164 }
1165
1166 #[rstest]
1167 #[tokio::test]
1168 async fn execute_simple_lifecycle_task_based_mem(
1169 #[values(LockingStrategy::ByFfqns, LockingStrategy::ByComponentDigest)]
1170 locking_strategy: LockingStrategy,
1171 ) {
1172 set_up();
1173 let created_at = Now.now();
1174 let clock_fn = Box::new(ConstClock(created_at));
1175 let (_guard, db_pool, db_close) = Database::Sqlite.set_up().await;
1176 let exec_config = ExecConfig {
1177 batch_size: 1,
1178 lock_expiry: Duration::from_secs(1),
1179 tick_sleep: Duration::ZERO,
1180 component_id: ComponentId::dummy_activity(),
1181 task_limiter_global: None,
1182 task_limiter_local: None,
1183 executor_id: ExecutorId::generate(),
1184 retry_config: ComponentRetryConfig::ZERO,
1185 locking_strategy,
1186 };
1187
1188 let worker = Arc::new(SimpleWorker::with_single_result(WorkerResult::Ok(
1189 WorkerResultOk::RunFinished(RunFinished {
1190 retval: SUPPORTED_RETURN_VALUE_OK_EMPTY,
1191 version: Version::new(2),
1192 http_client_traces: None,
1193 }),
1194 )));
1195 let db_connection = db_pool.connection_test().await.unwrap();
1196
1197 let execution_log = create_and_tick(
1198 CreateAndTickConfig {
1199 execution_id: ExecutionId::generate(),
1200 created_at,
1201 executed_at: created_at,
1202 },
1203 clock_fn,
1204 db_connection.as_ref(),
1205 db_pool,
1206 exec_config,
1207 worker,
1208 tick_fn,
1209 )
1210 .await;
1211 assert_matches!(
1212 execution_log.events.get(2).unwrap(),
1213 ExecutionEvent {
1214 event: ExecutionRequest::Finished {
1215 retval: SupportedFunctionReturnValue::Ok(None),
1216 http_client_traces: None
1217 },
1218 created_at: _,
1219 backtrace_id: None,
1220 version: Version(2),
1221 }
1222 );
1223 db_close.close().await;
1224 }
1225
1226 struct CreateAndTickConfig {
1227 execution_id: ExecutionId,
1228 created_at: DateTime<Utc>,
1229 executed_at: DateTime<Utc>,
1230 }
1231
1232 async fn create_and_tick<
1233 W: Worker,
1234 T: FnMut(ExecConfig, Box<dyn ClockFn>, Arc<dyn DbPool>, Arc<W>, DateTime<Utc>) -> F,
1235 F: Future<Output = Vec<ExecutionId>>,
1236 >(
1237 config: CreateAndTickConfig,
1238 clock_fn: Box<dyn ClockFn>,
1239 db_connection: &dyn DbConnectionTest,
1240 db_pool: Arc<dyn DbPool>,
1241 exec_config: ExecConfig,
1242 worker: Arc<W>,
1243 mut tick: T,
1244 ) -> ExecutionLog {
1245 db_connection
1247 .create(CreateRequest {
1248 created_at: config.created_at,
1249 execution_id: config.execution_id.clone(),
1250 ffqn: FFQN_SOME,
1251 params: Params::empty(),
1252 parent: None,
1253 metadata: concepts::ExecutionMetadata::empty(),
1254 scheduled_at: config.created_at,
1255 component_id: ComponentId::dummy_activity(),
1256 deployment_id: DEPLOYMENT_ID_DUMMY,
1257 scheduled_by: None,
1258 paused: false,
1259 })
1260 .await
1261 .unwrap();
1262 tick(exec_config, clock_fn, db_pool, worker, config.executed_at).await;
1264 let execution_log = db_connection.get(&config.execution_id).await.unwrap();
1265 debug!("Execution history after tick: {execution_log:?}");
1266 let actually_created_at = assert_matches!(
1268 execution_log.events.first().unwrap(),
1269 ExecutionEvent {
1270 event: ExecutionRequest::Created { .. },
1271 created_at: actually_created_at,
1272 backtrace_id: None,
1273 version: Version(0),
1274 }
1275 => *actually_created_at
1276 );
1277 assert_eq!(config.created_at, actually_created_at);
1278 let locked_at = assert_matches!(
1279 execution_log.events.get(1).unwrap(),
1280 ExecutionEvent {
1281 event: ExecutionRequest::Locked { .. },
1282 created_at: locked_at,
1283 backtrace_id: None,
1284 version: Version(1),
1285 } if config.created_at <= *locked_at
1286 => *locked_at
1287 );
1288 assert_matches!(execution_log.events.get(2).unwrap(), ExecutionEvent {
1289 event: _,
1290 created_at: executed_at,
1291 backtrace_id: None,
1292 version: Version(2),
1293 } if *executed_at >= locked_at);
1294 execution_log
1295 }
1296
1297 #[rstest]
1298 #[tokio::test]
1299 async fn activity_trap_should_trigger_an_execution_retry(
1300 #[values(LockingStrategy::ByFfqns, LockingStrategy::ByComponentDigest)]
1301 locking_strategy: LockingStrategy,
1302 ) {
1303 set_up();
1304 let sim_clock = SimClock::default();
1305 let (_guard, db_pool, db_close) = Database::Sqlite.set_up().await;
1306 let retry_exp_backoff = Duration::from_millis(100);
1307 let retry_config = ComponentRetryConfig {
1308 max_retries: Some(1),
1309 retry_exp_backoff,
1310 };
1311 let exec_config = ExecConfig {
1312 batch_size: 1,
1313 lock_expiry: Duration::from_secs(1),
1314 tick_sleep: Duration::ZERO,
1315 component_id: ComponentId::dummy_activity(),
1316 task_limiter_global: None,
1317 task_limiter_local: None,
1318 executor_id: ExecutorId::generate(),
1319 retry_config,
1320 locking_strategy,
1321 };
1322 let expected_reason = "error reason";
1323 let expected_detail = "error detail";
1324 let worker = Arc::new(SimpleWorker::with_single_result(WorkerResult::Err(
1325 WorkerError::ActivityTrap {
1326 reason: expected_reason.to_string(),
1327 trap_kind: concepts::TrapKind::Trap,
1328 detail: Some(expected_detail.to_string()),
1329 version: Version::new(2),
1330 http_client_traces: None,
1331 },
1332 )));
1333 debug!(now = %sim_clock.now(), "Creating an execution that should fail");
1334 let db_connection = db_pool.connection_test().await.unwrap();
1335 let execution_log = create_and_tick(
1336 CreateAndTickConfig {
1337 execution_id: ExecutionId::generate(),
1338 created_at: sim_clock.now(),
1339 executed_at: sim_clock.now(),
1340 },
1341 sim_clock.clone_box(),
1342 db_connection.as_ref(),
1343 db_pool.clone(),
1344 exec_config.clone(),
1345 worker,
1346 tick_fn,
1347 )
1348 .await;
1349 assert_eq!(3, execution_log.events.len());
1350 {
1351 let (reason, detail, at, expires_at) = assert_matches!(
1352 &execution_log.events.get(2).unwrap(),
1353 ExecutionEvent {
1354 event: ExecutionRequest::TemporarilyFailed {
1355 reason,
1356 detail,
1357 backoff_expires_at,
1358 http_client_traces: None,
1359 },
1360 created_at: at,
1361 backtrace_id: None,
1362 version: Version(2),
1363 }
1364 => (reason, detail, *at, *backoff_expires_at)
1365 );
1366 assert_eq!(format!("activity trap: {expected_reason}"), reason.deref());
1367 assert_eq!(Some(expected_detail), detail.as_deref());
1368 assert_eq!(at, sim_clock.now());
1369 assert_eq!(sim_clock.now() + retry_config.retry_exp_backoff, expires_at);
1370 }
1371 let worker = Arc::new(SimpleWorker::with_worker_results_rev(Arc::new(
1372 std::sync::Mutex::new(IndexMap::from([(
1373 Version::new(4),
1374 (
1375 vec![],
1376 WorkerResult::Ok(WorkerResultOk::RunFinished(RunFinished {
1377 retval: SUPPORTED_RETURN_VALUE_OK_EMPTY,
1378 version: Version::new(4),
1379 http_client_traces: None,
1380 })),
1381 ),
1382 )])),
1383 )));
1384 assert!(
1386 tick_fn(
1387 exec_config.clone(),
1388 sim_clock.clone_box(),
1389 db_pool.clone(),
1390 worker.clone(),
1391 sim_clock.now(),
1392 )
1393 .await
1394 .is_empty()
1395 );
1396 sim_clock.move_time_forward(retry_config.retry_exp_backoff);
1398 tick_fn(
1399 exec_config,
1400 sim_clock.clone_box(),
1401 db_pool.clone(),
1402 worker,
1403 sim_clock.now(),
1404 )
1405 .await;
1406 let execution_log = {
1407 let db_connection = db_pool.connection_test().await.unwrap();
1408 db_connection
1409 .get(&execution_log.execution_id)
1410 .await
1411 .unwrap()
1412 };
1413 debug!(now = %sim_clock.now(), "Execution history after second tick: {execution_log:?}");
1414 assert_matches!(
1415 execution_log.events.get(3).unwrap(),
1416 ExecutionEvent {
1417 event: ExecutionRequest::Locked { .. },
1418 created_at: at,
1419 backtrace_id: None,
1420 version: Version(3),
1421 } if *at == sim_clock.now()
1422 );
1423 assert_matches!(
1424 execution_log.events.get(4).unwrap(),
1425 ExecutionEvent {
1426 event: ExecutionRequest::Finished {
1427 retval: SupportedFunctionReturnValue::Ok(None),
1428 http_client_traces: None
1429 },
1430 created_at: finished_at,
1431 backtrace_id: None,
1432 version: Version(4),
1433 } if *finished_at == sim_clock.now()
1434 );
1435 db_close.close().await;
1436 }
1437
1438 #[rstest]
1439 #[tokio::test]
1440 async fn activity_trap_should_not_be_retried_if_no_retries_are_set(
1441 #[values(LockingStrategy::ByFfqns, LockingStrategy::ByComponentDigest)]
1442 locking_strategy: LockingStrategy,
1443 ) {
1444 set_up();
1445 let created_at = Now.now();
1446 let clock_fn = Box::new(ConstClock(created_at));
1447 let (_guard, db_pool, db_close) = Database::Sqlite.set_up().await;
1448 let exec_config = ExecConfig {
1449 batch_size: 1,
1450 lock_expiry: Duration::from_secs(1),
1451 tick_sleep: Duration::ZERO,
1452 component_id: ComponentId::dummy_activity(),
1453 task_limiter_global: None,
1454 task_limiter_local: None,
1455 executor_id: ExecutorId::generate(),
1456 retry_config: ComponentRetryConfig::ZERO,
1457 locking_strategy,
1458 };
1459
1460 let reason = "error reason";
1461 let expected_reason = format!("activity trap: {reason}");
1462 let expected_detail = "error detail";
1463 let worker = Arc::new(SimpleWorker::with_single_result(WorkerResult::Err(
1464 WorkerError::ActivityTrap {
1465 reason: reason.to_string(),
1466 trap_kind: concepts::TrapKind::Trap,
1467 detail: Some(expected_detail.to_string()),
1468 version: Version::new(2),
1469 http_client_traces: None,
1470 },
1471 )));
1472 let execution_log = create_and_tick(
1473 CreateAndTickConfig {
1474 execution_id: ExecutionId::generate(),
1475 created_at,
1476 executed_at: created_at,
1477 },
1478 clock_fn,
1479 db_pool.connection_test().await.unwrap().as_ref(),
1480 db_pool.clone(),
1481 exec_config.clone(),
1482 worker,
1483 tick_fn,
1484 )
1485 .await;
1486 assert_eq!(3, execution_log.events.len());
1487 let (reason, kind, detail) = assert_matches!(
1488 &execution_log.events.get(2).unwrap(),
1489 ExecutionEvent {
1490 event: ExecutionRequest::Finished{
1491 retval: SupportedFunctionReturnValue::ExecutionFailure(FinishedExecutionFailure{reason, kind, detail}),
1492 http_client_traces: None
1493 },
1494 created_at: at,
1495 backtrace_id: None,
1496 version: Version(2),
1497 } if *at == created_at
1498 => (reason, kind, detail)
1499 );
1500
1501 assert_eq!(Some(expected_reason), *reason);
1502 assert_eq!(Some(expected_detail), detail.as_deref());
1503 assert_eq!(ExecutionFailureKind::Uncategorized, *kind);
1504
1505 db_close.close().await;
1506 }
1507
1508 #[rstest]
1509 #[tokio::test]
1510 async fn child_execution_permanently_failed_should_notify_parent_permanent_failure(
1511 #[values(LockingStrategy::ByFfqns, LockingStrategy::ByComponentDigest)]
1512 locking_strategy: LockingStrategy,
1513 ) {
1514 let worker_error = WorkerError::ActivityTrap {
1515 reason: "error reason".to_string(),
1516 trap_kind: TrapKind::Trap,
1517 detail: Some("detail".to_string()),
1518 version: Version::new(2),
1519 http_client_traces: None,
1520 };
1521 let expected_child_err = FinishedExecutionFailure {
1522 kind: ExecutionFailureKind::Uncategorized,
1523 reason: Some("activity trap: error reason".to_string()),
1524 detail: Some("detail".to_string()),
1525 };
1526 child_execution_permanently_failed_should_notify_parent(
1527 WorkerResult::Err(worker_error),
1528 expected_child_err,
1529 locking_strategy,
1530 )
1531 .await;
1532 }
1533
1534 #[rstest]
1535 #[tokio::test]
1536 async fn child_execution_permanently_failed_handled_by_watcher_should_notify_parent_timeout(
1537 #[values(LockingStrategy::ByFfqns, LockingStrategy::ByComponentDigest)]
1538 locking_strategy: LockingStrategy,
1539 ) {
1540 let expected_child_err = FinishedExecutionFailure {
1541 kind: ExecutionFailureKind::TimedOut,
1542 reason: None,
1543 detail: None,
1544 };
1545 child_execution_permanently_failed_should_notify_parent(
1546 WorkerResult::Ok(WorkerResultOk::DbUpdatedByWorkerOrWatcher),
1547 expected_child_err,
1548 locking_strategy,
1549 )
1550 .await;
1551 }
1552
1553 async fn child_execution_permanently_failed_should_notify_parent(
1554 worker_result: WorkerResult,
1555 expected_child_err: FinishedExecutionFailure,
1556 locking_strategy: LockingStrategy,
1557 ) {
1558 use concepts::storage::JoinSetResponseEventOuter;
1559 const LOCK_EXPIRY: Duration = Duration::from_secs(1);
1560
1561 set_up();
1562 let sim_clock = SimClock::default();
1563 let (_guard, db_pool, db_close) = Database::Sqlite.set_up().await;
1564
1565 let parent_worker = Arc::new(SimpleWorker::with_single_result(WorkerResult::Ok(
1566 WorkerResultOk::DbUpdatedByWorkerOrWatcher,
1567 )));
1568 let parent_execution_id = ExecutionId::generate();
1569 db_pool
1570 .connection()
1571 .await
1572 .unwrap()
1573 .create(CreateRequest {
1574 created_at: sim_clock.now(),
1575 execution_id: parent_execution_id.clone(),
1576 ffqn: FFQN_SOME,
1577 params: Params::empty(),
1578 parent: None,
1579 metadata: concepts::ExecutionMetadata::empty(),
1580 scheduled_at: sim_clock.now(),
1581 component_id: ComponentId::dummy_activity(),
1582 deployment_id: DEPLOYMENT_ID_DUMMY,
1583 scheduled_by: None,
1584 paused: false,
1585 })
1586 .await
1587 .unwrap();
1588 let parent_executor_id = ExecutorId::generate();
1589 tick_fn(
1590 ExecConfig {
1591 batch_size: 1,
1592 lock_expiry: LOCK_EXPIRY,
1593 tick_sleep: Duration::ZERO,
1594 component_id: ComponentId::dummy_activity(),
1595 task_limiter_global: None,
1596 task_limiter_local: None,
1597 executor_id: parent_executor_id,
1598 retry_config: ComponentRetryConfig::ZERO,
1599 locking_strategy,
1600 },
1601 sim_clock.clone_box(),
1602 db_pool.clone(),
1603 parent_worker,
1604 sim_clock.now(),
1605 )
1606 .await;
1607
1608 let join_set_id = JoinSetId::new(JoinSetKind::OneOff, StrVariant::empty()).unwrap();
1609 let child_execution_id = parent_execution_id.next_level(&join_set_id);
1610 {
1612 let params = Params::empty();
1613 let child = CreateRequest {
1614 created_at: sim_clock.now(),
1615 execution_id: ExecutionId::Derived(child_execution_id.clone()),
1616 ffqn: FFQN_CHILD,
1617 params: params.clone(),
1618 parent: Some((parent_execution_id.clone(), join_set_id.clone())),
1619 metadata: concepts::ExecutionMetadata::empty(),
1620 scheduled_at: sim_clock.now(),
1621 component_id: ComponentId::dummy_activity(),
1622 deployment_id: DEPLOYMENT_ID_DUMMY,
1623 scheduled_by: None,
1624 paused: false,
1625 };
1626 let current_time = sim_clock.now();
1627 let join_set = AppendRequest {
1628 created_at: current_time,
1629 event: ExecutionRequest::HistoryEvent {
1630 event: HistoryEvent::JoinSetCreate {
1631 join_set_id: join_set_id.clone(),
1632 },
1633 },
1634 };
1635 let child_exec_req = AppendRequest {
1636 created_at: current_time,
1637 event: ExecutionRequest::HistoryEvent {
1638 event: HistoryEvent::JoinSetRequest {
1639 join_set_id: join_set_id.clone(),
1640 request: JoinSetRequest::ChildExecutionRequest {
1641 child_execution_id: child_execution_id.clone(),
1642 target_ffqn: FFQN_CHILD,
1643 params,
1644 result: Ok(()),
1645 },
1646 },
1647 },
1648 };
1649 let join_next = AppendRequest {
1650 created_at: current_time,
1651 event: ExecutionRequest::HistoryEvent {
1652 event: HistoryEvent::JoinNext {
1653 join_set_id: join_set_id.clone(),
1654 run_expires_at: sim_clock.now(),
1655 closing: false,
1656 requested_ffqn: Some(FFQN_CHILD),
1657 },
1658 },
1659 };
1660 db_pool
1661 .connection()
1662 .await
1663 .unwrap()
1664 .append_batch_create_new_execution(
1665 current_time,
1666 vec![join_set, child_exec_req, join_next],
1667 parent_execution_id.clone(),
1668 Version::new(2),
1669 vec![child],
1670 vec![],
1671 )
1672 .await
1673 .unwrap();
1674 }
1675
1676 let child_worker =
1677 Arc::new(SimpleWorker::with_single_result(worker_result).with_ffqn(FFQN_CHILD));
1678
1679 tick_fn(
1681 ExecConfig {
1682 batch_size: 1,
1683 lock_expiry: LOCK_EXPIRY,
1684 tick_sleep: Duration::ZERO,
1685 component_id: ComponentId::dummy_activity(),
1686 task_limiter_global: None,
1687 task_limiter_local: None,
1688 executor_id: ExecutorId::generate(),
1689 retry_config: ComponentRetryConfig::ZERO,
1690 locking_strategy,
1691 },
1692 sim_clock.clone_box(),
1693 db_pool.clone(),
1694 child_worker,
1695 sim_clock.now(),
1696 )
1697 .await;
1698 if matches!(expected_child_err.kind, ExecutionFailureKind::TimedOut) {
1699 sim_clock.move_time_forward(LOCK_EXPIRY);
1701 expired_timers_watcher::tick(
1702 db_pool.connection().await.unwrap().as_ref(),
1703 sim_clock.now(),
1704 )
1705 .await
1706 .unwrap();
1707 }
1708 let child_log = db_pool
1709 .connection_test()
1710 .await
1711 .unwrap()
1712 .get(&ExecutionId::Derived(child_execution_id.clone()))
1713 .await
1714 .unwrap();
1715 assert!(child_log.pending_state.is_finished());
1716 assert_eq!(
1717 Version(2),
1718 child_log.next_version,
1719 "created = 0, locked = 1, with_single_result = 2"
1720 );
1721 assert_eq!(
1722 ExecutionRequest::Finished {
1723 retval: SupportedFunctionReturnValue::ExecutionFailure(expected_child_err),
1724 http_client_traces: None
1725 },
1726 child_log.last_event().event
1727 );
1728 let parent_log = db_pool
1729 .connection_test()
1730 .await
1731 .unwrap()
1732 .get(&parent_execution_id)
1733 .await
1734 .unwrap();
1735 assert_matches!(
1736 parent_log.pending_state,
1737 PendingState::PendingAt(PendingStatePendingAt {
1738 scheduled_at,
1739 last_lock: Some(LockedBy { executor_id: found_executor_id, run_id: _}),
1740 }) if scheduled_at == sim_clock.now() && found_executor_id == parent_executor_id,
1741 "parent should be back to pending"
1742 );
1743 let (found_join_set_id, found_child_execution_id, child_finished_version, found_result) = assert_matches!(
1744 parent_log.responses.last().map(|resp| &resp.event),
1745 Some(JoinSetResponseEventOuter{
1746 created_at: at,
1747 event: JoinSetResponseEvent{
1748 join_set_id: found_join_set_id,
1749 event: JoinSetResponse::ChildExecutionFinished {
1750 child_execution_id: found_child_execution_id,
1751 finished_version,
1752 result: found_result,
1753 }
1754 }
1755 })
1756 if *at == sim_clock.now()
1757 => (found_join_set_id, found_child_execution_id, finished_version, found_result)
1758 );
1759 assert_eq!(join_set_id, *found_join_set_id);
1760 assert_eq!(child_execution_id, *found_child_execution_id);
1761 assert_eq!(child_log.next_version, *child_finished_version);
1762 assert_matches!(
1763 found_result,
1764 SupportedFunctionReturnValue::ExecutionFailure(_)
1765 );
1766
1767 db_close.close().await;
1768 }
1769
1770 #[derive(Clone, Debug)]
1771 struct SleepyWorker {
1772 duration: Duration,
1773 result: SupportedFunctionReturnValue,
1774 exported: [FunctionMetadata; 1],
1775 }
1776
1777 #[async_trait]
1778 impl Worker for SleepyWorker {
1779 async fn run(&self, ctx: WorkerContext) -> WorkerResult {
1780 tokio::time::sleep(self.duration).await;
1781 WorkerResult::Ok(WorkerResultOk::RunFinished(RunFinished {
1782 retval: self.result.clone(),
1783 version: ctx.version,
1784 http_client_traces: None,
1785 }))
1786 }
1787
1788 fn exported_functions_noext(&self) -> &[FunctionMetadata] {
1789 &self.exported
1790 }
1791 }
1792
1793 #[rstest]
1794 #[tokio::test]
1795 async fn hanging_lock_should_be_cleaned_and_execution_retried(
1796 #[values(LockingStrategy::ByFfqns, LockingStrategy::ByComponentDigest)]
1797 locking_strategy: LockingStrategy,
1798 ) {
1799 set_up();
1800 let sim_clock = SimClock::default();
1801 let (_guard, db_pool, db_close) = Database::Sqlite.set_up().await;
1802 let lock_expiry = Duration::from_millis(100);
1803 let timeout_duration = Duration::from_millis(300);
1804 let retry_config = ComponentRetryConfig {
1805 max_retries: Some(1),
1806 retry_exp_backoff: timeout_duration,
1807 };
1808 let exec_config = ExecConfig {
1809 batch_size: 1,
1810 lock_expiry,
1811 tick_sleep: Duration::ZERO,
1812 component_id: ComponentId::dummy_activity(),
1813 task_limiter_global: None,
1814 task_limiter_local: None,
1815 executor_id: ExecutorId::generate(),
1816 retry_config,
1817 locking_strategy,
1818 };
1819
1820 let worker = Arc::new(SleepyWorker {
1821 duration: lock_expiry + Duration::from_millis(1), result: SUPPORTED_RETURN_VALUE_OK_EMPTY,
1823 exported: [FunctionMetadata {
1824 ffqn: FFQN_SOME,
1825 parameter_types: ParameterTypes::default(),
1826 return_type: RETURN_TYPE_DUMMY,
1827 extension: None,
1828 submittable: true,
1829 }],
1830 });
1831 let execution_id = ExecutionId::generate();
1833 let db_connection = db_pool.connection_test().await.unwrap();
1834 db_connection
1835 .create(CreateRequest {
1836 created_at: sim_clock.now(),
1837 execution_id: execution_id.clone(),
1838 ffqn: FFQN_SOME,
1839 params: Params::empty(),
1840 parent: None,
1841 metadata: concepts::ExecutionMetadata::empty(),
1842 scheduled_at: sim_clock.now(),
1843 component_id: ComponentId::dummy_activity(),
1844 deployment_id: DEPLOYMENT_ID_DUMMY,
1845 scheduled_by: None,
1846 paused: false,
1847 })
1848 .await
1849 .unwrap();
1850
1851 let ffqns = super::extract_exported_ffqns_noext(worker.as_ref());
1852 let (executor, _close_tx) = ExecTask::new_test(
1853 exec_config.clone(),
1854 worker,
1855 sim_clock.clone_box(),
1856 db_pool.clone(),
1857 ffqns,
1858 );
1859 let db_exec = db_pool.db_exec_conn().await.unwrap();
1860 let mut first_execution_progress = executor
1861 .tick(
1862 db_exec.as_ref(),
1863 sim_clock.now(),
1864 RunId::generate(),
1865 DEPLOYMENT_ID_DUMMY,
1866 )
1867 .await
1868 .unwrap();
1869 assert_eq!(1, first_execution_progress.executions.len());
1870 sim_clock.move_time_forward(lock_expiry);
1872 let now_after_first_lock_expiry = sim_clock.now();
1874 {
1875 debug!(now = %now_after_first_lock_expiry, "Expecting an expired lock");
1876 let cleanup_progress = executor
1877 .tick(
1878 db_pool.db_exec_conn().await.unwrap().as_ref(),
1879 now_after_first_lock_expiry,
1880 RunId::generate(),
1881 DEPLOYMENT_ID_DUMMY,
1882 )
1883 .await
1884 .unwrap();
1885 assert!(cleanup_progress.executions.is_empty());
1886 }
1887 {
1888 let expired_locks = expired_timers_watcher::tick(
1889 db_pool.connection().await.unwrap().as_ref(),
1890 now_after_first_lock_expiry,
1891 )
1892 .await
1893 .unwrap()
1894 .expired_locks;
1895 assert_eq!(1, expired_locks);
1896 }
1897 assert!(
1898 !first_execution_progress
1899 .executions
1900 .pop()
1901 .unwrap()
1902 .1
1903 .is_finished()
1904 );
1905
1906 let execution_log = db_connection.get(&execution_id).await.unwrap();
1907 let expected_first_timeout_expiry = now_after_first_lock_expiry + timeout_duration;
1908 assert_matches!(
1909 &execution_log.events.get(2).unwrap(),
1910 ExecutionEvent {
1911 event: ExecutionRequest::TemporarilyTimedOut { backoff_expires_at, .. },
1912 created_at: at,
1913 backtrace_id: None,
1914 version: Version(2),
1915 } if *at == now_after_first_lock_expiry && *backoff_expires_at == expected_first_timeout_expiry
1916 );
1917 assert_matches!(
1918 execution_log.pending_state,
1919 PendingState::PendingAt(PendingStatePendingAt {
1920 scheduled_at: found_scheduled_by,
1921 last_lock: Some(LockedBy {
1922 executor_id: found_executor_id,
1923 run_id: _,
1924 }),
1925 }) if found_scheduled_by == expected_first_timeout_expiry && found_executor_id == exec_config.executor_id
1926 );
1927 sim_clock.move_time_forward(timeout_duration);
1928 let now_after_first_timeout = sim_clock.now();
1929 debug!(now = %now_after_first_timeout, "Second execution should hang again and result in a permanent timeout");
1930
1931 let mut second_execution_progress = executor
1932 .tick(
1933 db_pool.db_exec_conn().await.unwrap().as_ref(),
1934 now_after_first_timeout,
1935 RunId::generate(),
1936 DEPLOYMENT_ID_DUMMY,
1937 )
1938 .await
1939 .unwrap();
1940 assert_eq!(1, second_execution_progress.executions.len());
1941
1942 sim_clock.move_time_forward(lock_expiry);
1944 let now_after_second_lock_expiry = sim_clock.now();
1946 debug!(now = %now_after_second_lock_expiry, "Expecting the second lock to be expired");
1947 {
1948 let cleanup_progress = executor
1949 .tick(
1950 db_pool.db_exec_conn().await.unwrap().as_ref(),
1951 now_after_second_lock_expiry,
1952 RunId::generate(),
1953 DEPLOYMENT_ID_DUMMY,
1954 )
1955 .await
1956 .unwrap();
1957 assert!(cleanup_progress.executions.is_empty());
1958 }
1959 {
1960 let expired_locks = expired_timers_watcher::tick(
1961 db_pool.connection().await.unwrap().as_ref(),
1962 now_after_second_lock_expiry,
1963 )
1964 .await
1965 .unwrap()
1966 .expired_locks;
1967 assert_eq!(1, expired_locks);
1968 }
1969 assert!(
1970 !second_execution_progress
1971 .executions
1972 .pop()
1973 .unwrap()
1974 .1
1975 .is_finished()
1976 );
1977
1978 drop(db_connection);
1979 drop(executor);
1980 db_close.close().await;
1981 }
1982}