use std::sync::Arc;
use std::time::Duration;
use ironflow_engine::engine::Engine;
use ironflow_engine::escalation::{ApprovalEscalator, EscalationAction};
use tokio::time::interval;
use tokio_util::sync::CancellationToken;
use tracing::{error, info};
pub const DEFAULT_ESCALATOR_INTERVAL: Duration = Duration::from_secs(30);
pub const DEFAULT_ESCALATOR_BATCH_SIZE: u32 = 50;
pub struct Escalator {
escalator: ApprovalEscalator,
interval: Duration,
}
impl Escalator {
pub fn new(engine: Arc<Engine>) -> Self {
Self {
escalator: ApprovalEscalator::new(engine).batch_size(DEFAULT_ESCALATOR_BATCH_SIZE),
interval: DEFAULT_ESCALATOR_INTERVAL,
}
}
pub fn interval(mut self, interval: Duration) -> Self {
self.interval = interval;
self
}
pub fn batch_size(self, batch_size: u32) -> Self {
Self {
escalator: self.escalator.batch_size(batch_size),
interval: self.interval,
}
}
pub async fn run(self, shutdown: CancellationToken) {
let mut ticker = interval(self.interval);
ticker.tick().await;
info!(
interval_secs = self.interval.as_secs(),
"approval escalator started"
);
loop {
tokio::select! {
_ = shutdown.cancelled() => {
info!("approval escalator stopped");
return;
}
_ = ticker.tick() => {
self.tick().await;
}
}
}
}
pub async fn tick(&self) {
let records = match self.escalator.tick().await {
Ok(records) => records,
Err(err) => {
error!(error = %err, "failed to collect expired approval deadlines");
return;
}
};
for record in &records {
if record.action == EscalationAction::Stale {
continue;
}
info!(
run_id = %record.run_id,
step_id = %record.step_id,
stage = record.stage,
action = ?record.action,
reason = %record.reason,
"approval gate escalated"
);
}
}
}
#[cfg(test)]
mod tests {
use chrono::{TimeDelta, Utc};
use ironflow_core::providers::claude::ClaudeCodeProvider;
use ironflow_engine::config::{ApprovalConfig, EscalationPolicy};
use ironflow_store::entities::{
NewRun, NewStep, RunStatus, StepKind, StepStatus, StepUpdate, TriggerKind, step_trace_id,
};
use ironflow_store::memory::InMemoryStore;
use ironflow_store::store::{RunStore, Store};
use serde_json::json;
use std::collections::HashMap;
use uuid::Uuid;
use super::*;
async fn expired_gate(config: ApprovalConfig) -> (Arc<InMemoryStore>, Escalator, Uuid, Uuid) {
let store = Arc::new(InMemoryStore::new());
let store_dyn: Arc<dyn Store> = store.clone();
let engine = Arc::new(Engine::new(store_dyn, Arc::new(ClaudeCodeProvider::new())));
let run = store
.create_run(NewRun {
created_by: None,
workflow_name: "deploy".to_string(),
trigger: TriggerKind::Manual,
payload: json!({}),
max_retries: 0,
handler_version: None,
labels: HashMap::new(),
scheduled_at: None,
idempotency_key: None,
max_cost_usd: None,
})
.await
.expect("create run")
.into_run();
store
.update_run_status(run.id, RunStatus::Running)
.await
.expect("to running");
store
.update_run_status(run.id, RunStatus::AwaitingApproval)
.await
.expect("to awaiting approval");
let step = store
.create_step(NewStep {
run_id: run.id,
trace_id: step_trace_id(run.id, "prod-gate", 0),
name: "prod-gate".to_string(),
kind: StepKind::Approval,
position: 0,
input: Some(serde_json::to_value(&config).expect("serialize config")),
is_error_handler: false,
})
.await
.expect("create step");
store
.update_step(
step.id,
StepUpdate {
status: Some(StepStatus::Running),
..StepUpdate::default()
},
)
.await
.expect("to running");
store
.update_step(
step.id,
StepUpdate {
status: Some(StepStatus::AwaitingApproval),
approval_deadline_at: Some(Utc::now() - TimeDelta::seconds(1)),
..StepUpdate::default()
},
)
.await
.expect("arm an expired timer");
(store, Escalator::new(engine), run.id, step.id)
}
#[tokio::test]
async fn tick_auto_rejects_a_gate_past_its_deadline() {
let config = ApprovalConfig::new("Deploy?")
.with_deadline_secs(60)
.on_timeout(EscalationPolicy::AutoReject);
let (store, escalator, run_id, step_id) = expired_gate(config).await;
escalator.tick().await;
let run = store.get_run(run_id).await.unwrap().unwrap();
assert_eq!(run.status.state, RunStatus::Failed);
assert_eq!(run.error.as_deref(), Some("approval timeout"));
let step = store.get_step(step_id).await.unwrap().unwrap();
assert_eq!(step.status.state, StepStatus::Failed);
assert_eq!(step.error.as_deref(), Some("approval timeout"));
assert!(step.approval_deadline_at.is_none());
}
#[tokio::test]
async fn tick_leaves_a_gate_without_a_deadline_alone() {
let config = ApprovalConfig::new("Deploy?");
let (store, escalator, run_id, step_id) = expired_gate(config).await;
store
.update_step(
step_id,
StepUpdate {
clear_approval_deadline: true,
..StepUpdate::default()
},
)
.await
.expect("clear timer");
escalator.tick().await;
let run = store.get_run(run_id).await.unwrap().unwrap();
assert_eq!(run.status.state, RunStatus::AwaitingApproval);
let step = store.get_step(step_id).await.unwrap().unwrap();
assert_eq!(step.status.state, StepStatus::AwaitingApproval);
}
#[tokio::test]
async fn run_stops_on_shutdown() {
let config = ApprovalConfig::new("Deploy?");
let (_store, escalator, _run_id, _step_id) = expired_gate(config).await;
let shutdown = CancellationToken::new();
shutdown.cancel();
tokio::time::timeout(Duration::from_secs(5), escalator.run(shutdown))
.await
.expect("escalator stopped");
}
}