Skip to main content

obeli_sk_executor/
executor.rs

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>, // Used for obtaining current time when the execution finishes.
50    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    // Signals close() -> event loop to shut down
84    #[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    // Signals close() -> workers to shut down
92    executor_closing_signal_sender: tokio::sync::watch::Sender<bool>,
93    /// Tracks the number of worker tasks currently in-flight.
94    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    // Signals close() -> workers to shut down
102    executor_closing_signal_sender: tokio::sync::watch::Sender<bool>,
103    /// Tracks the number of worker tasks currently in-flight.
104    worker_count_rx: tokio::sync::watch::Receiver<usize>,
105}
106
107impl ExecutorTaskHandle {
108    /// Shut down the executor task, return worker tasks handle untouched.
109    #[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    /// Signal to worker tasks, then wait for all tasks to finish.
144    #[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)] // Stored permits limit semaphores until dropped.
201struct 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    /// Spawn new tokio worker locking executions and starting worker tasks.
277    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, // fetch expiring before now
465                            ffqns.clone(),
466                            executed_at, // created at
467                            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, // fetch expiring before now
481                            ffqns.clone(),
482                            executed_at, // created at
483                            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, // pending_at_or_sooner
497                            &self.config.component_id,
498                            deployment_id,
499                            executed_at, // created at
500                            self.config.executor_id,
501                            lock_expires_at,
502                            run_id,
503                            self.config.retry_config,
504                        )
505                        .await?
506                }
507            };
508            // Drop permits if too many were allocated.
509            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    /// Map the `WorkerError` to an optional append event
632    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                // Interpret returned `err` variant as a retry request, unless it is a permanent variant
651                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(); // Override with err's reason if no information is lost.
711
712                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: _, // reason_generic contains trap_kind + reason
769                        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, // Since self.primary_event is a finished event, the version will remain the same.
916                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                    // Ignore unsuccessful `Unlocked`
932                }
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        // Create an execution
1246        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        // execute!
1263        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        // check that DB contains Created and Locked events.
1267        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        // noop until `retry_exp_backoff` expires
1385        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        // tick again to finish the execution
1397        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        // executor does not append anything, this should have been written by the worker:
1611        {
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        // execute the child
1680        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            // In case of timeout, let the timers watcher handle it
1700            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), // sleep more than allowed by the lock expiry
1822            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        // Create an execution
1832        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        // Started hanging, wait for lock expiry.
1871        sim_clock.move_time_forward(lock_expiry);
1872        // cleanup should be called
1873        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        // Started hanging, wait for lock expiry.
1943        sim_clock.move_time_forward(lock_expiry);
1944        // cleanup should be called
1945        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}