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 if let Err(err) = self
484 .engine
485 .fail_ancestors(step.run_id, APPROVAL_TIMEOUT_ERROR)
486 .await
487 {
488 error!(
489 run_id = %step.run_id,
490 error = %err,
491 "failed to fail the ancestors of an auto-rejected run"
492 );
493 }
494
495 Ok(EscalationAction::Rejected)
496 }
497
498 async fn notify(
500 &self,
501 step: &Step,
502 config: &ApprovalConfig,
503 policy: &EscalationPolicy,
504 targets: &[NotificationTarget],
505 reason: &str,
506 ) -> Result<EscalationAction, EngineError> {
507 let event = self.escalated_event(
508 step,
509 step.approval_stage,
510 "notify",
511 EscalationAction::Notified(targets.len()).label(),
512 reason,
513 step.approval_assignee.clone(),
514 );
515 self.deliver_notifications(targets, &event).await;
516
517 self.rearm(step, config, policy, None).await?;
518 Ok(EscalationAction::Notified(targets.len()))
519 }
520
521 async fn reassign(
523 &self,
524 step: &Step,
525 config: &ApprovalConfig,
526 policy: &EscalationPolicy,
527 assignee: &Assignee,
528 ) -> Result<EscalationAction, EngineError> {
529 self.rearm(step, config, policy, Some(assignee)).await?;
530 Ok(EscalationAction::Reassigned(assignee.clone()))
531 }
532
533 async fn rearm(
539 &self,
540 step: &Step,
541 config: &ApprovalConfig,
542 policy: &EscalationPolicy,
543 assignee: Option<&Assignee>,
544 ) -> Result<(), EngineError> {
545 let next_stage = next_stage(policy, step.approval_stage);
546 let deadline = config
547 .effective_deadline_secs()
548 .map(|secs| Utc::now() + TimeDelta::seconds(secs as i64));
549
550 self.engine
551 .store()
552 .update_step(
553 step.id,
554 StepUpdate {
555 approval_deadline_at: deadline,
556 approval_stage: Some(next_stage),
557 approval_assignee: assignee.cloned(),
558 ..StepUpdate::default()
559 },
560 )
561 .await?;
562
563 Ok(())
564 }
565
566 fn escalated_event(
568 &self,
569 step: &Step,
570 stage: u32,
571 policy: &str,
572 action: &str,
573 reason: &str,
574 assignee: Option<Assignee>,
575 ) -> ApprovalEscalatedEvent {
576 ApprovalEscalatedEvent {
577 run_id: step.run_id,
578 step_id: step.id,
579 step_name: step.name.clone(),
580 stage,
581 policy: policy.to_string(),
582 action: action.to_string(),
583 reason: reason.to_string(),
584 assignee,
585 at: Utc::now(),
586 }
587 }
588
589 fn publish(
591 &self,
592 step: &Step,
593 stage: u32,
594 policy: &str,
595 action: &EscalationAction,
596 reason: &str,
597 assignee: Option<Assignee>,
598 ) {
599 let event = self.escalated_event(step, stage, policy, action.label(), reason, assignee);
600
601 info!(
602 run_id = %step.run_id,
603 step_id = %step.id,
604 stage,
605 policy,
606 action = action.label(),
607 reason,
608 "approval gate escalated"
609 );
610
611 self.engine
612 .event_publisher()
613 .publish(Event::ApprovalEscalated(event));
614 }
615
616 async fn deliver_notifications(
622 &self,
623 targets: &[NotificationTarget],
624 event: &ApprovalEscalatedEvent,
625 ) {
626 for target in targets {
627 match target {
628 NotificationTarget::Webhook { url } => {
629 deliver_with_retry(
630 &self.retry,
631 || self.client.post(url).json(event),
632 is_success_2xx,
633 "approval_escalation",
634 url,
635 )
636 .await;
637 }
638 NotificationTarget::Slack {
639 webhook_url,
640 channel,
641 } => {
642 let body = json!({
643 "channel": channel,
644 "text": format!(
645 "⏰ Approval gate `{}` on run {} hit its SLA: {}",
646 event.step_name, event.run_id, event.reason
647 ),
648 });
649 deliver_with_retry(
650 &self.retry,
651 || self.client.post(webhook_url).json(&body),
652 is_success_2xx,
653 "approval_escalation",
654 webhook_url,
655 )
656 .await;
657 }
658 }
659 }
660 }
661}
662
663fn resolve_stage(
668 policy: &EscalationPolicy,
669 stage: u32,
670) -> Result<Option<&EscalationPolicy>, EngineError> {
671 match policy.stage(stage as usize) {
672 Some(EscalationPolicy::Chain(_)) => Err(EngineError::StepConfig(
673 "nested escalation chains are not supported".to_string(),
674 )),
675 other => Ok(other),
676 }
677}
678
679fn next_stage(policy: &EscalationPolicy, stage: u32) -> u32 {
684 match policy {
685 EscalationPolicy::Chain(_) => stage.saturating_add(1),
686 _ => stage,
687 }
688}
689
690fn policy_label(policy: &EscalationPolicy) -> &'static str {
692 match policy {
693 EscalationPolicy::AutoApprove => "auto_approve",
694 EscalationPolicy::AutoReject => "auto_reject",
695 EscalationPolicy::Notify(_) => "notify",
696 EscalationPolicy::Escalate(_) => "escalate",
697 EscalationPolicy::Chain(_) => "chain",
698 }
699}
700
701#[cfg(test)]
702mod tests {
703 use super::*;
704
705 #[test]
706 fn actions_compare_by_value() {
707 assert_eq!(EscalationAction::Notified(2), EscalationAction::Notified(2));
708 assert_ne!(EscalationAction::Notified(2), EscalationAction::Notified(3));
709 assert_eq!(
710 EscalationAction::Reassigned(Assignee::group("sre")),
711 EscalationAction::Reassigned(Assignee::group("sre"))
712 );
713 assert_ne!(EscalationAction::Approved, EscalationAction::Rejected);
714 }
715
716 #[test]
717 fn action_labels_are_distinct() {
718 let labels = [
719 EscalationAction::Approved.label(),
720 EscalationAction::Rejected.label(),
721 EscalationAction::Notified(1).label(),
722 EscalationAction::Reassigned(Assignee::group("sre")).label(),
723 EscalationAction::Exhausted.label(),
724 EscalationAction::Stale.label(),
725 ];
726
727 for (i, a) in labels.iter().enumerate() {
728 for (j, b) in labels.iter().enumerate() {
729 if i != j {
730 assert_ne!(a, b, "labels {i} and {j} collide");
731 }
732 }
733 }
734 }
735
736 #[test]
737 fn a_bare_repeating_policy_stays_at_the_same_stage() {
738 let notify = EscalationPolicy::Notify(Vec::new());
739 assert_eq!(next_stage(¬ify, 0), 0);
740 assert_eq!(next_stage(¬ify, 7), 7);
741
742 let escalate = EscalationPolicy::Escalate(Assignee::group("sre"));
743 assert_eq!(next_stage(&escalate, 3), 3);
744 }
745
746 #[test]
747 fn a_chained_policy_advances_one_stage_per_expiry() {
748 let chain = EscalationPolicy::Chain(vec![
749 EscalationPolicy::Notify(Vec::new()),
750 EscalationPolicy::AutoReject,
751 ]);
752 assert_eq!(next_stage(&chain, 0), 1);
753 assert_eq!(next_stage(&chain, 1), 2);
754 assert_eq!(next_stage(&chain, u32::MAX), u32::MAX);
755 }
756
757 #[test]
758 fn resolve_stage_rejects_a_nested_chain() {
759 let policy = EscalationPolicy::Chain(vec![
760 EscalationPolicy::Chain(vec![EscalationPolicy::AutoReject]),
761 EscalationPolicy::AutoApprove,
762 ]);
763
764 let err = resolve_stage(&policy, 0).expect_err("nested chains are rejected");
765 assert!(
766 err.to_string().contains("nested escalation chains"),
767 "got {err}"
768 );
769
770 assert_eq!(
772 resolve_stage(&policy, 1).expect("valid stage"),
773 Some(&EscalationPolicy::AutoApprove)
774 );
775 }
776
777 #[test]
778 fn resolve_stage_reports_an_exhausted_chain_as_none() {
779 let policy = EscalationPolicy::Chain(vec![EscalationPolicy::AutoReject]);
780 assert!(resolve_stage(&policy, 1).expect("valid call").is_none());
781
782 let bare = EscalationPolicy::AutoApprove;
783 assert!(resolve_stage(&bare, 1).expect("valid call").is_none());
784 assert_eq!(
785 resolve_stage(&bare, 0).expect("valid call"),
786 Some(&EscalationPolicy::AutoApprove)
787 );
788 }
789
790 #[test]
791 fn policy_labels_match_the_wire_format() {
792 assert_eq!(policy_label(&EscalationPolicy::AutoApprove), "auto_approve");
793 assert_eq!(policy_label(&EscalationPolicy::AutoReject), "auto_reject");
794 assert_eq!(
795 policy_label(&EscalationPolicy::Notify(Vec::new())),
796 "notify"
797 );
798 assert_eq!(
799 policy_label(&EscalationPolicy::Escalate(Assignee::group("sre"))),
800 "escalate"
801 );
802 assert_eq!(policy_label(&EscalationPolicy::Chain(Vec::new())), "chain");
803 }
804}