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::notify::{WorkflowApprovalRequiredEvent, WorkflowEvent};
13use crate::plan::lock_plan;
14
15impl WorkflowContext {
16    /// Create a human approval gate.
17    ///
18    /// On first execution, records an approval step and returns
19    /// [`EngineError::ApprovalRequired`] to suspend the run. The engine
20    /// transitions the run to `AwaitingApproval`.
21    ///
22    /// On resume (after a human approved via the API), the approval step
23    /// is replayed: it is marked as `Completed` and execution continues
24    /// past it. Multiple approval gates in the same handler work -- each
25    /// one pauses and resumes independently.
26    ///
27    /// When the config carries an SLA
28    /// ([`ApprovalConfig::with_deadline`](crate::config::ApprovalConfig::with_deadline),
29    /// or the legacy `with_timeout_seconds`), the deadline is persisted on the
30    /// step so the API server's escalator can apply the configured
31    /// [`EscalationPolicy`](crate::config::EscalationPolicy) when it fires. The
32    /// timer is cleared as soon as the gate resolves.
33    ///
34    /// # Errors
35    ///
36    /// Returns [`EngineError::ApprovalRequired`] to pause the run on
37    /// first execution. Returns other [`EngineError`] variants on store
38    /// failures.
39    ///
40    /// # Examples
41    ///
42    /// ```no_run
43    /// use ironflow_engine::context::WorkflowContext;
44    /// use ironflow_engine::config::ApprovalConfig;
45    /// use ironflow_engine::error::EngineError;
46    ///
47    /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
48    /// ctx.approval("deploy-gate", ApprovalConfig::new("Approve deployment?")).await?;
49    /// // Execution continues here after approval
50    /// # Ok(())
51    /// # }
52    /// ```
53    pub async fn approval(
54        &mut self,
55        name: &str,
56        config: ApprovalConfig,
57    ) -> Result<(), EngineError> {
58        // Plan mode: record the gate and continue. Planning must never suspend,
59        // so this comes before the `ApprovalRequired` path below.
60        if let Some(plan) = self.plan().cloned() {
61            self.position += 1;
62            let mut recorder = lock_plan(&plan);
63            if recorder.record(name, StepKind::Approval, &self.workflow_name, None) {
64                recorder.set_last(vec![name.to_string()]);
65            }
66            return Ok(());
67        }
68
69        let position = self.position;
70        self.position += 1;
71
72        // Replay: if this approval step exists from a prior execution,
73        // the run was approved -- mark it completed (if not already) and continue.
74        if let Some(existing) = self.replay_steps.get(&position)
75            && existing.kind == StepKind::Approval
76        {
77            if existing.status.state == StepStatus::AwaitingApproval {
78                self.store
79                    .update_step(
80                        existing.id,
81                        StepUpdate {
82                            status: Some(StepStatus::Completed),
83                            completed_at: Some(Utc::now()),
84                            // An approved gate must never be escalated afterwards.
85                            clear_approval_deadline: true,
86                            ..StepUpdate::default()
87                        },
88                    )
89                    .await?;
90            }
91
92            self.last_step_ids = vec![existing.id];
93            info!(
94                run_id = %self.run_id,
95                step = %name,
96                position,
97                "approval step replayed (approved)"
98            );
99            return Ok(());
100        }
101
102        // Carried over: a human already approved this gate in an earlier
103        // attempt. Record a fresh step in the current attempt so that each
104        // attempt keeps a complete, self-contained DAG, and continue.
105        if let Some(&granted_in) = self.granted_approvals.get(&position) {
106            let trace_id = step_trace_id(self.run_id, name, position);
107            let step = self
108                .store
109                .create_step(NewStep {
110                    run_id: self.run_id,
111                    trace_id,
112                    name: name.to_string(),
113                    kind: StepKind::Approval,
114                    position,
115                    input: Some(to_value(&config)?),
116                    is_error_handler: false,
117                })
118                .await?;
119
120            let now = Utc::now();
121            self.start_step(step.id, now).await?;
122            self.store
123                .update_step(
124                    step.id,
125                    StepUpdate {
126                        status: Some(StepStatus::Completed),
127                        output: Some(json!({"approved_in_attempt": granted_in})),
128                        completed_at: Some(now),
129                        ..StepUpdate::default()
130                    },
131                )
132                .await?;
133
134            self.last_step_ids = vec![step.id];
135            info!(
136                run_id = %self.run_id,
137                step = %name,
138                position,
139                granted_in_attempt = granted_in,
140                attempt = self.attempt,
141                "approval carried over from a previous attempt"
142            );
143            return Ok(());
144        }
145
146        // First execution: create the approval step and suspend.
147        let trace_id = step_trace_id(self.run_id, name, position);
148        let step = self
149            .store
150            .create_step(NewStep {
151                run_id: self.run_id,
152                trace_id,
153                name: name.to_string(),
154                kind: StepKind::Approval,
155                position,
156                input: Some(to_value(&config)?),
157                is_error_handler: false,
158            })
159            .await?;
160
161        self.start_step(step.id, Utc::now()).await?;
162
163        // Transition the step to AwaitingApproval so it reflects the suspended
164        // state on the dashboard, and arm the SLA timer in the same update. The
165        // deadline lives in the store, so it survives an API or worker restart.
166        let deadline_at = config
167            .effective_deadline_secs()
168            .map(|secs| Utc::now() + TimeDelta::seconds(secs as i64));
169
170        self.store
171            .update_step(
172                step.id,
173                StepUpdate {
174                    status: Some(StepStatus::AwaitingApproval),
175                    approval_deadline_at: deadline_at,
176                    approval_stage: Some(0),
177                    approval_assignee: config.assignee().cloned(),
178                    ..StepUpdate::default()
179                },
180            )
181            .await?;
182
183        self.last_step_ids = vec![step.id];
184
185        if let Some(ref bus) = self.event_bus {
186            bus.publish(
187                self.run_id,
188                WorkflowEvent::ApprovalRequired(WorkflowApprovalRequiredEvent {
189                    step_name: name.to_string(),
190                    step_index: position,
191                    approval_id: step.id,
192                }),
193            );
194        }
195
196        Err(EngineError::ApprovalRequired {
197            run_id: self.run_id,
198            step_id: step.id,
199            message: config.message().to_string(),
200        })
201    }
202}