obeli-sk-executor 0.40.0

Internal package of obelisk
Documentation
use crate::AbortOnDropHandle;
use crate::executor::Append;
use crate::executor::ChildFinishedResponse;
use chrono::{DateTime, Utc};
use concepts::ExecutionFailureKind;
use concepts::ExecutionId;
use concepts::SupportedFunctionReturnValue;
use concepts::storage::AppendRequest;
use concepts::storage::DbConnection;
use concepts::storage::DbErrorWrite;
use concepts::storage::DbPool;
use concepts::storage::ExecutionLog;
use concepts::storage::ExpiredDelay;
use concepts::storage::Unlocked;
use concepts::time::ClockFn;
use concepts::{
    FinishedExecutionFailure,
    storage::{ExecutionRequest, ExpiredTimer},
};
use std::{sync::Arc, time::Duration};
use tracing::Instrument;
use tracing::Level;
use tracing::info_span;
use tracing::warn;
use tracing::{debug, info, instrument};

#[derive(derive_more::Debug)]
pub struct TimersWatcherConfig {
    pub tick_sleep: Duration,
    #[debug(skip)]
    pub clock_fn: Box<dyn ClockFn>,
    // A short duration that will be subtracted from now() so that:
    // a workflow that made progress (is blocked by join set) and is subscribed can either win or
    // the executor can lock it.
    // If watcher wins, the execution will be scheduled based on its backoff settings.
    pub leeway: Duration,
}

#[derive(Debug, PartialEq)]
pub struct TickProgress {
    pub expired_locks: usize,
    pub expired_async_timers: usize,
}

pub fn spawn_new(db_pool: Arc<dyn DbPool>, config: TimersWatcherConfig) -> AbortOnDropHandle {
    let tick_sleep = config.tick_sleep;
    AbortOnDropHandle::new(
        tokio::spawn(
            async move {
                debug!("Spawned expired timers watcher");
                let mut old_err = None;
                loop {
                    let executed_at = config.clock_fn.now() - config.leeway;
                    let res = match db_pool.connection().await {
                        Ok(conn) => tick(conn.as_ref(), executed_at).await,
                        Err(err) => Err(DbErrorWrite::from(err)),
                    };
                    log_err_if_new(res, &mut old_err);
                    tokio::time::sleep(tick_sleep).await;
                }
            }
            .instrument(info_span!(parent: None, "expired_timers_watcher")),
        )
        .abort_handle(),
    )
}

fn log_err_if_new(res: Result<TickProgress, DbErrorWrite>, old_err: &mut Option<DbErrorWrite>) {
    match (res, &old_err) {
        (Ok(_), _) => {
            *old_err = None;
        }
        (Err(err), Some(old)) if err == *old => {}
        (Err(err), _) => {
            warn!("Tick failed: {err:?}");
            *old_err = Some(err);
        }
    }
}

#[cfg(feature = "test")]
pub async fn tick_test(
    db_connection: &dyn DbConnection,
    executed_at: DateTime<Utc>,
) -> Result<TickProgress, DbErrorWrite> {
    tick(db_connection, executed_at).await
}

#[instrument(level = Level::TRACE, skip_all)]
pub(crate) async fn tick(
    db_connection: &dyn DbConnection,
    executed_at: DateTime<Utc>,
) -> Result<TickProgress, DbErrorWrite> {
    let mut expired_locks = 0;
    let mut expired_async_timers = 0;
    for expired_timer in db_connection.get_expired_timers(executed_at).await? {
        match expired_timer {
            ExpiredTimer::Lock(expired) => {
                let execution_id = expired.execution_id.clone();
                let append = if expired.max_retries.is_none()
                    && expired.locked_at_version.0 + 1 < expired.next_version.0
                {
                    // Workflow that made progress is unlocked and immediately available for locking.
                    info!(%execution_id, run_id = %expired.locked_by.run_id,
                        executor_id = %expired.locked_by.executor_id,
                        created_at = %executed_at,
                        "Unlocking workflow execution - {expired:?}");
                    Append {
                        created_at: executed_at,
                        primary_event: AppendRequest {
                            created_at: executed_at,
                            event: ExecutionRequest::Unlocked(Unlocked {
                                unlocked_at: executed_at,
                                reason: "made progress".into(),
                            }),
                        },
                        execution_id: execution_id.clone(),
                        version: expired.next_version,
                        child_finished: None,
                    }
                } else if let Some(duration) = ExecutionLog::can_be_retried_after(
                    expired.intermittent_event_count + 1,
                    expired.max_retries,
                    expired.retry_exp_backoff,
                ) {
                    let backoff_expires_at = executed_at + duration;
                    info!(%execution_id, run_id = %expired.locked_by.run_id, executor_id = %expired.locked_by.executor_id,
                        created_at = %executed_at,
                        "Retrying execution with expired lock after {duration:?} at {backoff_expires_at} - {expired:?}");
                    Append {
                        created_at: executed_at,
                        primary_event: AppendRequest {
                            created_at: executed_at,
                            event: ExecutionRequest::TemporarilyTimedOut {
                                backoff_expires_at,
                                http_client_traces: None,
                            },
                        },
                        execution_id: execution_id.clone(),
                        version: expired.next_version,
                        child_finished: None,
                    }
                } else {
                    info!(%execution_id, run_id = %expired.locked_by.run_id, executor_id = %expired.locked_by.executor_id,
                        created_at = %executed_at,
                        "Marking execution with expired lock as permanently timed out - {expired:?}");
                    let finished_exec_result =
                        SupportedFunctionReturnValue::ExecutionFailure(FinishedExecutionFailure {
                            kind: ExecutionFailureKind::TimedOut,
                            reason: None,
                            detail: None,
                        });
                    let parent = if let ExecutionId::Derived(derived) = &execution_id {
                        Some(derived.split_to_parts())
                    } else {
                        None
                    };
                    let child_finished = parent.map(|(parent_execution_id, parent_join_set)| {
                        ChildFinishedResponse {
                            parent_execution_id,
                            parent_join_set,
                            result: finished_exec_result.clone(),
                        }
                    });
                    Append {
                        created_at: executed_at,
                        primary_event: AppendRequest {
                            created_at: executed_at,
                            event: ExecutionRequest::Finished {
                                retval: finished_exec_result,
                                http_client_traces: None,
                            },
                        },
                        execution_id: execution_id.clone(),
                        version: expired.next_version,
                        child_finished,
                    }
                };
                let res = append.append(db_connection).await;
                if let Err(err) = res {
                    debug!(%execution_id, "Failed to update expired lock - {err:?}");
                } else {
                    expired_locks += 1;
                }
            }
            ExpiredTimer::Delay(ExpiredDelay {
                execution_id,
                join_set_id,
                delay_id,
            }) => {
                debug!(%execution_id, %join_set_id, %delay_id, created_at = %executed_at, "Appending delay response");
                let res = db_connection
                    .append_delay_response(
                        executed_at,
                        execution_id.clone(),
                        join_set_id,
                        delay_id,
                        Ok(()),
                    )
                    .await;
                if let Err(err) = res {
                    debug!(%execution_id, "Failed to append delay response - {err:?}");
                } else {
                    expired_async_timers += 1;
                }
            }
        }
    }
    Ok(TickProgress {
        expired_locks,
        expired_async_timers,
    })
}