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