ironflow_api/
escalator.rs1use std::sync::Arc;
13use std::time::Duration;
14
15use ironflow_engine::engine::Engine;
16use ironflow_engine::escalation::{ApprovalEscalator, EscalationAction};
17use tokio::time::interval;
18use tokio_util::sync::CancellationToken;
19use tracing::{error, info};
20
21pub const DEFAULT_ESCALATOR_INTERVAL: Duration = Duration::from_secs(30);
23
24pub const DEFAULT_ESCALATOR_BATCH_SIZE: u32 = 50;
30
31pub struct Escalator {
54 escalator: ApprovalEscalator,
55 interval: Duration,
56}
57
58impl Escalator {
59 pub fn new(engine: Arc<Engine>) -> Self {
61 Self {
62 escalator: ApprovalEscalator::new(engine).batch_size(DEFAULT_ESCALATOR_BATCH_SIZE),
63 interval: DEFAULT_ESCALATOR_INTERVAL,
64 }
65 }
66
67 pub fn interval(mut self, interval: Duration) -> Self {
72 self.interval = interval;
73 self
74 }
75
76 pub fn batch_size(self, batch_size: u32) -> Self {
78 Self {
79 escalator: self.escalator.batch_size(batch_size),
80 interval: self.interval,
81 }
82 }
83
84 pub async fn run(self, shutdown: CancellationToken) {
89 let mut ticker = interval(self.interval);
90 ticker.tick().await;
92
93 info!(
94 interval_secs = self.interval.as_secs(),
95 "approval escalator started"
96 );
97
98 loop {
99 tokio::select! {
100 _ = shutdown.cancelled() => {
101 info!("approval escalator stopped");
102 return;
103 }
104 _ = ticker.tick() => {
105 self.tick().await;
106 }
107 }
108 }
109 }
110
111 pub async fn tick(&self) {
115 let records = match self.escalator.tick().await {
116 Ok(records) => records,
117 Err(err) => {
118 error!(error = %err, "failed to collect expired approval deadlines");
119 return;
120 }
121 };
122
123 for record in &records {
124 if record.action == EscalationAction::Stale {
127 continue;
128 }
129
130 info!(
131 run_id = %record.run_id,
132 step_id = %record.step_id,
133 stage = record.stage,
134 action = ?record.action,
135 reason = %record.reason,
136 "approval gate escalated"
137 );
138 }
139 }
140}
141
142#[cfg(test)]
143mod tests {
144 use chrono::{TimeDelta, Utc};
145 use ironflow_core::providers::claude::ClaudeCodeProvider;
146 use ironflow_engine::config::{ApprovalConfig, EscalationPolicy};
147 use ironflow_store::entities::{
148 NewRun, NewStep, RunStatus, StepKind, StepStatus, StepUpdate, TriggerKind, step_trace_id,
149 };
150 use ironflow_store::memory::InMemoryStore;
151 use ironflow_store::store::{RunStore, Store};
152 use serde_json::json;
153 use std::collections::HashMap;
154 use uuid::Uuid;
155
156 use super::*;
157
158 async fn expired_gate(config: ApprovalConfig) -> (Arc<InMemoryStore>, Escalator, Uuid, Uuid) {
160 let store = Arc::new(InMemoryStore::new());
161 let store_dyn: Arc<dyn Store> = store.clone();
162 let engine = Arc::new(Engine::new(store_dyn, Arc::new(ClaudeCodeProvider::new())));
163
164 let run = store
165 .create_run(NewRun {
166 created_by: None,
167 workflow_name: "deploy".to_string(),
168 trigger: TriggerKind::Manual,
169 payload: json!({}),
170 max_retries: 0,
171 handler_version: None,
172 labels: HashMap::new(),
173 scheduled_at: None,
174 idempotency_key: None,
175 max_cost_usd: None,
176 })
177 .await
178 .expect("create run")
179 .into_run();
180 store
181 .update_run_status(run.id, RunStatus::Running)
182 .await
183 .expect("to running");
184 store
185 .update_run_status(run.id, RunStatus::AwaitingApproval)
186 .await
187 .expect("to awaiting approval");
188
189 let step = store
190 .create_step(NewStep {
191 run_id: run.id,
192 trace_id: step_trace_id(run.id, "prod-gate", 0),
193 name: "prod-gate".to_string(),
194 kind: StepKind::Approval,
195 position: 0,
196 input: Some(serde_json::to_value(&config).expect("serialize config")),
197 is_error_handler: false,
198 })
199 .await
200 .expect("create step");
201 store
202 .update_step(
203 step.id,
204 StepUpdate {
205 status: Some(StepStatus::Running),
206 ..StepUpdate::default()
207 },
208 )
209 .await
210 .expect("to running");
211 store
212 .update_step(
213 step.id,
214 StepUpdate {
215 status: Some(StepStatus::AwaitingApproval),
216 approval_deadline_at: Some(Utc::now() - TimeDelta::seconds(1)),
217 ..StepUpdate::default()
218 },
219 )
220 .await
221 .expect("arm an expired timer");
222
223 (store, Escalator::new(engine), run.id, step.id)
224 }
225
226 #[tokio::test]
227 async fn tick_auto_rejects_a_gate_past_its_deadline() {
228 let config = ApprovalConfig::new("Deploy?")
229 .with_deadline_secs(60)
230 .on_timeout(EscalationPolicy::AutoReject);
231 let (store, escalator, run_id, step_id) = expired_gate(config).await;
232
233 escalator.tick().await;
234
235 let run = store.get_run(run_id).await.unwrap().unwrap();
236 assert_eq!(run.status.state, RunStatus::Failed);
237 assert_eq!(run.error.as_deref(), Some("approval timeout"));
238
239 let step = store.get_step(step_id).await.unwrap().unwrap();
240 assert_eq!(step.status.state, StepStatus::Failed);
241 assert_eq!(step.error.as_deref(), Some("approval timeout"));
242 assert!(step.approval_deadline_at.is_none());
243 }
244
245 #[tokio::test]
246 async fn tick_leaves_a_gate_without_a_deadline_alone() {
247 let config = ApprovalConfig::new("Deploy?");
248 let (store, escalator, run_id, step_id) = expired_gate(config).await;
249
250 store
252 .update_step(
253 step_id,
254 StepUpdate {
255 clear_approval_deadline: true,
256 ..StepUpdate::default()
257 },
258 )
259 .await
260 .expect("clear timer");
261
262 escalator.tick().await;
263
264 let run = store.get_run(run_id).await.unwrap().unwrap();
265 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
266 let step = store.get_step(step_id).await.unwrap().unwrap();
267 assert_eq!(step.status.state, StepStatus::AwaitingApproval);
268 }
269
270 #[tokio::test]
271 async fn run_stops_on_shutdown() {
272 let config = ApprovalConfig::new("Deploy?");
273 let (_store, escalator, _run_id, _step_id) = expired_gate(config).await;
274 let shutdown = CancellationToken::new();
275 shutdown.cancel();
276
277 tokio::time::timeout(Duration::from_secs(5), escalator.run(shutdown))
279 .await
280 .expect("escalator stopped");
281 }
282}