1use std::sync::Arc;
21
22use chrono::{TimeDelta, Utc};
23use reqwest::Client;
24use serde_json::json;
25use tracing::{error, info, warn};
26use uuid::Uuid;
27
28use ironflow_store::models::{
29 Assignee, RunStatus, RunUpdate, Step, StepKind, StepStatus, StepUpdate,
30};
31use strum::IntoStaticStr;
32
33use crate::config::{ApprovalConfig, EscalationPolicy, NotificationTarget};
34use crate::engine::{Engine, ExecutionMode};
35use crate::error::EngineError;
36use crate::notify::{
37 ApprovalEscalatedEvent, ApprovalGrantedEvent, ApprovalRejectedEvent, Event, RetryConfig,
38 deliver_with_retry, is_success_2xx,
39};
40
41pub const SYSTEM_TIMEOUT_ACTOR: &str = "system:timeout";
43
44pub const APPROVAL_TIMEOUT_ERROR: &str = "approval timeout";
46
47pub const DEFAULT_ESCALATION_BATCH_SIZE: u32 = 50;
49
50#[derive(Debug, Clone, PartialEq, Eq, IntoStaticStr)]
61#[strum(serialize_all = "snake_case")]
62pub enum EscalationAction {
63 Approved,
65 Rejected,
67 Notified(usize),
69 Reassigned(Assignee),
71 Exhausted,
73 Stale,
75}
76
77impl EscalationAction {
78 pub fn label(&self) -> &'static str {
89 self.into()
90 }
91}
92
93#[derive(Debug, Clone)]
111pub struct EscalationRecord {
112 pub run_id: Uuid,
114 pub step_id: Uuid,
116 pub stage: u32,
118 pub action: EscalationAction,
120 pub reason: String,
122}
123
124pub struct ApprovalEscalator {
150 engine: Arc<Engine>,
151 client: Client,
152 retry: RetryConfig,
153 batch_size: u32,
154}
155
156impl ApprovalEscalator {
157 pub fn new(engine: Arc<Engine>) -> Self {
179 let retry = RetryConfig::default();
180 Self {
181 engine,
182 client: retry.build_client(),
183 retry,
184 batch_size: DEFAULT_ESCALATION_BATCH_SIZE,
185 }
186 }
187
188 pub fn batch_size(mut self, batch_size: u32) -> Self {
205 self.batch_size = batch_size;
206 self
207 }
208
209 pub async fn tick(&self) -> Result<Vec<EscalationRecord>, EngineError> {
236 let steps = self
237 .engine
238 .store()
239 .claim_due_approval_deadlines(self.batch_size)
240 .await?;
241
242 let mut records = Vec::with_capacity(steps.len());
243 for step in &steps {
244 match self.escalate(step).await {
245 Ok(record) => records.push(record),
246 Err(err) => {
247 error!(
248 run_id = %step.run_id,
249 step_id = %step.id,
250 error = %err,
251 "failed to escalate an expired approval gate"
252 );
253 records.push(EscalationRecord {
254 run_id: step.run_id,
255 step_id: step.id,
256 stage: step.approval_stage,
257 action: EscalationAction::Stale,
258 reason: err.to_string(),
259 });
260 }
261 }
262 }
263
264 Ok(records)
265 }
266
267 async fn escalate(&self, step: &Step) -> Result<EscalationRecord, EngineError> {
269 let stage = step.approval_stage;
270
271 let Some(input) = step.input.as_ref() else {
272 error!(step_id = %step.id, "approval step has no stored configuration");
273 return Ok(self.stale(step, "approval step has no stored configuration"));
274 };
275 let config: ApprovalConfig = match serde_json::from_value(input.clone()) {
276 Ok(config) => config,
277 Err(err) => {
278 error!(step_id = %step.id, error = %err, "unreadable approval configuration");
279 return Ok(self.stale(step, "unreadable approval configuration"));
280 }
281 };
282
283 let run = self.engine.store().get_run(step.run_id).await?;
285 let Some(run) = run else {
286 return Ok(self.stale(step, "run no longer exists"));
287 };
288 if run.status.state != RunStatus::AwaitingApproval
289 || step.status.state != StepStatus::AwaitingApproval
290 {
291 return Ok(self.stale(step, "gate already resolved"));
292 }
293
294 let policy = config.effective_policy();
295 let reason = format!(
296 "approval deadline of {}s expired",
297 config.effective_deadline_secs().unwrap_or(0)
298 );
299
300 let Some(current) = resolve_stage(&policy, stage)? else {
301 warn!(
302 run_id = %step.run_id,
303 step_id = %step.id,
304 stage,
305 "escalation chain exhausted; the approval gate stays open with no timer"
306 );
307 self.publish(
308 step,
309 stage,
310 "chain",
311 &EscalationAction::Exhausted,
312 &reason,
313 None,
314 );
315 return Ok(EscalationRecord {
316 run_id: step.run_id,
317 step_id: step.id,
318 stage,
319 action: EscalationAction::Exhausted,
320 reason,
321 });
322 };
323
324 let action = match current {
325 EscalationPolicy::AutoApprove if step.kind == StepKind::HumanInput => {
330 self.auto_reject(step).await?
331 }
332 EscalationPolicy::AutoApprove => self.auto_approve(step, &reason).await?,
333 EscalationPolicy::AutoReject => self.auto_reject(step).await?,
334 EscalationPolicy::Notify(targets) => {
335 self.notify(step, &config, &policy, targets, &reason)
336 .await?
337 }
338 EscalationPolicy::Escalate(assignee) => {
339 self.reassign(step, &config, &policy, assignee).await?
340 }
341 EscalationPolicy::Chain(_) => unreachable!("nested chains are rejected upfront"),
343 };
344
345 let assignee = match &action {
346 EscalationAction::Reassigned(to) => Some(to.clone()),
347 _ => step.approval_assignee.clone(),
348 };
349 self.publish(
350 step,
351 stage,
352 policy_label(current),
353 &action,
354 &reason,
355 assignee,
356 );
357
358 Ok(EscalationRecord {
359 run_id: step.run_id,
360 step_id: step.id,
361 stage,
362 action,
363 reason,
364 })
365 }
366
367 fn stale(&self, step: &Step, reason: &str) -> EscalationRecord {
369 EscalationRecord {
370 run_id: step.run_id,
371 step_id: step.id,
372 stage: step.approval_stage,
373 action: EscalationAction::Stale,
374 reason: reason.to_string(),
375 }
376 }
377
378 async fn auto_approve(
380 &self,
381 step: &Step,
382 reason: &str,
383 ) -> Result<EscalationAction, EngineError> {
384 let now = Utc::now();
385 let store = self.engine.store();
386
387 store
388 .update_step(
389 step.id,
390 StepUpdate {
391 status: Some(StepStatus::Completed),
392 output: Some(json!({
393 "approved_by": SYSTEM_TIMEOUT_ACTOR,
394 "escalated_at": now,
395 "reason": reason,
396 })),
397 completed_at: Some(now),
398 clear_approval_deadline: true,
399 ..StepUpdate::default()
400 },
401 )
402 .await?;
403 let resume_status = match self.engine.execution_mode() {
404 ExecutionMode::Local => RunStatus::Running,
405 ExecutionMode::Workers => RunStatus::Pending,
406 };
407 store.update_run_status(step.run_id, resume_status).await?;
408
409 self.engine
410 .event_publisher()
411 .publish(Event::ApprovalGranted(ApprovalGrantedEvent {
412 run_id: step.run_id,
413 step_id: Some(step.id),
414 approved_by: SYSTEM_TIMEOUT_ACTOR.to_string(),
415 approvals_received: step.approvals.len() as u32,
416 approvals_required: step
417 .approval_requirement
418 .as_ref()
419 .map_or(1, |r| r.required_approvers),
420 requirement: step.approval_requirement.clone(),
421 at: now,
422 }));
423
424 if matches!(self.engine.execution_mode(), ExecutionMode::Local)
430 && let Err(err) = self.engine.resume_run(step.run_id).await
431 {
432 error!(
433 run_id = %step.run_id,
434 error = %err,
435 "failed to resume run after an auto-approved gate"
436 );
437 }
438
439 Ok(EscalationAction::Approved)
440 }
441
442 async fn auto_reject(&self, step: &Step) -> Result<EscalationAction, EngineError> {
444 let now = Utc::now();
445 let store = self.engine.store();
446
447 store
448 .update_step(
449 step.id,
450 StepUpdate {
451 status: Some(StepStatus::Failed),
452 error: Some(APPROVAL_TIMEOUT_ERROR.to_string()),
453 completed_at: Some(now),
454 clear_approval_deadline: true,
455 ..StepUpdate::default()
456 },
457 )
458 .await?;
459 store
460 .update_run(
461 step.run_id,
462 RunUpdate {
463 status: Some(RunStatus::Failed),
464 error: Some(APPROVAL_TIMEOUT_ERROR.to_string()),
465 completed_at: Some(now),
466 ..RunUpdate::default()
467 },
468 )
469 .await?;
470
471 self.engine
472 .event_publisher()
473 .publish(Event::ApprovalRejected(ApprovalRejectedEvent {
474 run_id: step.run_id,
475 step_id: Some(step.id),
476 rejected_by: SYSTEM_TIMEOUT_ACTOR.to_string(),
477 requirement: step.approval_requirement.clone(),
478 at: now,
479 }));
480
481 Ok(EscalationAction::Rejected)
482 }
483
484 async fn notify(
486 &self,
487 step: &Step,
488 config: &ApprovalConfig,
489 policy: &EscalationPolicy,
490 targets: &[NotificationTarget],
491 reason: &str,
492 ) -> Result<EscalationAction, EngineError> {
493 let event = self.escalated_event(
494 step,
495 step.approval_stage,
496 "notify",
497 EscalationAction::Notified(targets.len()).label(),
498 reason,
499 step.approval_assignee.clone(),
500 );
501 self.deliver_notifications(targets, &event).await;
502
503 self.rearm(step, config, policy, None).await?;
504 Ok(EscalationAction::Notified(targets.len()))
505 }
506
507 async fn reassign(
509 &self,
510 step: &Step,
511 config: &ApprovalConfig,
512 policy: &EscalationPolicy,
513 assignee: &Assignee,
514 ) -> Result<EscalationAction, EngineError> {
515 self.rearm(step, config, policy, Some(assignee)).await?;
516 Ok(EscalationAction::Reassigned(assignee.clone()))
517 }
518
519 async fn rearm(
525 &self,
526 step: &Step,
527 config: &ApprovalConfig,
528 policy: &EscalationPolicy,
529 assignee: Option<&Assignee>,
530 ) -> Result<(), EngineError> {
531 let next_stage = next_stage(policy, step.approval_stage);
532 let deadline = config
533 .effective_deadline_secs()
534 .map(|secs| Utc::now() + TimeDelta::seconds(secs as i64));
535
536 self.engine
537 .store()
538 .update_step(
539 step.id,
540 StepUpdate {
541 approval_deadline_at: deadline,
542 approval_stage: Some(next_stage),
543 approval_assignee: assignee.cloned(),
544 ..StepUpdate::default()
545 },
546 )
547 .await?;
548
549 Ok(())
550 }
551
552 fn escalated_event(
554 &self,
555 step: &Step,
556 stage: u32,
557 policy: &str,
558 action: &str,
559 reason: &str,
560 assignee: Option<Assignee>,
561 ) -> ApprovalEscalatedEvent {
562 ApprovalEscalatedEvent {
563 run_id: step.run_id,
564 step_id: step.id,
565 step_name: step.name.clone(),
566 stage,
567 policy: policy.to_string(),
568 action: action.to_string(),
569 reason: reason.to_string(),
570 assignee,
571 at: Utc::now(),
572 }
573 }
574
575 fn publish(
577 &self,
578 step: &Step,
579 stage: u32,
580 policy: &str,
581 action: &EscalationAction,
582 reason: &str,
583 assignee: Option<Assignee>,
584 ) {
585 let event = self.escalated_event(step, stage, policy, action.label(), reason, assignee);
586
587 info!(
588 run_id = %step.run_id,
589 step_id = %step.id,
590 stage,
591 policy,
592 action = action.label(),
593 reason,
594 "approval gate escalated"
595 );
596
597 self.engine
598 .event_publisher()
599 .publish(Event::ApprovalEscalated(event));
600 }
601
602 async fn deliver_notifications(
608 &self,
609 targets: &[NotificationTarget],
610 event: &ApprovalEscalatedEvent,
611 ) {
612 for target in targets {
613 match target {
614 NotificationTarget::Webhook { url } => {
615 deliver_with_retry(
616 &self.retry,
617 || self.client.post(url).json(event),
618 is_success_2xx,
619 "approval_escalation",
620 url,
621 )
622 .await;
623 }
624 NotificationTarget::Slack {
625 webhook_url,
626 channel,
627 } => {
628 let body = json!({
629 "channel": channel,
630 "text": format!(
631 "⏰ Approval gate `{}` on run {} hit its SLA: {}",
632 event.step_name, event.run_id, event.reason
633 ),
634 });
635 deliver_with_retry(
636 &self.retry,
637 || self.client.post(webhook_url).json(&body),
638 is_success_2xx,
639 "approval_escalation",
640 webhook_url,
641 )
642 .await;
643 }
644 }
645 }
646 }
647}
648
649fn resolve_stage(
654 policy: &EscalationPolicy,
655 stage: u32,
656) -> Result<Option<&EscalationPolicy>, EngineError> {
657 match policy.stage(stage as usize) {
658 Some(EscalationPolicy::Chain(_)) => Err(EngineError::StepConfig(
659 "nested escalation chains are not supported".to_string(),
660 )),
661 other => Ok(other),
662 }
663}
664
665fn next_stage(policy: &EscalationPolicy, stage: u32) -> u32 {
670 match policy {
671 EscalationPolicy::Chain(_) => stage.saturating_add(1),
672 _ => stage,
673 }
674}
675
676fn policy_label(policy: &EscalationPolicy) -> &'static str {
678 match policy {
679 EscalationPolicy::AutoApprove => "auto_approve",
680 EscalationPolicy::AutoReject => "auto_reject",
681 EscalationPolicy::Notify(_) => "notify",
682 EscalationPolicy::Escalate(_) => "escalate",
683 EscalationPolicy::Chain(_) => "chain",
684 }
685}
686
687#[cfg(test)]
688mod tests {
689 use super::*;
690
691 #[test]
692 fn actions_compare_by_value() {
693 assert_eq!(EscalationAction::Notified(2), EscalationAction::Notified(2));
694 assert_ne!(EscalationAction::Notified(2), EscalationAction::Notified(3));
695 assert_eq!(
696 EscalationAction::Reassigned(Assignee::group("sre")),
697 EscalationAction::Reassigned(Assignee::group("sre"))
698 );
699 assert_ne!(EscalationAction::Approved, EscalationAction::Rejected);
700 }
701
702 #[test]
703 fn action_labels_are_distinct() {
704 let labels = [
705 EscalationAction::Approved.label(),
706 EscalationAction::Rejected.label(),
707 EscalationAction::Notified(1).label(),
708 EscalationAction::Reassigned(Assignee::group("sre")).label(),
709 EscalationAction::Exhausted.label(),
710 EscalationAction::Stale.label(),
711 ];
712
713 for (i, a) in labels.iter().enumerate() {
714 for (j, b) in labels.iter().enumerate() {
715 if i != j {
716 assert_ne!(a, b, "labels {i} and {j} collide");
717 }
718 }
719 }
720 }
721
722 #[test]
723 fn a_bare_repeating_policy_stays_at_the_same_stage() {
724 let notify = EscalationPolicy::Notify(Vec::new());
725 assert_eq!(next_stage(¬ify, 0), 0);
726 assert_eq!(next_stage(¬ify, 7), 7);
727
728 let escalate = EscalationPolicy::Escalate(Assignee::group("sre"));
729 assert_eq!(next_stage(&escalate, 3), 3);
730 }
731
732 #[test]
733 fn a_chained_policy_advances_one_stage_per_expiry() {
734 let chain = EscalationPolicy::Chain(vec![
735 EscalationPolicy::Notify(Vec::new()),
736 EscalationPolicy::AutoReject,
737 ]);
738 assert_eq!(next_stage(&chain, 0), 1);
739 assert_eq!(next_stage(&chain, 1), 2);
740 assert_eq!(next_stage(&chain, u32::MAX), u32::MAX);
741 }
742
743 #[test]
744 fn resolve_stage_rejects_a_nested_chain() {
745 let policy = EscalationPolicy::Chain(vec![
746 EscalationPolicy::Chain(vec![EscalationPolicy::AutoReject]),
747 EscalationPolicy::AutoApprove,
748 ]);
749
750 let err = resolve_stage(&policy, 0).expect_err("nested chains are rejected");
751 assert!(
752 err.to_string().contains("nested escalation chains"),
753 "got {err}"
754 );
755
756 assert_eq!(
758 resolve_stage(&policy, 1).expect("valid stage"),
759 Some(&EscalationPolicy::AutoApprove)
760 );
761 }
762
763 #[test]
764 fn resolve_stage_reports_an_exhausted_chain_as_none() {
765 let policy = EscalationPolicy::Chain(vec![EscalationPolicy::AutoReject]);
766 assert!(resolve_stage(&policy, 1).expect("valid call").is_none());
767
768 let bare = EscalationPolicy::AutoApprove;
769 assert!(resolve_stage(&bare, 1).expect("valid call").is_none());
770 assert_eq!(
771 resolve_stage(&bare, 0).expect("valid call"),
772 Some(&EscalationPolicy::AutoApprove)
773 );
774 }
775
776 #[test]
777 fn policy_labels_match_the_wire_format() {
778 assert_eq!(policy_label(&EscalationPolicy::AutoApprove), "auto_approve");
779 assert_eq!(policy_label(&EscalationPolicy::AutoReject), "auto_reject");
780 assert_eq!(
781 policy_label(&EscalationPolicy::Notify(Vec::new())),
782 "notify"
783 );
784 assert_eq!(
785 policy_label(&EscalationPolicy::Escalate(Assignee::group("sre"))),
786 "escalate"
787 );
788 assert_eq!(policy_label(&EscalationPolicy::Chain(Vec::new())), "chain");
789 }
790}