use std::sync::{Arc, Mutex};
use chrono::Utc;
use crate::EngineError;
use crate::durability::ActiveWorkflowRecoverySeam;
use super::api::Engine;
use super::startup::{
StartupRecoveryContext, recover_active_workflows_on_startup, recover_timers_on_startup,
};
pub(crate) struct DeferredStartupRecovery {
pub(crate) recovery: Option<Arc<dyn ActiveWorkflowRecoverySeam>>,
pub(crate) bootstrap_schedule_coordinator: bool,
}
pub(super) enum DeferredRecoverySlot {
NotDeferred,
Pending(DeferredStartupRecovery),
WorkflowsRecovered,
Completed,
}
impl DeferredRecoverySlot {
pub(super) fn from_build(deferred: Option<DeferredStartupRecovery>) -> Mutex<Self> {
Mutex::new(deferred.map_or(Self::NotDeferred, Self::Pending))
}
}
impl Engine {
pub async fn run_startup_recovery(&self) -> Result<(), EngineError> {
self.recover_workflows_on_startup().await?;
self.run_startup_catchup().await
}
pub async fn recover_workflows_on_startup(&self) -> Result<(), EngineError> {
let deferred = {
let mut slot = self
.deferred_startup_recovery
.lock()
.map_err(|_| EngineError::StartupRecoverySlotPoisoned)?;
match std::mem::replace(&mut *slot, DeferredRecoverySlot::WorkflowsRecovered) {
DeferredRecoverySlot::Pending(deferred) => deferred,
DeferredRecoverySlot::NotDeferred => {
*slot = DeferredRecoverySlot::NotDeferred;
return Err(EngineError::StartupRecoveryNotDeferred);
}
DeferredRecoverySlot::WorkflowsRecovered => {
*slot = DeferredRecoverySlot::WorkflowsRecovered;
return Err(EngineError::StartupRecoveryAlreadyRan);
}
DeferredRecoverySlot::Completed => {
*slot = DeferredRecoverySlot::Completed;
return Err(EngineError::StartupRecoveryAlreadyRan);
}
}
};
recover_active_workflows_on_startup(StartupRecoveryContext {
store: Arc::clone(&self.store),
visibility_store: Arc::clone(&self.visibility_store),
runtime: Arc::clone(&self.runtime),
catalog: Arc::clone(&self.catalog),
registry: Arc::clone(&self.registry),
supervision: Arc::clone(&self.supervision),
recovery: deferred.recovery,
search_attribute_schema: Arc::clone(&self.search_attribute_schema),
bootstrap_schedule_coordinator: deferred.bootstrap_schedule_coordinator,
})
.await
}
pub async fn run_startup_catchup(&self) -> Result<(), EngineError> {
{
let mut slot = self
.deferred_startup_recovery
.lock()
.map_err(|_| EngineError::StartupRecoverySlotPoisoned)?;
match std::mem::replace(&mut *slot, DeferredRecoverySlot::Completed) {
DeferredRecoverySlot::WorkflowsRecovered => {}
DeferredRecoverySlot::Pending(deferred) => {
*slot = DeferredRecoverySlot::Pending(deferred);
return Err(EngineError::StartupCatchupBeforeWorkflowRecovery);
}
DeferredRecoverySlot::NotDeferred => {
*slot = DeferredRecoverySlot::NotDeferred;
return Err(EngineError::StartupRecoveryNotDeferred);
}
DeferredRecoverySlot::Completed => {
return Err(EngineError::StartupRecoveryAlreadyRan);
}
}
}
recover_timers_on_startup(self.runtime.nif_state(), Arc::clone(&self.store)).await?;
self.catchup_schedule_coordinator().await?;
self.recover_schedules_on_startup(Utc::now()).await?;
Ok(())
}
}
#[cfg(test)]
#[path = "startup_deferred_tests.rs"]
mod tests;