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 }