oxedyne/fe2o3/fe2o3_steel/src/srv/stop.rs
10.1 KiB, 17 runs
created by r1870400018:21670, which is this file's identity for as long as the history lasts, whatever it is later renamed to
download · who wrote it · its history
| 1 | //! Steel hearing the operating system ask it to stop. |
| 2 | //! |
| 3 | //! The listening itself belongs to [`oxedyne_fe2o3_core::stop::on_stop_request`], |
| 4 | //! which answers `SIGINT` and `SIGTERM` on unix and the three console events on |
| 5 | //! Windows. What is here is the other half: the state a caught signal leaves |
| 6 | //! behind, and the two ways the rest of Steel reads it. |
| 7 | //! |
| 8 | //! # Why a flag is not enough |
| 9 | //! |
| 10 | //! A server asked to stop is almost always sitting in `accept`, waiting for a |
| 11 | //! connection that may never come. Setting a flag it will look at *next time |
| 12 | //! round* stops nothing, because there is no next time round until somebody |
| 13 | //! opens a socket. So the ask is published twice over: as a counter anything may |
| 14 | //! read at any moment ([`asked`]), and as a bell an asynchronous task can wait |
| 15 | //! on ([`wait`]) inside a `tokio::select!` alongside the accept itself. |
| 16 | //! |
| 17 | //! The order in [`wait`] is load-bearing. Interest is registered with the bell |
| 18 | //! *before* the counter is read, so an ask that lands between the two is heard |
| 19 | //! by the wait rather than slept through. `notify_waiters` stores no permit for |
| 20 | //! a waiter that is not yet there, which is why the counter is written first and |
| 21 | //! the bell rung second. |
| 22 | //! |
| 23 | //! # What a stop means for a process with no server in it |
| 24 | //! |
| 25 | //! Steel is a shell with a server in it, not the other way round: `run` may be |
| 26 | //! executing `wallet --list` or sitting at a prompt, with no accept loop |
| 27 | //! anywhere. Catching a signal there and merely noting it would be worse than |
| 28 | //! not catching it at all -- the process would ignore a `SIGTERM` it used to |
| 29 | //! obey. So [`listen`] ends such a process itself, once the log has been |
| 30 | //! flushed. [`serving`] is what tells the two cases apart. |
| 31 | //! |
| 32 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 33 | //! Anthropic Claude |
| 34 | |
| 35 | use oxedyne_fe2o3_core::prelude::*; |
| 36 | |
| 37 | use std::{ |
| 38 | sync::{ |
| 39 | atomic::{ |
| 40 | AtomicBool, |
| 41 | AtomicUsize, |
| 42 | Ordering, |
| 43 | }, |
| 44 | Arc, |
| 45 | OnceLock, |
| 46 | }, |
| 47 | time::{ |
| 48 | Duration, |
| 49 | Instant, |
| 50 | }, |
| 51 | }; |
| 52 | |
| 53 | use tokio::sync::Notify; |
| 54 | |
| 55 | static ASKS: AtomicUsize = AtomicUsize::new(0); |
| 56 | |
| 57 | static SERVING: AtomicBool = AtomicBool::new(false); |
| 58 | |
| 59 | /// The bell an accept loop waits on, rung once per ask. |
| 60 | /// |
| 61 | /// Built on first use rather than declared, because `Notify::new` is not a |
| 62 | /// constant. `OnceLock` rather than a lazy cell so the crate takes no dependency |
| 63 | /// for one static. |
| 64 | fn bell() -> &'static Notify { |
| 65 | static BELL: OnceLock<Notify> = OnceLock::new(); |
| 66 | BELL.get_or_init(Notify::new) |
| 67 | } |
| 68 | |
| 69 | pub fn asked() -> bool { |
| 70 | ASKS.load(Ordering::SeqCst) > 0 |
| 71 | } |
| 72 | |
| 73 | /// More than one means somebody asked again while the first ask was still being |
| 74 | /// obeyed, which is worth saying out loud even though Steel answers both the |
| 75 | /// same way; see [`listen`]. |
| 76 | pub fn asks() -> usize { |
| 77 | ASKS.load(Ordering::SeqCst) |
| 78 | } |
| 79 | |
| 80 | /// Records an ask and wakes whatever is waiting on one, returning the number of |
| 81 | /// asks so far -- `1` for the first. |
| 82 | /// |
| 83 | /// Safe to call from any thread, and from more than one at once: the counter |
| 84 | /// moves first, so a waiter that misses the bell still sees the count. |
| 85 | pub fn ask() -> usize { |
| 86 | let n = ASKS.fetch_add(1, Ordering::SeqCst) + 1; |
| 87 | bell().notify_waiters(); |
| 88 | n |
| 89 | } |
| 90 | |
| 91 | /// Resolves as soon as this process has been asked to stop, and at once if it |
| 92 | /// already has. |
| 93 | /// |
| 94 | /// Cancel-safe, which is what allows it to sit in a `tokio::select!` opposite an |
| 95 | /// `accept`: dropped part way through, it has consumed nothing, and the next |
| 96 | /// call registers afresh and re-reads the counter. |
| 97 | pub async fn wait() { |
| 98 | loop { |
| 99 | // Registered before the counter is read. The other order drops an ask |
| 100 | // that lands in the gap between the two. |
| 101 | let ringing = bell().notified(); |
| 102 | if asked() { |
| 103 | return; |
| 104 | } |
| 105 | ringing.await; |
| 106 | if asked() { |
| 107 | return; |
| 108 | } |
| 109 | } |
| 110 | } |
| 111 | |
| 112 | /// Declares whether something is running that will wind itself up when asked. |
| 113 | /// |
| 114 | /// Set around the server's whole life -- from before its databases start opening |
| 115 | /// to after they have been closed -- rather than around the accept loop alone, |
| 116 | /// so that a signal arriving while a store is being opened is still answered by |
| 117 | /// the orderly path. |
| 118 | pub fn serving(now: bool) { |
| 119 | SERVING.store(now, Ordering::SeqCst); |
| 120 | } |
| 121 | |
| 122 | pub fn is_serving() -> bool { |
| 123 | SERVING.load(Ordering::SeqCst) |
| 124 | } |
| 125 | |
| 126 | /// Installs this process's stop-request listener. |
| 127 | /// |
| 128 | /// One to a process, because a signal arrives at a process rather than at an |
| 129 | /// object; the underlying listener refuses a second. Called from |
| 130 | /// [`crate::app::tui::run_with_extension`], which is the real `main` of every |
| 131 | /// Steel application, stock or extended. |
| 132 | /// |
| 133 | /// A listener already installed, or a thread that cannot be spawned, is not |
| 134 | /// fatal to a server -- it will simply be killed rather than asked when the |
| 135 | /// machine goes -- so the caller logs and carries on. |
| 136 | pub fn listen() -> Outcome<()> { |
| 137 | res!(oxedyne_fe2o3_core::stop::on_stop_request(|| { |
| 138 | let n = ask(); |
| 139 | if n > 1 { |
| 140 | // Said, and not acted on. A program whose store is a cache can |
| 141 | // read a second Ctrl-C as "now" and leave where it stands, because |
| 142 | // the worst it costs is a rebuild. A Steel store is the site's own |
| 143 | // data, so the firmer ask is refused rather than obeyed: what it |
| 144 | // would interrupt is a database part way through closing, and the |
| 145 | // wait it is impatient with is bounded already -- see |
| 146 | // `srv::server::DRAIN_SECS`. |
| 147 | warn!("Asked to stop {} times. The first ask is being obeyed and \ |
| 148 | the wind-up is bounded; a store is not abandoned part way \ |
| 149 | through closing.", n); |
| 150 | return; |
| 151 | } |
| 152 | info!("Asked to stop."); |
| 153 | if !is_serving() { |
| 154 | info!("Nothing is serving, so there is nothing to wind up."); |
| 155 | if let Err(e) = flush_log() { |
| 156 | error!(e, "Flushing the log on the way out."); |
| 157 | } |
| 158 | // Nought, not 130. This process was not felled: it was asked, and |
| 159 | // it did as it was asked. A service manager reads anything else as |
| 160 | // a unit that failed. |
| 161 | std::process::exit(0); |
| 162 | } |
| 163 | })); |
| 164 | Ok(()) |
| 165 | } |
| 166 | |
| 167 | /// Waits for the logger to finish, so nothing written on the way out is lost. |
| 168 | /// |
| 169 | /// The macro can return an error, and one that can has to sit in a function that |
| 170 | /// can, which the listener's closure is not. |
| 171 | fn flush_log() -> Outcome<()> { |
| 172 | log_finish_wait!(); |
| 173 | Ok(()) |
| 174 | } |
| 175 | |
| 176 | |
| 177 | // ┌───────────────────────────────────────────────────────────────────────────┐ |
| 178 | // │ WORK IN FLIGHT │ |
| 179 | // └───────────────────────────────────────────────────────────────────────────┘ |
| 180 | |
| 181 | /// One piece of work, counted for exactly as long as it lasts. |
| 182 | /// |
| 183 | /// Held by each connection task in the accept loop so that a wind-up can tell |
| 184 | /// the difference between a server with nothing left to do and one part way |
| 185 | /// through a reply. The count falls in `Drop`, so a task that panics is not |
| 186 | /// counted for ever after. |
| 187 | #[derive(Debug)] |
| 188 | pub struct InFlight(Arc<AtomicUsize>); |
| 189 | |
| 190 | impl InFlight { |
| 191 | /// Counts one more piece of work, until the returned value is dropped. |
| 192 | pub fn begin(count: &Arc<AtomicUsize>) -> Self { |
| 193 | count.fetch_add(1, Ordering::SeqCst); |
| 194 | Self(count.clone()) |
| 195 | } |
| 196 | } |
| 197 | |
| 198 | impl Drop for InFlight { |
| 199 | fn drop(&mut self) { |
| 200 | self.0.fetch_sub(1, Ordering::SeqCst); |
| 201 | } |
| 202 | } |
| 203 | |
| 204 | /// Waits for work in flight to finish, for no longer than `bound`, and reports |
| 205 | /// how much of it was still going when the wait ended. |
| 206 | /// |
| 207 | /// A bound rather than a wait, because some of what is counted here does not |
| 208 | /// end on its own: a WebSocket held open by a browser tab is in flight until the |
| 209 | /// tab closes, which may be tomorrow. Everything that does end -- a response |
| 210 | /// being written, a file being read off disk -- ends in milliseconds, so the |
| 211 | /// bound is generous for the case it is for and short enough that a service |
| 212 | /// manager's own patience is never the thing that runs out first. |
| 213 | pub async fn drain(count: &Arc<AtomicUsize>, bound: Duration) -> usize { |
| 214 | let began = Instant::now(); |
| 215 | loop { |
| 216 | let left = count.load(Ordering::SeqCst); |
| 217 | if left == 0 || began.elapsed() >= bound { |
| 218 | return left; |
| 219 | } |
| 220 | tokio::time::sleep(Duration::from_millis(50)).await; |
| 221 | } |
| 222 | } |
| 223 | |
| 224 | |
| 225 | #[cfg(test)] |
| 226 | mod tests { |
| 227 | use super::*; |
| 228 | |
| 229 | /// The count and the bell agree, and a wait that arrives late still returns. |
| 230 | /// |
| 231 | /// What can be checked in-process, and no more. A test cannot send itself a |
| 232 | /// signal and go on being a test; that claim is `tests/stopping_signal.rs`, |
| 233 | /// which starts a real server and sends it a real `SIGTERM`. |
| 234 | #[test] |
| 235 | fn test_an_ask_is_heard_late_00() -> Outcome<()> { |
| 236 | req!(asked(), false, "Nothing has asked yet."); |
| 237 | req!(ask(), 1, "The first ask is the first."); |
| 238 | req!(asked(), true, "An ask was made and not heard."); |
| 239 | req!(ask(), 2, "A second ask counts separately."); |
| 240 | |
| 241 | // Registered after the fact, and returns anyway. |
| 242 | let rt = res!(tokio::runtime::Builder::new_current_thread() |
| 243 | .enable_all() |
| 244 | .build(), Init, System); |
| 245 | let returned = rt.block_on(async { |
| 246 | tokio::time::timeout(Duration::from_secs(5), wait()).await |
| 247 | }); |
| 248 | req!(returned.is_ok(), true, |
| 249 | "A wait begun after the ask never returned."); |
| 250 | Ok(()) |
| 251 | } |
| 252 | |
| 253 | /// Work in flight is counted while it lasts and not afterwards. |
| 254 | #[test] |
| 255 | fn test_work_in_flight_is_counted_01() -> Outcome<()> { |
| 256 | let count = Arc::new(AtomicUsize::new(0)); |
| 257 | { |
| 258 | let _one = InFlight::begin(&count); |
| 259 | let _two = InFlight::begin(&count); |
| 260 | req!(count.load(Ordering::SeqCst), 2, "Two pieces of work."); |
| 261 | } |
| 262 | req!(count.load(Ordering::SeqCst), 0, "Both have finished."); |
| 263 | Ok(()) |
| 264 | } |
| 265 | } |