Skip to main content

asupersync/signal/
shutdown.rs

1//! Coordinated shutdown controller using sync primitives.
2//!
3//! Provides a centralized mechanism for initiating and propagating shutdown
4//! signals throughout an application. Uses our sync primitives (Notify) to
5//! coordinate without external dependencies.
6
7use std::future::Future;
8use std::io;
9use std::sync::Arc;
10use std::sync::atomic::{AtomicBool, Ordering};
11
12use super::{SignalKind, signal};
13use crate::sync::Notify;
14use crate::tracing_compat::{info, warn};
15
16/// Internal state shared between controller and receivers.
17#[derive(Debug)]
18struct ShutdownState {
19    /// Tracks whether shutdown has been initiated.
20    initiated: AtomicBool,
21    /// Ensures signal listeners are only installed once per controller.
22    signal_listeners_started: AtomicBool,
23    /// Notifier for broadcast notifications.
24    notify: Notify,
25}
26
27/// Controller for coordinated graceful shutdown.
28///
29/// This provides a clean way to propagate shutdown signals through an application.
30/// Multiple receivers can subscribe to receive shutdown notifications.
31///
32/// # Example
33///
34/// ```ignore
35/// use asupersync::signal::ShutdownController;
36///
37/// async fn run_server() {
38///     let controller = ShutdownController::new();
39///     let mut receiver = controller.subscribe();
40///
41///     // Spawn a task that will receive the shutdown signal
42///     let handle = async move {
43///         receiver.wait().await;
44///         println!("Shutting down...");
45///     };
46///
47///     // Later, initiate shutdown
48///     controller.shutdown();
49/// }
50/// ```
51#[derive(Debug)]
52pub struct ShutdownController {
53    /// Shared state between controller and receivers.
54    state: Arc<ShutdownState>,
55}
56
57/// Internal state shared between reload controller and receivers.
58#[derive(Debug)]
59struct ReloadState {
60    /// Monotone reload request sequence.
61    requests: std::sync::atomic::AtomicU64,
62    /// Ensures the SIGHUP listener is installed at most once per controller.
63    signal_listener_started: AtomicBool,
64    /// Notifier for reload request broadcasts.
65    notify: Notify,
66}
67
68/// Controller for SIGHUP-style configuration reload notifications.
69///
70/// Reloads are independent from shutdown. Calling [`request_reload`](Self::request_reload)
71/// or receiving SIGHUP through [`listen_for_sighup`](Self::listen_for_sighup)
72/// increments a monotone reload sequence and wakes subscribed receivers, but it
73/// never marks a [`ShutdownController`] as shutting down.
74#[derive(Debug)]
75pub struct ReloadController {
76    /// Shared state between controller and receivers.
77    state: Arc<ReloadState>,
78}
79
80impl ReloadController {
81    /// Creates a new reload controller.
82    #[must_use]
83    pub fn new() -> Self {
84        Self {
85            state: Arc::new(ReloadState {
86                requests: std::sync::atomic::AtomicU64::new(0),
87                signal_listener_started: AtomicBool::new(false),
88                notify: Notify::new(),
89            }),
90        }
91    }
92
93    /// Gets a receiver for future reload notifications.
94    ///
95    /// The receiver starts after the current reload sequence, so it observes
96    /// reloads requested after subscription rather than replaying historical
97    /// requests.
98    #[must_use]
99    pub fn subscribe(&self) -> ReloadReceiver {
100        ReloadReceiver {
101            state: Arc::clone(&self.state),
102            seen_requests: self.reload_count(),
103        }
104    }
105
106    /// Requests a reload and returns the new reload sequence number.
107    ///
108    /// This wakes all receivers waiting for a reload notification. The request
109    /// is event-like, not sticky: receivers created after this call start after
110    /// the returned sequence.
111    pub fn request_reload(&self) -> u64 {
112        Self::trigger_reload_state(&self.state)
113    }
114
115    /// Returns the number of reload requests recorded by this controller.
116    #[must_use]
117    pub fn reload_count(&self) -> u64 {
118        self.state.requests.load(Ordering::Acquire)
119    }
120
121    /// Installs an opt-in SIGHUP listener for this reload controller.
122    ///
123    /// The listener is Unix-only because SIGHUP has no portable Windows
124    /// equivalent. Unsupported platforms return a deterministic
125    /// [`io::ErrorKind::Unsupported`] error.
126    ///
127    /// Calling this method more than once is idempotent.
128    pub fn listen_for_sighup(self: &Arc<Self>) -> io::Result<()> {
129        if self
130            .state
131            .signal_listener_started
132            .swap(true, Ordering::AcqRel)
133        {
134            return Ok(());
135        }
136
137        match Self::spawn_sighup_listener(Arc::downgrade(&self.state)) {
138            Ok(()) => Ok(()),
139            Err(err) => {
140                self.state
141                    .signal_listener_started
142                    .store(false, Ordering::Release);
143                Err(err)
144            }
145        }
146    }
147
148    fn trigger_reload_state(state: &ReloadState) -> u64 {
149        let sequence = state.requests.fetch_add(1, Ordering::AcqRel) + 1;
150        info!(reload_sequence = sequence, "reload requested");
151        state.notify.notify_waiters();
152        sequence
153    }
154
155    #[cfg(unix)]
156    fn spawn_sighup_listener(state: std::sync::Weak<ReloadState>) -> io::Result<()> {
157        let mut stream = signal(SignalKind::hangup())?;
158        std::thread::Builder::new()
159            .name("asupersync-reload-sighup".to_string())
160            .spawn(move || {
161                while futures_lite::future::block_on(stream.recv()).is_some() {
162                    let Some(state) = state.upgrade() else {
163                        break;
164                    };
165                    Self::trigger_reload_state(&state);
166                }
167            })
168            .map(|_| ())
169    }
170
171    #[cfg(not(unix))]
172    fn spawn_sighup_listener(_state: std::sync::Weak<ReloadState>) -> io::Result<()> {
173        Err(io::Error::new(
174            io::ErrorKind::Unsupported,
175            "SIGHUP reload listener is only supported on Unix",
176        ))
177    }
178}
179
180impl Default for ReloadController {
181    fn default() -> Self {
182        Self::new()
183    }
184}
185
186impl Clone for ReloadController {
187    fn clone(&self) -> Self {
188        Self {
189            state: Arc::clone(&self.state),
190        }
191    }
192}
193
194/// Outcome of handling one reload request.
195#[derive(Debug, Clone, PartialEq, Eq)]
196pub enum ReloadOutcome<E> {
197    /// The reload handler completed successfully.
198    Completed {
199        /// Reload sequence handled by the caller.
200        sequence: u64,
201    },
202    /// The reload handler returned an error.
203    Failed {
204        /// Reload sequence handled by the caller.
205        sequence: u64,
206        /// Error returned by the handler.
207        error: E,
208    },
209    /// Reload handling was cancelled before completion.
210    Cancelled {
211        /// Reload sequence that was cancelled, or `None` if shutdown won
212        /// before a reload request was selected.
213        sequence: Option<u64>,
214    },
215}
216
217/// Receiver for reload notifications.
218///
219/// A receiver observes a monotone sequence of reload requests. If several
220/// reloads arrive before the receiver is polled again, repeated calls to
221/// [`wait`](Self::wait) drain the pending sequence numbers without losing
222/// notifications.
223#[derive(Debug)]
224pub struct ReloadReceiver {
225    /// Shared state with the controller.
226    state: Arc<ReloadState>,
227    /// Last reload sequence observed by this receiver.
228    seen_requests: u64,
229}
230
231impl ReloadReceiver {
232    /// Waits for the next reload request and returns its sequence number.
233    pub async fn wait(&mut self) -> u64 {
234        let state = Arc::clone(&self.state);
235        loop {
236            let current = state.requests.load(Ordering::Acquire);
237            if current > self.seen_requests {
238                self.seen_requests = self.seen_requests.saturating_add(1);
239                return self.seen_requests;
240            }
241
242            let mut notified = std::pin::pin!(state.notify.notified());
243            std::future::poll_fn(|cx| {
244                let current = state.requests.load(Ordering::Acquire);
245                if current > self.seen_requests
246                    || std::future::Future::poll(notified.as_mut(), cx).is_ready()
247                {
248                    return std::task::Poll::Ready(());
249                }
250                std::task::Poll::Pending
251            })
252            .await;
253        }
254    }
255
256    /// Returns the last reload sequence this receiver has observed.
257    #[must_use]
258    pub fn seen_reload_count(&self) -> u64 {
259        self.seen_requests
260    }
261
262    /// Waits for one reload request and runs the supplied async handler.
263    ///
264    /// Structured tracing records the requested sequence, successful
265    /// completion, handler failure, and cancellation if the returned future is
266    /// dropped while the handler is still running.
267    pub async fn handle_next_reload<F, Fut, E>(&mut self, handler: F) -> ReloadOutcome<E>
268    where
269        F: FnOnce(u64) -> Fut,
270        Fut: Future<Output = Result<(), E>>,
271    {
272        let sequence = self.wait().await;
273        let mut guard = ReloadAttemptGuard::new(sequence);
274        match handler(sequence).await {
275            Ok(()) => {
276                guard.finish();
277                info!(reload_sequence = sequence, "reload completed");
278                ReloadOutcome::Completed { sequence }
279            }
280            Err(error) => {
281                guard.finish();
282                warn!(reload_sequence = sequence, "reload failed");
283                ReloadOutcome::Failed { sequence, error }
284            }
285        }
286    }
287
288    /// Waits for the next reload request unless shutdown is signalled first.
289    ///
290    /// Returns `Some(sequence)` when a reload request is selected and `None`
291    /// when shutdown wins the race. Both underlying waits are cancel-safe, so
292    /// dropping this future does not consume either notification.
293    pub async fn wait_or_shutdown(&mut self, shutdown: &mut ShutdownReceiver) -> Option<u64> {
294        let mut reload_wait = std::pin::pin!(self.wait());
295        let mut shutdown_wait = std::pin::pin!(shutdown.wait());
296
297        std::future::poll_fn(|cx| {
298            if let std::task::Poll::Ready(sequence) = Future::poll(reload_wait.as_mut(), cx) {
299                return std::task::Poll::Ready(Some(sequence));
300            }
301
302            if Future::poll(shutdown_wait.as_mut(), cx).is_ready() {
303                return std::task::Poll::Ready(None);
304            }
305
306            std::task::Poll::Pending
307        })
308        .await
309    }
310
311    /// Waits for one reload request and runs the handler unless shutdown wins.
312    ///
313    /// Shutdown before a reload request returns [`ReloadOutcome::Cancelled`]
314    /// with no sequence. Shutdown while a handler is running cancels that
315    /// handler future, records a structured cancellation event, and returns
316    /// [`ReloadOutcome::Cancelled`] for the selected sequence.
317    pub async fn handle_next_reload_or_shutdown<F, Fut, E>(
318        &mut self,
319        shutdown: &mut ShutdownReceiver,
320        handler: F,
321    ) -> ReloadOutcome<E>
322    where
323        F: FnOnce(u64) -> Fut,
324        Fut: Future<Output = Result<(), E>>,
325    {
326        let Some(sequence) = self.wait_or_shutdown(shutdown).await else {
327            warn!(reload_sequence = 0_u64, "reload cancelled");
328            return ReloadOutcome::Cancelled { sequence: None };
329        };
330
331        let mut guard = ReloadAttemptGuard::new(sequence);
332        let mut handler = std::pin::pin!(handler(sequence));
333        let mut shutdown_wait = std::pin::pin!(shutdown.wait());
334
335        match std::future::poll_fn(|cx| {
336            if let std::task::Poll::Ready(result) = Future::poll(handler.as_mut(), cx) {
337                return std::task::Poll::Ready(Some(result));
338            }
339
340            if Future::poll(shutdown_wait.as_mut(), cx).is_ready() {
341                return std::task::Poll::Ready(None);
342            }
343
344            std::task::Poll::Pending
345        })
346        .await
347        {
348            Some(Ok(())) => {
349                guard.finish();
350                info!(reload_sequence = sequence, "reload completed");
351                ReloadOutcome::Completed { sequence }
352            }
353            Some(Err(error)) => {
354                guard.finish();
355                warn!(reload_sequence = sequence, "reload failed");
356                ReloadOutcome::Failed { sequence, error }
357            }
358            None => {
359                guard.finish();
360                warn!(reload_sequence = sequence, "reload cancelled");
361                ReloadOutcome::Cancelled {
362                    sequence: Some(sequence),
363                }
364            }
365        }
366    }
367}
368
369impl Clone for ReloadReceiver {
370    fn clone(&self) -> Self {
371        Self {
372            state: Arc::clone(&self.state),
373            seen_requests: self.seen_requests,
374        }
375    }
376}
377
378#[derive(Debug)]
379struct ReloadAttemptGuard {
380    sequence: u64,
381    finished: bool,
382}
383
384impl ReloadAttemptGuard {
385    fn new(sequence: u64) -> Self {
386        Self {
387            sequence,
388            finished: false,
389        }
390    }
391
392    fn finish(&mut self) {
393        self.finished = true;
394    }
395}
396
397impl Drop for ReloadAttemptGuard {
398    fn drop(&mut self) {
399        if !self.finished && self.sequence > 0 {
400            warn!(reload_sequence = self.sequence, "reload cancelled");
401        }
402    }
403}
404
405impl ShutdownController {
406    /// Creates a new shutdown controller.
407    #[must_use]
408    pub fn new() -> Self {
409        Self {
410            state: Arc::new(ShutdownState {
411                initiated: AtomicBool::new(false),
412                signal_listeners_started: AtomicBool::new(false),
413                notify: Notify::new(),
414            }),
415        }
416    }
417
418    /// Gets a handle for receiving shutdown notifications.
419    ///
420    /// Multiple receivers can be created and they will all be notified
421    /// when shutdown is initiated.
422    #[must_use]
423    pub fn subscribe(&self) -> ShutdownReceiver {
424        ShutdownReceiver {
425            state: Arc::clone(&self.state),
426        }
427    }
428
429    /// Initiates shutdown.
430    ///
431    /// This wakes all receivers that are currently waiting for shutdown.
432    /// The shutdown state is persistent - once initiated, it cannot be reset.
433    pub fn shutdown(&self) {
434        Self::trigger_shutdown_state(&self.state);
435    }
436
437    /// Checks if shutdown has been initiated.
438    #[must_use]
439    pub fn is_shutting_down(&self) -> bool {
440        self.state.initiated.load(Ordering::Acquire)
441    }
442
443    /// Spawns a background task to listen for shutdown signals.
444    ///
445    /// This is a convenience method that sets up signal handling
446    /// (when available) to automatically trigger shutdown.
447    ///
448    /// # Note
449    ///
450    /// The listeners are installed at most once per controller. When a watched
451    /// signal arrives, the controller transitions to shutdown just as if
452    /// [`ShutdownController::shutdown`] had been called manually.
453    pub fn listen_for_signals(self: &Arc<Self>) {
454        if self
455            .state
456            .signal_listeners_started
457            .swap(true, Ordering::AcqRel)
458        {
459            return;
460        }
461
462        let state = Arc::downgrade(&self.state);
463        let mut installed = false;
464
465        for kind in watched_signal_kinds() {
466            if Self::spawn_signal_listener(state.clone(), kind).is_ok() {
467                installed = true;
468            }
469        }
470
471        if !installed {
472            self.state
473                .signal_listeners_started
474                .store(false, Ordering::Release);
475        }
476    }
477
478    fn trigger_shutdown_state(state: &ShutdownState) {
479        if state
480            .initiated
481            .compare_exchange(false, true, Ordering::Release, Ordering::Relaxed)
482            .is_ok()
483        {
484            state.notify.notify_waiters();
485        }
486    }
487
488    fn spawn_signal_listener(
489        state: std::sync::Weak<ShutdownState>,
490        kind: SignalKind,
491    ) -> std::io::Result<()> {
492        let mut stream = signal(kind)?;
493        std::thread::Builder::new()
494            .name(format!(
495                "asupersync-shutdown-{}",
496                kind.name().to_ascii_lowercase()
497            ))
498            .spawn(move || {
499                if futures_lite::future::block_on(stream.recv()).is_some()
500                    && let Some(state) = state.upgrade()
501                {
502                    Self::trigger_shutdown_state(&state);
503                }
504            })
505            .map(|_| ())
506    }
507}
508
509#[cfg(unix)]
510fn watched_signal_kinds() -> [SignalKind; 2] {
511    [SignalKind::interrupt(), SignalKind::terminate()]
512}
513
514#[cfg(windows)]
515fn watched_signal_kinds() -> [SignalKind; 3] {
516    [
517        SignalKind::interrupt(),
518        SignalKind::terminate(),
519        SignalKind::quit(),
520    ]
521}
522
523#[cfg(not(any(unix, windows)))]
524fn watched_signal_kinds() -> [SignalKind; 0] {
525    []
526}
527
528impl Default for ShutdownController {
529    fn default() -> Self {
530        Self::new()
531    }
532}
533
534impl Clone for ShutdownController {
535    fn clone(&self) -> Self {
536        Self {
537            state: Arc::clone(&self.state),
538        }
539    }
540}
541
542/// Receiver for shutdown notifications.
543///
544/// This is a handle that can wait for shutdown to be initiated.
545/// Multiple receivers can be created from a single controller.
546#[derive(Debug)]
547pub struct ShutdownReceiver {
548    /// Shared state with the controller.
549    state: Arc<ShutdownState>,
550}
551
552impl ShutdownReceiver {
553    /// Waits for shutdown to be initiated.
554    ///
555    /// This method returns immediately if shutdown has already been initiated.
556    /// Otherwise, it waits until the controller's `shutdown()` method is called.
557    pub async fn wait(&mut self) {
558        let state = Arc::clone(&self.state);
559        loop {
560            if state.initiated.load(Ordering::Acquire) {
561                return;
562            }
563
564            let mut notified = std::pin::pin!(state.notify.notified());
565            std::future::poll_fn(|cx| {
566                if std::future::Future::poll(notified.as_mut(), cx).is_ready()
567                    || state.initiated.load(Ordering::Acquire)
568                {
569                    return std::task::Poll::Ready(());
570                }
571                std::task::Poll::Pending
572            })
573            .await;
574
575            if state.initiated.load(Ordering::Acquire) {
576                return;
577            }
578        }
579    }
580
581    /// Checks if shutdown has been initiated.
582    #[must_use]
583    pub fn is_shutting_down(&self) -> bool {
584        self.state.initiated.load(Ordering::Acquire)
585    }
586}
587
588impl Clone for ShutdownReceiver {
589    fn clone(&self) -> Self {
590        Self {
591            state: Arc::clone(&self.state),
592        }
593    }
594}
595
596#[cfg(test)]
597mod tests {
598    #![allow(
599        clippy::pedantic,
600        clippy::nursery,
601        clippy::expect_fun_call,
602        clippy::map_unwrap_or,
603        clippy::cast_possible_wrap,
604        clippy::future_not_send
605    )]
606    use super::super::SignalKind;
607    use super::super::signal::inject_test_signal;
608    use super::*;
609    use serde_json::json;
610    use std::sync::Arc;
611    use std::task::{Context, Poll, Waker};
612    use std::thread;
613    use std::time::Duration;
614    #[cfg(unix)]
615    use std::time::Instant;
616
617    fn noop_waker() -> Waker {
618        std::task::Waker::noop().clone()
619    }
620
621    fn poll_once<F: std::future::Future + Unpin>(fut: &mut F) -> Poll<F::Output> {
622        let waker = noop_waker();
623        let mut cx = Context::from_waker(&waker);
624        std::pin::Pin::new(fut).poll(&mut cx)
625    }
626
627    fn init_test(name: &str) {
628        crate::test_utils::init_test_logging();
629        crate::test_phase!(name);
630    }
631
632    #[cfg(unix)]
633    fn wait_until(mut condition: impl FnMut() -> bool) -> bool {
634        let deadline = Instant::now() + Duration::from_secs(5);
635        while Instant::now() < deadline {
636            if condition() {
637                return true;
638            }
639            thread::sleep(Duration::from_millis(10));
640        }
641        condition()
642    }
643
644    #[test]
645    fn shutdown_controller_initial_state() {
646        init_test("shutdown_controller_initial_state");
647        let controller = ShutdownController::new();
648        let shutting_down = controller.is_shutting_down();
649        crate::assert_with_log!(
650            !shutting_down,
651            "controller not shutting down",
652            false,
653            shutting_down
654        );
655
656        let receiver = controller.subscribe();
657        let rx_shutdown = receiver.is_shutting_down();
658        crate::assert_with_log!(
659            !rx_shutdown,
660            "receiver not shutting down",
661            false,
662            rx_shutdown
663        );
664        crate::test_complete!("shutdown_controller_initial_state");
665    }
666
667    #[test]
668    fn shutdown_controller_initiates() {
669        init_test("shutdown_controller_initiates");
670        let controller = ShutdownController::new();
671        let receiver = controller.subscribe();
672
673        controller.shutdown();
674
675        let ctrl_shutdown = controller.is_shutting_down();
676        crate::assert_with_log!(
677            ctrl_shutdown,
678            "controller shutting down",
679            true,
680            ctrl_shutdown
681        );
682        let rx_shutdown = receiver.is_shutting_down();
683        crate::assert_with_log!(rx_shutdown, "receiver shutting down", true, rx_shutdown);
684        crate::test_complete!("shutdown_controller_initiates");
685    }
686
687    #[test]
688    fn shutdown_only_once() {
689        init_test("shutdown_only_once");
690        let controller = ShutdownController::new();
691
692        // Multiple shutdown calls should be idempotent.
693        controller.shutdown();
694        controller.shutdown();
695        controller.shutdown();
696
697        let shutting_down = controller.is_shutting_down();
698        crate::assert_with_log!(shutting_down, "shutting down", true, shutting_down);
699        crate::test_complete!("shutdown_only_once");
700    }
701
702    #[test]
703    fn multiple_receivers() {
704        init_test("multiple_receivers");
705        let controller = ShutdownController::new();
706        let rx1 = controller.subscribe();
707        let rx2 = controller.subscribe();
708        let rx3 = controller.subscribe();
709
710        let rx1_shutdown = rx1.is_shutting_down();
711        crate::assert_with_log!(!rx1_shutdown, "rx1 not shutting down", false, rx1_shutdown);
712        let rx2_shutdown = rx2.is_shutting_down();
713        crate::assert_with_log!(!rx2_shutdown, "rx2 not shutting down", false, rx2_shutdown);
714        let rx3_shutdown = rx3.is_shutting_down();
715        crate::assert_with_log!(!rx3_shutdown, "rx3 not shutting down", false, rx3_shutdown);
716
717        controller.shutdown();
718
719        let rx1_shutdown = rx1.is_shutting_down();
720        crate::assert_with_log!(rx1_shutdown, "rx1 shutting down", true, rx1_shutdown);
721        let rx2_shutdown = rx2.is_shutting_down();
722        crate::assert_with_log!(rx2_shutdown, "rx2 shutting down", true, rx2_shutdown);
723        let rx3_shutdown = rx3.is_shutting_down();
724        crate::assert_with_log!(rx3_shutdown, "rx3 shutting down", true, rx3_shutdown);
725        crate::test_complete!("multiple_receivers");
726    }
727
728    #[test]
729    fn receiver_wait_after_shutdown() {
730        init_test("receiver_wait_after_shutdown");
731        let controller = ShutdownController::new();
732        let mut receiver = controller.subscribe();
733
734        controller.shutdown();
735
736        // Wait should return immediately.
737        let mut fut = Box::pin(receiver.wait());
738        let ready = poll_once(&mut fut).is_ready();
739        crate::assert_with_log!(ready, "wait ready", true, ready);
740        crate::test_complete!("receiver_wait_after_shutdown");
741    }
742
743    #[test]
744    fn receiver_wait_before_shutdown() {
745        init_test("receiver_wait_before_shutdown");
746        let controller = Arc::new(ShutdownController::new());
747        let controller2 = Arc::clone(&controller);
748        let mut receiver = controller.subscribe();
749
750        let handle = thread::spawn(move || {
751            thread::sleep(Duration::from_millis(50));
752            controller2.shutdown();
753        });
754
755        // First poll should be pending.
756        let mut fut = Box::pin(receiver.wait());
757        let pending = poll_once(&mut fut).is_pending();
758        crate::assert_with_log!(pending, "wait pending", true, pending);
759
760        // Wait for shutdown.
761        handle.join().expect("thread panicked");
762
763        // Now should be ready.
764        let ready = poll_once(&mut fut).is_ready();
765        crate::assert_with_log!(ready, "wait ready", true, ready);
766        crate::test_complete!("receiver_wait_before_shutdown");
767    }
768
769    #[test]
770    fn receiver_clone() {
771        init_test("receiver_clone");
772        let controller = ShutdownController::new();
773        let rx1 = controller.subscribe();
774        let rx2 = rx1.clone();
775
776        let rx1_shutdown = rx1.is_shutting_down();
777        crate::assert_with_log!(!rx1_shutdown, "rx1 not shutting down", false, rx1_shutdown);
778        let rx2_shutdown = rx2.is_shutting_down();
779        crate::assert_with_log!(!rx2_shutdown, "rx2 not shutting down", false, rx2_shutdown);
780
781        controller.shutdown();
782
783        let rx1_shutdown = rx1.is_shutting_down();
784        crate::assert_with_log!(rx1_shutdown, "rx1 shutting down", true, rx1_shutdown);
785        let rx2_shutdown = rx2.is_shutting_down();
786        crate::assert_with_log!(rx2_shutdown, "rx2 shutting down", true, rx2_shutdown);
787        crate::test_complete!("receiver_clone");
788    }
789
790    #[test]
791    fn receiver_clone_preserves_state() {
792        init_test("receiver_clone_preserves_state");
793        let controller = ShutdownController::new();
794        controller.shutdown();
795
796        let rx1 = controller.subscribe();
797        let rx2 = rx1.clone();
798
799        // Both should see shutdown already initiated.
800        let rx1_shutdown = rx1.is_shutting_down();
801        crate::assert_with_log!(rx1_shutdown, "rx1 shutting down", true, rx1_shutdown);
802        let rx2_shutdown = rx2.is_shutting_down();
803        crate::assert_with_log!(rx2_shutdown, "rx2 shutting down", true, rx2_shutdown);
804        crate::test_complete!("receiver_clone_preserves_state");
805    }
806
807    #[test]
808    fn controller_clone() {
809        init_test("controller_clone");
810        let controller1 = ShutdownController::new();
811        let controller2 = controller1.clone();
812        let receiver = controller1.subscribe();
813
814        // Shutdown via clone.
815        controller2.shutdown();
816
817        // All should see it.
818        let ctrl1 = controller1.is_shutting_down();
819        crate::assert_with_log!(ctrl1, "controller1 shutting down", true, ctrl1);
820        let ctrl2 = controller2.is_shutting_down();
821        crate::assert_with_log!(ctrl2, "controller2 shutting down", true, ctrl2);
822        let rx_shutdown = receiver.is_shutting_down();
823        crate::assert_with_log!(rx_shutdown, "receiver shutting down", true, rx_shutdown);
824        crate::test_complete!("controller_clone");
825    }
826
827    #[cfg(any(unix, windows))]
828    #[test]
829    fn listen_for_signals_triggers_shutdown() {
830        init_test("listen_for_signals_triggers_shutdown");
831        let controller = Arc::new(ShutdownController::new());
832        let mut receiver = controller.subscribe();
833
834        controller.listen_for_signals();
835        inject_test_signal(SignalKind::terminate()).expect("test signal injection");
836
837        let mut fut = Box::pin(receiver.wait());
838        for _ in 0..50 {
839            if poll_once(&mut fut).is_ready() {
840                let shutting_down = controller.is_shutting_down();
841                crate::assert_with_log!(
842                    shutting_down,
843                    "controller shutting down via signal listener",
844                    true,
845                    shutting_down
846                );
847                crate::test_complete!("listen_for_signals_triggers_shutdown");
848                return;
849            }
850            thread::sleep(Duration::from_millis(10));
851        }
852
853        crate::assert_with_log!(
854            false,
855            "signal listener triggered shutdown before timeout",
856            true,
857            false
858        );
859    }
860
861    #[cfg(any(unix, windows))]
862    #[test]
863    fn listen_for_signals_is_idempotent() {
864        init_test("listen_for_signals_is_idempotent");
865        let controller = Arc::new(ShutdownController::new());
866
867        controller.listen_for_signals();
868        controller.listen_for_signals();
869
870        let started = controller
871            .state
872            .signal_listeners_started
873            .load(Ordering::Acquire);
874        crate::assert_with_log!(started, "signal listeners installed once", true, started);
875
876        controller.shutdown();
877        let shutting_down = controller.is_shutting_down();
878        crate::assert_with_log!(
879            shutting_down,
880            "manual shutdown still works",
881            true,
882            shutting_down
883        );
884        crate::test_complete!("listen_for_signals_is_idempotent");
885    }
886
887    #[test]
888    fn reload_controller_request_wakes_receiver_without_shutdown() {
889        init_test("reload_controller_request_wakes_receiver_without_shutdown");
890        let reload = ReloadController::new();
891        let shutdown = ShutdownController::new();
892        let mut reload_rx = reload.subscribe();
893        let shutdown_rx = shutdown.subscribe();
894
895        let sequence = reload.request_reload();
896        crate::assert_with_log!(sequence == 1, "reload sequence", 1, sequence);
897
898        let mut fut = Box::pin(reload_rx.wait());
899        let observed = futures_lite::future::block_on(fut.as_mut());
900        crate::assert_with_log!(observed == 1, "receiver observed sequence", 1, observed);
901        crate::assert_with_log!(
902            !shutdown.is_shutting_down(),
903            "reload does not trigger shutdown controller",
904            false,
905            shutdown.is_shutting_down()
906        );
907        crate::assert_with_log!(
908            !shutdown_rx.is_shutting_down(),
909            "reload does not trigger shutdown receiver",
910            false,
911            shutdown_rx.is_shutting_down()
912        );
913        crate::test_complete!("reload_controller_request_wakes_receiver_without_shutdown");
914    }
915
916    #[test]
917    fn reload_receiver_drains_queued_sequences() {
918        init_test("reload_receiver_drains_queued_sequences");
919        let reload = ReloadController::new();
920        let mut receiver = reload.subscribe();
921
922        reload.request_reload();
923        reload.request_reload();
924
925        let first = futures_lite::future::block_on(receiver.wait());
926        let second = futures_lite::future::block_on(receiver.wait());
927        crate::assert_with_log!(first == 1, "first reload sequence", 1, first);
928        crate::assert_with_log!(second == 2, "second reload sequence", 2, second);
929        crate::assert_with_log!(
930            receiver.seen_reload_count() == 2,
931            "receiver seen sequence",
932            2,
933            receiver.seen_reload_count()
934        );
935        crate::test_complete!("reload_receiver_drains_queued_sequences");
936    }
937
938    #[test]
939    fn reload_receiver_invokes_handler_and_reports_outcome() {
940        init_test("reload_receiver_invokes_handler_and_reports_outcome");
941        let reload = ReloadController::new();
942        let mut receiver = reload.subscribe();
943
944        reload.request_reload();
945        let completed =
946            futures_lite::future::block_on(receiver.handle_next_reload(|sequence| async move {
947                crate::assert_with_log!(sequence == 1, "handler sequence", 1, sequence);
948                Ok::<(), &'static str>(())
949            }));
950        crate::assert_with_log!(
951            completed == ReloadOutcome::Completed { sequence: 1 },
952            "handler completed",
953            ReloadOutcome::<&'static str>::Completed { sequence: 1 },
954            completed
955        );
956
957        reload.request_reload();
958        let failed =
959            futures_lite::future::block_on(receiver.handle_next_reload(|_sequence| async {
960                Err::<(), &'static str>("reload failed")
961            }));
962        crate::assert_with_log!(
963            failed
964                == ReloadOutcome::Failed {
965                    sequence: 2,
966                    error: "reload failed"
967                },
968            "handler failed",
969            ReloadOutcome::Failed {
970                sequence: 2,
971                error: "reload failed"
972            },
973            failed
974        );
975        crate::test_complete!("reload_receiver_invokes_handler_and_reports_outcome");
976    }
977
978    #[test]
979    fn reload_wait_or_shutdown_returns_none_when_shutdown_wins() {
980        init_test("reload_wait_or_shutdown_returns_none_when_shutdown_wins");
981        let reload = ReloadController::new();
982        let shutdown = ShutdownController::new();
983        let mut reload_rx = reload.subscribe();
984        let mut shutdown_rx = shutdown.subscribe();
985
986        shutdown.shutdown();
987
988        let observed = futures_lite::future::block_on(reload_rx.wait_or_shutdown(&mut shutdown_rx));
989        crate::assert_with_log!(
990            observed.is_none(),
991            "shutdown wins before reload request",
992            None::<u64>,
993            observed
994        );
995        crate::assert_with_log!(
996            reload_rx.seen_reload_count() == 0,
997            "no reload sequence consumed",
998            0,
999            reload_rx.seen_reload_count()
1000        );
1001        crate::test_complete!("reload_wait_or_shutdown_returns_none_when_shutdown_wins");
1002    }
1003
1004    #[test]
1005    fn reload_handler_reports_cancelled_when_shutdown_wins() {
1006        init_test("reload_handler_reports_cancelled_when_shutdown_wins");
1007        let reload = ReloadController::new();
1008        let shutdown = ShutdownController::new();
1009        let mut reload_rx = reload.subscribe();
1010        let mut shutdown_rx = shutdown.subscribe();
1011
1012        reload.request_reload();
1013        shutdown.shutdown();
1014
1015        let outcome = futures_lite::future::block_on(reload_rx.handle_next_reload_or_shutdown(
1016            &mut shutdown_rx,
1017            |sequence| {
1018                crate::assert_with_log!(sequence == 1, "handler sequence", 1, sequence);
1019                std::future::pending::<Result<(), &'static str>>()
1020            },
1021        ));
1022        crate::assert_with_log!(
1023            outcome == ReloadOutcome::<&'static str>::Cancelled { sequence: Some(1) },
1024            "shutdown cancels pending reload handler",
1025            ReloadOutcome::<&'static str>::Cancelled { sequence: Some(1) },
1026            outcome
1027        );
1028        crate::test_complete!("reload_handler_reports_cancelled_when_shutdown_wins");
1029    }
1030
1031    #[cfg(unix)]
1032    #[test]
1033    fn sighup_triggers_reload_only_and_sigterm_triggers_shutdown() {
1034        init_test("sighup_triggers_reload_only_and_sigterm_triggers_shutdown");
1035        let reload = Arc::new(ReloadController::new());
1036        let shutdown = Arc::new(ShutdownController::new());
1037        let mut reload_rx = reload.subscribe();
1038        let mut shutdown_rx = shutdown.subscribe();
1039
1040        let listener_installed = reload.listen_for_sighup().is_ok();
1041        crate::assert_with_log!(
1042            listener_installed,
1043            "install SIGHUP listener",
1044            true,
1045            listener_installed
1046        );
1047        if !listener_installed {
1048            return;
1049        }
1050        let shutdown_watches_sighup = watched_signal_kinds().contains(&SignalKind::hangup());
1051        crate::assert_with_log!(
1052            !shutdown_watches_sighup,
1053            "shutdown listener excludes SIGHUP",
1054            false,
1055            shutdown_watches_sighup
1056        );
1057
1058        let sighup_injected = inject_test_signal(SignalKind::hangup()).is_ok();
1059        crate::assert_with_log!(sighup_injected, "inject SIGHUP", true, sighup_injected);
1060        if !sighup_injected {
1061            return;
1062        }
1063        if !wait_until(|| reload.reload_count() > 0) {
1064            crate::assert_with_log!(false, "SIGHUP triggered reload before timeout", true, false);
1065            return;
1066        }
1067
1068        let mut reload_fut = Box::pin(reload_rx.wait());
1069        let reload_sequence = match poll_once(&mut reload_fut) {
1070            Poll::Ready(sequence) => sequence,
1071            Poll::Pending => {
1072                crate::assert_with_log!(
1073                    false,
1074                    "SIGHUP triggered reload before timeout",
1075                    true,
1076                    false
1077                );
1078                return;
1079            }
1080        };
1081        crate::assert_with_log!(
1082            reload_sequence == 1,
1083            "SIGHUP triggers reload sequence",
1084            1,
1085            reload_sequence
1086        );
1087        crate::assert_with_log!(
1088            !shutdown.is_shutting_down(),
1089            "SIGHUP does not trigger shutdown",
1090            false,
1091            shutdown.is_shutting_down()
1092        );
1093
1094        shutdown.listen_for_signals();
1095        let sigterm_injected = inject_test_signal(SignalKind::terminate()).is_ok();
1096        crate::assert_with_log!(sigterm_injected, "inject SIGTERM", true, sigterm_injected);
1097        if !sigterm_injected {
1098            return;
1099        }
1100        let mut fut = Box::pin(shutdown_rx.wait());
1101        if wait_until(|| poll_once(&mut fut).is_ready()) {
1102            crate::assert_with_log!(
1103                shutdown.is_shutting_down(),
1104                "SIGTERM triggers shutdown",
1105                true,
1106                shutdown.is_shutting_down()
1107            );
1108            crate::test_complete!("sighup_triggers_reload_only_and_sigterm_triggers_shutdown");
1109            return;
1110        }
1111
1112        crate::assert_with_log!(
1113            false,
1114            "SIGTERM triggered shutdown before timeout",
1115            true,
1116            false
1117        );
1118    }
1119
1120    #[test]
1121    fn shutdown_sequence_snapshot_scrubbed() {
1122        let controller = ShutdownController::new();
1123        let rx_a = controller.subscribe();
1124        let rx_b = controller.subscribe();
1125
1126        let before = json!({
1127            "controller": controller.is_shutting_down(),
1128            "receivers": [
1129                {"receiver": "[RX_A]", "shutting_down": rx_a.is_shutting_down()},
1130                {"receiver": "[RX_B]", "shutting_down": rx_b.is_shutting_down()},
1131            ],
1132        });
1133
1134        controller.shutdown();
1135
1136        insta::assert_json_snapshot!(
1137            "shutdown_sequence_scrubbed",
1138            json!({
1139                "before": before,
1140                "after": {
1141                    "controller": controller.is_shutting_down(),
1142                    "receivers": [
1143                        {"receiver": "[RX_A]", "shutting_down": rx_a.is_shutting_down()},
1144                        {"receiver": "[RX_B]", "shutting_down": rx_b.is_shutting_down()},
1145                    ],
1146                }
1147            })
1148        );
1149    }
1150}