Skip to main content

runifold_workflow/worker/
supervisor.rs

1use 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    /// Creates a continuous supervisor around one shareable worker.
10    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    /// Overrides supervisor sleeping for deterministic runtimes and tests.
20    #[must_use]
21    pub fn with_sleeper(mut self, sleeper: Arc<dyn WorkflowWorkerSleeper>) -> Self {
22        self.sleeper = sleeper;
23        self
24    }
25
26    /// Uses the supplied cumulative metric set.
27    #[must_use]
28    pub fn with_metrics(mut self, metrics: WorkflowSupervisorMetrics) -> Self {
29        self.metrics = metrics;
30        self
31    }
32
33    /// Returns the cumulative metric set used by this supervisor.
34    pub const fn metrics(&self) -> &WorkflowSupervisorMetrics {
35        &self.metrics
36    }
37
38    /// Continuously polls and executes work until shutdown, then drains started cycles.
39    ///
40    /// Infrastructure errors are counted and retried with backoff. Once shutdown is
41    /// observed, no replacement cycles are scheduled. Already-started claim or
42    /// execution operations are awaited so their lease protocol can finish safely.
43    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}