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