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}