Skip to main content

ironflow_api/
escalator.rs

1//! Resolution of approval gates that missed their SLA deadline.
2//!
3//! An approval gate can carry a deadline. The deadline lives on the step row, so
4//! it survives an API or worker restart. The [`Escalator`] periodically hands
5//! every expired gate to the engine's
6//! [`ApprovalEscalator`](ironflow_engine::escalation::ApprovalEscalator), which
7//! applies the gate's configured escalation policy.
8//!
9//! Gates without a deadline are never touched, and a deadline fires at most once
10//! even when several API instances run this loop.
11
12use 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
21/// How often expired approval deadlines are collected.
22pub const DEFAULT_ESCALATOR_INTERVAL: Duration = Duration::from_secs(30);
23
24/// How many gates a single tick resolves.
25///
26/// Bounded so that a burst of expiries (a long outage with many open gates)
27/// resorbs progressively instead of holding a long transaction on the steps
28/// table.
29pub const DEFAULT_ESCALATOR_BATCH_SIZE: u32 = 50;
30
31/// Periodic task that escalates approval gates past their deadline.
32///
33/// # Examples
34///
35/// ```no_run
36/// use std::sync::Arc;
37/// use std::time::Duration;
38/// use ironflow_api::escalator::Escalator;
39/// use ironflow_core::providers::claude::ClaudeCodeProvider;
40/// use ironflow_engine::engine::Engine;
41/// use ironflow_store::memory::InMemoryStore;
42/// use ironflow_store::store::Store;
43/// use tokio_util::sync::CancellationToken;
44///
45/// # async fn example() {
46/// let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
47/// let engine = Arc::new(Engine::new(store, Arc::new(ClaudeCodeProvider::new())));
48///
49/// let escalator = Escalator::new(engine).interval(Duration::from_secs(15));
50/// tokio::spawn(escalator.run(CancellationToken::new()));
51/// # }
52/// ```
53pub struct Escalator {
54    escalator: ApprovalEscalator,
55    interval: Duration,
56}
57
58impl Escalator {
59    /// Create an escalator with the default interval and batch size.
60    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    /// Set how often expired deadlines are collected.
68    ///
69    /// Keep it well below the shortest SLA in use, otherwise a gate overshoots
70    /// its deadline by up to one interval.
71    pub fn interval(mut self, interval: Duration) -> Self {
72        self.interval = interval;
73        self
74    }
75
76    /// Set how many gates a single tick resolves.
77    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    /// Run the escalation loop until `shutdown` is cancelled.
85    ///
86    /// Store errors are logged and the loop keeps going: a transient database
87    /// failure must not silently stop escalation.
88    pub async fn run(self, shutdown: CancellationToken) {
89        let mut ticker = interval(self.interval);
90        // The first tick fires immediately; skip it so startup is not a burst.
91        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    /// Escalate one batch of expired gates.
112    ///
113    /// Exposed for tests and for callers that drive the schedule themselves.
114    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            // A stale gate resolved on its own between the claim and the
125            // escalation; that is routine, not news.
126            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    /// Build an escalator over a store holding one run stuck on an expired gate.
159    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        // Drop the timer the fixture armed: this gate carries no SLA.
251        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        // Returns instead of looping forever.
278        tokio::time::timeout(Duration::from_secs(5), escalator.run(shutdown))
279            .await
280            .expect("escalator stopped");
281    }
282}