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