saluki_core/runtime/
supervisable.rs

1//! The contract that a supervised process implements.
2//!
3//! A supervisor drives a [`Supervisable`]. The [`Supervisable`] names the process, tells how to shut it down, and
4//! builds the future of the process each time the process starts. This module also contains the types in its signature.
5
6use std::{future::Future, pin::Pin, time::Duration};
7
8use async_trait::async_trait;
9use saluki_common::sync::shutdown::ShutdownHandle;
10use saluki_error::GenericError;
11use snafu::Snafu;
12
13/// A `Future` that represents the execution of a supervised process.
14pub type SupervisorFuture = Pin<Box<dyn Future<Output = Result<(), GenericError>> + Send>>;
15
16/// Initialization errors.
17///
18/// Initialization errors are distinct from runtime errors: they indicate that a process couldn't be started at all
19/// (for example, failed to bind a port, missing configuration). These errors don't trigger restart logic; instead, they
20/// immediately propagate up and fail the supervisor.
21#[derive(Debug, Snafu)]
22#[snafu(context(suffix(false)))]
23pub enum InitializationError {
24    /// The process couldn't be initialized due to an error.
25    #[snafu(display("Process failed to initialize: {}", source))]
26    Failed {
27        /// The underlying error that caused initialization to fail.
28        source: GenericError,
29    },
30}
31
32impl From<GenericError> for InitializationError {
33    fn from(source: GenericError) -> Self {
34        Self::Failed { source }
35    }
36}
37
38/// Strategy for shutting down a process.
39#[derive(Clone, Copy, Debug)]
40pub enum ShutdownStrategy {
41    /// Waits for the configured duration for the process to exit, and then forcefully aborts it otherwise.
42    Graceful(Duration),
43
44    /// Forcefully aborts the process without waiting.
45    Brutal,
46}
47
48/// A supervisable process.
49#[async_trait]
50pub trait Supervisable: Send + Sync {
51    /// Returns the name of the process.
52    fn name(&self) -> &str;
53
54    /// Returns the shutdown strategy for the process.
55    fn shutdown_strategy(&self) -> ShutdownStrategy {
56        ShutdownStrategy::Graceful(Duration::from_secs(5))
57    }
58
59    /// Returns whether this process observes the shutdown signal it is given.
60    ///
61    /// Shutting a subtree down is a _trigger_, not an enforcement: many workers ignore the signal entirely and stop
62    /// only when they reach their own terminal condition, such as an input channel closing. Reporting `false` lets the
63    /// supervisor skip creating a shutdown coordinator it would never usefully fire, and hand the process a
64    /// [`ShutdownHandle::noop`] instead.
65    ///
66    /// This says nothing about _whether_ the supervisor waits for the process -- that's
67    /// [`shutdown_strategy`][Self::shutdown_strategy]. A process that ignores the signal is still waited for, up to
68    /// whatever deadline applies to it.
69    ///
70    /// Defaults to `true`.
71    fn wants_shutdown_signal(&self) -> bool {
72        true
73    }
74
75    /// Initializes the process asynchronously.
76    ///
77    /// During initialization, any resources or configuration for the process can be created asynchronously, and the
78    /// same runtime that's used for running the process is used for initialization. The resulting future is expected to
79    /// complete as soon as reasonably possible after `shutdown` resolves.
80    ///
81    /// **Important:** The `process_shutdown` signal must be moved into the returned [`SupervisorFuture`] so the worker
82    /// can respond to supervisor-initiated shutdown. If `process_shutdown` is dropped during initialization, the worker
83    /// will be unable to shut down gracefully and will be forcefully aborted after the shutdown timeout.
84    ///
85    /// # Errors
86    ///
87    /// If the process can't be initialized, an error is returned.
88    async fn initialize(&self, process_shutdown: ShutdownHandle) -> Result<SupervisorFuture, InitializationError>;
89}
90
91#[cfg(test)]
92mod tests {
93    use super::*;
94
95    /// The smallest possible implementation, existing only to be erased below.
96    struct Noop;
97
98    #[async_trait]
99    impl Supervisable for Noop {
100        fn name(&self) -> &str {
101            "noop"
102        }
103
104        async fn initialize(&self, process_shutdown: ShutdownHandle) -> Result<SupervisorFuture, InitializationError> {
105            Ok(Box::pin(async move {
106                process_shutdown.await;
107                Ok(())
108            }))
109        }
110    }
111
112    #[test]
113    fn trait_is_object_safe() {
114        // Producers hand their background work over as `Vec<Box<dyn Supervisable>>`, so object safety is load-bearing
115        // here rather than incidental: a default method taking `self` by value, or a generic one, would break every
116        // one of those call sites.
117        let worker: Box<dyn Supervisable> = Box::new(Noop);
118
119        assert_eq!(worker.name(), "noop");
120        assert!(worker.wants_shutdown_signal());
121        assert!(matches!(
122            worker.shutdown_strategy(),
123            ShutdownStrategy::Graceful(timeout) if timeout == Duration::from_secs(5)
124        ));
125    }
126}