runifold_workflow/worker/
supervisor.rs1use super::{
2 Arc, CancellationToken, Duration, Either, Future, FuturesUnordered, Pin, StreamExt,
3 SystemWorkflowWorkerSleeper, WorkflowSupervisor, WorkflowSupervisorConfig,
4 WorkflowSupervisorMetrics, WorkflowSupervisorReport, WorkflowWorker, WorkflowWorkerError,
5 WorkflowWorkerOutcome, WorkflowWorkerSleeper, select,
6};
7
8impl WorkflowSupervisor {
9 pub fn new(worker: Arc<WorkflowWorker>, config: WorkflowSupervisorConfig) -> Self {
11 Self {
12 worker,
13 config,
14 sleeper: Arc::new(SystemWorkflowWorkerSleeper),
15 metrics: WorkflowSupervisorMetrics::default(),
16 }
17 }
18
19 #[must_use]
21 pub fn with_sleeper(mut self, sleeper: Arc<dyn WorkflowWorkerSleeper>) -> Self {
22 self.sleeper = sleeper;
23 self
24 }
25
26 #[must_use]
28 pub fn with_metrics(mut self, metrics: WorkflowSupervisorMetrics) -> Self {
29 self.metrics = metrics;
30 self
31 }
32
33 pub const fn metrics(&self) -> &WorkflowSupervisorMetrics {
35 &self.metrics
36 }
37
38 pub async fn run(&self, shutdown: &CancellationToken) -> WorkflowSupervisorReport {
44 let mut cycles = FuturesUnordered::new();
45 for _ in 0..self.config.max_concurrency {
46 self.schedule(&mut cycles, Duration::ZERO, shutdown.clone());
47 }
48
49 let mut report = WorkflowSupervisorReport::default();
50 let mut next_backoff = self.config.initial_backoff;
51 while let Either::Right((Some(result), _)) =
52 select(Box::pin(shutdown.cancelled()), Box::pin(cycles.next())).await
53 {
54 self.metrics.cycle_stopped();
55 let delay = self.observe_cycle(&mut report, &result, &mut next_backoff);
56 if !matches!(result, WorkflowSupervisorCycleResult::Stopped) {
57 self.schedule(&mut cycles, delay, shutdown.clone());
58 }
59 }
60
61 while let Some(result) = cycles.next().await {
62 self.metrics.cycle_stopped();
63 self.observe_cycle(&mut report, &result, &mut next_backoff);
64 }
65 report
66 }
67
68 fn schedule(
69 &self,
70 cycles: &mut FuturesUnordered<WorkflowSupervisorCycleFuture>,
71 delay: Duration,
72 shutdown: CancellationToken,
73 ) {
74 if !delay.is_zero() {
75 self.metrics.record_backoff();
76 }
77 self.metrics.cycle_started();
78 cycles.push(Box::pin(run_supervisor_cycle(
79 self.worker.clone(),
80 self.sleeper.clone(),
81 delay,
82 shutdown,
83 )));
84 }
85
86 fn observe_cycle(
87 &self,
88 report: &mut WorkflowSupervisorReport,
89 result: &WorkflowSupervisorCycleResult,
90 next_backoff: &mut Duration,
91 ) -> Duration {
92 let WorkflowSupervisorCycleResult::Finished(result) = result else {
93 return Duration::ZERO;
94 };
95 self.metrics.record_result(result);
96 report.record(result);
97 if matches!(result, Ok(WorkflowWorkerOutcome::Idle) | Err(_)) {
98 let delay = *next_backoff;
99 *next_backoff = next_backoff.saturating_mul(2).min(self.config.max_backoff);
100 delay
101 } else {
102 *next_backoff = self.config.initial_backoff;
103 Duration::ZERO
104 }
105 }
106}
107
108impl std::fmt::Debug for WorkflowSupervisor {
109 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
110 formatter
111 .debug_struct("WorkflowSupervisor")
112 .field("worker", &self.worker)
113 .field("config", &self.config)
114 .field("metrics", &self.metrics.snapshot())
115 .finish_non_exhaustive()
116 }
117}
118
119impl WorkflowSupervisorReport {
120 fn record(&mut self, result: &Result<WorkflowWorkerOutcome, WorkflowWorkerError>) {
121 match result {
122 Ok(WorkflowWorkerOutcome::Idle) => self.idle_polls += 1,
123 Ok(WorkflowWorkerOutcome::Completed { .. }) => self.completed += 1,
124 Ok(WorkflowWorkerOutcome::Retried { .. }) => self.retried += 1,
125 Ok(WorkflowWorkerOutcome::Suspended { .. }) => self.suspended += 1,
126 Ok(WorkflowWorkerOutcome::Failed { .. }) => self.failed += 1,
127 Ok(WorkflowWorkerOutcome::DefinitionUnavailable { .. }) => {
128 self.definitions_unavailable += 1;
129 }
130 Ok(WorkflowWorkerOutcome::LeaseLost { .. }) => self.leases_lost += 1,
131 Err(_) => self.infrastructure_errors += 1,
132 }
133 }
134}
135
136#[cfg(not(target_arch = "wasm32"))]
137type WorkflowSupervisorCycleFuture =
138 Pin<Box<dyn Future<Output = WorkflowSupervisorCycleResult> + Send>>;
139
140#[cfg(target_arch = "wasm32")]
141type WorkflowSupervisorCycleFuture = Pin<Box<dyn Future<Output = WorkflowSupervisorCycleResult>>>;
142
143enum WorkflowSupervisorCycleResult {
144 Finished(Result<WorkflowWorkerOutcome, WorkflowWorkerError>),
145 Stopped,
146}
147
148async fn run_supervisor_cycle(
149 worker: Arc<WorkflowWorker>,
150 sleeper: Arc<dyn WorkflowWorkerSleeper>,
151 delay: Duration,
152 shutdown: CancellationToken,
153) -> WorkflowSupervisorCycleResult {
154 if !delay.is_zero() {
155 match select(
156 Box::pin(shutdown.cancelled()),
157 Box::pin(sleeper.sleep(delay)),
158 )
159 .await
160 {
161 Either::Left(_) => return WorkflowSupervisorCycleResult::Stopped,
162 Either::Right(_) => {}
163 }
164 }
165 if shutdown.is_cancelled() {
166 return WorkflowSupervisorCycleResult::Stopped;
167 }
168 WorkflowSupervisorCycleResult::Finished(worker.run_once().await)
169}