use std::sync::Arc;
use chrono::Utc;
use tracing::{debug, info};
use uuid::Uuid;
use ironflow_store::error::StoreError;
use ironflow_store::models::{Run, RunFilter, RunStatus, RunUpdate, WorkflowPause};
use crate::engine::{Engine, ExecutionMode, chain_root};
use crate::error::EngineError;
use crate::notify::{Event, RunStatusChangedEvent};
const HELD_RUNS_PAGE_SIZE: u32 = 100;
#[derive(Debug, Clone)]
pub struct RunPause {
pub run: Run,
pub paused_descendants: Vec<Uuid>,
}
#[derive(Debug, Clone)]
pub struct RunResume {
pub run: Run,
pub resumed_descendants: Vec<Uuid>,
}
impl Engine {
pub async fn pause_run(&self, run_id: Uuid) -> Result<RunPause, EngineError> {
let run = self.load_pausable_root(run_id).await?;
let from = run.status.state;
if !from.can_transition_to(&RunStatus::Paused) {
return Err(EngineError::Store(StoreError::InvalidTransition {
from,
to: RunStatus::Paused,
}));
}
self.store()
.update_run_status(run_id, RunStatus::Paused)
.await?;
if from == RunStatus::Running {
self.interrupt_running_steps(run_id).await?;
}
self.publish_transition(&run, RunStatus::Paused);
info!(run_id = %run_id, from = %from, "run paused");
let mut paused_descendants = Vec::new();
for descendant in self.store().list_active_descendants(run_id).await? {
if !descendant
.status
.state
.can_transition_to(&RunStatus::Paused)
{
continue;
}
match self
.store()
.update_run_status(descendant.id, RunStatus::Paused)
.await
{
Ok(()) => {}
Err(StoreError::InvalidTransition { .. }) => {
debug!(
run_id = %descendant.id,
"descendant run no longer pausable, skipped"
);
continue;
}
Err(err) => return Err(err.into()),
}
if descendant.status.state == RunStatus::Running {
self.interrupt_running_steps(descendant.id).await?;
}
self.publish_transition(&descendant, RunStatus::Paused);
paused_descendants.push(descendant.id);
}
if !paused_descendants.is_empty() {
info!(
run_id = %run_id,
count = paused_descendants.len(),
"descendant runs paused"
);
}
Ok(RunPause {
run: self.load_run(run_id).await?,
paused_descendants,
})
}
pub async fn resume_paused_run(
self: &Arc<Self>,
run_id: Uuid,
) -> Result<RunResume, EngineError> {
let run = self.load_pausable_root(run_id).await?;
if run.status.state != RunStatus::Paused {
return Err(EngineError::Store(StoreError::InvalidTransition {
from: run.status.state,
to: RunStatus::Pending,
}));
}
let mut resumed_descendants = Vec::new();
let mut queued_descendants = Vec::new();
for descendant in self.store().list_active_descendants(run_id).await? {
if descendant.status.state != RunStatus::Paused {
continue;
}
let target = descendant.resume_status.unwrap_or(RunStatus::Pending);
self.store()
.update_run_status(descendant.id, target)
.await?;
self.publish_transition(&descendant, target);
if target == RunStatus::Pending {
queued_descendants.push(descendant.id);
}
resumed_descendants.push(descendant.id);
}
let target = match run.resume_status {
Some(RunStatus::Running) | None => {
self.interrupt_running_steps(run_id).await?;
RunStatus::Pending
}
Some(RunStatus::Sleeping) if run.scheduled_at.is_some_and(|at| at <= Utc::now()) => {
RunStatus::Pending
}
Some(status) => status,
};
self.store().update_run_status(run_id, target).await?;
self.publish_transition(&run, target);
info!(run_id = %run_id, to = %target, "run resumed");
if self.execution_mode() == ExecutionMode::Local {
if target == RunStatus::Pending {
self.continue_in_flight_execution(run_id).await?;
self.spawn_local_resume(run_id);
} else {
for child_id in queued_descendants {
self.continue_in_flight_execution(child_id).await?;
self.spawn_local_resume(child_id);
}
}
}
Ok(RunResume {
run: self.load_run(run_id).await?,
resumed_descendants,
})
}
async fn continue_in_flight_execution(&self, run_id: Uuid) -> Result<(), EngineError> {
if self.is_executing(run_id) {
self.store()
.update_run_status(run_id, RunStatus::Running)
.await?;
debug!(run_id = %run_id, "execution still in flight, run continues");
}
Ok(())
}
pub async fn pause_workflow(
&self,
workflow_name: &str,
paused_by: Option<Uuid>,
) -> Result<WorkflowPause, EngineError> {
self.require_handler(workflow_name)?;
let pause = self
.store()
.pause_workflow(workflow_name, paused_by)
.await?;
info!(workflow = %workflow_name, "workflow paused");
Ok(pause)
}
pub async fn resume_workflow(
self: &Arc<Self>,
workflow_name: &str,
) -> Result<bool, EngineError> {
self.require_handler(workflow_name)?;
let was_paused = self.store().resume_workflow(workflow_name).await?;
if was_paused {
info!(workflow = %workflow_name, "workflow resumed");
if self.execution_mode() == ExecutionMode::Local {
self.start_held_runs(workflow_name).await?;
}
}
Ok(was_paused)
}
async fn start_held_runs(self: &Arc<Self>, workflow_name: &str) -> Result<(), EngineError> {
let mut held = Vec::new();
let mut page = 1;
loop {
let filter = RunFilter {
workflow_name: Some(workflow_name.to_string()),
status: Some(RunStatus::Pending),
..RunFilter::default()
};
let result = self
.store()
.list_runs(filter, page, HELD_RUNS_PAGE_SIZE)
.await?;
if result.items.is_empty() {
break;
}
held.extend(result.items);
if held.len() as u64 >= result.total {
break;
}
page += 1;
}
let now = Utc::now();
for run in held {
if chain_root(&run).is_some() || run.scheduled_at.is_some_and(|at| at > now) {
continue;
}
self.spawn_local_resume(run.id);
}
Ok(())
}
pub(crate) async fn is_workflow_paused(
&self,
workflow_name: &str,
) -> Result<bool, EngineError> {
let pauses = self.store().list_workflow_pauses().await?;
Ok(pauses
.iter()
.any(|pause| pause.workflow_name == workflow_name))
}
fn require_handler(&self, workflow_name: &str) -> Result<(), EngineError> {
match self.get_handler(workflow_name) {
Some(_) => Ok(()),
None => Err(EngineError::InvalidWorkflow(format!(
"no handler registered for workflow '{workflow_name}'"
))),
}
}
pub(crate) async fn requeue_paused_root(&self, run: &Run) -> Result<(), EngineError> {
let Some(root_id) = chain_root(run) else {
return Ok(());
};
let root = self.load_run(root_id).await?;
let suspended = matches!(
root.resume_status,
Some(RunStatus::AwaitingApproval | RunStatus::Sleeping)
);
if root.status.state != RunStatus::Paused || !suspended {
return Ok(());
}
let update = RunUpdate {
resume_status: Some(RunStatus::Pending),
..RunUpdate::default()
};
self.store().update_run(root_id, update).await?;
info!(run_id = %run.id, root_run_id = %root_id, "paused root run will resume to observe its child");
Ok(())
}
async fn load_pausable_root(&self, run_id: Uuid) -> Result<Run, EngineError> {
let run = self.load_run(run_id).await?;
match chain_root(&run) {
Some(root_run_id) => Err(EngineError::ChildRunNotPausable {
run_id,
root_run_id,
}),
None => Ok(run),
}
}
fn publish_transition(&self, run: &Run, to: RunStatus) {
self.event_publisher()
.publish(Event::RunStatusChanged(RunStatusChangedEvent {
run_id: run.id,
workflow_name: run.workflow_name.clone(),
from: run.status.state,
to,
error: None,
cost_usd: run.cost_usd,
duration_ms: run.duration_ms,
labels: run.labels.clone(),
at: Utc::now(),
}));
}
}