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