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