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 }