Skip to main content

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}