Skip to main content

ironflow_engine/context/steps/
approval.rs

1//! Human approval gate for [`WorkflowContext`].
2
3use chrono::{TimeDelta, Utc};
4use serde_json::{Map, Value, json, to_value};
5use tracing::info;
6use uuid::Uuid;
7
8use ironflow_store::error::StoreError;
9use ironflow_store::models::{NewStep, Run, Step, StepKind, StepStatus, StepUpdate, step_trace_id};
10
11use crate::config::ApprovalConfig;
12use crate::context::WorkflowContext;
13use crate::error::EngineError;
14use crate::executor::ApprovalOutcome;
15use crate::notify::{WorkflowApprovalRequiredEvent, WorkflowEvent};
16use crate::plan::lock_plan;
17
18impl WorkflowContext {
19    /// Create a human approval gate.
20    ///
21    /// On first execution, records an approval step and returns
22    /// [`EngineError::ApprovalRequired`] to suspend the run. The engine
23    /// transitions the run to `AwaitingApproval`.
24    ///
25    /// On resume (after a human approved via the API), the approval step
26    /// is replayed: it is marked as `Completed` and execution continues
27    /// past it. Multiple approval gates in the same handler work -- each
28    /// one pauses and resumes independently.
29    ///
30    /// When the config carries an SLA
31    /// ([`ApprovalConfig::with_deadline`](crate::config::ApprovalConfig::with_deadline),
32    /// or the legacy `with_timeout_seconds`), the deadline is persisted on the
33    /// step so the API server's escalator can apply the configured
34    /// [`EscalationPolicy`](crate::config::EscalationPolicy) when it fires. The
35    /// timer is cleared as soon as the gate resolves.
36    ///
37    /// A [`StepInterceptor`](crate::executor::StepInterceptor) wired into the
38    /// context resolves the gate inline instead of suspending: the step is
39    /// recorded, then completed or rejected without waiting for a human. This
40    /// is what [`crate::testing::TestEngine`] uses to run gated handlers end to
41    /// end.
42    ///
43    /// When the config carries approval rules
44    /// ([`ApprovalConfig::with_rule`](crate::config::ApprovalConfig::with_rule)),
45    /// they are evaluated once, when the gate opens, against a context holding
46    /// the `output` of the previous step (the last one of a parallel batch),
47    /// the run `payload` and `labels`, run `metadata` (`run_id`,
48    /// `workflow_name`, `trigger`, `attempt`, `handler_version`) and the
49    /// completed steps of the current attempt under `steps.<name>` (`output`,
50    /// `kind`, `status`). The resulting
51    /// [`ApprovalRequirement`](crate::config::ApprovalRequirement) is stored on
52    /// the step and stays the source of truth on replay and resume: rules are
53    /// never re-evaluated. A config without rules stores no requirement.
54    ///
55    /// # Errors
56    ///
57    /// Returns [`EngineError::ApprovalRequired`] to pause the run on
58    /// first execution. Returns [`EngineError::ApprovalRejected`] when an
59    /// interceptor refuses the gate. Returns other [`EngineError`] variants on
60    /// store failures.
61    ///
62    /// # Examples
63    ///
64    /// ```no_run
65    /// use ironflow_engine::context::WorkflowContext;
66    /// use ironflow_engine::config::ApprovalConfig;
67    /// use ironflow_engine::error::EngineError;
68    ///
69    /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
70    /// ctx.approval("deploy-gate", ApprovalConfig::new("Approve deployment?")).await?;
71    /// // Execution continues here after approval
72    /// # Ok(())
73    /// # }
74    /// ```
75    pub async fn approval(
76        &mut self,
77        name: &str,
78        config: ApprovalConfig,
79    ) -> Result<(), EngineError> {
80        // Plan mode: record the gate and continue. Planning must never suspend,
81        // so this comes before the `ApprovalRequired` path below.
82        if let Some(plan) = self.plan().cloned() {
83            self.position += 1;
84            let mut recorder = lock_plan(&plan);
85            if recorder.record(name, StepKind::Approval, &self.workflow_name, None) {
86                recorder.set_last(vec![name.to_string()]);
87            }
88            return Ok(());
89        }
90
91        let position = self.position;
92        self.position += 1;
93
94        // Replay: if this approval step exists from a prior execution,
95        // the run was approved -- mark it completed (if not already) and continue.
96        if let Some(existing) = self.replay_steps.get(&position)
97            && existing.kind == StepKind::Approval
98        {
99            if existing.status.state == StepStatus::AwaitingApproval {
100                self.store
101                    .update_step(
102                        existing.id,
103                        StepUpdate {
104                            status: Some(StepStatus::Completed),
105                            completed_at: Some(Utc::now()),
106                            // An approved gate must never be escalated afterwards.
107                            clear_approval_deadline: true,
108                            ..StepUpdate::default()
109                        },
110                    )
111                    .await?;
112            }
113
114            self.last_step_ids = vec![existing.id];
115            info!(
116                run_id = %self.run_id,
117                step = %name,
118                position,
119                "approval step replayed (approved)"
120            );
121            return Ok(());
122        }
123
124        // Carried over: a human already approved this gate in an earlier
125        // attempt. Record a fresh step in the current attempt so that each
126        // attempt keeps a complete, self-contained DAG, and continue.
127        if let Some(&granted_in) = self.granted_approvals.get(&position) {
128            let trace_id = step_trace_id(self.run_id, name, position);
129            let step = self
130                .store
131                .create_step(NewStep {
132                    run_id: self.run_id,
133                    trace_id,
134                    name: name.to_string(),
135                    kind: StepKind::Approval,
136                    position,
137                    input: Some(to_value(&config)?),
138                    is_error_handler: false,
139                })
140                .await?;
141
142            let now = Utc::now();
143            self.start_step(step.id, now).await?;
144            self.store
145                .update_step(
146                    step.id,
147                    StepUpdate {
148                        status: Some(StepStatus::Completed),
149                        output: Some(json!({"approved_in_attempt": granted_in})),
150                        completed_at: Some(now),
151                        ..StepUpdate::default()
152                    },
153                )
154                .await?;
155
156            self.last_step_ids = vec![step.id];
157            info!(
158                run_id = %self.run_id,
159                step = %name,
160                position,
161                granted_in_attempt = granted_in,
162                attempt = self.attempt,
163                "approval carried over from a previous attempt"
164            );
165            return Ok(());
166        }
167
168        // Rules are evaluated only when the gate opens, never on replay or
169        // carry-over: the stored requirement is the source of truth from here.
170        let requirement = if config.rules().is_empty() {
171            None
172        } else {
173            let run = self
174                .store
175                .get_run(self.run_id)
176                .await?
177                .ok_or(EngineError::Store(StoreError::RunNotFound(self.run_id)))?;
178            let steps = self.store.list_steps(self.run_id).await?;
179            let ctx = approval_expression_context(
180                &run,
181                &steps,
182                self.attempt,
183                position,
184                self.last_step_ids.last().copied(),
185            );
186            let requirement = config.evaluate_rules(&ctx);
187            info!(
188                run_id = %self.run_id,
189                step = %name,
190                rule_index = ?requirement.rule_index,
191                required_approvers = requirement.required_approvers,
192                approver_groups = ?requirement.approver_groups,
193                "approval rules evaluated"
194            );
195            Some(requirement)
196        };
197
198        // An interceptor resolves the gate inline: the run neither suspends nor
199        // waits for a human.
200        if let Some(interceptor) = self.interceptor.clone()
201            && let Some(outcome) = interceptor.intercept_approval(name, &config)
202        {
203            let trace_id = step_trace_id(self.run_id, name, position);
204            let step = self
205                .store
206                .create_step(NewStep {
207                    run_id: self.run_id,
208                    trace_id,
209                    name: name.to_string(),
210                    kind: StepKind::Approval,
211                    position,
212                    input: Some(to_value(&config)?),
213                    is_error_handler: false,
214                })
215                .await?;
216
217            let now = Utc::now();
218            self.start_step(step.id, now).await?;
219            self.last_step_ids = vec![step.id];
220
221            return match outcome {
222                ApprovalOutcome::Approved => {
223                    self.store
224                        .update_step(
225                            step.id,
226                            StepUpdate {
227                                status: Some(StepStatus::Completed),
228                                output: Some(json!({"approved_by": "step-interceptor"})),
229                                completed_at: Some(now),
230                                approval_requirement: requirement.clone(),
231                                ..StepUpdate::default()
232                            },
233                        )
234                        .await?;
235                    info!(
236                        run_id = %self.run_id,
237                        step = %name,
238                        position,
239                        "approval granted by the step interceptor"
240                    );
241                    Ok(())
242                }
243                ApprovalOutcome::Rejected { reason } => {
244                    // The step FSM only reaches Rejected from AwaitingApproval.
245                    self.store
246                        .update_step(
247                            step.id,
248                            StepUpdate {
249                                status: Some(StepStatus::AwaitingApproval),
250                                approval_requirement: requirement.clone(),
251                                ..StepUpdate::default()
252                            },
253                        )
254                        .await?;
255                    self.store
256                        .update_step(
257                            step.id,
258                            StepUpdate {
259                                status: Some(StepStatus::Rejected),
260                                error: Some(reason.clone()),
261                                completed_at: Some(Utc::now()),
262                                ..StepUpdate::default()
263                            },
264                        )
265                        .await?;
266                    info!(
267                        run_id = %self.run_id,
268                        step = %name,
269                        position,
270                        %reason,
271                        "approval rejected by the step interceptor"
272                    );
273                    Err(EngineError::ApprovalRejected {
274                        run_id: self.run_id,
275                        step_id: step.id,
276                        reason,
277                    })
278                }
279            };
280        }
281
282        // First execution: create the approval step and suspend.
283        let trace_id = step_trace_id(self.run_id, name, position);
284        let step = self
285            .store
286            .create_step(NewStep {
287                run_id: self.run_id,
288                trace_id,
289                name: name.to_string(),
290                kind: StepKind::Approval,
291                position,
292                input: Some(to_value(&config)?),
293                is_error_handler: false,
294            })
295            .await?;
296
297        self.start_step(step.id, Utc::now()).await?;
298
299        // Transition the step to AwaitingApproval so it reflects the suspended
300        // state on the dashboard, and arm the SLA timer in the same update. The
301        // deadline lives in the store, so it survives an API or worker restart.
302        let deadline_at = config
303            .effective_deadline_secs()
304            .map(|secs| Utc::now() + TimeDelta::seconds(secs as i64));
305
306        self.store
307            .update_step(
308                step.id,
309                StepUpdate {
310                    status: Some(StepStatus::AwaitingApproval),
311                    approval_deadline_at: deadline_at,
312                    approval_stage: Some(0),
313                    approval_assignee: config.assignee().cloned(),
314                    approval_requirement: requirement,
315                    ..StepUpdate::default()
316                },
317            )
318            .await?;
319
320        self.last_step_ids = vec![step.id];
321
322        if let Some(ref bus) = self.event_bus {
323            bus.publish(
324                self.run_id,
325                WorkflowEvent::ApprovalRequired(WorkflowApprovalRequiredEvent {
326                    step_name: name.to_string(),
327                    step_index: position,
328                    approval_id: step.id,
329                }),
330            );
331        }
332
333        Err(EngineError::ApprovalRequired {
334            run_id: self.run_id,
335            step_id: step.id,
336            message: config.message().to_string(),
337        })
338    }
339}
340
341/// Build the JSON context approval rule conditions are evaluated against.
342///
343/// - `output`: output of the step `last_step_id` (the previous step, or the
344///   last one of a parallel batch), `null` when absent.
345/// - `payload`, `labels`: taken from the run.
346/// - `metadata`: `run_id`, `workflow_name`, `trigger`, `attempt`,
347///   `handler_version`.
348/// - `steps.<name>`: `output`, `kind` and `status` of every completed step of
349///   `attempt` positioned before `before_position`. When two steps share a
350///   name, the one with the higher position wins.
351pub(crate) fn approval_expression_context(
352    run: &Run,
353    steps: &[Step],
354    attempt: u32,
355    before_position: u32,
356    last_step_id: Option<Uuid>,
357) -> Value {
358    let output = last_step_id
359        .and_then(|id| steps.iter().find(|s| s.id == id))
360        .and_then(|s| s.output.clone())
361        .unwrap_or(Value::Null);
362
363    let mut completed: Vec<&Step> = steps
364        .iter()
365        .filter(|s| {
366            s.attempt == attempt
367                && s.position < before_position
368                && s.status.state == StepStatus::Completed
369        })
370        .collect();
371    completed.sort_by_key(|s| s.position);
372
373    let mut by_name = Map::new();
374    for step in completed {
375        by_name.insert(
376            step.name.clone(),
377            json!({
378                "output": step.output,
379                "kind": step.kind,
380                "status": step.status.state,
381            }),
382        );
383    }
384
385    json!({
386        "output": output,
387        "payload": run.payload,
388        "labels": run.labels,
389        "metadata": {
390            "run_id": run.id,
391            "workflow_name": run.workflow_name,
392            "trigger": run.trigger,
393            "attempt": attempt,
394            "handler_version": run.handler_version,
395        },
396        "steps": by_name,
397    })
398}
399
400#[cfg(test)]
401mod tests {
402    use std::collections::HashMap;
403    use std::slice;
404
405    use ironflow_store::memory::InMemoryStore;
406    use ironflow_store::models::{NewRun, TriggerKind};
407    use ironflow_store::store::RunStore;
408
409    use super::*;
410
411    async fn run_with_labels(store: &InMemoryStore) -> Run {
412        store
413            .create_run(NewRun {
414                created_by: None,
415                workflow_name: "payments".to_string(),
416                trigger: TriggerKind::Manual,
417                payload: json!({"amount": 15000}),
418                max_retries: 0,
419                handler_version: Some("v2".to_string()),
420                labels: HashMap::from([("env".to_string(), "production".to_string())]),
421                scheduled_at: None,
422                idempotency_key: None,
423                max_cost_usd: None,
424            })
425            .await
426            .expect("create run")
427            .into_run()
428    }
429
430    /// Create a step at `position`, optionally driving it to `Completed`.
431    async fn step(
432        store: &InMemoryStore,
433        run_id: Uuid,
434        name: &str,
435        position: u32,
436        output: Option<Value>,
437    ) -> Step {
438        let step = store
439            .create_step(NewStep {
440                run_id,
441                trace_id: step_trace_id(run_id, name, position),
442                name: name.to_string(),
443                kind: StepKind::Shell,
444                position,
445                input: None,
446                is_error_handler: false,
447            })
448            .await
449            .expect("create step");
450        store
451            .update_step(
452                step.id,
453                StepUpdate {
454                    status: Some(StepStatus::Running),
455                    ..StepUpdate::default()
456                },
457            )
458            .await
459            .expect("to running");
460        if output.is_some() {
461            store
462                .update_step(
463                    step.id,
464                    StepUpdate {
465                        status: Some(StepStatus::Completed),
466                        output,
467                        ..StepUpdate::default()
468                    },
469                )
470                .await
471                .expect("to completed");
472        }
473        store.get_step(step.id).await.expect("get").expect("exists")
474    }
475
476    #[tokio::test]
477    async fn context_holds_run_data_and_the_previous_output() {
478        let store = InMemoryStore::new();
479        let run = run_with_labels(&store).await;
480        let risk = step(&store, run.id, "risk", 0, Some(json!({"level": "high"}))).await;
481
482        let steps = slice::from_ref(&risk);
483        let ctx = approval_expression_context(&run, steps, 1, 1, Some(risk.id));
484
485        assert_eq!(ctx["output"], json!({"level": "high"}));
486        assert_eq!(ctx["payload"], json!({"amount": 15000}));
487        assert_eq!(ctx["labels"], json!({"env": "production"}));
488        assert_eq!(ctx["metadata"]["run_id"], json!(run.id));
489        assert_eq!(ctx["metadata"]["workflow_name"], json!("payments"));
490        assert_eq!(ctx["metadata"]["trigger"], json!(run.trigger));
491        assert_eq!(ctx["metadata"]["attempt"], json!(1));
492        assert_eq!(ctx["metadata"]["handler_version"], json!("v2"));
493        assert_eq!(ctx["steps"]["risk"]["output"], json!({"level": "high"}));
494        assert_eq!(ctx["steps"]["risk"]["kind"], json!("shell"));
495        assert_eq!(ctx["steps"]["risk"]["status"], json!("completed"));
496    }
497
498    #[tokio::test]
499    async fn output_is_null_without_a_previous_step() {
500        let store = InMemoryStore::new();
501        let run = run_with_labels(&store).await;
502
503        let ctx = approval_expression_context(&run, &[], 1, 0, None);
504
505        assert_eq!(ctx["output"], Value::Null);
506        assert_eq!(ctx["steps"], json!({}));
507    }
508
509    #[tokio::test]
510    async fn steps_keep_only_completed_steps_of_the_attempt_before_the_gate() {
511        let store = InMemoryStore::new();
512        let run = run_with_labels(&store).await;
513        let done = step(&store, run.id, "done", 0, Some(json!(1))).await;
514        let running = step(&store, run.id, "running", 1, None).await;
515        let later = step(&store, run.id, "later", 5, Some(json!(2))).await;
516        let mut previous_attempt = step(&store, run.id, "old", 2, Some(json!(3))).await;
517        previous_attempt.attempt = 2;
518
519        let steps = [done.clone(), running, later, previous_attempt];
520        let ctx = approval_expression_context(&run, &steps, 1, 3, Some(done.id));
521
522        let names: Vec<&String> = ctx["steps"]
523            .as_object()
524            .expect("steps object")
525            .keys()
526            .collect();
527        assert_eq!(names, vec!["done"]);
528        assert_eq!(ctx["output"], json!(1));
529    }
530
531    #[tokio::test]
532    async fn duplicate_names_keep_the_highest_position() {
533        let store = InMemoryStore::new();
534        let run = run_with_labels(&store).await;
535        let second = step(&store, run.id, "check", 1, Some(json!("second"))).await;
536        let first = step(&store, run.id, "check", 0, Some(json!("first"))).await;
537
538        let ctx = approval_expression_context(&run, &[second, first], 1, 2, None);
539
540        assert_eq!(ctx["steps"]["check"]["output"], json!("second"));
541    }
542}