lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

shutdown.rs (28649B)


      1 //! Bounded graceful-shutdown orchestration without signal installation.
      2 
      3 use core::{fmt, future::Future, pin::Pin, time::Duration};
      4 use std::error::Error;
      5 
      6 use crate::{HostError, MonotonicClock, MonotonicClockError, MonotonicDeadline};
      7 
      8 use super::{ShutdownPhase, SupervisionFailure, SupervisionFailureKind, TaskSupervisor};
      9 
     10 const ORDERED_PHASES: [ShutdownPhase; 7] = [
     11     ShutdownPhase::RejectNewMutations,
     12     ShutdownPhase::CancelIngress,
     13     ShutdownPhase::DrainOperations,
     14     ShutdownPhase::PersistRecoverableWork,
     15     ShutdownPhase::CloseNetwork,
     16     ShutdownPhase::CloseSqlite,
     17     ShutdownPhase::CloseSockets,
     18 ];
     19 
     20 /// Service-owned asynchronous work performed when entering one shutdown phase.
     21 pub type ShutdownPhaseFuture<'a> = Pin<Box<dyn Future<Output = Result<(), HostError>> + Send + 'a>>;
     22 
     23 /// Executes service-specific phase work without transferring lifecycle ownership.
     24 ///
     25 /// If the caller cancels [`GracefulShutdown::run`] before one `enter` future
     26 /// completes, a later call re-enters that incomplete phase under the original
     27 /// deadline. Implementations must therefore make each phase idempotent and
     28 /// cancellation safe. A completed phase is never re-entered.
     29 pub trait ShutdownPhaseHandler: Send {
     30     fn enter(&mut self, phase: ShutdownPhase) -> ShutdownPhaseFuture<'_>;
     31 }
     32 
     33 /// Reusable, idempotent bounded shutdown coordinator.
     34 pub struct GracefulShutdown {
     35     grace: Duration,
     36     progress: Option<ShutdownProgress>,
     37     completed: Option<ShutdownSummary>,
     38     phase_failure: Option<ShutdownPhaseFailure>,
     39     task_failure: Option<SupervisionFailure>,
     40 }
     41 
     42 #[derive(Clone, Copy)]
     43 struct ShutdownProgress {
     44     deadline: MonotonicDeadline,
     45     runtime_deadline: tokio::time::Instant,
     46     phase_index: usize,
     47     stage: ShutdownPhaseStage,
     48     disposition: ShutdownDisposition,
     49 }
     50 
     51 #[derive(Clone, Copy, PartialEq, Eq)]
     52 enum ShutdownPhaseStage {
     53     Enter,
     54     Drain,
     55     Abort(ShutdownDisposition),
     56 }
     57 
     58 impl GracefulShutdown {
     59     pub fn new(grace: Duration) -> Result<Self, ShutdownConfigError> {
     60         if grace.is_zero() {
     61             return Err(ShutdownConfigError::ZeroGrace);
     62         }
     63         Ok(Self {
     64             grace,
     65             progress: None,
     66             completed: None,
     67             phase_failure: None,
     68             task_failure: None,
     69         })
     70     }
     71 
     72     #[must_use]
     73     pub const fn grace(&self) -> Duration {
     74         self.grace
     75     }
     76 
     77     /// Returns the trusted phase failure retained by the completed run, if any.
     78     #[must_use]
     79     pub const fn phase_failure(&self) -> Option<&ShutdownPhaseFailure> {
     80         self.phase_failure.as_ref()
     81     }
     82 
     83     /// Returns the trusted task failure retained by the completed run, if any.
     84     #[must_use]
     85     pub const fn task_failure(&self) -> Option<&SupervisionFailure> {
     86         self.task_failure.as_ref()
     87     }
     88 
     89     /// Runs or resumes the exact shutdown sequence under one retained absolute deadline.
     90     ///
     91     /// Cancelling this future retains the last completed handler/drain boundary.
     92     /// A retry resumes that boundary with the original remaining duration.
     93     /// Completed runs return the first summary unchanged.
     94     pub async fn run<C, F>(
     95         &mut self,
     96         clock: &C,
     97         supervisor: &mut TaskSupervisor,
     98         handler: &mut dyn ShutdownPhaseHandler,
     99         force: F,
    100     ) -> Result<ShutdownSummary, ShutdownStartError>
    101     where
    102         C: MonotonicClock,
    103         F: Future<Output = ()> + Send,
    104     {
    105         if let Some(completed) = self.completed {
    106             return Ok(completed);
    107         }
    108         if self.progress.is_none() {
    109             let deadline = clock
    110                 .deadline_after(self.grace)
    111                 .map_err(ShutdownStartError::Deadline)?;
    112             let runtime_deadline = tokio::time::Instant::now().checked_add(self.grace).ok_or(
    113                 ShutdownStartError::Deadline(MonotonicClockError::DeadlineOverflow),
    114             )?;
    115             self.progress = Some(ShutdownProgress {
    116                 deadline,
    117                 runtime_deadline,
    118                 phase_index: 0,
    119                 stage: ShutdownPhaseStage::Enter,
    120                 disposition: ShutdownDisposition::Completed,
    121             });
    122         }
    123         tokio::pin!(force);
    124 
    125         loop {
    126             let progress = self
    127                 .progress
    128                 .expect("shutdown progress must be initialized");
    129             if let ShutdownPhaseStage::Abort(disposition) = progress.stage {
    130                 supervisor.abort_and_drain().await;
    131                 return Ok(self.complete(progress.deadline, disposition));
    132             }
    133             let Some(&phase) = ORDERED_PHASES.get(progress.phase_index) else {
    134                 return Ok(self.complete(progress.deadline, progress.disposition));
    135             };
    136             let deadline_wait = tokio::time::sleep_until(progress.runtime_deadline);
    137             tokio::pin!(deadline_wait);
    138 
    139             let outcome = match progress.stage {
    140                 ShutdownPhaseStage::Enter => wait_bounded(
    141                     || async {
    142                         supervisor.request_phase_cancellation(phase);
    143                         handler.enter(phase).await
    144                     },
    145                     force.as_mut(),
    146                     deadline_wait.as_mut(),
    147                 )
    148                 .await
    149                 .map(ShutdownPhaseOutcome::Entered),
    150                 ShutdownPhaseStage::Drain => wait_bounded(
    151                     || supervisor.supervise_phase(phase),
    152                     force.as_mut(),
    153                     deadline_wait.as_mut(),
    154                 )
    155                 .await
    156                 .map(ShutdownPhaseOutcome::Drained),
    157                 ShutdownPhaseStage::Abort(_) => unreachable!("abort stage handled before phase"),
    158             };
    159 
    160             match outcome {
    161                 BoundedWait::Completed(ShutdownPhaseOutcome::Entered(result)) => {
    162                     if let Err(error) = result {
    163                         let unfinished = supervisor.unfinished_work();
    164                         if self.phase_failure.is_none() {
    165                             self.phase_failure = Some(ShutdownPhaseFailure { phase, error });
    166                         }
    167                         self.retain_first_disposition(ShutdownDisposition::PhaseFailed {
    168                             phase,
    169                             unfinished,
    170                         });
    171                     }
    172                     self.progress.as_mut().expect("shutdown progress").stage =
    173                         ShutdownPhaseStage::Drain;
    174                 }
    175                 BoundedWait::Completed(ShutdownPhaseOutcome::Drained(result)) => match result {
    176                     Ok(_exits) => {
    177                         let progress = self.progress.as_mut().expect("shutdown progress");
    178                         progress.phase_index += 1;
    179                         progress.stage = ShutdownPhaseStage::Enter;
    180                     }
    181                     Err(failure) => {
    182                         let kind = failure.kind();
    183                         if self.task_failure.is_none() {
    184                             self.task_failure = Some(failure);
    185                         }
    186                         self.retain_first_disposition(ShutdownDisposition::TaskFailed { kind });
    187                     }
    188                 },
    189                 BoundedWait::Forced => {
    190                     let unfinished = supervisor.unfinished_work();
    191                     self.progress.as_mut().expect("shutdown progress").stage =
    192                         ShutdownPhaseStage::Abort(ShutdownDisposition::Forced { unfinished });
    193                 }
    194                 BoundedWait::GraceExpired => {
    195                     let unfinished = supervisor.unfinished_work();
    196                     self.progress.as_mut().expect("shutdown progress").stage =
    197                         ShutdownPhaseStage::Abort(ShutdownDisposition::GraceExpired { unfinished });
    198                 }
    199             }
    200         }
    201     }
    202 
    203     fn retain_first_disposition(&mut self, disposition: ShutdownDisposition) {
    204         let progress = self.progress.as_mut().expect("shutdown progress");
    205         if progress.disposition == ShutdownDisposition::Completed {
    206             progress.disposition = disposition;
    207         }
    208     }
    209 
    210     fn complete(
    211         &mut self,
    212         deadline: MonotonicDeadline,
    213         disposition: ShutdownDisposition,
    214     ) -> ShutdownSummary {
    215         let summary = ShutdownSummary {
    216             deadline,
    217             disposition,
    218         };
    219         self.completed = Some(summary);
    220         self.progress = None;
    221         summary
    222     }
    223 }
    224 
    225 enum ShutdownPhaseOutcome {
    226     Entered(Result<(), HostError>),
    227     Drained(Result<Vec<super::SupervisedTaskExit>, SupervisionFailure>),
    228 }
    229 
    230 enum BoundedWait<T> {
    231     Completed(T),
    232     Forced,
    233     GraceExpired,
    234 }
    235 
    236 impl<T> BoundedWait<T> {
    237     fn map<U>(self, map: impl FnOnce(T) -> U) -> BoundedWait<U> {
    238         match self {
    239             Self::Completed(value) => BoundedWait::Completed(map(value)),
    240             Self::Forced => BoundedWait::Forced,
    241             Self::GraceExpired => BoundedWait::GraceExpired,
    242         }
    243     }
    244 }
    245 
    246 async fn wait_bounded<T, MakeWork, Work, Force>(
    247     make_work: MakeWork,
    248     mut force: Pin<&mut Force>,
    249     mut deadline: Pin<&mut tokio::time::Sleep>,
    250 ) -> BoundedWait<T>
    251 where
    252     MakeWork: FnOnce() -> Work,
    253     Work: Future<Output = T>,
    254     Force: Future<Output = ()>,
    255 {
    256     tokio::select! {
    257         biased;
    258         () = force.as_mut() => BoundedWait::Forced,
    259         () = deadline.as_mut() => BoundedWait::GraceExpired,
    260         completed = async move { make_work().await } => BoundedWait::Completed(completed),
    261     }
    262 }
    263 
    264 /// Trusted phase failure retained separately from the stable shutdown summary.
    265 pub struct ShutdownPhaseFailure {
    266     phase: ShutdownPhase,
    267     error: HostError,
    268 }
    269 
    270 impl ShutdownPhaseFailure {
    271     #[must_use]
    272     pub const fn phase(&self) -> ShutdownPhase {
    273         self.phase
    274     }
    275 
    276     /// Returns the original host error for trusted internal inspection.
    277     #[must_use]
    278     pub const fn error(&self) -> &HostError {
    279         &self.error
    280     }
    281 }
    282 
    283 impl fmt::Debug for ShutdownPhaseFailure {
    284     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    285         formatter
    286             .debug_struct("ShutdownPhaseFailure")
    287             .field("phase", &self.phase)
    288             .field("error", &self.error.safe_error())
    289             .field("source", &self.error.source().map(|_| "<redacted>"))
    290             .finish()
    291     }
    292 }
    293 
    294 impl fmt::Display for ShutdownPhaseFailure {
    295     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    296         formatter.write_str("service shutdown phase failed")
    297     }
    298 }
    299 
    300 impl Error for ShutdownPhaseFailure {
    301     fn source(&self) -> Option<&(dyn Error + 'static)> {
    302         Some(&self.error)
    303     }
    304 }
    305 
    306 /// Immutable result of one shutdown run.
    307 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    308 pub struct ShutdownSummary {
    309     deadline: MonotonicDeadline,
    310     disposition: ShutdownDisposition,
    311 }
    312 
    313 impl ShutdownSummary {
    314     #[must_use]
    315     pub const fn deadline(self) -> MonotonicDeadline {
    316         self.deadline
    317     }
    318 
    319     #[must_use]
    320     pub const fn disposition(self) -> ShutdownDisposition {
    321         self.disposition
    322     }
    323 }
    324 
    325 /// Final bounded shutdown classification.
    326 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    327 pub enum ShutdownDisposition {
    328     Completed,
    329     TaskFailed {
    330         kind: SupervisionFailureKind,
    331     },
    332     PhaseFailed {
    333         phase: ShutdownPhase,
    334         unfinished: UnfinishedWork,
    335     },
    336     GraceExpired {
    337         unfinished: UnfinishedWork,
    338     },
    339     Forced {
    340         unfinished: UnfinishedWork,
    341     },
    342 }
    343 
    344 /// Whether forcibly stopped work may safely recover after restart.
    345 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    346 pub enum UnfinishedWork {
    347     None,
    348     RecoverableOptional,
    349     FatalAuthoritative,
    350 }
    351 
    352 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    353 pub enum ShutdownConfigError {
    354     ZeroGrace,
    355 }
    356 
    357 impl fmt::Display for ShutdownConfigError {
    358     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    359         formatter.write_str("shutdown grace must be greater than zero")
    360     }
    361 }
    362 
    363 impl Error for ShutdownConfigError {}
    364 
    365 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    366 pub enum ShutdownStartError {
    367     Deadline(MonotonicClockError),
    368 }
    369 
    370 impl fmt::Display for ShutdownStartError {
    371     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    372         formatter.write_str("shutdown grace deadline could not be represented")
    373     }
    374 }
    375 
    376 impl Error for ShutdownStartError {
    377     fn source(&self) -> Option<&(dyn Error + 'static)> {
    378         match self {
    379             Self::Deadline(error) => Some(error),
    380         }
    381     }
    382 }
    383 
    384 #[cfg(test)]
    385 mod tests {
    386     use core::future::{pending, ready};
    387     use std::{
    388         error::Error,
    389         sync::{Arc, Mutex},
    390     };
    391 
    392     use crate::{HostErrorKind, MonotonicTime, TaskClassification, TaskMetadata, TaskName};
    393 
    394     use super::*;
    395 
    396     struct FakeClock {
    397         now: MonotonicTime,
    398     }
    399 
    400     impl MonotonicClock for FakeClock {
    401         fn now_monotonic(&self) -> MonotonicTime {
    402             self.now
    403         }
    404     }
    405 
    406     struct RecordingHandler {
    407         phases: Arc<Mutex<Vec<ShutdownPhase>>>,
    408         fail_at: Option<ShutdownPhase>,
    409     }
    410 
    411     impl ShutdownPhaseHandler for RecordingHandler {
    412         fn enter(&mut self, phase: ShutdownPhase) -> ShutdownPhaseFuture<'_> {
    413             self.phases.lock().unwrap().push(phase);
    414             let fail = self.fail_at == Some(phase);
    415             Box::pin(async move {
    416                 if fail {
    417                     Err(HostError::with_source(
    418                         HostErrorKind::Lifecycle,
    419                         SensitiveCause,
    420                     ))
    421                 } else {
    422                     Ok(())
    423                 }
    424             })
    425         }
    426     }
    427 
    428     #[derive(Debug)]
    429     struct SensitiveCause;
    430 
    431     impl fmt::Display for SensitiveCause {
    432         fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    433             formatter.write_str("sensitive shutdown detail")
    434         }
    435     }
    436 
    437     impl Error for SensitiveCause {}
    438 
    439     struct BlockingHandler {
    440         phases: Arc<Mutex<Vec<ShutdownPhase>>>,
    441         block_at: ShutdownPhase,
    442         entered: Arc<tokio::sync::Notify>,
    443     }
    444 
    445     impl ShutdownPhaseHandler for BlockingHandler {
    446         fn enter(&mut self, phase: ShutdownPhase) -> ShutdownPhaseFuture<'_> {
    447             self.phases.lock().unwrap().push(phase);
    448             if phase == self.block_at {
    449                 self.entered.notify_one();
    450                 Box::pin(pending())
    451             } else {
    452                 Box::pin(ready(Ok(())))
    453             }
    454         }
    455     }
    456 
    457     fn clock() -> FakeClock {
    458         FakeClock {
    459             now: MonotonicTime::from_duration_since_origin(Duration::from_secs(5)),
    460         }
    461     }
    462 
    463     fn handler(
    464         fail_at: Option<ShutdownPhase>,
    465     ) -> (RecordingHandler, Arc<Mutex<Vec<ShutdownPhase>>>) {
    466         let phases = Arc::new(Mutex::new(Vec::new()));
    467         (
    468             RecordingHandler {
    469                 phases: Arc::clone(&phases),
    470                 fail_at,
    471             },
    472             phases,
    473         )
    474     }
    475 
    476     fn task_metadata(name: &str, classification: TaskClassification) -> TaskMetadata {
    477         let shutdown_phase = classification
    478             .requires_shutdown_phase()
    479             .then_some(ShutdownPhase::CancelIngress);
    480         TaskMetadata::new(TaskName::new(name).unwrap(), classification, shutdown_phase).unwrap()
    481     }
    482 
    483     #[tokio::test(start_paused = true)]
    484     async fn phases_are_ordered_and_repeated_run_is_idempotent() {
    485         let mut supervisor = TaskSupervisor::new();
    486         supervisor
    487             .spawn(
    488                 task_metadata("critical_worker", TaskClassification::Critical),
    489                 |token| async move {
    490                     token.cancelled().await;
    491                     Ok(())
    492                 },
    493             )
    494             .unwrap();
    495         let (mut handler, phases) = handler(None);
    496         let mut shutdown = GracefulShutdown::new(Duration::from_secs(30)).unwrap();
    497         assert_eq!(shutdown.grace(), Duration::from_secs(30));
    498 
    499         let first = shutdown
    500             .run(&clock(), &mut supervisor, &mut handler, pending())
    501             .await
    502             .unwrap();
    503         let second = shutdown
    504             .run(&clock(), &mut supervisor, &mut handler, ready(()))
    505             .await
    506             .unwrap();
    507 
    508         assert_eq!(first, second);
    509         assert_eq!(first.disposition(), ShutdownDisposition::Completed);
    510         assert_eq!(
    511             first.deadline().time().duration_since_origin(),
    512             Duration::from_secs(35)
    513         );
    514         assert_eq!(*phases.lock().unwrap(), ORDERED_PHASES);
    515         assert!(supervisor.is_empty());
    516     }
    517 
    518     #[tokio::test(start_paused = true)]
    519     async fn grace_timeout_classifies_and_drains_recoverable_and_fatal_work() {
    520         for (classification, expected) in [
    521             (
    522                 TaskClassification::Optional,
    523                 UnfinishedWork::RecoverableOptional,
    524             ),
    525             (
    526                 TaskClassification::Critical,
    527                 UnfinishedWork::FatalAuthoritative,
    528             ),
    529         ] {
    530             let mut supervisor = TaskSupervisor::new();
    531             supervisor
    532                 .spawn(task_metadata("stuck_worker", classification), |_| async {
    533                     pending::<()>().await;
    534                     Ok(())
    535                 })
    536                 .unwrap();
    537             let (mut handler, _) = handler(None);
    538             let mut shutdown = GracefulShutdown::new(Duration::from_secs(1)).unwrap();
    539 
    540             let summary = shutdown
    541                 .run(&clock(), &mut supervisor, &mut handler, pending())
    542                 .await
    543                 .unwrap();
    544             assert_eq!(
    545                 summary.disposition(),
    546                 ShutdownDisposition::GraceExpired {
    547                     unfinished: expected
    548                 }
    549             );
    550             assert!(supervisor.is_empty());
    551         }
    552     }
    553 
    554     #[tokio::test]
    555     async fn force_input_aborts_and_classifies_active_authoritative_work() {
    556         let mut supervisor = TaskSupervisor::new();
    557         supervisor
    558             .spawn(
    559                 task_metadata("stuck_worker", TaskClassification::Critical),
    560                 |_| async {
    561                     pending::<()>().await;
    562                     Ok(())
    563                 },
    564             )
    565             .unwrap();
    566         let (mut handler, phases) = handler(None);
    567         let mut shutdown = GracefulShutdown::new(Duration::from_secs(30)).unwrap();
    568 
    569         let summary = shutdown
    570             .run(&clock(), &mut supervisor, &mut handler, ready(()))
    571             .await
    572             .unwrap();
    573         assert_eq!(
    574             summary.disposition(),
    575             ShutdownDisposition::Forced {
    576                 unfinished: UnfinishedWork::FatalAuthoritative
    577             }
    578         );
    579         assert!(phases.lock().unwrap().is_empty());
    580         assert!(supervisor.is_empty());
    581     }
    582 
    583     #[tokio::test]
    584     async fn phase_failure_is_fatal_and_records_exact_phase() {
    585         let mut supervisor = TaskSupervisor::new();
    586         let (mut handler, phases) = handler(Some(ShutdownPhase::PersistRecoverableWork));
    587         let mut shutdown = GracefulShutdown::new(Duration::from_secs(30)).unwrap();
    588 
    589         let summary = shutdown
    590             .run(&clock(), &mut supervisor, &mut handler, pending())
    591             .await
    592             .unwrap();
    593         assert_eq!(
    594             summary.disposition(),
    595             ShutdownDisposition::PhaseFailed {
    596                 phase: ShutdownPhase::PersistRecoverableWork,
    597                 unfinished: UnfinishedWork::None,
    598             }
    599         );
    600         assert_eq!(*phases.lock().unwrap(), ORDERED_PHASES);
    601         let failure = shutdown.phase_failure().unwrap();
    602         assert_eq!(failure.phase(), ShutdownPhase::PersistRecoverableWork);
    603         assert_eq!(
    604             failure.error().source().map(ToString::to_string).as_deref(),
    605             Some("sensitive shutdown detail")
    606         );
    607         assert!(!failure.to_string().contains("sensitive"));
    608         assert!(!format!("{failure:?}").contains("sensitive shutdown detail"));
    609         assert!(failure.source().is_some());
    610     }
    611 
    612     #[tokio::test]
    613     async fn fatal_task_result_still_runs_the_remaining_close_phases() {
    614         let mut supervisor = TaskSupervisor::new();
    615         supervisor
    616             .spawn(
    617                 task_metadata("fatal_worker", TaskClassification::Critical),
    618                 |_| async { Err(HostError::new(HostErrorKind::TaskFailure)) },
    619             )
    620             .unwrap();
    621         let (mut handler, phases) = handler(None);
    622         let mut shutdown = GracefulShutdown::new(Duration::from_secs(30)).unwrap();
    623 
    624         let summary = shutdown
    625             .run(&clock(), &mut supervisor, &mut handler, pending())
    626             .await
    627             .unwrap();
    628         assert_eq!(
    629             summary.disposition(),
    630             ShutdownDisposition::TaskFailed {
    631                 kind: SupervisionFailureKind::TaskReturnedError
    632             }
    633         );
    634         assert_eq!(*phases.lock().unwrap(), ORDERED_PHASES);
    635         assert!(supervisor.is_empty());
    636         let failure = shutdown.task_failure().unwrap();
    637         assert_eq!(failure.metadata().unwrap().name().as_str(), "fatal_worker");
    638         assert!(failure.source().is_some());
    639     }
    640 
    641     #[tokio::test(start_paused = true)]
    642     async fn cancelled_run_retains_its_absolute_deadline_and_completed_phase_progress() {
    643         let mut supervisor = TaskSupervisor::new();
    644         let phases = Arc::new(Mutex::new(Vec::new()));
    645         let entered = Arc::new(tokio::sync::Notify::new());
    646         let mut handler = BlockingHandler {
    647             phases: Arc::clone(&phases),
    648             block_at: ShutdownPhase::PersistRecoverableWork,
    649             entered: Arc::clone(&entered),
    650         };
    651         let mut shutdown = GracefulShutdown::new(Duration::from_secs(10)).unwrap();
    652         let shutdown_clock = clock();
    653 
    654         {
    655             let mut run =
    656                 Box::pin(shutdown.run(&shutdown_clock, &mut supervisor, &mut handler, pending()));
    657             tokio::select! {
    658                 () = entered.notified() => {}
    659                 result = &mut run => panic!("blocked phase completed unexpectedly: {result:?}"),
    660             }
    661         }
    662         assert_eq!(*phases.lock().unwrap(), ORDERED_PHASES[..=3].to_vec());
    663 
    664         tokio::time::advance(Duration::from_secs(10)).await;
    665         let summary = tokio::time::timeout(
    666             Duration::from_millis(1),
    667             shutdown.run(&shutdown_clock, &mut supervisor, &mut handler, pending()),
    668         )
    669         .await
    670         .expect("retry must use the original elapsed runtime deadline")
    671         .unwrap();
    672         assert_eq!(
    673             summary.disposition(),
    674             ShutdownDisposition::GraceExpired {
    675                 unfinished: UnfinishedWork::None
    676             }
    677         );
    678         assert_eq!(
    679             summary.deadline().time().duration_since_origin(),
    680             Duration::from_secs(15)
    681         );
    682         assert_eq!(
    683             *phases.lock().unwrap(),
    684             ORDERED_PHASES[..=3].to_vec(),
    685             "an elapsed retry must not reconstruct the incomplete phase"
    686         );
    687     }
    688 
    689     #[tokio::test]
    690     async fn cancellation_during_phase_drain_resumes_without_reentering_completed_handler() {
    691         let (release, wait_for_release) = tokio::sync::oneshot::channel();
    692         let cleanup_started = Arc::new(tokio::sync::Notify::new());
    693         let mut supervisor = TaskSupervisor::new();
    694         let started = Arc::clone(&cleanup_started);
    695         let metadata = TaskMetadata::new(
    696             TaskName::new("mutation_gate").unwrap(),
    697             TaskClassification::Critical,
    698             Some(ShutdownPhase::RejectNewMutations),
    699         )
    700         .unwrap();
    701         supervisor
    702             .spawn(metadata, move |token| async move {
    703                 token.cancelled().await;
    704                 started.notify_one();
    705                 let _ = wait_for_release.await;
    706                 Ok(())
    707             })
    708             .unwrap();
    709         let (mut handler, phases) = handler(None);
    710         let mut shutdown = GracefulShutdown::new(Duration::from_secs(30)).unwrap();
    711         let shutdown_clock = clock();
    712 
    713         {
    714             let mut run =
    715                 Box::pin(shutdown.run(&shutdown_clock, &mut supervisor, &mut handler, pending()));
    716             tokio::select! {
    717                 () = cleanup_started.notified() => {}
    718                 result = &mut run => panic!("phase drain completed unexpectedly: {result:?}"),
    719             }
    720         }
    721         assert_eq!(
    722             *phases.lock().unwrap(),
    723             vec![ShutdownPhase::RejectNewMutations]
    724         );
    725 
    726         release.send(()).unwrap();
    727         let summary = shutdown
    728             .run(&shutdown_clock, &mut supervisor, &mut handler, pending())
    729             .await
    730             .unwrap();
    731         assert_eq!(summary.disposition(), ShutdownDisposition::Completed);
    732         assert_eq!(*phases.lock().unwrap(), ORDERED_PHASES);
    733         assert!(supervisor.is_empty());
    734     }
    735 
    736     #[tokio::test(start_paused = true)]
    737     async fn cancellation_while_fatal_peers_drain_retains_the_first_task_failure() {
    738         let (release, wait_for_release) = tokio::sync::oneshot::channel();
    739         let peer_waiting = Arc::new(tokio::sync::Notify::new());
    740         let mut supervisor = TaskSupervisor::new();
    741         let phase_metadata = |name| {
    742             TaskMetadata::new(
    743                 TaskName::new(name).unwrap(),
    744                 TaskClassification::Critical,
    745                 Some(ShutdownPhase::RejectNewMutations),
    746             )
    747             .unwrap()
    748         };
    749         supervisor
    750             .spawn(phase_metadata("failing_gate"), |_| async {
    751                 Err(HostError::new(HostErrorKind::TaskFailure))
    752             })
    753             .unwrap();
    754         let waiting = Arc::clone(&peer_waiting);
    755         supervisor
    756             .spawn(phase_metadata("draining_peer"), move |token| async move {
    757                 token.cancelled().await;
    758                 tokio::time::sleep(Duration::from_secs(1)).await;
    759                 waiting.notify_one();
    760                 let _ = wait_for_release.await;
    761                 Ok(())
    762             })
    763             .unwrap();
    764         let (mut handler, phases) = handler(None);
    765         let mut shutdown = GracefulShutdown::new(Duration::from_secs(30)).unwrap();
    766         let shutdown_clock = clock();
    767 
    768         {
    769             let mut run =
    770                 Box::pin(shutdown.run(&shutdown_clock, &mut supervisor, &mut handler, pending()));
    771             tokio::select! {
    772                 () = peer_waiting.notified() => {}
    773                 result = &mut run => panic!("fatal peer drain completed unexpectedly: {result:?}"),
    774             }
    775         }
    776         assert_eq!(
    777             shutdown.task_failure().map(SupervisionFailure::kind),
    778             Some(SupervisionFailureKind::TaskReturnedError)
    779         );
    780 
    781         release.send(()).unwrap();
    782         let summary = shutdown
    783             .run(&shutdown_clock, &mut supervisor, &mut handler, pending())
    784             .await
    785             .unwrap();
    786         assert_eq!(
    787             summary.disposition(),
    788             ShutdownDisposition::TaskFailed {
    789                 kind: SupervisionFailureKind::TaskReturnedError
    790             }
    791         );
    792         assert_eq!(*phases.lock().unwrap(), ORDERED_PHASES);
    793         assert!(supervisor.is_empty());
    794     }
    795 
    796     #[test]
    797     fn zero_grace_fails_closed() {
    798         let error = GracefulShutdown::new(Duration::ZERO).err().unwrap();
    799         assert_eq!(error, ShutdownConfigError::ZeroGrace);
    800         assert_eq!(
    801             error.to_string(),
    802             "shutdown grace must be greater than zero"
    803         );
    804     }
    805 
    806     #[tokio::test]
    807     async fn unrepresentable_deadline_fails_before_side_effects() {
    808         let overflow_clock = FakeClock {
    809             now: MonotonicTime::from_duration_since_origin(Duration::MAX),
    810         };
    811         let mut supervisor = TaskSupervisor::new();
    812         let cancellation = supervisor.cancellation_token();
    813         let (mut handler, phases) = handler(None);
    814         let mut shutdown = GracefulShutdown::new(Duration::from_nanos(1)).unwrap();
    815 
    816         let error = shutdown
    817             .run(&overflow_clock, &mut supervisor, &mut handler, pending())
    818             .await
    819             .unwrap_err();
    820         assert_eq!(
    821             error,
    822             ShutdownStartError::Deadline(MonotonicClockError::DeadlineOverflow)
    823         );
    824         assert_eq!(
    825             error.to_string(),
    826             "shutdown grace deadline could not be represented"
    827         );
    828         assert!(error.source().is_some());
    829         assert!(phases.lock().unwrap().is_empty());
    830         assert!(!cancellation.is_cancelled());
    831         assert!(shutdown.phase_failure().is_none());
    832         assert!(shutdown.task_failure().is_none());
    833     }
    834 }