1use std::sync::Arc;
27
28use chrono::{TimeDelta, Utc};
29use reqwest::Client;
30use serde_json::json;
31use tracing::{error, info, warn};
32use uuid::Uuid;
33
34use ironflow_store::models::{
35 Assignee, Run, RunStatus, RunUpdate, Step, StepKind, StepStatus, StepUpdate,
36};
37use strum::IntoStaticStr;
38
39use crate::config::{ApprovalConfig, EscalationPolicy, NotificationTarget};
40use crate::engine::{Engine, ExecutionMode};
41use crate::error::EngineError;
42use crate::notify::{
43 ApprovalEscalatedEvent, ApprovalGrantedEvent, ApprovalRejectedEvent, Event, RetryConfig,
44 deliver_with_retry, is_success_2xx,
45};
46
47pub const SYSTEM_TIMEOUT_ACTOR: &str = "system:timeout";
49
50pub const APPROVAL_TIMEOUT_ERROR: &str = "approval timeout";
52
53pub const DEFAULT_ESCALATION_BATCH_SIZE: u32 = 50;
55
56#[derive(Debug, Clone, PartialEq, Eq, IntoStaticStr)]
67#[strum(serialize_all = "snake_case")]
68pub enum EscalationAction {
69 Approved,
71 Rejected,
73 Notified(usize),
75 Reassigned(Assignee),
77 Exhausted,
79 Stale,
81}
82
83impl EscalationAction {
84 pub fn label(&self) -> &'static str {
95 self.into()
96 }
97}
98
99#[derive(Debug, Clone)]
117pub struct EscalationRecord {
118 pub run_id: Uuid,
120 pub step_id: Uuid,
122 pub stage: u32,
124 pub action: EscalationAction,
126 pub reason: String,
128}
129
130pub struct ApprovalEscalator {
156 engine: Arc<Engine>,
157 client: Client,
158 retry: RetryConfig,
159 batch_size: u32,
160}
161
162impl ApprovalEscalator {
163 pub fn new(engine: Arc<Engine>) -> Self {
185 let retry = RetryConfig::default();
186 Self {
187 engine,
188 client: retry.build_client(),
189 retry,
190 batch_size: DEFAULT_ESCALATION_BATCH_SIZE,
191 }
192 }
193
194 pub fn batch_size(mut self, batch_size: u32) -> Self {
211 self.batch_size = batch_size;
212 self
213 }
214
215 pub async fn tick(&self) -> Result<Vec<EscalationRecord>, EngineError> {
245 let steps = self
246 .engine
247 .store()
248 .claim_due_approval_deadlines(self.batch_size)
249 .await?;
250
251 let mut records = Vec::with_capacity(steps.len());
252 for step in &steps {
253 match self.escalate(step).await {
254 Ok(record) => records.push(record),
255 Err(err) => {
256 error!(
257 run_id = %step.run_id,
258 step_id = %step.id,
259 error = %err,
260 "failed to escalate an expired approval gate"
261 );
262 records.push(EscalationRecord {
263 run_id: step.run_id,
264 step_id: step.id,
265 stage: step.approval_stage,
266 action: EscalationAction::Stale,
267 reason: err.to_string(),
268 });
269 }
270 }
271 }
272
273 Ok(records)
274 }
275
276 async fn escalate(&self, step: &Step) -> Result<EscalationRecord, EngineError> {
278 let stage = step.approval_stage;
279
280 let Some(input) = step.input.as_ref() else {
281 error!(step_id = %step.id, "approval step has no stored configuration");
282 return Ok(self.stale(step, "approval step has no stored configuration"));
283 };
284 let config: ApprovalConfig = match serde_json::from_value(input.clone()) {
285 Ok(config) => config,
286 Err(err) => {
287 error!(step_id = %step.id, error = %err, "unreadable approval configuration");
288 return Ok(self.stale(step, "unreadable approval configuration"));
289 }
290 };
291
292 let run = self.engine.store().get_run(step.run_id).await?;
294 let Some(run) = run else {
295 return Ok(self.stale(step, "run no longer exists"));
296 };
297 if !awaits_gate(&run) || step.status.state != StepStatus::AwaitingApproval {
298 return Ok(self.stale(step, "gate already resolved"));
299 }
300 let paused = run.status.state == RunStatus::Paused;
301
302 let policy = config.effective_policy();
303 let reason = format!(
304 "approval deadline of {}s expired",
305 config.effective_deadline_secs().unwrap_or(0)
306 );
307
308 let Some(current) = resolve_stage(&policy, stage)? else {
309 warn!(
310 run_id = %step.run_id,
311 step_id = %step.id,
312 stage,
313 "escalation chain exhausted; the approval gate stays open with no timer"
314 );
315 self.publish(
316 step,
317 stage,
318 "chain",
319 &EscalationAction::Exhausted,
320 &reason,
321 None,
322 );
323 return Ok(EscalationRecord {
324 run_id: step.run_id,
325 step_id: step.id,
326 stage,
327 action: EscalationAction::Exhausted,
328 reason,
329 });
330 };
331
332 let action = match current {
333 EscalationPolicy::AutoApprove if step.kind == StepKind::HumanInput => {
338 self.auto_reject(step).await?
339 }
340 EscalationPolicy::AutoApprove => self.auto_approve(step, &reason, paused).await?,
341 EscalationPolicy::AutoReject => self.auto_reject(step).await?,
342 EscalationPolicy::Notify(targets) => {
343 self.notify(step, &config, &policy, targets, &reason)
344 .await?
345 }
346 EscalationPolicy::Escalate(assignee) => {
347 self.reassign(step, &config, &policy, assignee).await?
348 }
349 EscalationPolicy::Chain(_) => unreachable!("nested chains are rejected upfront"),
351 };
352
353 let assignee = match &action {
354 EscalationAction::Reassigned(to) => Some(to.clone()),
355 _ => step.approval_assignee.clone(),
356 };
357 self.publish(
358 step,
359 stage,
360 policy_label(current),
361 &action,
362 &reason,
363 assignee,
364 );
365
366 Ok(EscalationRecord {
367 run_id: step.run_id,
368 step_id: step.id,
369 stage,
370 action,
371 reason,
372 })
373 }
374
375 fn stale(&self, step: &Step, reason: &str) -> EscalationRecord {
377 EscalationRecord {
378 run_id: step.run_id,
379 step_id: step.id,
380 stage: step.approval_stage,
381 action: EscalationAction::Stale,
382 reason: reason.to_string(),
383 }
384 }
385
386 async fn auto_approve(
388 &self,
389 step: &Step,
390 reason: &str,
391 paused: bool,
392 ) -> Result<EscalationAction, EngineError> {
393 let now = Utc::now();
394 let store = self.engine.store();
395
396 store
397 .update_step(
398 step.id,
399 StepUpdate {
400 status: Some(StepStatus::Completed),
401 output: Some(json!({
402 "approved_by": SYSTEM_TIMEOUT_ACTOR,
403 "escalated_at": now,
404 "reason": reason,
405 })),
406 completed_at: Some(now),
407 clear_approval_deadline: true,
408 ..StepUpdate::default()
409 },
410 )
411 .await?;
412 if paused {
414 store
415 .update_run(
416 step.run_id,
417 RunUpdate {
418 resume_status: Some(RunStatus::Pending),
419 ..RunUpdate::default()
420 },
421 )
422 .await?;
423 } else {
424 let resume_status = match self.engine.execution_mode() {
425 ExecutionMode::Local => RunStatus::Running,
426 ExecutionMode::Workers => RunStatus::Pending,
427 };
428 store.update_run_status(step.run_id, resume_status).await?;
429 }
430
431 self.engine
432 .event_publisher()
433 .publish(Event::ApprovalGranted(ApprovalGrantedEvent {
434 run_id: step.run_id,
435 step_id: Some(step.id),
436 approved_by: SYSTEM_TIMEOUT_ACTOR.to_string(),
437 approvals_received: step.approvals.len() as u32,
438 approvals_required: step
439 .approval_requirement
440 .as_ref()
441 .map_or(1, |r| r.required_approvers),
442 requirement: step.approval_requirement.clone(),
443 at: now,
444 }));
445
446 if paused {
452 info!(
453 run_id = %step.run_id,
454 "paused run will resume after its auto-approved gate"
455 );
456 } else if matches!(self.engine.execution_mode(), ExecutionMode::Local)
457 && let Err(err) = self.engine.resume_run(step.run_id).await
458 {
459 error!(
460 run_id = %step.run_id,
461 error = %err,
462 "failed to resume run after an auto-approved gate"
463 );
464 }
465
466 Ok(EscalationAction::Approved)
467 }
468
469 async fn auto_reject(&self, step: &Step) -> Result<EscalationAction, EngineError> {
472 let now = Utc::now();
473 let store = self.engine.store();
474
475 store
476 .update_step(
477 step.id,
478 StepUpdate {
479 status: Some(StepStatus::Failed),
480 error: Some(APPROVAL_TIMEOUT_ERROR.to_string()),
481 completed_at: Some(now),
482 clear_approval_deadline: true,
483 ..StepUpdate::default()
484 },
485 )
486 .await?;
487 store
488 .update_run(
489 step.run_id,
490 RunUpdate {
491 status: Some(RunStatus::Failed),
492 error: Some(APPROVAL_TIMEOUT_ERROR.to_string()),
493 completed_at: Some(now),
494 ..RunUpdate::default()
495 },
496 )
497 .await?;
498
499 self.engine
500 .event_publisher()
501 .publish(Event::ApprovalRejected(ApprovalRejectedEvent {
502 run_id: step.run_id,
503 step_id: Some(step.id),
504 rejected_by: SYSTEM_TIMEOUT_ACTOR.to_string(),
505 requirement: step.approval_requirement.clone(),
506 at: now,
507 }));
508
509 if let Err(err) = self
512 .engine
513 .fail_ancestors(step.run_id, APPROVAL_TIMEOUT_ERROR)
514 .await
515 {
516 error!(
517 run_id = %step.run_id,
518 error = %err,
519 "failed to fail the ancestors of an auto-rejected run"
520 );
521 }
522
523 Ok(EscalationAction::Rejected)
524 }
525
526 async fn notify(
528 &self,
529 step: &Step,
530 config: &ApprovalConfig,
531 policy: &EscalationPolicy,
532 targets: &[NotificationTarget],
533 reason: &str,
534 ) -> Result<EscalationAction, EngineError> {
535 let event = self.escalated_event(
536 step,
537 step.approval_stage,
538 "notify",
539 EscalationAction::Notified(targets.len()).label(),
540 reason,
541 step.approval_assignee.clone(),
542 );
543 self.deliver_notifications(targets, &event).await;
544
545 self.rearm(step, config, policy, None).await?;
546 Ok(EscalationAction::Notified(targets.len()))
547 }
548
549 async fn reassign(
551 &self,
552 step: &Step,
553 config: &ApprovalConfig,
554 policy: &EscalationPolicy,
555 assignee: &Assignee,
556 ) -> Result<EscalationAction, EngineError> {
557 self.rearm(step, config, policy, Some(assignee)).await?;
558 Ok(EscalationAction::Reassigned(assignee.clone()))
559 }
560
561 async fn rearm(
567 &self,
568 step: &Step,
569 config: &ApprovalConfig,
570 policy: &EscalationPolicy,
571 assignee: Option<&Assignee>,
572 ) -> Result<(), EngineError> {
573 let next_stage = next_stage(policy, step.approval_stage);
574 let deadline = config
575 .effective_deadline_secs()
576 .map(|secs| Utc::now() + TimeDelta::seconds(secs as i64));
577
578 self.engine
579 .store()
580 .update_step(
581 step.id,
582 StepUpdate {
583 approval_deadline_at: deadline,
584 approval_stage: Some(next_stage),
585 approval_assignee: assignee.cloned(),
586 ..StepUpdate::default()
587 },
588 )
589 .await?;
590
591 Ok(())
592 }
593
594 fn escalated_event(
596 &self,
597 step: &Step,
598 stage: u32,
599 policy: &str,
600 action: &str,
601 reason: &str,
602 assignee: Option<Assignee>,
603 ) -> ApprovalEscalatedEvent {
604 ApprovalEscalatedEvent {
605 run_id: step.run_id,
606 step_id: step.id,
607 step_name: step.name.clone(),
608 stage,
609 policy: policy.to_string(),
610 action: action.to_string(),
611 reason: reason.to_string(),
612 assignee,
613 at: Utc::now(),
614 }
615 }
616
617 fn publish(
619 &self,
620 step: &Step,
621 stage: u32,
622 policy: &str,
623 action: &EscalationAction,
624 reason: &str,
625 assignee: Option<Assignee>,
626 ) {
627 let event = self.escalated_event(step, stage, policy, action.label(), reason, assignee);
628
629 info!(
630 run_id = %step.run_id,
631 step_id = %step.id,
632 stage,
633 policy,
634 action = action.label(),
635 reason,
636 "approval gate escalated"
637 );
638
639 self.engine
640 .event_publisher()
641 .publish(Event::ApprovalEscalated(event));
642 }
643
644 async fn deliver_notifications(
650 &self,
651 targets: &[NotificationTarget],
652 event: &ApprovalEscalatedEvent,
653 ) {
654 for target in targets {
655 match target {
656 NotificationTarget::Webhook { url } => {
657 deliver_with_retry(
658 &self.retry,
659 || self.client.post(url).json(event),
660 is_success_2xx,
661 "approval_escalation",
662 url,
663 )
664 .await;
665 }
666 NotificationTarget::Slack {
667 webhook_url,
668 channel,
669 } => {
670 let body = json!({
671 "channel": channel,
672 "text": format!(
673 "⏰ Approval gate `{}` on run {} hit its SLA: {}",
674 event.step_name, event.run_id, event.reason
675 ),
676 });
677 deliver_with_retry(
678 &self.retry,
679 || self.client.post(webhook_url).json(&body),
680 is_success_2xx,
681 "approval_escalation",
682 webhook_url,
683 )
684 .await;
685 }
686 }
687 }
688 }
689}
690
691fn awaits_gate(run: &Run) -> bool {
694 run.status.state == RunStatus::AwaitingApproval
695 || (run.status.state == RunStatus::Paused
696 && run.resume_status == Some(RunStatus::AwaitingApproval))
697}
698
699fn resolve_stage(
704 policy: &EscalationPolicy,
705 stage: u32,
706) -> Result<Option<&EscalationPolicy>, EngineError> {
707 match policy.stage(stage as usize) {
708 Some(EscalationPolicy::Chain(_)) => Err(EngineError::StepConfig(
709 "nested escalation chains are not supported".to_string(),
710 )),
711 other => Ok(other),
712 }
713}
714
715fn next_stage(policy: &EscalationPolicy, stage: u32) -> u32 {
720 match policy {
721 EscalationPolicy::Chain(_) => stage.saturating_add(1),
722 _ => stage,
723 }
724}
725
726fn policy_label(policy: &EscalationPolicy) -> &'static str {
728 match policy {
729 EscalationPolicy::AutoApprove => "auto_approve",
730 EscalationPolicy::AutoReject => "auto_reject",
731 EscalationPolicy::Notify(_) => "notify",
732 EscalationPolicy::Escalate(_) => "escalate",
733 EscalationPolicy::Chain(_) => "chain",
734 }
735}
736
737#[cfg(test)]
738mod tests {
739 use super::*;
740
741 #[test]
742 fn actions_compare_by_value() {
743 assert_eq!(EscalationAction::Notified(2), EscalationAction::Notified(2));
744 assert_ne!(EscalationAction::Notified(2), EscalationAction::Notified(3));
745 assert_eq!(
746 EscalationAction::Reassigned(Assignee::group("sre")),
747 EscalationAction::Reassigned(Assignee::group("sre"))
748 );
749 assert_ne!(EscalationAction::Approved, EscalationAction::Rejected);
750 }
751
752 #[test]
753 fn action_labels_are_distinct() {
754 let labels = [
755 EscalationAction::Approved.label(),
756 EscalationAction::Rejected.label(),
757 EscalationAction::Notified(1).label(),
758 EscalationAction::Reassigned(Assignee::group("sre")).label(),
759 EscalationAction::Exhausted.label(),
760 EscalationAction::Stale.label(),
761 ];
762
763 for (i, a) in labels.iter().enumerate() {
764 for (j, b) in labels.iter().enumerate() {
765 if i != j {
766 assert_ne!(a, b, "labels {i} and {j} collide");
767 }
768 }
769 }
770 }
771
772 #[test]
773 fn a_bare_repeating_policy_stays_at_the_same_stage() {
774 let notify = EscalationPolicy::Notify(Vec::new());
775 assert_eq!(next_stage(¬ify, 0), 0);
776 assert_eq!(next_stage(¬ify, 7), 7);
777
778 let escalate = EscalationPolicy::Escalate(Assignee::group("sre"));
779 assert_eq!(next_stage(&escalate, 3), 3);
780 }
781
782 #[test]
783 fn a_chained_policy_advances_one_stage_per_expiry() {
784 let chain = EscalationPolicy::Chain(vec![
785 EscalationPolicy::Notify(Vec::new()),
786 EscalationPolicy::AutoReject,
787 ]);
788 assert_eq!(next_stage(&chain, 0), 1);
789 assert_eq!(next_stage(&chain, 1), 2);
790 assert_eq!(next_stage(&chain, u32::MAX), u32::MAX);
791 }
792
793 #[test]
794 fn resolve_stage_rejects_a_nested_chain() {
795 let policy = EscalationPolicy::Chain(vec![
796 EscalationPolicy::Chain(vec![EscalationPolicy::AutoReject]),
797 EscalationPolicy::AutoApprove,
798 ]);
799
800 let err = resolve_stage(&policy, 0).expect_err("nested chains are rejected");
801 assert!(
802 err.to_string().contains("nested escalation chains"),
803 "got {err}"
804 );
805
806 assert_eq!(
808 resolve_stage(&policy, 1).expect("valid stage"),
809 Some(&EscalationPolicy::AutoApprove)
810 );
811 }
812
813 #[test]
814 fn resolve_stage_reports_an_exhausted_chain_as_none() {
815 let policy = EscalationPolicy::Chain(vec![EscalationPolicy::AutoReject]);
816 assert!(resolve_stage(&policy, 1).expect("valid call").is_none());
817
818 let bare = EscalationPolicy::AutoApprove;
819 assert!(resolve_stage(&bare, 1).expect("valid call").is_none());
820 assert_eq!(
821 resolve_stage(&bare, 0).expect("valid call"),
822 Some(&EscalationPolicy::AutoApprove)
823 );
824 }
825
826 #[test]
827 fn policy_labels_match_the_wire_format() {
828 assert_eq!(policy_label(&EscalationPolicy::AutoApprove), "auto_approve");
829 assert_eq!(policy_label(&EscalationPolicy::AutoReject), "auto_reject");
830 assert_eq!(
831 policy_label(&EscalationPolicy::Notify(Vec::new())),
832 "notify"
833 );
834 assert_eq!(
835 policy_label(&EscalationPolicy::Escalate(Assignee::group("sre"))),
836 "escalate"
837 );
838 assert_eq!(policy_label(&EscalationPolicy::Chain(Vec::new())), "chain");
839 }
840}