Oregami
Repositories/oxedyne/fe2o3

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
35use oxedyne_fe2o3_core::prelude::*;
36
37use 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
53use tokio::sync::Notify;
54
55static ASKS: AtomicUsize = AtomicUsize::new(0);
56
57static 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.
64fn bell() -> &'static Notify {
65 static BELL: OnceLock<Notify> = OnceLock::new();
66 BELL.get_or_init(Notify::new)
67}
68
69pub 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`].
76pub 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.
85pub 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.
97pub 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.
118pub fn serving(now: bool) {
119 SERVING.store(now, Ordering::SeqCst);
120}
121
122pub 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.
136pub 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.
171fn 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)]
188pub struct InFlight(Arc<AtomicUsize>);
189
190impl 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
198impl 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.
213pub 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)]
226mod 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}