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 approved_by: SYSTEM_TIMEOUT_ACTOR.to_string(),
403 at: now,
404 }));
405
406 if let Err(err) = self.engine.resume_run(step.run_id).await {
410 error!(
411 run_id = %step.run_id,
412 error = %err,
413 "failed to resume run after an auto-approved gate"
414 );
415 }
416
417 Ok(EscalationAction::Approved)
418 }
419
420 async fn auto_reject(&self, step: &Step) -> Result<EscalationAction, EngineError> {
422 let now = Utc::now();
423 let store = self.engine.store();
424
425 store
426 .update_step(
427 step.id,
428 StepUpdate {
429 status: Some(StepStatus::Failed),
430 error: Some(APPROVAL_TIMEOUT_ERROR.to_string()),
431 completed_at: Some(now),
432 clear_approval_deadline: true,
433 ..StepUpdate::default()
434 },
435 )
436 .await?;
437 store
438 .update_run(
439 step.run_id,
440 RunUpdate {
441 status: Some(RunStatus::Failed),
442 error: Some(APPROVAL_TIMEOUT_ERROR.to_string()),
443 completed_at: Some(now),
444 ..RunUpdate::default()
445 },
446 )
447 .await?;
448
449 self.engine
450 .event_publisher()
451 .publish(Event::ApprovalRejected(ApprovalRejectedEvent {
452 run_id: step.run_id,
453 rejected_by: SYSTEM_TIMEOUT_ACTOR.to_string(),
454 at: now,
455 }));
456
457 Ok(EscalationAction::Rejected)
458 }
459
460 async fn notify(
462 &self,
463 step: &Step,
464 config: &ApprovalConfig,
465 policy: &EscalationPolicy,
466 targets: &[NotificationTarget],
467 reason: &str,
468 ) -> Result<EscalationAction, EngineError> {
469 let event = self.escalated_event(
470 step,
471 step.approval_stage,
472 "notify",
473 EscalationAction::Notified(targets.len()).label(),
474 reason,
475 step.approval_assignee.clone(),
476 );
477 self.deliver_notifications(targets, &event).await;
478
479 self.rearm(step, config, policy, None).await?;
480 Ok(EscalationAction::Notified(targets.len()))
481 }
482
483 async fn reassign(
485 &self,
486 step: &Step,
487 config: &ApprovalConfig,
488 policy: &EscalationPolicy,
489 assignee: &Assignee,
490 ) -> Result<EscalationAction, EngineError> {
491 self.rearm(step, config, policy, Some(assignee)).await?;
492 Ok(EscalationAction::Reassigned(assignee.clone()))
493 }
494
495 async fn rearm(
501 &self,
502 step: &Step,
503 config: &ApprovalConfig,
504 policy: &EscalationPolicy,
505 assignee: Option<&Assignee>,
506 ) -> Result<(), EngineError> {
507 let next_stage = next_stage(policy, step.approval_stage);
508 let deadline = config
509 .effective_deadline_secs()
510 .map(|secs| Utc::now() + TimeDelta::seconds(secs as i64));
511
512 self.engine
513 .store()
514 .update_step(
515 step.id,
516 StepUpdate {
517 approval_deadline_at: deadline,
518 approval_stage: Some(next_stage),
519 approval_assignee: assignee.cloned(),
520 ..StepUpdate::default()
521 },
522 )
523 .await?;
524
525 Ok(())
526 }
527
528 fn escalated_event(
530 &self,
531 step: &Step,
532 stage: u32,
533 policy: &str,
534 action: &str,
535 reason: &str,
536 assignee: Option<Assignee>,
537 ) -> ApprovalEscalatedEvent {
538 ApprovalEscalatedEvent {
539 run_id: step.run_id,
540 step_id: step.id,
541 step_name: step.name.clone(),
542 stage,
543 policy: policy.to_string(),
544 action: action.to_string(),
545 reason: reason.to_string(),
546 assignee,
547 at: Utc::now(),
548 }
549 }
550
551 fn publish(
553 &self,
554 step: &Step,
555 stage: u32,
556 policy: &str,
557 action: &EscalationAction,
558 reason: &str,
559 assignee: Option<Assignee>,
560 ) {
561 let event = self.escalated_event(step, stage, policy, action.label(), reason, assignee);
562
563 info!(
564 run_id = %step.run_id,
565 step_id = %step.id,
566 stage,
567 policy,
568 action = action.label(),
569 reason,
570 "approval gate escalated"
571 );
572
573 self.engine
574 .event_publisher()
575 .publish(Event::ApprovalEscalated(event));
576 }
577
578 async fn deliver_notifications(
584 &self,
585 targets: &[NotificationTarget],
586 event: &ApprovalEscalatedEvent,
587 ) {
588 for target in targets {
589 match target {
590 NotificationTarget::Webhook { url } => {
591 deliver_with_retry(
592 &self.retry,
593 || self.client.post(url).json(event),
594 is_success_2xx,
595 "approval_escalation",
596 url,
597 )
598 .await;
599 }
600 NotificationTarget::Slack {
601 webhook_url,
602 channel,
603 } => {
604 let body = json!({
605 "channel": channel,
606 "text": format!(
607 "⏰ Approval gate `{}` on run {} hit its SLA: {}",
608 event.step_name, event.run_id, event.reason
609 ),
610 });
611 deliver_with_retry(
612 &self.retry,
613 || self.client.post(webhook_url).json(&body),
614 is_success_2xx,
615 "approval_escalation",
616 webhook_url,
617 )
618 .await;
619 }
620 }
621 }
622 }
623}
624
625fn resolve_stage(
630 policy: &EscalationPolicy,
631 stage: u32,
632) -> Result<Option<&EscalationPolicy>, EngineError> {
633 match policy.stage(stage as usize) {
634 Some(EscalationPolicy::Chain(_)) => Err(EngineError::StepConfig(
635 "nested escalation chains are not supported".to_string(),
636 )),
637 other => Ok(other),
638 }
639}
640
641fn next_stage(policy: &EscalationPolicy, stage: u32) -> u32 {
646 match policy {
647 EscalationPolicy::Chain(_) => stage.saturating_add(1),
648 _ => stage,
649 }
650}
651
652fn policy_label(policy: &EscalationPolicy) -> &'static str {
654 match policy {
655 EscalationPolicy::AutoApprove => "auto_approve",
656 EscalationPolicy::AutoReject => "auto_reject",
657 EscalationPolicy::Notify(_) => "notify",
658 EscalationPolicy::Escalate(_) => "escalate",
659 EscalationPolicy::Chain(_) => "chain",
660 }
661}
662
663#[cfg(test)]
664mod tests {
665 use super::*;
666
667 #[test]
668 fn actions_compare_by_value() {
669 assert_eq!(EscalationAction::Notified(2), EscalationAction::Notified(2));
670 assert_ne!(EscalationAction::Notified(2), EscalationAction::Notified(3));
671 assert_eq!(
672 EscalationAction::Reassigned(Assignee::group("sre")),
673 EscalationAction::Reassigned(Assignee::group("sre"))
674 );
675 assert_ne!(EscalationAction::Approved, EscalationAction::Rejected);
676 }
677
678 #[test]
679 fn action_labels_are_distinct() {
680 let labels = [
681 EscalationAction::Approved.label(),
682 EscalationAction::Rejected.label(),
683 EscalationAction::Notified(1).label(),
684 EscalationAction::Reassigned(Assignee::group("sre")).label(),
685 EscalationAction::Exhausted.label(),
686 EscalationAction::Stale.label(),
687 ];
688
689 for (i, a) in labels.iter().enumerate() {
690 for (j, b) in labels.iter().enumerate() {
691 if i != j {
692 assert_ne!(a, b, "labels {i} and {j} collide");
693 }
694 }
695 }
696 }
697
698 #[test]
699 fn a_bare_repeating_policy_stays_at_the_same_stage() {
700 let notify = EscalationPolicy::Notify(Vec::new());
701 assert_eq!(next_stage(¬ify, 0), 0);
702 assert_eq!(next_stage(¬ify, 7), 7);
703
704 let escalate = EscalationPolicy::Escalate(Assignee::group("sre"));
705 assert_eq!(next_stage(&escalate, 3), 3);
706 }
707
708 #[test]
709 fn a_chained_policy_advances_one_stage_per_expiry() {
710 let chain = EscalationPolicy::Chain(vec![
711 EscalationPolicy::Notify(Vec::new()),
712 EscalationPolicy::AutoReject,
713 ]);
714 assert_eq!(next_stage(&chain, 0), 1);
715 assert_eq!(next_stage(&chain, 1), 2);
716 assert_eq!(next_stage(&chain, u32::MAX), u32::MAX);
717 }
718
719 #[test]
720 fn resolve_stage_rejects_a_nested_chain() {
721 let policy = EscalationPolicy::Chain(vec![
722 EscalationPolicy::Chain(vec![EscalationPolicy::AutoReject]),
723 EscalationPolicy::AutoApprove,
724 ]);
725
726 let err = resolve_stage(&policy, 0).expect_err("nested chains are rejected");
727 assert!(
728 err.to_string().contains("nested escalation chains"),
729 "got {err}"
730 );
731
732 assert_eq!(
734 resolve_stage(&policy, 1).expect("valid stage"),
735 Some(&EscalationPolicy::AutoApprove)
736 );
737 }
738
739 #[test]
740 fn resolve_stage_reports_an_exhausted_chain_as_none() {
741 let policy = EscalationPolicy::Chain(vec![EscalationPolicy::AutoReject]);
742 assert!(resolve_stage(&policy, 1).expect("valid call").is_none());
743
744 let bare = EscalationPolicy::AutoApprove;
745 assert!(resolve_stage(&bare, 1).expect("valid call").is_none());
746 assert_eq!(
747 resolve_stage(&bare, 0).expect("valid call"),
748 Some(&EscalationPolicy::AutoApprove)
749 );
750 }
751
752 #[test]
753 fn policy_labels_match_the_wire_format() {
754 assert_eq!(policy_label(&EscalationPolicy::AutoApprove), "auto_approve");
755 assert_eq!(policy_label(&EscalationPolicy::AutoReject), "auto_reject");
756 assert_eq!(
757 policy_label(&EscalationPolicy::Notify(Vec::new())),
758 "notify"
759 );
760 assert_eq!(
761 policy_label(&EscalationPolicy::Escalate(Assignee::group("sre"))),
762 "escalate"
763 );
764 assert_eq!(policy_label(&EscalationPolicy::Chain(Vec::new())), "chain");
765 }
766}