nmbrs_runtime/session_signals.rs
1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Session-wide signal handling — the THREE-LEVEL shutdown ladder.
5//!
6//! A Ctrl-C (SIGINT or SIGTERM, or the raw-mode key-watcher's
7//! translated keystroke) advances one rung at a time:
8//!
9//! - **Level 1 — graceful, cooperative.** The session-stop flag is
10//! set; active fiber loops observe it at their cycle boundary and
11//! exit cleanly; end-of-run cleanup runs in the normal order
12//! (profiler flush, cadence reporter shutdown, metrics.db WAL
13//! consolidation, summary writes). A visible 10-second countdown
14//! starts: if the drain hasn't finished when it expires, the
15//! ladder advances to level 2 automatically.
16//! - **Level 2 — cancel in-flight ops, keep process cleanup.**
17//! Entered by the countdown expiring or a second Ctrl-C. Ops
18//! parked inside a hung adapter call (a request that will only
19//! ever end by client timeout) are CANCELLED — their futures are
20//! dropped at the fiber's dispatch point — so the drain completes
21//! and the process-level graceful shutdown (WAL consolidation,
22//! summaries) still runs. This is the rung that used to not
23//! exist: previously a second Ctrl-C hard-exited, skipping WAL
24//! cleanup exactly when hung ops made it matter.
25//! - **Level 3 — force-exit.** A further Ctrl-C once level 2 is in
26//! force exits immediately (`process::exit(130)`); metrics and
27//! profiler output may be incomplete.
28//!
29//! The state is intentionally global: there is one session per
30//! process by construction. Tests shouldn't need to install or
31//! consult it (no `RunObserver` test sets up signals).
32//!
33//! ## SRD-93 M7 — the detached-console signal contract
34//!
35//! The ladder MUST be reachable from outside the process even when
36//! the async runtime is wedged (SRD-93 A8: the stop door never
37//! depends on runtime health — the 2026-08-03 incident proved a
38//! tokio-task watcher is unreachable exactly when it matters). So
39//! signal watching runs on a dedicated OS thread: `main` calls
40//! [`block_shutdown_signals`] before any thread spawns (every later
41//! thread inherits the block) and [`spawn_signal_dispatcher`] owns
42//! the signals via `sigwait`, advancing the ladder synchronously.
43//!
44//! - `SIGINT` / `SIGTERM` — one rung per signal; past the cancel
45//! rung the dispatcher force-exits `128 + signo` (130 / 143) on
46//! its own thread, so the hard floor works even when nothing
47//! else does. Unarmed (no run in progress) they exit
48//! `128 + signo` directly — CLI commands keep their default
49//! interrupt semantics.
50//! - `SIGHUP` — NEVER a stop. Armed: console-loss (log, run the
51//! [`set_console_loss_hook`] if installed, continue headless).
52//! Unarmed: exit 129.
53//! - `SIGQUIT` — diagnostics dump (ladder level, stop flags, plus
54//! the [`set_diag_dump_hook`] inventory if installed); no state
55//! change. Unarmed: exit 131.
56//!
57//! The tokio `ctrl_c` watcher in [`install_signal_handler`] remains
58//! as a redundant door for embeddings whose `main` never blocked
59//! the signals; when the dispatcher owns them, that watcher simply
60//! never fires (the two are mutually exclusive by construction).
61
62use std::sync::Arc;
63use std::sync::OnceLock;
64use std::sync::atomic::{AtomicBool, Ordering};
65
66/// Shared session-stop flag. Initialized lazily on first read or
67/// the call to [`install_signal_handler`].
68static SESSION_STOP: OnceLock<Arc<AtomicBool>> = OnceLock::new();
69
70fn flag() -> &'static Arc<AtomicBool> {
71 SESSION_STOP.get_or_init(|| Arc::new(AtomicBool::new(false)))
72}
73
74/// Returns `true` once a stop has been requested for the current execution.
75/// Cheap relaxed atomic load — safe to call from a hot fiber loop.
76///
77/// SRD-88: a fiber running inside an [`ExecutionContext`](crate::execution_context)
78/// stops when EITHER the process-global session stop (Ctrl-C, which halts
79/// *every* execution) OR its own per-execution stop flag is set — so a stop
80/// scoped to one execution isolates to that execution while Ctrl-C still stops
81/// all. Outside any execution scope (single-run / CLI / tests) only the global
82/// flag applies — behavior is identical to before the seam (axiom A1).
83#[inline]
84pub fn stop_requested() -> bool {
85 let global = SESSION_STOP
86 .get()
87 .map(|f| f.load(Ordering::Relaxed))
88 .unwrap_or(false);
89 let local = crate::execution_context::current_stop()
90 .map(|f| f.load(Ordering::Relaxed))
91 .unwrap_or(false);
92 global || local
93}
94
95/// Programmatically request a session-wide stop. Used by the
96/// signal handler, but also available to other lifecycle code
97/// that wants to short-circuit the run.
98pub fn request_stop() {
99 flag().store(true, Ordering::Relaxed);
100}
101
102// ─── The shutdown ladder (module doc) ────────────────────────────────
103
104/// Countdown from level 1 (graceful) to level 2 (cancel in-flight ops).
105const SHUTDOWN_COUNTDOWN_SECS: u64 = 10;
106
107/// Ladder level, published through a `watch` channel so per-fiber op
108/// dispatch can race a pending adapter call against the cancel rung
109/// without polling. 0 = running, 1 = graceful, 2 = cancel-ops.
110/// (Level 3 — force-exit — never publishes; it exits.)
111static SHUTDOWN: OnceLock<tokio::sync::watch::Sender<u8>> = OnceLock::new();
112
113/// Set once the runner's process-level shutdown has completed — the
114/// countdown thread goes quiet instead of announcing op cancellation
115/// for a run that has already drained.
116static SHUTDOWN_DONE: OnceLock<Arc<AtomicBool>> = OnceLock::new();
117
118fn shutdown_tx() -> &'static tokio::sync::watch::Sender<u8> {
119 SHUTDOWN.get_or_init(|| tokio::sync::watch::channel(0u8).0)
120}
121
122fn done_flag() -> &'static Arc<AtomicBool> {
123 SHUTDOWN_DONE.get_or_init(|| Arc::new(AtomicBool::new(false)))
124}
125
126/// The ladder level currently in force.
127#[inline]
128pub fn shutdown_level() -> u8 {
129 SHUTDOWN.get().map(|tx| *tx.borrow()).unwrap_or(0)
130}
131
132/// True once level 2 has been entered — in-flight ops should cancel.
133#[inline]
134pub fn cancel_ops_requested() -> bool {
135 shutdown_level() >= 2
136}
137
138/// Subscribe to ladder-level changes. Each fiber holds one receiver
139/// and races its op dispatch against [`ops_cancelled`].
140pub fn subscribe_shutdown() -> tokio::sync::watch::Receiver<u8> {
141 shutdown_tx().subscribe()
142}
143
144/// Resolves when the cancel rung (level ≥ 2) is in force. Ready
145/// immediately if it already is. Used in a `select!` against the op
146/// future at the fiber's dispatch point — dropping the op future is
147/// the cancellation.
148pub async fn ops_cancelled(rx: &mut tokio::sync::watch::Receiver<u8>) {
149 loop {
150 if *rx.borrow() >= 2 {
151 return;
152 }
153 if rx.changed().await.is_err() {
154 // Sender can't drop (static), but never spin if it did.
155 std::future::pending::<()>().await;
156 }
157 }
158}
159
160/// Mark the runner's process-level shutdown complete: the countdown
161/// (if still running) goes quiet, and further escalations are moot.
162pub fn mark_shutdown_complete() {
163 done_flag().store(true, Ordering::Relaxed);
164}
165
166/// What advanced the shutdown ladder — used only to phrase the level-1
167/// announcement accurately. A programmatic `action: abort` trip drives the
168/// SAME ladder as Ctrl-C, so it must not be reported as a Ctrl-C.
169#[derive(Debug, Clone, Copy, PartialEq, Eq)]
170pub enum ShutdownOrigin {
171 /// The SIGINT / Ctrl-C watcher (`install_signal_handler`) or the
172 /// dispatcher thread receiving SIGINT (SRD-93 M7).
173 CtrlC,
174 /// The dispatcher thread receiving SIGTERM (SRD-93 M7) — the
175 /// canonical orchestrator graceful-terminate (systemd / Docker /
176 /// Kubernetes send it ahead of their grace-period SIGKILL).
177 Term,
178 /// A `stop_when` condition with `action: abort` (SRD-83 follow-up).
179 StopAction,
180}
181
182impl ShutdownOrigin {
183 /// The leading clause of the level-1 log line, naming the trigger. The
184 /// trailing `(Ctrl-C: cancel them now …)` guidance stays fixed — Ctrl-C
185 /// escalates the ladder regardless of what first advanced it.
186 fn lead(self) -> &'static str {
187 match self {
188 ShutdownOrigin::CtrlC => "session: graceful shutdown requested (Ctrl-C).",
189 ShutdownOrigin::Term => "session: graceful shutdown requested (SIGTERM).",
190 ShutdownOrigin::StopAction => "session: shutdown requested by stop action `abort`.",
191 }
192 }
193
194 /// Short noun phrase naming the trigger, for inline use in a sentence.
195 fn trigger(self) -> &'static str {
196 match self {
197 ShutdownOrigin::CtrlC => "Ctrl-C",
198 ShutdownOrigin::Term => "SIGTERM",
199 ShutdownOrigin::StopAction => "stop action `abort`",
200 }
201 }
202}
203
204/// `action: abort` (SRD-83 follow-up) — jump the shutdown ladder STRAIGHT
205/// to the cancel-ops rung (level 2), skipping the level-1 cooperative-drain
206/// countdown. A stop driven by errors must not wait for in-flight ops or
207/// remaining phases: their futures are dropped at the fiber dispatch point
208/// NOW, the walk unwinds, and control passes straight to the runner's
209/// graceful SESSION shutdown. Cleanup still runs (metrics flush, WAL
210/// consolidation, summaries via the RAII guard) — this is NOT a force-exit;
211/// a further Ctrl-C is. Also raises the global session stop so the walk
212/// halts and no new phase starts. Idempotent; a no-op if the ladder is
213/// already at level ≥ 2. Returns the level now in force.
214pub fn abort_shutdown(origin: ShutdownOrigin) -> u8 {
215 request_stop();
216 let modified = shutdown_tx().send_if_modified(|level| {
217 if *level < 2 {
218 *level = 2;
219 true
220 } else {
221 false
222 }
223 });
224 if modified {
225 crate::diag!(
226 crate::observer::LogLevel::Warn,
227 "session: aborting on {} — cancelling in-flight ops now \
228 (skipping the cooperative drain); remaining phases skipped. \
229 Process-level cleanup (metrics flush, WAL consolidation, \
230 summaries) still runs. Ctrl-C to force-exit.",
231 origin.trigger()
232 );
233 }
234 2
235}
236
237/// Advance the ladder ONE rung: 0 → 1 (graceful + countdown),
238/// 1 → 2 (cancel in-flight ops). At ≥ 2 this is a no-op returning the
239/// current level — the FORCE-EXIT decision (level 3) stays with the
240/// caller, which owns its own terminal hygiene (the raw-mode
241/// key-watcher must restore the terminal before exiting; the SIGINT
242/// watcher just exits). `origin` phrases the level-1 line only. Returns
243/// the level now in force.
244pub fn escalate_shutdown(origin: ShutdownOrigin) -> u8 {
245 let tx = shutdown_tx();
246 let mut entered: u8 = 0;
247 tx.send_if_modified(|level| {
248 if *level >= 2 {
249 entered = *level;
250 false
251 } else {
252 *level += 1;
253 entered = *level;
254 true
255 }
256 });
257 match entered {
258 1 => {
259 request_stop();
260 crate::diag!(
261 crate::observer::LogLevel::Info,
262 "{} Active fibers exit at the next cycle boundary; profiler / \
263 metrics / summaries will flush. In-flight ops will be \
264 CANCELLED in {SHUTDOWN_COUNTDOWN_SECS}s (Ctrl-C: cancel \
265 them now; a further Ctrl-C force-exits).",
266 origin.lead()
267 );
268 // Not under `cfg(test)`: the in-crate unit tests exercise the
269 // ladder's transitions against process-global state; a live
270 // countdown escalating that state seconds later would race
271 // every concurrently-running test. Integration tests (their
272 // own processes) compile the lib without `test` and get the
273 // real countdown.
274 #[cfg(not(test))]
275 spawn_cancel_countdown();
276 }
277 2 => announce_cancel_ops(),
278 _ => {}
279 }
280 entered
281}
282
283/// Enter the cancel rung directly (countdown expiry). Idempotent.
284/// (Only the countdown calls this, and the countdown is compiled out
285/// of the in-crate unit-test build — see `escalate_shutdown`.)
286#[cfg_attr(test, allow(dead_code))]
287fn escalate_cancel_ops() {
288 let tx = shutdown_tx();
289 let modified = tx.send_if_modified(|level| {
290 if *level < 2 {
291 *level = 2;
292 true
293 } else {
294 false
295 }
296 });
297 if modified {
298 announce_cancel_ops();
299 }
300}
301
302fn announce_cancel_ops() {
303 crate::diag!(
304 crate::observer::LogLevel::Warn,
305 "session: cancelling in-flight ops — process-level cleanup \
306 (metrics flush, WAL consolidation, summaries) continues. \
307 Ctrl-C again to force-exit."
308 );
309}
310
311/// The visible level-1 → level-2 countdown. A plain thread (both
312/// entry points can spawn it — the raw-mode key-watcher runs outside
313/// the tokio runtime): one line per second, silenced the moment the
314/// run drains ([`mark_shutdown_complete`]) or the ladder advances by
315/// keypress; escalates to the cancel rung on expiry.
316#[cfg_attr(test, allow(dead_code))]
317fn spawn_cancel_countdown() {
318 let done = done_flag().clone();
319 std::thread::Builder::new()
320 .name("shutdown-countdown".into())
321 .spawn(move || {
322 for remaining in (1..=SHUTDOWN_COUNTDOWN_SECS).rev() {
323 if done.load(Ordering::Relaxed) || shutdown_level() >= 2 {
324 return;
325 }
326 crate::diag!(
327 crate::observer::LogLevel::Info,
328 "session: cancelling in-flight ops in {remaining}s \
329 (Ctrl-C to cancel now)"
330 );
331 std::thread::sleep(std::time::Duration::from_secs(1));
332 }
333 if !done.load(Ordering::Relaxed) {
334 escalate_cancel_ops();
335 }
336 })
337 .expect("spawn shutdown-countdown thread");
338}
339
340/// Shared "a stop condition gracefully halted the walk" flag. Distinct
341/// from [`SESSION_STOP`] (Ctrl-C): set when a workload-shell stop
342/// condition (SRD-83) intentionally halts the remaining walk. The
343/// end-of-run unreached-phase check consults it to distinguish an
344/// *intended* early stop (later phases deliberately skipped, not an
345/// error) from the genuine "a phase failed and stranded downstream
346/// phases" case.
347static GRACEFUL_STOP: OnceLock<Arc<AtomicBool>> = OnceLock::new();
348
349/// True once a stop condition has gracefully halted the walk.
350#[inline]
351pub fn graceful_stop_requested() -> bool {
352 GRACEFUL_STOP
353 .get()
354 .map(|f| f.load(Ordering::Relaxed))
355 .unwrap_or(false)
356}
357
358/// Record that a stop condition gracefully halted the remaining walk
359/// (SRD-83 workload shell). Idempotent.
360pub fn request_graceful_stop() {
361 GRACEFUL_STOP
362 .get_or_init(|| Arc::new(AtomicBool::new(false)))
363 .store(true, Ordering::Relaxed);
364}
365
366/// Why a shell-driven stop halted the walk (SRD-82 Part 4). A
367/// [`StopCause::Fault`] is a `fail`-effect trip — a child phase failed,
368/// so the run's validity is `Failed` and the halt records
369/// `Interrupted + Failed`. A [`StopCause::Interrupt`] is a clean `stop`
370/// (a graceful condition or user Ctrl-C) — later phases are deliberately
371/// skipped and the result is re-usable.
372#[derive(Debug, Clone, Copy, PartialEq, Eq)]
373pub enum StopCause {
374 Interrupt,
375 Fault,
376}
377
378/// Shared "a `fail`-effect stop condition halted the walk" flag (SRD-82
379/// Part 4). Distinct from [`GRACEFUL_STOP`]: a fault halt skips the tail
380/// phases (so the unreached-phase check stays quiet, like graceful) but
381/// the run still exits non-zero — the failing phase's own `Err` carries
382/// that via the `run_result` path. The flag records the *cause* so the
383/// end-of-run accounting and the trip log read "fault", not "graceful".
384static FAULT_STOP: OnceLock<Arc<AtomicBool>> = OnceLock::new();
385
386/// True once a `fail`-effect stop condition has halted the walk.
387#[inline]
388pub fn fault_stop_requested() -> bool {
389 FAULT_STOP
390 .get()
391 .map(|f| f.load(Ordering::Relaxed))
392 .unwrap_or(false)
393}
394
395/// Record that a `fail`-effect stop condition halted the remaining walk
396/// (SRD-82 Part 4, `StopCause::Fault`). Idempotent.
397pub fn request_fault_stop() {
398 FAULT_STOP
399 .get_or_init(|| Arc::new(AtomicBool::new(false)))
400 .store(true, Ordering::Relaxed);
401}
402
403/// Route a shell stop to the right session signal by its
404/// [`StopCause`]. Both halt the tail; `Fault` additionally marks the
405/// run failed (non-zero exit via the failing phase's `Err`).
406pub fn request_shell_stop(cause: StopCause) {
407 match cause {
408 StopCause::Fault => request_fault_stop(),
409 StopCause::Interrupt => request_graceful_stop(),
410 }
411}
412
413/// SRD-92 Step 0 — one cooperative-stop view, consulted at every boundary.
414/// Bundles the per-execution stop sources so a boundary check is a single
415/// call, and so a unit that previously held only ONE flag (the `while:`
416/// wrapper held only the activity `stop_flag`) observes ALL of them. The
417/// global / per-execution session stop ([`stop_requested`]) is always
418/// folded in; a `fail`-effect global ([`fault_stop_requested`]) is folded
419/// into [`StopView::poll`].
420///
421/// The `daemon` source (the SRD-82 Part 6 daemon-group completion) is a
422/// CLEAN termination: it ends a loop ([`StopView::stopped`]) but is NOT a
423/// fault, so it is excluded from [`StopView::abnormal`] (the
424/// failure-determining set). Fiber-pool scale-down
425/// ([`crate::fiber_pool::StopFlag`]) is a per-fiber concern, deliberately
426/// NOT part of this view.
427#[derive(Clone, Default)]
428pub struct StopView {
429 activity: Option<Arc<AtomicBool>>,
430 walk: Option<Arc<AtomicBool>>,
431 daemon: Option<Arc<AtomicBool>>,
432}
433
434impl StopView {
435 /// Build from the per-execution stop sources (any may be absent).
436 pub fn new(
437 activity: Option<Arc<AtomicBool>>,
438 walk: Option<Arc<AtomicBool>>,
439 daemon: Option<Arc<AtomicBool>>,
440 ) -> Self {
441 Self {
442 activity,
443 walk,
444 daemon,
445 }
446 }
447
448 #[inline]
449 fn on(f: &Option<Arc<AtomicBool>>) -> bool {
450 f.as_ref().is_some_and(|b| b.load(Ordering::Relaxed))
451 }
452
453 /// Any cooperative stop — incl. the clean daemon-group completion and
454 /// the global / per-execution session stop. Use at loop BREAK boundaries.
455 #[inline]
456 pub fn stopped(&self) -> bool {
457 Self::on(&self.activity)
458 || stop_requested()
459 || Self::on(&self.walk)
460 || Self::on(&self.daemon)
461 }
462
463 /// A stop that marks the unit FAILED / abnormal — EXCLUDES the clean
464 /// daemon-group stop. Use for the failure-determining return.
465 #[inline]
466 pub fn abnormal(&self) -> bool {
467 Self::on(&self.activity) || stop_requested() || Self::on(&self.walk)
468 }
469
470 /// The stop CAUSE for shell-level recording (SRD-82 Part 4): a fault
471 /// (the activity error-handler `stop_flag`, or a `fail`-effect global)
472 /// outranks a clean interrupt.
473 #[inline]
474 pub fn poll(&self) -> Option<StopCause> {
475 if fault_stop_requested() || Self::on(&self.activity) {
476 Some(StopCause::Fault)
477 } else if self.stopped() {
478 Some(StopCause::Interrupt)
479 } else {
480 None
481 }
482 }
483}
484
485/// Install a tokio task that watches `ctrl_c()` and drives SIGINT
486/// through the three-level ladder described in the module doc.
487/// Idempotent — only the first call wins; subsequent calls are
488/// no-ops. Must be called from inside a tokio runtime context.
489pub fn install_signal_handler() {
490 // SRD-93 M7 — arm the dispatcher's ladder routing (idempotent,
491 // and meaningful on every call: a run is now in progress).
492 LADDER_ARMED.store(true, Ordering::Relaxed);
493 static INSTALLED: OnceLock<()> = OnceLock::new();
494 if INSTALLED.set(()).is_err() {
495 return;
496 }
497 // Touch the flag to ensure it's initialized before any
498 // observer or fiber checks `stop_requested()`.
499 let _ = flag();
500 tokio::spawn(async move {
501 // Every Ctrl-C advances one rung. The messages route through
502 // `crate::diag!` so they reach every sink the rest of the
503 // runtime uses (session.log via the async sink, plus the
504 // registered RunObserver — the TUI log panel in TUI mode, the
505 // stderr fallback otherwise). The leading-newline cosmetics
506 // for the terminal-echoed `^C` live in [`StderrObserver::log`]
507 // so the structured log isn't littered with blank lines.
508 loop {
509 if tokio::signal::ctrl_c().await.is_err() {
510 return;
511 }
512 if shutdown_level() >= 2 {
513 // Level 3: force-exit.
514 crate::diag!(
515 crate::observer::LogLevel::Warn,
516 "session: force-exit (Ctrl-C past the cancel rung) — \
517 profiler output and metrics may be incomplete."
518 );
519 std::process::exit(130);
520 }
521 escalate_shutdown(ShutdownOrigin::CtrlC);
522 }
523 });
524}
525
526// ─── SRD-93 M7 — dedicated signal dispatcher (the A8 stop door) ──────
527
528/// True once a run has armed the ladder ([`install_signal_handler`]).
529/// Unarmed, the dispatcher preserves default CLI interrupt semantics
530/// (exit `128 + signo`); armed, signals drive the ladder.
531static LADDER_ARMED: AtomicBool = AtomicBool::new(false);
532
533#[cfg(unix)]
534#[inline]
535fn ladder_armed() -> bool {
536 LADDER_ARMED.load(Ordering::Relaxed)
537}
538
539/// Console-loss hook (SIGHUP while armed). The TUI installs a closure
540/// that restores the terminal and switches rendering headless; without
541/// one, the dispatcher just logs and continues — the run survives
542/// terminal loss either way.
543static CONSOLE_LOSS_HOOK: OnceLock<Box<dyn Fn() + Send + Sync>> = OnceLock::new();
544
545/// Install the console-loss hook. First call wins (one console).
546pub fn set_console_loss_hook(hook: Box<dyn Fn() + Send + Sync>) {
547 let _ = CONSOLE_LOSS_HOOK.set(hook);
548}
549
550/// Diagnostics-dump hook (SIGQUIT while armed). The runner installs a
551/// closure that logs the per-activity / component inventory; the
552/// dispatcher always logs the ladder + stop-flag line first.
553static DIAG_DUMP_HOOK: OnceLock<Box<dyn Fn() + Send + Sync>> = OnceLock::new();
554
555/// Install the diagnostics-dump hook. First call wins.
556pub fn set_diag_dump_hook(hook: Box<dyn Fn() + Send + Sync>) {
557 let _ = DIAG_DUMP_HOOK.set(hook);
558}
559
560/// The signal set the dispatcher owns. SIGINT/SIGTERM drive the
561/// ladder; SIGHUP is console-loss (never a stop); SIGQUIT dumps
562/// diagnostics.
563#[cfg(unix)]
564const DISPATCHED_SIGNALS: [libc::c_int; 4] =
565 [libc::SIGINT, libc::SIGTERM, libc::SIGHUP, libc::SIGQUIT];
566
567/// What the dispatcher does with one received signal. Pure decision,
568/// separated from the thread loop so the routing is unit-testable
569/// without raising real signals.
570#[cfg(unix)]
571#[derive(Debug, Clone, Copy, PartialEq, Eq)]
572enum SignalAction {
573 /// Exit `128 + signo` now (unarmed default semantics).
574 Exit(i32),
575 /// Advance the ladder one rung.
576 Escalate(ShutdownOrigin),
577 /// Past the cancel rung: announce and exit `128 + signo`.
578 ForceExit(i32),
579 /// SIGHUP armed: log, run the console-loss hook, continue.
580 ConsoleLoss,
581 /// SIGQUIT armed: log ladder + stop flags, run the dump hook,
582 /// continue.
583 DiagDump,
584}
585
586/// Route one signal by (signo, armed, current ladder level).
587#[cfg(unix)]
588fn dispatch_decision(signo: libc::c_int, armed: bool, level: u8) -> SignalAction {
589 let exit_code = 128 + signo as i32;
590 match signo {
591 libc::SIGINT | libc::SIGTERM => {
592 let origin = if signo == libc::SIGTERM {
593 ShutdownOrigin::Term
594 } else {
595 ShutdownOrigin::CtrlC
596 };
597 if !armed {
598 SignalAction::Exit(exit_code)
599 } else if level >= 2 {
600 SignalAction::ForceExit(exit_code)
601 } else {
602 SignalAction::Escalate(origin)
603 }
604 }
605 libc::SIGHUP if armed => SignalAction::ConsoleLoss,
606 libc::SIGQUIT if armed => SignalAction::DiagDump,
607 _ => SignalAction::Exit(exit_code),
608 }
609}
610
611/// Block the dispatched signals in the calling thread. MUST run first
612/// thing in `main`, before any thread spawns, so every later thread
613/// inherits the block and delivery can only land in the dispatcher's
614/// `sigwait`. This also makes the ladder immune to a third-party
615/// library later installing its own `sigaction` for these signals —
616/// `sigwait` consumes a blocked signal before any handler would run.
617/// (Child processes are unaffected: `std::process::Command` resets
618/// the child's signal mask and dispositions.)
619#[cfg(unix)]
620pub fn block_shutdown_signals() {
621 unsafe {
622 let mut set: libc::sigset_t = std::mem::zeroed();
623 libc::sigemptyset(&mut set);
624 for s in DISPATCHED_SIGNALS {
625 libc::sigaddset(&mut set, s);
626 }
627 libc::pthread_sigmask(libc::SIG_BLOCK, &set, std::ptr::null_mut());
628 }
629}
630
631#[cfg(not(unix))]
632pub fn block_shutdown_signals() {}
633
634/// Spawn the dedicated dispatcher thread (SRD-93 M7 / A8). Pair with
635/// [`block_shutdown_signals`] — without the block, `sigwait` never
636/// receives anything and the legacy tokio watcher keeps the door.
637/// Idempotent.
638#[cfg(unix)]
639pub fn spawn_signal_dispatcher() {
640 static SPAWNED: OnceLock<()> = OnceLock::new();
641 if SPAWNED.set(()).is_err() {
642 return;
643 }
644 std::thread::Builder::new()
645 .name("signal-dispatch".into())
646 .spawn(|| {
647 let mut set: libc::sigset_t = unsafe { std::mem::zeroed() };
648 unsafe {
649 libc::sigemptyset(&mut set);
650 for s in DISPATCHED_SIGNALS {
651 libc::sigaddset(&mut set, s);
652 }
653 }
654 loop {
655 let mut signo: libc::c_int = 0;
656 if unsafe { libc::sigwait(&set, &mut signo) } != 0 {
657 // EINVAL on the fixed set can't happen; if it
658 // somehow does, retiring the thread just reverts
659 // to the legacy watcher door.
660 return;
661 }
662 match dispatch_decision(signo, ladder_armed(), shutdown_level()) {
663 SignalAction::Exit(code) => std::process::exit(code),
664 SignalAction::Escalate(origin) => {
665 escalate_shutdown(origin);
666 }
667 SignalAction::ForceExit(code) => {
668 crate::diag!(
669 crate::observer::LogLevel::Warn,
670 "session: force-exit (signal past the cancel rung) — \
671 profiler output and metrics may be incomplete."
672 );
673 std::process::exit(code);
674 }
675 SignalAction::ConsoleLoss => {
676 // Hook FIRST: it suppresses terminal-bound
677 // writes (the terminal is gone — a subsequent
678 // stderr write could fail hard), so the diag
679 // below reaches only the file sinks.
680 if let Some(hook) = CONSOLE_LOSS_HOOK.get() {
681 hook();
682 }
683 crate::diag!(
684 crate::observer::LogLevel::Warn,
685 "session: SIGHUP — controlling terminal lost; \
686 continuing headless (the run is unaffected)."
687 );
688 }
689 SignalAction::DiagDump => {
690 crate::diag!(
691 crate::observer::LogLevel::Info,
692 "session: SIGQUIT diagnostics — shutdown ladder level {}, \
693 session stop {}, graceful stop {}, fault stop {}.",
694 shutdown_level(),
695 stop_requested(),
696 graceful_stop_requested(),
697 fault_stop_requested()
698 );
699 if let Some(hook) = DIAG_DUMP_HOOK.get() {
700 hook();
701 }
702 }
703 }
704 }
705 })
706 .expect("spawn signal-dispatch thread");
707}
708
709#[cfg(not(unix))]
710pub fn spawn_signal_dispatcher() {}
711
712/// Test-only: serialize the tests that touch the process-global
713/// `SESSION_STOP` flag. `SESSION_STOP` is a never-reset
714/// `OnceLock<Arc<AtomicBool>>`, so a test that sets it would
715/// otherwise leak into every sibling test in the same binary
716/// (notably the per-execution isolation test, which requires the
717/// global to be clear). Tests acquire this lock and clear the flag
718/// before asserting.
719#[cfg(test)]
720pub(crate) static STOP_GLOBAL_TEST_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
721
722/// Test-only: reset the process-global stop flag to unset. Paired
723/// with [`STOP_GLOBAL_TEST_LOCK`] so global-flag tests don't leak
724/// state into each other.
725#[cfg(test)]
726pub(crate) fn clear_session_stop_for_test() {
727 flag().store(false, Ordering::Relaxed);
728}
729
730/// Test-only: reset the shutdown ladder to level 0 (and clear the
731/// done flag). Same locking discipline as
732/// [`clear_session_stop_for_test`].
733#[cfg(test)]
734pub(crate) fn reset_shutdown_ladder_for_test() {
735 let _ = shutdown_tx().send_replace(0);
736 done_flag().store(false, Ordering::Relaxed);
737}
738
739#[cfg(test)]
740mod tests {
741 use super::*;
742
743 #[test]
744 fn flag_starts_unset_and_responds_to_request() {
745 // Serialize with the other global-flag test and reset the
746 // process-global flag around the assertions so this test
747 // neither sees nor leaks a stale stop.
748 let _guard = STOP_GLOBAL_TEST_LOCK
749 .lock()
750 .unwrap_or_else(|e| e.into_inner());
751 clear_session_stop_for_test();
752 assert!(!stop_requested());
753 request_stop();
754 assert!(stop_requested());
755 clear_session_stop_for_test();
756 }
757
758 /// The ladder advances one rung per escalation — 0 → 1 (graceful,
759 /// session stop set) → 2 (cancel ops) — and holds at 2:
760 /// `escalate_shutdown` never enters level 3 itself (force-exit is
761 /// the CALLER's decision, with its own terminal hygiene).
762 #[test]
763 fn ladder_advances_one_rung_per_escalation_and_holds_at_cancel() {
764 let _guard = STOP_GLOBAL_TEST_LOCK
765 .lock()
766 .unwrap_or_else(|e| e.into_inner());
767 clear_session_stop_for_test();
768 reset_shutdown_ladder_for_test();
769
770 assert_eq!(shutdown_level(), 0);
771 assert!(!cancel_ops_requested());
772
773 assert_eq!(
774 escalate_shutdown(ShutdownOrigin::CtrlC),
775 1,
776 "first rung: graceful"
777 );
778 assert!(stop_requested(), "graceful rung sets the session stop");
779 assert!(!cancel_ops_requested());
780
781 assert_eq!(
782 escalate_shutdown(ShutdownOrigin::CtrlC),
783 2,
784 "second rung: cancel ops"
785 );
786 assert!(cancel_ops_requested());
787
788 assert_eq!(
789 escalate_shutdown(ShutdownOrigin::CtrlC),
790 2,
791 "ladder holds at cancel"
792 );
793 assert!(cancel_ops_requested());
794
795 clear_session_stop_for_test();
796 reset_shutdown_ladder_for_test();
797 }
798
799 /// `abort_shutdown` jumps STRAIGHT from level 0 to the cancel-ops
800 /// rung (level 2) — no intermediate graceful rung — and raises the
801 /// session stop, so an `action: abort` trip cancels in-flight ops
802 /// without the cooperative-drain wait.
803 #[test]
804 fn abort_jumps_straight_to_cancel_rung() {
805 let _guard = STOP_GLOBAL_TEST_LOCK
806 .lock()
807 .unwrap_or_else(|e| e.into_inner());
808 clear_session_stop_for_test();
809 reset_shutdown_ladder_for_test();
810
811 assert_eq!(shutdown_level(), 0);
812 assert!(!stop_requested());
813
814 assert_eq!(
815 abort_shutdown(ShutdownOrigin::StopAction),
816 2,
817 "abort skips the graceful rung and lands on cancel-ops"
818 );
819 assert!(cancel_ops_requested(), "in-flight ops cancel immediately");
820 assert!(stop_requested(), "abort also raises the session stop");
821
822 // Idempotent, and never regresses the rung.
823 assert_eq!(abort_shutdown(ShutdownOrigin::StopAction), 2);
824
825 clear_session_stop_for_test();
826 reset_shutdown_ladder_for_test();
827 }
828
829 /// `ops_cancelled` resolves the moment the cancel rung is in force
830 /// — including when it already was at subscribe time — and does
831 /// NOT resolve for the graceful rung alone.
832 #[tokio::test]
833 async fn ops_cancelled_resolves_at_cancel_rung_only() {
834 let _guard = STOP_GLOBAL_TEST_LOCK
835 .lock()
836 .unwrap_or_else(|e| e.into_inner());
837 clear_session_stop_for_test();
838 reset_shutdown_ladder_for_test();
839
840 // Graceful rung alone must NOT resolve the cancel future.
841 let mut rx = subscribe_shutdown();
842 escalate_shutdown(ShutdownOrigin::CtrlC); // → 1
843 let pending =
844 tokio::time::timeout(std::time::Duration::from_millis(50), ops_cancelled(&mut rx))
845 .await;
846 assert!(pending.is_err(), "graceful rung must not cancel ops");
847
848 // Cancel rung resolves it — and resolves immediately for a
849 // subscriber that arrives after the fact.
850 escalate_shutdown(ShutdownOrigin::CtrlC); // → 2
851 tokio::time::timeout(
852 std::time::Duration::from_millis(200),
853 ops_cancelled(&mut rx),
854 )
855 .await
856 .expect("cancel rung resolves the in-flight race");
857 let mut late = subscribe_shutdown();
858 tokio::time::timeout(
859 std::time::Duration::from_millis(200),
860 ops_cancelled(&mut late),
861 )
862 .await
863 .expect("already-cancelled resolves immediately");
864
865 clear_session_stop_for_test();
866 reset_shutdown_ladder_for_test();
867 }
868
869 /// SRD-93 M7 — the dispatcher's pure routing. Unarmed, every
870 /// signal keeps its default CLI exit semantics (`128 + signo`);
871 /// armed, INT/TERM climb the ladder and force-exit past the
872 /// cancel rung with the origin-correct code (130 / 143), HUP is
873 /// console-loss (never a stop), QUIT is a diagnostics dump.
874 #[cfg(unix)]
875 #[test]
876 fn dispatch_routes_by_signal_arming_and_rung() {
877 use SignalAction::*;
878
879 // Unarmed: default exit codes, any level.
880 assert_eq!(dispatch_decision(libc::SIGINT, false, 0), Exit(130));
881 assert_eq!(dispatch_decision(libc::SIGTERM, false, 0), Exit(143));
882 assert_eq!(dispatch_decision(libc::SIGHUP, false, 0), Exit(129));
883 assert_eq!(dispatch_decision(libc::SIGQUIT, false, 0), Exit(131));
884
885 // Armed, below the cancel rung: one rung per signal, with
886 // the origin carrying the trigger name.
887 assert_eq!(
888 dispatch_decision(libc::SIGINT, true, 0),
889 Escalate(ShutdownOrigin::CtrlC)
890 );
891 assert_eq!(
892 dispatch_decision(libc::SIGTERM, true, 1),
893 Escalate(ShutdownOrigin::Term)
894 );
895
896 // Armed, at/past the cancel rung: the hard floor, exit code
897 // by the signal that delivered the final blow.
898 assert_eq!(dispatch_decision(libc::SIGINT, true, 2), ForceExit(130));
899 assert_eq!(dispatch_decision(libc::SIGTERM, true, 2), ForceExit(143));
900
901 // Armed HUP/QUIT never stop and never escalate.
902 assert_eq!(dispatch_decision(libc::SIGHUP, true, 0), ConsoleLoss);
903 assert_eq!(dispatch_decision(libc::SIGHUP, true, 2), ConsoleLoss);
904 assert_eq!(dispatch_decision(libc::SIGQUIT, true, 0), DiagDump);
905 assert_eq!(dispatch_decision(libc::SIGQUIT, true, 2), DiagDump);
906 }
907}