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;
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    /// # Errors
42    ///
43    /// Returns [`EngineError::ApprovalRequired`] to pause the run on
44    /// first execution. Returns [`EngineError::ApprovalRejected`] when an
45    /// interceptor refuses the gate. Returns other [`EngineError`] variants on
46    /// store failures.
47    ///
48    /// # Examples
49    ///
50    /// ```no_run
51    /// use ironflow_engine::context::WorkflowContext;
52    /// use ironflow_engine::config::ApprovalConfig;
53    /// use ironflow_engine::error::EngineError;
54    ///
55    /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
56    /// ctx.approval("deploy-gate", ApprovalConfig::new("Approve deployment?")).await?;
57    /// // Execution continues here after approval
58    /// # Ok(())
59    /// # }
60    /// ```
61    pub async fn approval(
62        &mut self,
63        name: &str,
64        config: ApprovalConfig,
65    ) -> Result<(), EngineError> {
66        // Plan mode: record the gate and continue. Planning must never suspend,
67        // so this comes before the `ApprovalRequired` path below.
68        if let Some(plan) = self.plan().cloned() {
69            self.position += 1;
70            let mut recorder = lock_plan(&plan);
71            if recorder.record(name, StepKind::Approval, &self.workflow_name, None) {
72                recorder.set_last(vec![name.to_string()]);
73            }
74            return Ok(());
75        }
76
77        let position = self.position;
78        self.position += 1;
79
80        // Replay: if this approval step exists from a prior execution,
81        // the run was approved -- mark it completed (if not already) and continue.
82        if let Some(existing) = self.replay_steps.get(&position)
83            && existing.kind == StepKind::Approval
84        {
85            if existing.status.state == StepStatus::AwaitingApproval {
86                self.store
87                    .update_step(
88                        existing.id,
89                        StepUpdate {
90                            status: Some(StepStatus::Completed),
91                            completed_at: Some(Utc::now()),
92                            // An approved gate must never be escalated afterwards.
93                            clear_approval_deadline: true,
94                            ..StepUpdate::default()
95                        },
96                    )
97                    .await?;
98            }
99
100            self.last_step_ids = vec![existing.id];
101            info!(
102                run_id = %self.run_id,
103                step = %name,
104                position,
105                "approval step replayed (approved)"
106            );
107            return Ok(());
108        }
109
110        // Carried over: a human already approved this gate in an earlier
111        // attempt. Record a fresh step in the current attempt so that each
112        // attempt keeps a complete, self-contained DAG, and continue.
113        if let Some(&granted_in) = self.granted_approvals.get(&position) {
114            let trace_id = step_trace_id(self.run_id, name, position);
115            let step = self
116                .store
117                .create_step(NewStep {
118                    run_id: self.run_id,
119                    trace_id,
120                    name: name.to_string(),
121                    kind: StepKind::Approval,
122                    position,
123                    input: Some(to_value(&config)?),
124                    is_error_handler: false,
125                })
126                .await?;
127
128            let now = Utc::now();
129            self.start_step(step.id, now).await?;
130            self.store
131                .update_step(
132                    step.id,
133                    StepUpdate {
134                        status: Some(StepStatus::Completed),
135                        output: Some(json!({"approved_in_attempt": granted_in})),
136                        completed_at: Some(now),
137                        ..StepUpdate::default()
138                    },
139                )
140                .await?;
141
142            self.last_step_ids = vec![step.id];
143            info!(
144                run_id = %self.run_id,
145                step = %name,
146                position,
147                granted_in_attempt = granted_in,
148                attempt = self.attempt,
149                "approval carried over from a previous attempt"
150            );
151            return Ok(());
152        }
153
154        // An interceptor resolves the gate inline: the run neither suspends nor
155        // waits for a human.
156        if let Some(interceptor) = self.interceptor.clone()
157            && let Some(outcome) = interceptor.intercept_approval(name, &config)
158        {
159            let trace_id = step_trace_id(self.run_id, name, position);
160            let step = self
161                .store
162                .create_step(NewStep {
163                    run_id: self.run_id,
164                    trace_id,
165                    name: name.to_string(),
166                    kind: StepKind::Approval,
167                    position,
168                    input: Some(to_value(&config)?),
169                    is_error_handler: false,
170                })
171                .await?;
172
173            let now = Utc::now();
174            self.start_step(step.id, now).await?;
175            self.last_step_ids = vec![step.id];
176
177            return match outcome {
178                ApprovalOutcome::Approved => {
179                    self.store
180                        .update_step(
181                            step.id,
182                            StepUpdate {
183                                status: Some(StepStatus::Completed),
184                                output: Some(json!({"approved_by": "step-interceptor"})),
185                                completed_at: Some(now),
186                                ..StepUpdate::default()
187                            },
188                        )
189                        .await?;
190                    info!(
191                        run_id = %self.run_id,
192                        step = %name,
193                        position,
194                        "approval granted by the step interceptor"
195                    );
196                    Ok(())
197                }
198                ApprovalOutcome::Rejected { reason } => {
199                    // The step FSM only reaches Rejected from AwaitingApproval.
200                    self.store
201                        .update_step(
202                            step.id,
203                            StepUpdate {
204                                status: Some(StepStatus::AwaitingApproval),
205                                ..StepUpdate::default()
206                            },
207                        )
208                        .await?;
209                    self.store
210                        .update_step(
211                            step.id,
212                            StepUpdate {
213                                status: Some(StepStatus::Rejected),
214                                error: Some(reason.clone()),
215                                completed_at: Some(Utc::now()),
216                                ..StepUpdate::default()
217                            },
218                        )
219                        .await?;
220                    info!(
221                        run_id = %self.run_id,
222                        step = %name,
223                        position,
224                        %reason,
225                        "approval rejected by the step interceptor"
226                    );
227                    Err(EngineError::ApprovalRejected {
228                        run_id: self.run_id,
229                        step_id: step.id,
230                        reason,
231                    })
232                }
233            };
234        }
235
236        // First execution: create the approval step and suspend.
237        let trace_id = step_trace_id(self.run_id, name, position);
238        let step = self
239            .store
240            .create_step(NewStep {
241                run_id: self.run_id,
242                trace_id,
243                name: name.to_string(),
244                kind: StepKind::Approval,
245                position,
246                input: Some(to_value(&config)?),
247                is_error_handler: false,
248            })
249            .await?;
250
251        self.start_step(step.id, Utc::now()).await?;
252
253        // Transition the step to AwaitingApproval so it reflects the suspended
254        // state on the dashboard, and arm the SLA timer in the same update. The
255        // deadline lives in the store, so it survives an API or worker restart.
256        let deadline_at = config
257            .effective_deadline_secs()
258            .map(|secs| Utc::now() + TimeDelta::seconds(secs as i64));
259
260        self.store
261            .update_step(
262                step.id,
263                StepUpdate {
264                    status: Some(StepStatus::AwaitingApproval),
265                    approval_deadline_at: deadline_at,
266                    approval_stage: Some(0),
267                    approval_assignee: config.assignee().cloned(),
268                    ..StepUpdate::default()
269                },
270            )
271            .await?;
272
273        self.last_step_ids = vec![step.id];
274
275        if let Some(ref bus) = self.event_bus {
276            bus.publish(
277                self.run_id,
278                WorkflowEvent::ApprovalRequired(WorkflowApprovalRequiredEvent {
279                    step_name: name.to_string(),
280                    step_index: position,
281                    approval_id: step.id,
282                }),
283            );
284        }
285
286        Err(EngineError::ApprovalRequired {
287            run_id: self.run_id,
288            step_id: step.id,
289            message: config.message().to_string(),
290        })
291    }
292}