lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

signal.rs (7601B)


      1 //! Process-signal sequencing over a binary-owned, injected signal source.
      2 
      3 use core::{fmt, future::Future, pin::Pin};
      4 use std::error::Error;
      5 
      6 /// One normalized process signal supplied by a consuming binary.
      7 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
      8 pub enum ProcessSignal {
      9     /// Interactive interrupt, conventionally produced by Ctrl-C.
     10     Interrupt,
     11     /// Unix termination request, conventionally produced by SIGTERM.
     12     #[cfg(unix)]
     13     Terminate,
     14 }
     15 
     16 impl ProcessSignal {
     17     #[must_use]
     18     pub const fn as_str(self) -> &'static str {
     19         match self {
     20             Self::Interrupt => "interrupt",
     21             #[cfg(unix)]
     22             Self::Terminate => "terminate",
     23         }
     24     }
     25 }
     26 
     27 impl fmt::Display for ProcessSignal {
     28     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
     29         formatter.write_str(self.as_str())
     30     }
     31 }
     32 
     33 /// Binary-owned asynchronous source of normalized process signals.
     34 ///
     35 /// Implementations install any operating-system handlers at the binary
     36 /// boundary. The service-host library only consumes the resulting events.
     37 pub type ProcessSignalFuture<'a> = Pin<Box<dyn Future<Output = Option<ProcessSignal>> + Send + 'a>>;
     38 
     39 pub trait ProcessSignalSource: Send {
     40     fn next_signal(&mut self) -> ProcessSignalFuture<'_>;
     41 }
     42 
     43 /// Action assigned to an observed signal by the shared sequencing contract.
     44 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
     45 pub enum ProcessSignalAction {
     46     BeginGracefulShutdown { signal: ProcessSignal },
     47     ForceTermination { signal: ProcessSignal },
     48 }
     49 
     50 impl ProcessSignalAction {
     51     #[must_use]
     52     pub const fn signal(self) -> ProcessSignal {
     53         match self {
     54             Self::BeginGracefulShutdown { signal } | Self::ForceTermination { signal } => signal,
     55         }
     56     }
     57 
     58     #[must_use]
     59     pub const fn forces_termination(self) -> bool {
     60         matches!(self, Self::ForceTermination { .. })
     61     }
     62 }
     63 
     64 /// Which event the adapter was awaiting when its source closed.
     65 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
     66 pub enum ProcessSignalStage {
     67     FirstCancellation,
     68     ForceTermination,
     69 }
     70 
     71 /// Typed closure of an injected signal source before the expected event.
     72 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
     73 pub struct ProcessSignalSourceClosed {
     74     stage: ProcessSignalStage,
     75 }
     76 
     77 impl ProcessSignalSourceClosed {
     78     #[must_use]
     79     pub const fn stage(self) -> ProcessSignalStage {
     80         self.stage
     81     }
     82 }
     83 
     84 impl fmt::Display for ProcessSignalSourceClosed {
     85     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
     86         formatter.write_str("process signal source closed before the expected event")
     87     }
     88 }
     89 
     90 impl Error for ProcessSignalSourceClosed {}
     91 
     92 /// Maps the first injected signal to graceful cancellation and later signals to force.
     93 pub struct ProcessSignalAdapter<S> {
     94     source: S,
     95     first_observed: bool,
     96 }
     97 
     98 impl<S> ProcessSignalAdapter<S>
     99 where
    100     S: ProcessSignalSource,
    101 {
    102     #[must_use]
    103     pub const fn new(source: S) -> Self {
    104         Self {
    105             source,
    106             first_observed: false,
    107         }
    108     }
    109 
    110     #[must_use]
    111     pub const fn stage(&self) -> ProcessSignalStage {
    112         if self.first_observed {
    113             ProcessSignalStage::ForceTermination
    114         } else {
    115             ProcessSignalStage::FirstCancellation
    116         }
    117     }
    118 
    119     /// Waits for and classifies the next binary-supplied process signal.
    120     ///
    121     /// Cancelling this future does not advance adapter state before the source
    122     /// returns an event. Source-specific event cancellation semantics remain
    123     /// the responsibility of the binary-owned source implementation.
    124     pub async fn next_action(&mut self) -> Result<ProcessSignalAction, ProcessSignalSourceClosed> {
    125         let stage = self.stage();
    126         let signal = self
    127             .source
    128             .next_signal()
    129             .await
    130             .ok_or(ProcessSignalSourceClosed { stage })?;
    131         if self.first_observed {
    132             Ok(ProcessSignalAction::ForceTermination { signal })
    133         } else {
    134             self.first_observed = true;
    135             Ok(ProcessSignalAction::BeginGracefulShutdown { signal })
    136         }
    137     }
    138 
    139     #[must_use]
    140     pub fn into_inner(self) -> S {
    141         self.source
    142     }
    143 }
    144 
    145 #[cfg(test)]
    146 mod tests {
    147     use std::collections::VecDeque;
    148 
    149     use super::*;
    150 
    151     struct InjectedSignals {
    152         events: VecDeque<ProcessSignal>,
    153     }
    154 
    155     impl InjectedSignals {
    156         fn new(events: impl IntoIterator<Item = ProcessSignal>) -> Self {
    157             Self {
    158                 events: events.into_iter().collect(),
    159             }
    160         }
    161     }
    162 
    163     impl ProcessSignalSource for InjectedSignals {
    164         fn next_signal(&mut self) -> ProcessSignalFuture<'_> {
    165             Box::pin(async move { self.events.pop_front() })
    166         }
    167     }
    168 
    169     #[tokio::test]
    170     async fn first_event_cancels_and_every_later_event_forces() {
    171         #[cfg(unix)]
    172         let events = [
    173             ProcessSignal::Interrupt,
    174             ProcessSignal::Terminate,
    175             ProcessSignal::Interrupt,
    176         ];
    177         #[cfg(not(unix))]
    178         let events = [
    179             ProcessSignal::Interrupt,
    180             ProcessSignal::Interrupt,
    181             ProcessSignal::Interrupt,
    182         ];
    183         let mut adapter = ProcessSignalAdapter::new(InjectedSignals::new(events));
    184 
    185         assert_eq!(adapter.stage(), ProcessSignalStage::FirstCancellation);
    186         assert_eq!(
    187             adapter.next_action().await.unwrap(),
    188             ProcessSignalAction::BeginGracefulShutdown {
    189                 signal: ProcessSignal::Interrupt,
    190             }
    191         );
    192         assert_eq!(adapter.stage(), ProcessSignalStage::ForceTermination);
    193 
    194         let second = adapter.next_action().await.unwrap();
    195         assert!(second.forces_termination());
    196         #[cfg(unix)]
    197         assert_eq!(second.signal(), ProcessSignal::Terminate);
    198         #[cfg(not(unix))]
    199         assert_eq!(second.signal(), ProcessSignal::Interrupt);
    200 
    201         assert!(adapter.next_action().await.unwrap().forces_termination());
    202     }
    203 
    204     #[tokio::test]
    205     async fn source_closure_is_typed_and_does_not_advance_the_stage() {
    206         let mut before_first = ProcessSignalAdapter::new(InjectedSignals::new([]));
    207         let error = before_first.next_action().await.unwrap_err();
    208         assert_eq!(error.stage(), ProcessSignalStage::FirstCancellation);
    209         assert_eq!(
    210             error.to_string(),
    211             "process signal source closed before the expected event"
    212         );
    213         assert_eq!(before_first.stage(), ProcessSignalStage::FirstCancellation);
    214 
    215         let mut before_force =
    216             ProcessSignalAdapter::new(InjectedSignals::new([ProcessSignal::Interrupt]));
    217         assert!(
    218             !before_force
    219                 .next_action()
    220                 .await
    221                 .unwrap()
    222                 .forces_termination()
    223         );
    224         let error = before_force.next_action().await.unwrap_err();
    225         assert_eq!(error.stage(), ProcessSignalStage::ForceTermination);
    226         assert_eq!(before_force.stage(), ProcessSignalStage::ForceTermination);
    227     }
    228 
    229     #[test]
    230     fn normalized_names_are_stable_and_platform_bounded() {
    231         assert_eq!(ProcessSignal::Interrupt.as_str(), "interrupt");
    232         assert_eq!(ProcessSignal::Interrupt.to_string(), "interrupt");
    233         #[cfg(unix)]
    234         {
    235             assert_eq!(ProcessSignal::Terminate.as_str(), "terminate");
    236             assert_eq!(ProcessSignal::Terminate.to_string(), "terminate");
    237         }
    238 
    239         let adapter = ProcessSignalAdapter::new(InjectedSignals::new([ProcessSignal::Interrupt]));
    240         assert_eq!(adapter.into_inner().events.len(), 1);
    241     }
    242 }