Skip to main content

ironflow_engine/context/steps/
approval.rs

1//! Human approval gate for [`WorkflowContext`].
2
3use chrono::{TimeDelta, Utc};
4use serde_json::{json, to_value};
5use tracing::info;
6
7use ironflow_store::models::{NewStep, StepKind, StepStatus, StepUpdate, step_trace_id};
8
9use crate::config::{ApprovalConfig, Approvers};
10use crate::context::WorkflowContext;
11use crate::error::EngineError;
12use crate::executor::ApprovalOutcome;
13use crate::notify::{WorkflowApprovalRequiredEvent, WorkflowEvent};
14use crate::plan::lock_plan;
15
16impl WorkflowContext {
17    /// Create a human approval gate.
18    ///
19    /// On first execution, records an approval step and returns
20    /// [`EngineError::ApprovalRequired`] to suspend the run. The engine
21    /// transitions the run to `AwaitingApproval`.
22    ///
23    /// On resume (after a human approved via the API), the approval step
24    /// is replayed: it is marked as `Completed` and execution continues
25    /// past it. Multiple approval gates in the same handler work -- each
26    /// one pauses and resumes independently.
27    ///
28    /// When the config carries an SLA
29    /// ([`ApprovalConfig::with_deadline`](crate::config::ApprovalConfig::with_deadline),
30    /// or the legacy `with_timeout_seconds`), the deadline is persisted on the
31    /// step so the API server's escalator can apply the configured
32    /// [`EscalationPolicy`](crate::config::EscalationPolicy) when it fires. The
33    /// timer is cleared as soon as the gate resolves.
34    ///
35    /// A [`StepInterceptor`](crate::executor::StepInterceptor) wired into the
36    /// context resolves the gate inline instead of suspending: the step is
37    /// recorded, then completed or rejected without waiting for a human. This
38    /// is what [`crate::testing::TestEngine`] uses to run gated handlers end to
39    /// end.
40    ///
41    /// When the config carries [`Approvers`]
42    /// ([`ApprovalConfig::requiring`]),
43    /// they are recorded on the step as an
44    /// [`ApprovalRequirement`](crate::config::ApprovalRequirement) when the gate
45    /// opens. That record stays the source of truth on replay and resume: the
46    /// approvers the handler computes on a later execution are ignored. A config
47    /// without approvers stores no requirement.
48    ///
49    /// # Errors
50    ///
51    /// Returns [`EngineError::ApprovalRequired`] to pause the run on
52    /// first execution. Returns [`EngineError::ApprovalRejected`] when an
53    /// interceptor refuses the gate. Returns other [`EngineError`] variants on
54    /// store failures.
55    ///
56    /// # Examples
57    ///
58    /// ```no_run
59    /// use ironflow_engine::context::WorkflowContext;
60    /// use ironflow_engine::config::{ApprovalConfig, Approvers};
61    /// use ironflow_engine::error::EngineError;
62    /// use serde::Deserialize;
63    ///
64    /// #[derive(Deserialize)]
65    /// struct Payment {
66    ///     amount: u64,
67    /// }
68    ///
69    /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
70    /// let payment: Payment = ctx.input().await?;
71    /// let approvers = match payment.amount {
72    ///     a if a > 10_000 => Approvers::at_least(2).from_groups(["finance"]).because("amount > 10k"),
73    ///     _ => Approvers::any(),
74    /// };
75    /// ctx.approval("payment-gate", ApprovalConfig::new("Release the payment?").requiring(approvers))
76    ///     .await?;
77    /// // Execution continues here after approval
78    /// # Ok(())
79    /// # }
80    /// ```
81    pub async fn approval(
82        &mut self,
83        name: &str,
84        config: ApprovalConfig,
85    ) -> Result<(), EngineError> {
86        // Plan mode: record the gate and continue. Planning must never suspend,
87        // so this comes before the `ApprovalRequired` path below.
88        if let Some(plan) = self.plan().cloned() {
89            self.position += 1;
90            let mut recorder = lock_plan(&plan);
91            if recorder.record(name, StepKind::Approval, &self.workflow_name, None) {
92                recorder.set_last(vec![name.to_string()]);
93            }
94            return Ok(());
95        }
96
97        let position = self.position;
98        self.position += 1;
99
100        // Replay: if this approval step exists from a prior execution,
101        // the run was approved -- mark it completed (if not already) and continue.
102        if let Some(existing) = self.replay_steps.get(&position)
103            && existing.kind == StepKind::Approval
104        {
105            if existing.status.state == StepStatus::AwaitingApproval {
106                self.store
107                    .update_step(
108                        existing.id,
109                        StepUpdate {
110                            status: Some(StepStatus::Completed),
111                            completed_at: Some(Utc::now()),
112                            // An approved gate must never be escalated afterwards.
113                            clear_approval_deadline: true,
114                            ..StepUpdate::default()
115                        },
116                    )
117                    .await?;
118            }
119
120            self.last_step_ids = vec![existing.id];
121            info!(
122                run_id = %self.run_id,
123                step = %name,
124                position,
125                "approval step replayed (approved)"
126            );
127            return Ok(());
128        }
129
130        // Carried over: a human already approved this gate in an earlier
131        // attempt. Record a fresh step in the current attempt so that each
132        // attempt keeps a complete, self-contained DAG, and continue.
133        if let Some(&granted_in) = self.granted_approvals.get(&position) {
134            let trace_id = step_trace_id(self.run_id, name, position);
135            let step = self
136                .store
137                .create_step(NewStep {
138                    run_id: self.run_id,
139                    trace_id,
140                    name: name.to_string(),
141                    kind: StepKind::Approval,
142                    position,
143                    input: Some(to_value(&config)?),
144                    is_error_handler: false,
145                })
146                .await?;
147
148            let now = Utc::now();
149            self.start_step(step.id, now).await?;
150            self.store
151                .update_step(
152                    step.id,
153                    StepUpdate {
154                        status: Some(StepStatus::Completed),
155                        output: Some(json!({"approved_in_attempt": granted_in})),
156                        completed_at: Some(now),
157                        ..StepUpdate::default()
158                    },
159                )
160                .await?;
161
162            self.last_step_ids = vec![step.id];
163            info!(
164                run_id = %self.run_id,
165                step = %name,
166                position,
167                granted_in_attempt = granted_in,
168                attempt = self.attempt,
169                "approval carried over from a previous attempt"
170            );
171            return Ok(());
172        }
173
174        // The approvers are recorded only when the gate opens, never on replay
175        // or carry-over: the stored requirement is the source of truth from here.
176        let requirement = config.approvers().map(Approvers::to_requirement);
177        if let Some(requirement) = &requirement {
178            info!(
179                run_id = %self.run_id,
180                step = %name,
181                reason = ?requirement.reason,
182                required_approvers = requirement.required_approvers,
183                approver_groups = ?requirement.approver_groups,
184                "approval requirement recorded"
185            );
186        }
187
188        // An interceptor resolves the gate inline: the run neither suspends nor
189        // waits for a human.
190        if let Some(interceptor) = self.interceptor.clone()
191            && let Some(outcome) = interceptor.intercept_approval(name, &config)
192        {
193            let trace_id = step_trace_id(self.run_id, name, position);
194            let step = self
195                .store
196                .create_step(NewStep {
197                    run_id: self.run_id,
198                    trace_id,
199                    name: name.to_string(),
200                    kind: StepKind::Approval,
201                    position,
202                    input: Some(to_value(&config)?),
203                    is_error_handler: false,
204                })
205                .await?;
206
207            let now = Utc::now();
208            self.start_step(step.id, now).await?;
209            self.last_step_ids = vec![step.id];
210
211            return match outcome {
212                ApprovalOutcome::Approved => {
213                    self.store
214                        .update_step(
215                            step.id,
216                            StepUpdate {
217                                status: Some(StepStatus::Completed),
218                                output: Some(json!({"approved_by": "step-interceptor"})),
219                                completed_at: Some(now),
220                                approval_requirement: requirement.clone(),
221                                ..StepUpdate::default()
222                            },
223                        )
224                        .await?;
225                    info!(
226                        run_id = %self.run_id,
227                        step = %name,
228                        position,
229                        "approval granted by the step interceptor"
230                    );
231                    Ok(())
232                }
233                ApprovalOutcome::Rejected { reason } => {
234                    // The step FSM only reaches Rejected from AwaitingApproval.
235                    self.store
236                        .update_step(
237                            step.id,
238                            StepUpdate {
239                                status: Some(StepStatus::AwaitingApproval),
240                                approval_requirement: requirement.clone(),
241                                ..StepUpdate::default()
242                            },
243                        )
244                        .await?;
245                    self.store
246                        .update_step(
247                            step.id,
248                            StepUpdate {
249                                status: Some(StepStatus::Rejected),
250                                error: Some(reason.clone()),
251                                completed_at: Some(Utc::now()),
252                                ..StepUpdate::default()
253                            },
254                        )
255                        .await?;
256                    info!(
257                        run_id = %self.run_id,
258                        step = %name,
259                        position,
260                        %reason,
261                        "approval rejected by the step interceptor"
262                    );
263                    Err(EngineError::ApprovalRejected {
264                        run_id: self.run_id,
265                        step_id: step.id,
266                        reason,
267                    })
268                }
269            };
270        }
271
272        // First execution: create the approval step and suspend.
273        let trace_id = step_trace_id(self.run_id, name, position);
274        let step = self
275            .store
276            .create_step(NewStep {
277                run_id: self.run_id,
278                trace_id,
279                name: name.to_string(),
280                kind: StepKind::Approval,
281                position,
282                input: Some(to_value(&config)?),
283                is_error_handler: false,
284            })
285            .await?;
286
287        self.start_step(step.id, Utc::now()).await?;
288
289        // Transition the step to AwaitingApproval so it reflects the suspended
290        // state on the dashboard, and arm the SLA timer in the same update. The
291        // deadline lives in the store, so it survives an API or worker restart.
292        let deadline_at = config
293            .effective_deadline_secs()
294            .map(|secs| Utc::now() + TimeDelta::seconds(secs as i64));
295
296        self.store
297            .update_step(
298                step.id,
299                StepUpdate {
300                    status: Some(StepStatus::AwaitingApproval),
301                    approval_deadline_at: deadline_at,
302                    approval_stage: Some(0),
303                    approval_assignee: config.assignee().cloned(),
304                    approval_requirement: requirement,
305                    ..StepUpdate::default()
306                },
307            )
308            .await?;
309
310        self.last_step_ids = vec![step.id];
311
312        if let Some(ref bus) = self.event_bus {
313            bus.publish(
314                self.run_id,
315                WorkflowEvent::ApprovalRequired(WorkflowApprovalRequiredEvent {
316                    step_name: name.to_string(),
317                    step_index: position,
318                    approval_id: step.id,
319                }),
320            );
321        }
322
323        Err(EngineError::ApprovalRequired {
324            run_id: self.run_id,
325            step_id: step.id,
326            message: config.message().to_string(),
327        })
328    }
329}