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}