1use chrono::{TimeDelta, Utc};
4use serde_json::{Map, Value, json, to_value};
5use tracing::info;
6use uuid::Uuid;
7
8use ironflow_store::error::StoreError;
9use ironflow_store::models::{NewStep, Run, Step, StepKind, StepStatus, StepUpdate, step_trace_id};
10
11use crate::config::ApprovalConfig;
12use crate::context::WorkflowContext;
13use crate::error::EngineError;
14use crate::executor::ApprovalOutcome;
15use crate::notify::{WorkflowApprovalRequiredEvent, WorkflowEvent};
16use crate::plan::lock_plan;
17
18impl WorkflowContext {
19 pub async fn approval(
76 &mut self,
77 name: &str,
78 config: ApprovalConfig,
79 ) -> Result<(), EngineError> {
80 if let Some(plan) = self.plan().cloned() {
83 self.position += 1;
84 let mut recorder = lock_plan(&plan);
85 if recorder.record(name, StepKind::Approval, &self.workflow_name, None) {
86 recorder.set_last(vec![name.to_string()]);
87 }
88 return Ok(());
89 }
90
91 let position = self.position;
92 self.position += 1;
93
94 if let Some(existing) = self.replay_steps.get(&position)
97 && existing.kind == StepKind::Approval
98 {
99 if existing.status.state == StepStatus::AwaitingApproval {
100 self.store
101 .update_step(
102 existing.id,
103 StepUpdate {
104 status: Some(StepStatus::Completed),
105 completed_at: Some(Utc::now()),
106 clear_approval_deadline: true,
108 ..StepUpdate::default()
109 },
110 )
111 .await?;
112 }
113
114 self.last_step_ids = vec![existing.id];
115 info!(
116 run_id = %self.run_id,
117 step = %name,
118 position,
119 "approval step replayed (approved)"
120 );
121 return Ok(());
122 }
123
124 if let Some(&granted_in) = self.granted_approvals.get(&position) {
128 let trace_id = step_trace_id(self.run_id, name, position);
129 let step = self
130 .store
131 .create_step(NewStep {
132 run_id: self.run_id,
133 trace_id,
134 name: name.to_string(),
135 kind: StepKind::Approval,
136 position,
137 input: Some(to_value(&config)?),
138 is_error_handler: false,
139 })
140 .await?;
141
142 let now = Utc::now();
143 self.start_step(step.id, now).await?;
144 self.store
145 .update_step(
146 step.id,
147 StepUpdate {
148 status: Some(StepStatus::Completed),
149 output: Some(json!({"approved_in_attempt": granted_in})),
150 completed_at: Some(now),
151 ..StepUpdate::default()
152 },
153 )
154 .await?;
155
156 self.last_step_ids = vec![step.id];
157 info!(
158 run_id = %self.run_id,
159 step = %name,
160 position,
161 granted_in_attempt = granted_in,
162 attempt = self.attempt,
163 "approval carried over from a previous attempt"
164 );
165 return Ok(());
166 }
167
168 let requirement = if config.rules().is_empty() {
171 None
172 } else {
173 let run = self
174 .store
175 .get_run(self.run_id)
176 .await?
177 .ok_or(EngineError::Store(StoreError::RunNotFound(self.run_id)))?;
178 let steps = self.store.list_steps(self.run_id).await?;
179 let ctx = approval_expression_context(
180 &run,
181 &steps,
182 self.attempt,
183 position,
184 self.last_step_ids.last().copied(),
185 );
186 let requirement = config.evaluate_rules(&ctx);
187 info!(
188 run_id = %self.run_id,
189 step = %name,
190 rule_index = ?requirement.rule_index,
191 required_approvers = requirement.required_approvers,
192 approver_groups = ?requirement.approver_groups,
193 "approval rules evaluated"
194 );
195 Some(requirement)
196 };
197
198 if let Some(interceptor) = self.interceptor.clone()
201 && let Some(outcome) = interceptor.intercept_approval(name, &config)
202 {
203 let trace_id = step_trace_id(self.run_id, name, position);
204 let step = self
205 .store
206 .create_step(NewStep {
207 run_id: self.run_id,
208 trace_id,
209 name: name.to_string(),
210 kind: StepKind::Approval,
211 position,
212 input: Some(to_value(&config)?),
213 is_error_handler: false,
214 })
215 .await?;
216
217 let now = Utc::now();
218 self.start_step(step.id, now).await?;
219 self.last_step_ids = vec![step.id];
220
221 return match outcome {
222 ApprovalOutcome::Approved => {
223 self.store
224 .update_step(
225 step.id,
226 StepUpdate {
227 status: Some(StepStatus::Completed),
228 output: Some(json!({"approved_by": "step-interceptor"})),
229 completed_at: Some(now),
230 approval_requirement: requirement.clone(),
231 ..StepUpdate::default()
232 },
233 )
234 .await?;
235 info!(
236 run_id = %self.run_id,
237 step = %name,
238 position,
239 "approval granted by the step interceptor"
240 );
241 Ok(())
242 }
243 ApprovalOutcome::Rejected { reason } => {
244 self.store
246 .update_step(
247 step.id,
248 StepUpdate {
249 status: Some(StepStatus::AwaitingApproval),
250 approval_requirement: requirement.clone(),
251 ..StepUpdate::default()
252 },
253 )
254 .await?;
255 self.store
256 .update_step(
257 step.id,
258 StepUpdate {
259 status: Some(StepStatus::Rejected),
260 error: Some(reason.clone()),
261 completed_at: Some(Utc::now()),
262 ..StepUpdate::default()
263 },
264 )
265 .await?;
266 info!(
267 run_id = %self.run_id,
268 step = %name,
269 position,
270 %reason,
271 "approval rejected by the step interceptor"
272 );
273 Err(EngineError::ApprovalRejected {
274 run_id: self.run_id,
275 step_id: step.id,
276 reason,
277 })
278 }
279 };
280 }
281
282 let trace_id = step_trace_id(self.run_id, name, position);
284 let step = self
285 .store
286 .create_step(NewStep {
287 run_id: self.run_id,
288 trace_id,
289 name: name.to_string(),
290 kind: StepKind::Approval,
291 position,
292 input: Some(to_value(&config)?),
293 is_error_handler: false,
294 })
295 .await?;
296
297 self.start_step(step.id, Utc::now()).await?;
298
299 let deadline_at = config
303 .effective_deadline_secs()
304 .map(|secs| Utc::now() + TimeDelta::seconds(secs as i64));
305
306 self.store
307 .update_step(
308 step.id,
309 StepUpdate {
310 status: Some(StepStatus::AwaitingApproval),
311 approval_deadline_at: deadline_at,
312 approval_stage: Some(0),
313 approval_assignee: config.assignee().cloned(),
314 approval_requirement: requirement,
315 ..StepUpdate::default()
316 },
317 )
318 .await?;
319
320 self.last_step_ids = vec![step.id];
321
322 if let Some(ref bus) = self.event_bus {
323 bus.publish(
324 self.run_id,
325 WorkflowEvent::ApprovalRequired(WorkflowApprovalRequiredEvent {
326 step_name: name.to_string(),
327 step_index: position,
328 approval_id: step.id,
329 }),
330 );
331 }
332
333 Err(EngineError::ApprovalRequired {
334 run_id: self.run_id,
335 step_id: step.id,
336 message: config.message().to_string(),
337 })
338 }
339}
340
341pub(crate) fn approval_expression_context(
352 run: &Run,
353 steps: &[Step],
354 attempt: u32,
355 before_position: u32,
356 last_step_id: Option<Uuid>,
357) -> Value {
358 let output = last_step_id
359 .and_then(|id| steps.iter().find(|s| s.id == id))
360 .and_then(|s| s.output.clone())
361 .unwrap_or(Value::Null);
362
363 let mut completed: Vec<&Step> = steps
364 .iter()
365 .filter(|s| {
366 s.attempt == attempt
367 && s.position < before_position
368 && s.status.state == StepStatus::Completed
369 })
370 .collect();
371 completed.sort_by_key(|s| s.position);
372
373 let mut by_name = Map::new();
374 for step in completed {
375 by_name.insert(
376 step.name.clone(),
377 json!({
378 "output": step.output,
379 "kind": step.kind,
380 "status": step.status.state,
381 }),
382 );
383 }
384
385 json!({
386 "output": output,
387 "payload": run.payload,
388 "labels": run.labels,
389 "metadata": {
390 "run_id": run.id,
391 "workflow_name": run.workflow_name,
392 "trigger": run.trigger,
393 "attempt": attempt,
394 "handler_version": run.handler_version,
395 },
396 "steps": by_name,
397 })
398}
399
400#[cfg(test)]
401mod tests {
402 use std::collections::HashMap;
403 use std::slice;
404
405 use ironflow_store::memory::InMemoryStore;
406 use ironflow_store::models::{NewRun, TriggerKind};
407 use ironflow_store::store::RunStore;
408
409 use super::*;
410
411 async fn run_with_labels(store: &InMemoryStore) -> Run {
412 store
413 .create_run(NewRun {
414 created_by: None,
415 workflow_name: "payments".to_string(),
416 trigger: TriggerKind::Manual,
417 payload: json!({"amount": 15000}),
418 max_retries: 0,
419 handler_version: Some("v2".to_string()),
420 labels: HashMap::from([("env".to_string(), "production".to_string())]),
421 scheduled_at: None,
422 idempotency_key: None,
423 max_cost_usd: None,
424 })
425 .await
426 .expect("create run")
427 .into_run()
428 }
429
430 async fn step(
432 store: &InMemoryStore,
433 run_id: Uuid,
434 name: &str,
435 position: u32,
436 output: Option<Value>,
437 ) -> Step {
438 let step = store
439 .create_step(NewStep {
440 run_id,
441 trace_id: step_trace_id(run_id, name, position),
442 name: name.to_string(),
443 kind: StepKind::Shell,
444 position,
445 input: None,
446 is_error_handler: false,
447 })
448 .await
449 .expect("create step");
450 store
451 .update_step(
452 step.id,
453 StepUpdate {
454 status: Some(StepStatus::Running),
455 ..StepUpdate::default()
456 },
457 )
458 .await
459 .expect("to running");
460 if output.is_some() {
461 store
462 .update_step(
463 step.id,
464 StepUpdate {
465 status: Some(StepStatus::Completed),
466 output,
467 ..StepUpdate::default()
468 },
469 )
470 .await
471 .expect("to completed");
472 }
473 store.get_step(step.id).await.expect("get").expect("exists")
474 }
475
476 #[tokio::test]
477 async fn context_holds_run_data_and_the_previous_output() {
478 let store = InMemoryStore::new();
479 let run = run_with_labels(&store).await;
480 let risk = step(&store, run.id, "risk", 0, Some(json!({"level": "high"}))).await;
481
482 let steps = slice::from_ref(&risk);
483 let ctx = approval_expression_context(&run, steps, 1, 1, Some(risk.id));
484
485 assert_eq!(ctx["output"], json!({"level": "high"}));
486 assert_eq!(ctx["payload"], json!({"amount": 15000}));
487 assert_eq!(ctx["labels"], json!({"env": "production"}));
488 assert_eq!(ctx["metadata"]["run_id"], json!(run.id));
489 assert_eq!(ctx["metadata"]["workflow_name"], json!("payments"));
490 assert_eq!(ctx["metadata"]["trigger"], json!(run.trigger));
491 assert_eq!(ctx["metadata"]["attempt"], json!(1));
492 assert_eq!(ctx["metadata"]["handler_version"], json!("v2"));
493 assert_eq!(ctx["steps"]["risk"]["output"], json!({"level": "high"}));
494 assert_eq!(ctx["steps"]["risk"]["kind"], json!("shell"));
495 assert_eq!(ctx["steps"]["risk"]["status"], json!("completed"));
496 }
497
498 #[tokio::test]
499 async fn output_is_null_without_a_previous_step() {
500 let store = InMemoryStore::new();
501 let run = run_with_labels(&store).await;
502
503 let ctx = approval_expression_context(&run, &[], 1, 0, None);
504
505 assert_eq!(ctx["output"], Value::Null);
506 assert_eq!(ctx["steps"], json!({}));
507 }
508
509 #[tokio::test]
510 async fn steps_keep_only_completed_steps_of_the_attempt_before_the_gate() {
511 let store = InMemoryStore::new();
512 let run = run_with_labels(&store).await;
513 let done = step(&store, run.id, "done", 0, Some(json!(1))).await;
514 let running = step(&store, run.id, "running", 1, None).await;
515 let later = step(&store, run.id, "later", 5, Some(json!(2))).await;
516 let mut previous_attempt = step(&store, run.id, "old", 2, Some(json!(3))).await;
517 previous_attempt.attempt = 2;
518
519 let steps = [done.clone(), running, later, previous_attempt];
520 let ctx = approval_expression_context(&run, &steps, 1, 3, Some(done.id));
521
522 let names: Vec<&String> = ctx["steps"]
523 .as_object()
524 .expect("steps object")
525 .keys()
526 .collect();
527 assert_eq!(names, vec!["done"]);
528 assert_eq!(ctx["output"], json!(1));
529 }
530
531 #[tokio::test]
532 async fn duplicate_names_keep_the_highest_position() {
533 let store = InMemoryStore::new();
534 let run = run_with_labels(&store).await;
535 let second = step(&store, run.id, "check", 1, Some(json!("second"))).await;
536 let first = step(&store, run.id, "check", 0, Some(json!("first"))).await;
537
538 let ctx = approval_expression_context(&run, &[second, first], 1, 2, None);
539
540 assert_eq!(ctx["steps"]["check"]["output"], json!("second"));
541 }
542}