commit 62fd018c1d76d4d894af02e0d3edc9d6fd242ee3
parent e1354dc3456d8465ff5e2e6a0cf77fea222221e6
Author: triesap <tyson@radroots.org>
Date: Sun, 23 Aug 2026 05:33:08 +0000
service-host: make shutdown phase aware
- cancel only tasks assigned to each orderly shutdown phase
- retain one absolute deadline and completed boundaries across retry
- preserve the first failure while later close phases continue
- verify lifecycle, package, API, Clippy, doctest, and Rustdoc gates
Diffstat:
4 files changed, 479 insertions(+), 75 deletions(-)
diff --git a/crates/service_host/README.md b/crates/service_host/README.md
@@ -121,6 +121,17 @@ or own service-domain policy. Process binaries inject clocks and entropy,
normalize signals through [`ProcessSignalAdapter`], register every authoritative
task with [`TaskSupervisor`], and execute [`GracefulShutdown`] explicitly.
+Shutdown task cancellation is phase aware. Entering a phase cancels only tasks
+assigned to that phase and does not advance until their joins are observed;
+bounded one-shot work drains without cancellation during `DrainOperations`,
+while a fatal task outcome still
+cancels and joins the complete graph. `GracefulShutdown` retains one absolute
+deadline plus the completed handler/drain boundary across cancellation and
+retry, so no phase or cleanup attempt receives a fresh grace period. A caller-
+cancelled incomplete handler may be entered again and must be idempotent and
+cancellation safe. The first phase or task failure is retained while later
+close phases continue as long as the original deadline remains.
+
## Supported targets and publication
The generic host/config/status/lifecycle surfaces support the workspace's
diff --git a/crates/service_host/src/lifecycle/shutdown.rs b/crates/service_host/src/lifecycle/shutdown.rs
@@ -21,6 +21,11 @@ const ORDERED_PHASES: [ShutdownPhase; 7] = [
pub type ShutdownPhaseFuture<'a> = Pin<Box<dyn Future<Output = Result<(), HostError>> + Send + 'a>>;
/// Executes service-specific phase work without transferring lifecycle ownership.
+///
+/// If the caller cancels [`GracefulShutdown::run`] before one `enter` future
+/// completes, a later call re-enters that incomplete phase under the original
+/// deadline. Implementations must therefore make each phase idempotent and
+/// cancellation safe. A completed phase is never re-entered.
pub trait ShutdownPhaseHandler: Send {
fn enter(&mut self, phase: ShutdownPhase) -> ShutdownPhaseFuture<'_>;
}
@@ -28,11 +33,28 @@ pub trait ShutdownPhaseHandler: Send {
/// Reusable, idempotent bounded shutdown coordinator.
pub struct GracefulShutdown {
grace: Duration,
+ progress: Option<ShutdownProgress>,
completed: Option<ShutdownSummary>,
phase_failure: Option<ShutdownPhaseFailure>,
task_failure: Option<SupervisionFailure>,
}
+#[derive(Clone, Copy)]
+struct ShutdownProgress {
+ deadline: MonotonicDeadline,
+ runtime_deadline: tokio::time::Instant,
+ phase_index: usize,
+ stage: ShutdownPhaseStage,
+ disposition: ShutdownDisposition,
+}
+
+#[derive(Clone, Copy, PartialEq, Eq)]
+enum ShutdownPhaseStage {
+ Enter,
+ Drain,
+ Abort(ShutdownDisposition),
+}
+
impl GracefulShutdown {
pub fn new(grace: Duration) -> Result<Self, ShutdownConfigError> {
if grace.is_zero() {
@@ -40,6 +62,7 @@ impl GracefulShutdown {
}
Ok(Self {
grace,
+ progress: None,
completed: None,
phase_failure: None,
task_failure: None,
@@ -63,7 +86,11 @@ impl GracefulShutdown {
self.task_failure.as_ref()
}
- /// Runs the exact shutdown sequence once; later calls return the first summary unchanged.
+ /// Runs or resumes the exact shutdown sequence under one retained absolute deadline.
+ ///
+ /// Cancelling this future retains the last completed handler/drain boundary.
+ /// A retry resumes that boundary with the original remaining duration.
+ /// Completed runs return the first summary unchanged.
pub async fn run<C, F>(
&mut self,
clock: &C,
@@ -78,82 +105,106 @@ impl GracefulShutdown {
if let Some(completed) = self.completed {
return Ok(completed);
}
- let deadline = clock
- .deadline_after(self.grace)
- .map_err(ShutdownStartError::Deadline)?;
- let deadline_wait = tokio::time::sleep(self.grace);
- tokio::pin!(deadline_wait);
+ if self.progress.is_none() {
+ let deadline = clock
+ .deadline_after(self.grace)
+ .map_err(ShutdownStartError::Deadline)?;
+ let runtime_deadline = tokio::time::Instant::now().checked_add(self.grace).ok_or(
+ ShutdownStartError::Deadline(MonotonicClockError::DeadlineOverflow),
+ )?;
+ self.progress = Some(ShutdownProgress {
+ deadline,
+ runtime_deadline,
+ phase_index: 0,
+ stage: ShutdownPhaseStage::Enter,
+ disposition: ShutdownDisposition::Completed,
+ });
+ }
tokio::pin!(force);
- let mut disposition = ShutdownDisposition::Completed;
- for phase in ORDERED_PHASES {
- if phase == ShutdownPhase::CancelIngress {
- supervisor.request_cancellation();
- }
-
- match wait_bounded(
- || handler.enter(phase),
- force.as_mut(),
- deadline_wait.as_mut(),
- )
- .await
- {
- BoundedWait::Completed(Ok(())) => {}
- BoundedWait::Completed(Err(error)) => {
- let unfinished = supervisor.unfinished_work();
- self.phase_failure = Some(ShutdownPhaseFailure { phase, error });
- supervisor.abort_and_drain().await;
- return Ok(self.complete(
- deadline,
- ShutdownDisposition::PhaseFailed { phase, unfinished },
- ));
- }
- BoundedWait::Forced => {
- let unfinished = supervisor.unfinished_work();
- supervisor.abort_and_drain().await;
- return Ok(self.complete(deadline, ShutdownDisposition::Forced { unfinished }));
- }
- BoundedWait::GraceExpired => {
- let unfinished = supervisor.unfinished_work();
- supervisor.abort_and_drain().await;
- return Ok(
- self.complete(deadline, ShutdownDisposition::GraceExpired { unfinished })
- );
- }
+ loop {
+ let progress = self
+ .progress
+ .expect("shutdown progress must be initialized");
+ if let ShutdownPhaseStage::Abort(disposition) = progress.stage {
+ supervisor.abort_and_drain().await;
+ return Ok(self.complete(progress.deadline, disposition));
}
-
- if phase == ShutdownPhase::DrainOperations {
- match wait_bounded(
- || supervisor.supervise(),
+ let Some(&phase) = ORDERED_PHASES.get(progress.phase_index) else {
+ return Ok(self.complete(progress.deadline, progress.disposition));
+ };
+ let deadline_wait = tokio::time::sleep_until(progress.runtime_deadline);
+ tokio::pin!(deadline_wait);
+
+ let outcome = match progress.stage {
+ ShutdownPhaseStage::Enter => wait_bounded(
+ || async {
+ supervisor.request_phase_cancellation(phase);
+ handler.enter(phase).await
+ },
force.as_mut(),
deadline_wait.as_mut(),
)
.await
- {
- BoundedWait::Completed(Ok(_exits)) => {}
- BoundedWait::Completed(Err(failure)) => {
- let kind = failure.kind();
- self.task_failure = Some(failure);
- disposition = ShutdownDisposition::TaskFailed { kind };
- }
- BoundedWait::Forced => {
+ .map(ShutdownPhaseOutcome::Entered),
+ ShutdownPhaseStage::Drain => wait_bounded(
+ || supervisor.supervise_phase(phase),
+ force.as_mut(),
+ deadline_wait.as_mut(),
+ )
+ .await
+ .map(ShutdownPhaseOutcome::Drained),
+ ShutdownPhaseStage::Abort(_) => unreachable!("abort stage handled before phase"),
+ };
+
+ match outcome {
+ BoundedWait::Completed(ShutdownPhaseOutcome::Entered(result)) => {
+ if let Err(error) = result {
let unfinished = supervisor.unfinished_work();
- supervisor.abort_and_drain().await;
- return Ok(
- self.complete(deadline, ShutdownDisposition::Forced { unfinished })
- );
+ if self.phase_failure.is_none() {
+ self.phase_failure = Some(ShutdownPhaseFailure { phase, error });
+ }
+ self.retain_first_disposition(ShutdownDisposition::PhaseFailed {
+ phase,
+ unfinished,
+ });
}
- BoundedWait::GraceExpired => {
- let unfinished = supervisor.unfinished_work();
- supervisor.abort_and_drain().await;
- return Ok(self
- .complete(deadline, ShutdownDisposition::GraceExpired { unfinished }));
+ self.progress.as_mut().expect("shutdown progress").stage =
+ ShutdownPhaseStage::Drain;
+ }
+ BoundedWait::Completed(ShutdownPhaseOutcome::Drained(result)) => match result {
+ Ok(_exits) => {
+ let progress = self.progress.as_mut().expect("shutdown progress");
+ progress.phase_index += 1;
+ progress.stage = ShutdownPhaseStage::Enter;
+ }
+ Err(failure) => {
+ let kind = failure.kind();
+ if self.task_failure.is_none() {
+ self.task_failure = Some(failure);
+ }
+ self.retain_first_disposition(ShutdownDisposition::TaskFailed { kind });
}
+ },
+ BoundedWait::Forced => {
+ let unfinished = supervisor.unfinished_work();
+ self.progress.as_mut().expect("shutdown progress").stage =
+ ShutdownPhaseStage::Abort(ShutdownDisposition::Forced { unfinished });
+ }
+ BoundedWait::GraceExpired => {
+ let unfinished = supervisor.unfinished_work();
+ self.progress.as_mut().expect("shutdown progress").stage =
+ ShutdownPhaseStage::Abort(ShutdownDisposition::GraceExpired { unfinished });
}
}
}
+ }
- Ok(self.complete(deadline, disposition))
+ fn retain_first_disposition(&mut self, disposition: ShutdownDisposition) {
+ let progress = self.progress.as_mut().expect("shutdown progress");
+ if progress.disposition == ShutdownDisposition::Completed {
+ progress.disposition = disposition;
+ }
}
fn complete(
@@ -166,16 +217,32 @@ impl GracefulShutdown {
disposition,
};
self.completed = Some(summary);
+ self.progress = None;
summary
}
}
+enum ShutdownPhaseOutcome {
+ Entered(Result<(), HostError>),
+ Drained(Result<Vec<super::SupervisedTaskExit>, SupervisionFailure>),
+}
+
enum BoundedWait<T> {
Completed(T),
Forced,
GraceExpired,
}
+impl<T> BoundedWait<T> {
+ fn map<U>(self, map: impl FnOnce(T) -> U) -> BoundedWait<U> {
+ match self {
+ Self::Completed(value) => BoundedWait::Completed(map(value)),
+ Self::Forced => BoundedWait::Forced,
+ Self::GraceExpired => BoundedWait::GraceExpired,
+ }
+ }
+}
+
async fn wait_bounded<T, MakeWork, Work, Force>(
make_work: MakeWork,
mut force: Pin<&mut Force>,
@@ -369,6 +436,24 @@ mod tests {
impl Error for SensitiveCause {}
+ struct BlockingHandler {
+ phases: Arc<Mutex<Vec<ShutdownPhase>>>,
+ block_at: ShutdownPhase,
+ entered: Arc<tokio::sync::Notify>,
+ }
+
+ impl ShutdownPhaseHandler for BlockingHandler {
+ fn enter(&mut self, phase: ShutdownPhase) -> ShutdownPhaseFuture<'_> {
+ self.phases.lock().unwrap().push(phase);
+ if phase == self.block_at {
+ self.entered.notify_one();
+ Box::pin(pending())
+ } else {
+ Box::pin(ready(Ok(())))
+ }
+ }
+ }
+
fn clock() -> FakeClock {
FakeClock {
now: MonotonicTime::from_duration_since_origin(Duration::from_secs(5)),
@@ -512,7 +597,7 @@ mod tests {
unfinished: UnfinishedWork::None,
}
);
- assert_eq!(*phases.lock().unwrap(), ORDERED_PHASES[..=3].to_vec());
+ assert_eq!(*phases.lock().unwrap(), ORDERED_PHASES);
let failure = shutdown.phase_failure().unwrap();
assert_eq!(failure.phase(), ShutdownPhase::PersistRecoverableWork);
assert_eq!(
@@ -553,6 +638,161 @@ mod tests {
assert!(failure.source().is_some());
}
+ #[tokio::test(start_paused = true)]
+ async fn cancelled_run_retains_its_absolute_deadline_and_completed_phase_progress() {
+ let mut supervisor = TaskSupervisor::new();
+ let phases = Arc::new(Mutex::new(Vec::new()));
+ let entered = Arc::new(tokio::sync::Notify::new());
+ let mut handler = BlockingHandler {
+ phases: Arc::clone(&phases),
+ block_at: ShutdownPhase::PersistRecoverableWork,
+ entered: Arc::clone(&entered),
+ };
+ let mut shutdown = GracefulShutdown::new(Duration::from_secs(10)).unwrap();
+ let shutdown_clock = clock();
+
+ {
+ let mut run =
+ Box::pin(shutdown.run(&shutdown_clock, &mut supervisor, &mut handler, pending()));
+ tokio::select! {
+ () = entered.notified() => {}
+ result = &mut run => panic!("blocked phase completed unexpectedly: {result:?}"),
+ }
+ }
+ assert_eq!(*phases.lock().unwrap(), ORDERED_PHASES[..=3].to_vec());
+
+ tokio::time::advance(Duration::from_secs(10)).await;
+ let summary = tokio::time::timeout(
+ Duration::from_millis(1),
+ shutdown.run(&shutdown_clock, &mut supervisor, &mut handler, pending()),
+ )
+ .await
+ .expect("retry must use the original elapsed runtime deadline")
+ .unwrap();
+ assert_eq!(
+ summary.disposition(),
+ ShutdownDisposition::GraceExpired {
+ unfinished: UnfinishedWork::None
+ }
+ );
+ assert_eq!(
+ summary.deadline().time().duration_since_origin(),
+ Duration::from_secs(15)
+ );
+ assert_eq!(
+ *phases.lock().unwrap(),
+ ORDERED_PHASES[..=3].to_vec(),
+ "an elapsed retry must not reconstruct the incomplete phase"
+ );
+ }
+
+ #[tokio::test]
+ async fn cancellation_during_phase_drain_resumes_without_reentering_completed_handler() {
+ let (release, wait_for_release) = tokio::sync::oneshot::channel();
+ let cleanup_started = Arc::new(tokio::sync::Notify::new());
+ let mut supervisor = TaskSupervisor::new();
+ let started = Arc::clone(&cleanup_started);
+ let metadata = TaskMetadata::new(
+ TaskName::new("mutation_gate").unwrap(),
+ TaskClassification::Critical,
+ Some(ShutdownPhase::RejectNewMutations),
+ )
+ .unwrap();
+ supervisor
+ .spawn(metadata, move |token| async move {
+ token.cancelled().await;
+ started.notify_one();
+ let _ = wait_for_release.await;
+ Ok(())
+ })
+ .unwrap();
+ let (mut handler, phases) = handler(None);
+ let mut shutdown = GracefulShutdown::new(Duration::from_secs(30)).unwrap();
+ let shutdown_clock = clock();
+
+ {
+ let mut run =
+ Box::pin(shutdown.run(&shutdown_clock, &mut supervisor, &mut handler, pending()));
+ tokio::select! {
+ () = cleanup_started.notified() => {}
+ result = &mut run => panic!("phase drain completed unexpectedly: {result:?}"),
+ }
+ }
+ assert_eq!(
+ *phases.lock().unwrap(),
+ vec![ShutdownPhase::RejectNewMutations]
+ );
+
+ release.send(()).unwrap();
+ let summary = shutdown
+ .run(&shutdown_clock, &mut supervisor, &mut handler, pending())
+ .await
+ .unwrap();
+ assert_eq!(summary.disposition(), ShutdownDisposition::Completed);
+ assert_eq!(*phases.lock().unwrap(), ORDERED_PHASES);
+ assert!(supervisor.is_empty());
+ }
+
+ #[tokio::test(start_paused = true)]
+ async fn cancellation_while_fatal_peers_drain_retains_the_first_task_failure() {
+ let (release, wait_for_release) = tokio::sync::oneshot::channel();
+ let peer_waiting = Arc::new(tokio::sync::Notify::new());
+ let mut supervisor = TaskSupervisor::new();
+ let phase_metadata = |name| {
+ TaskMetadata::new(
+ TaskName::new(name).unwrap(),
+ TaskClassification::Critical,
+ Some(ShutdownPhase::RejectNewMutations),
+ )
+ .unwrap()
+ };
+ supervisor
+ .spawn(phase_metadata("failing_gate"), |_| async {
+ Err(HostError::new(HostErrorKind::TaskFailure))
+ })
+ .unwrap();
+ let waiting = Arc::clone(&peer_waiting);
+ supervisor
+ .spawn(phase_metadata("draining_peer"), move |token| async move {
+ token.cancelled().await;
+ tokio::time::sleep(Duration::from_secs(1)).await;
+ waiting.notify_one();
+ let _ = wait_for_release.await;
+ Ok(())
+ })
+ .unwrap();
+ let (mut handler, phases) = handler(None);
+ let mut shutdown = GracefulShutdown::new(Duration::from_secs(30)).unwrap();
+ let shutdown_clock = clock();
+
+ {
+ let mut run =
+ Box::pin(shutdown.run(&shutdown_clock, &mut supervisor, &mut handler, pending()));
+ tokio::select! {
+ () = peer_waiting.notified() => {}
+ result = &mut run => panic!("fatal peer drain completed unexpectedly: {result:?}"),
+ }
+ }
+ assert_eq!(
+ shutdown.task_failure().map(SupervisionFailure::kind),
+ Some(SupervisionFailureKind::TaskReturnedError)
+ );
+
+ release.send(()).unwrap();
+ let summary = shutdown
+ .run(&shutdown_clock, &mut supervisor, &mut handler, pending())
+ .await
+ .unwrap();
+ assert_eq!(
+ summary.disposition(),
+ ShutdownDisposition::TaskFailed {
+ kind: SupervisionFailureKind::TaskReturnedError
+ }
+ );
+ assert_eq!(*phases.lock().unwrap(), ORDERED_PHASES);
+ assert!(supervisor.is_empty());
+ }
+
#[test]
fn zero_grace_fails_closed() {
let error = GracefulShutdown::new(Duration::ZERO).err().unwrap();
diff --git a/crates/service_host/src/lifecycle/supervisor.rs b/crates/service_host/src/lifecycle/supervisor.rs
@@ -7,14 +7,19 @@ use tokio::task::{Id, JoinError, JoinSet};
use crate::HostError;
-use super::{CancellationToken, TaskClassification, TaskMetadata, UnfinishedWork};
+use super::{CancellationToken, ShutdownPhase, TaskClassification, TaskMetadata, UnfinishedWork};
/// Owns every spawned service task until its join result is observed.
#[must_use = "a task supervisor must be run or drained so authoritative tasks are joined"]
pub struct TaskSupervisor {
cancellation: CancellationToken,
tasks: JoinSet<TaskCompletion>,
- metadata: HashMap<Id, TaskMetadata>,
+ controls: HashMap<Id, TaskControl>,
+}
+
+struct TaskControl {
+ metadata: TaskMetadata,
+ cancellation: CancellationToken,
}
impl TaskSupervisor {
@@ -22,7 +27,7 @@ impl TaskSupervisor {
Self {
cancellation: CancellationToken::new(),
tasks: JoinSet::new(),
- metadata: HashMap::new(),
+ controls: HashMap::new(),
}
}
@@ -58,9 +63,9 @@ impl TaskSupervisor {
Fut: Future<Output = Result<(), HostError>> + Send + 'static,
{
if self
- .metadata
+ .controls
.values()
- .any(|active| active.name() == metadata.name())
+ .any(|active| active.metadata.name() == metadata.name())
{
return Err(TaskRegistrationError::DuplicateName);
}
@@ -69,6 +74,7 @@ impl TaskSupervisor {
let task_metadata = metadata.clone();
let child = self.cancellation.child_token();
let completion_observer = child.clone();
+ let phase_cancellation = child.clone();
let abort = self.tasks.spawn_on(
async move {
let result = task(child).await;
@@ -80,10 +86,42 @@ impl TaskSupervisor {
},
&runtime,
);
- self.metadata.insert(abort.id(), metadata);
+ self.controls.insert(
+ abort.id(),
+ TaskControl {
+ metadata,
+ cancellation: phase_cancellation,
+ },
+ );
Ok(())
}
+ pub(crate) fn request_phase_cancellation(&self, phase: ShutdownPhase) {
+ for control in self.controls.values() {
+ if control.metadata.shutdown_phase() == Some(phase) {
+ control.cancellation.cancel();
+ }
+ }
+ }
+
+ pub(crate) async fn supervise_phase(
+ &mut self,
+ phase: ShutdownPhase,
+ ) -> Result<Vec<SupervisedTaskExit>, SupervisionFailure> {
+ let mut exits = Vec::new();
+ while self.has_phase_work(phase) {
+ let outcome = self
+ .join_next()
+ .await
+ .expect("phase work must retain a join-owned task");
+ match outcome {
+ Ok(exit) => exits.push(exit),
+ Err(failure) => return Err(failure),
+ }
+ }
+ Ok(exits)
+ }
+
/// Observes all task exits, cancels peers on the first fatal outcome, and drains every join.
pub async fn supervise(&mut self) -> Result<Vec<SupervisedTaskExit>, SupervisionFailure> {
let mut exits = Vec::with_capacity(self.tasks.len());
@@ -125,11 +163,14 @@ impl TaskSupervisor {
) -> Result<SupervisedTaskExit, SupervisionFailure> {
match joined {
Ok((id, completion)) => {
- self.metadata.remove(&id);
+ self.controls.remove(&id);
classify_completion(completion)
}
Err(error) => {
- let metadata = self.metadata.remove(&error.id());
+ let metadata = self
+ .controls
+ .remove(&error.id())
+ .map(|control| control.metadata);
let cancelled_during_shutdown =
error.is_cancelled() && self.cancellation.is_cancelled();
if cancelled_during_shutdown && let Some(metadata) = metadata {
@@ -157,12 +198,12 @@ impl TaskSupervisor {
}
pub(crate) fn unfinished_work(&self) -> UnfinishedWork {
- if self.metadata.is_empty() {
+ if self.controls.is_empty() {
UnfinishedWork::None
} else if self
- .metadata
+ .controls
.values()
- .any(|metadata| metadata.classification().failure_is_fatal())
+ .any(|control| control.metadata.classification().failure_is_fatal())
{
UnfinishedWork::FatalAuthoritative
} else {
@@ -175,6 +216,14 @@ impl TaskSupervisor {
self.tasks.abort_all();
while self.join_next().await.is_some() {}
}
+
+ fn has_phase_work(&self, phase: ShutdownPhase) -> bool {
+ self.controls.values().any(|control| {
+ control.metadata.shutdown_phase() == Some(phase)
+ || (phase == ShutdownPhase::DrainOperations
+ && control.metadata.classification() == TaskClassification::OneShot)
+ })
+ }
}
impl Default for TaskSupervisor {
@@ -380,6 +429,14 @@ mod tests {
TaskMetadata::new(TaskName::new(name).unwrap(), classification, shutdown_phase).unwrap()
}
+ fn metadata_at(
+ name: &str,
+ classification: TaskClassification,
+ shutdown_phase: Option<ShutdownPhase>,
+ ) -> TaskMetadata {
+ TaskMetadata::new(TaskName::new(name).unwrap(), classification, shutdown_phase).unwrap()
+ }
+
#[test]
fn registration_without_a_runtime_fails_before_spawning() {
let mut supervisor = TaskSupervisor::new();
@@ -425,6 +482,88 @@ mod tests {
}
#[tokio::test]
+ async fn phase_cancellation_stops_and_joins_only_the_assigned_tasks() {
+ let ingress_stopped = Arc::new(AtomicUsize::new(0));
+ let network_stopped = Arc::new(AtomicUsize::new(0));
+ let mut supervisor = TaskSupervisor::new();
+ let ingress = Arc::clone(&ingress_stopped);
+ supervisor
+ .spawn(
+ metadata_at(
+ "ingress_worker",
+ TaskClassification::Critical,
+ Some(ShutdownPhase::CancelIngress),
+ ),
+ move |token| async move {
+ token.cancelled().await;
+ ingress.fetch_add(1, Ordering::SeqCst);
+ Ok(())
+ },
+ )
+ .unwrap();
+ let network = Arc::clone(&network_stopped);
+ supervisor
+ .spawn(
+ metadata_at(
+ "network_worker",
+ TaskClassification::Critical,
+ Some(ShutdownPhase::CloseNetwork),
+ ),
+ move |token| async move {
+ token.cancelled().await;
+ network.fetch_add(1, Ordering::SeqCst);
+ Ok(())
+ },
+ )
+ .unwrap();
+
+ supervisor.request_phase_cancellation(ShutdownPhase::CancelIngress);
+ let ingress_exits = supervisor
+ .supervise_phase(ShutdownPhase::CancelIngress)
+ .await
+ .unwrap();
+ assert_eq!(ingress_exits.len(), 1);
+ assert_eq!(ingress_stopped.load(Ordering::SeqCst), 1);
+ assert_eq!(network_stopped.load(Ordering::SeqCst), 0);
+ assert_eq!(supervisor.task_count(), 1);
+
+ supervisor.request_phase_cancellation(ShutdownPhase::CloseNetwork);
+ let network_exits = supervisor
+ .supervise_phase(ShutdownPhase::CloseNetwork)
+ .await
+ .unwrap();
+ assert_eq!(network_exits.len(), 1);
+ assert_eq!(network_stopped.load(Ordering::SeqCst), 1);
+ assert!(supervisor.is_empty());
+ }
+
+ #[tokio::test]
+ async fn drain_operations_joins_one_shot_work_without_cancelling_it() {
+ let mut supervisor = TaskSupervisor::new();
+ supervisor
+ .spawn(
+ metadata_at("bounded_operation", TaskClassification::OneShot, None),
+ |token| async move {
+ assert!(!token.is_cancelled());
+ Ok(())
+ },
+ )
+ .unwrap();
+
+ supervisor.request_phase_cancellation(ShutdownPhase::DrainOperations);
+ let exits = supervisor
+ .supervise_phase(ShutdownPhase::DrainOperations)
+ .await
+ .unwrap();
+ assert_eq!(exits.len(), 1);
+ assert_eq!(
+ exits[0].status(),
+ SupervisedTaskExitStatus::ExpectedCompletion
+ );
+ assert!(supervisor.is_empty());
+ }
+
+ #[tokio::test]
async fn critical_panic_and_early_success_are_fatal() {
for (name, expected, panic_task) in [
("panic_worker", SupervisionFailureKind::TaskPanicked, true),
diff --git a/crates/service_host/tests/package_boundary.rs b/crates/service_host/tests/package_boundary.rs
@@ -178,6 +178,20 @@ fn documentation_and_reviewed_public_api_are_complete_and_dependency_safe() {
"Ordinary `Debug` is redacted and no",
"`Display` implementation exposes the retained value",
"explicit borrowed `as_str` accessor for serialization",
+ "Shutdown task cancellation is phase aware",
+ "Entering a phase cancels only tasks",
+ "assigned to that phase and does not advance until their joins are observed",
+ "bounded one-shot work drains without",
+ "cancellation during `DrainOperations`",
+ "a fatal task outcome still",
+ "cancels and joins the complete graph",
+ "retains one absolute",
+ "deadline plus the completed handler/drain boundary across cancellation and",
+ "no phase or cleanup attempt receives a fresh grace period",
+ "incomplete handler may be entered again and must be idempotent and",
+ "cancellation safe",
+ "first phase or task failure is retained while later",
+ "close phases continue as long as the original deadline remains",
] {
assert!(README.contains(required), "README is missing `{required}`");
}