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}