From c136a167e19a7ae072980ba03196c62f8c7a0e94 Mon Sep 17 00:00:00 2001 From: Abhisek Behera <123497213+abhisek343@users.noreply.github.com> Date: Mon, 28 Sep 2026 04:08:01 +0530 Subject: [PATCH 1/2] Skip graceful worker retirement during cluster teardown --- distributed/deploy/spec.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/distributed/deploy/spec.py b/distributed/deploy/spec.py index 49d51490d75..82ba4e66072 100644 --- a/distributed/deploy/spec.py +++ b/distributed/deploy/spec.py @@ -355,7 +355,10 @@ async def _correct_state_internal(self) -> None: to_close = set(self.workers) - set(self.worker_spec) if to_close: - if self.scheduler.status == Status.running: + if ( + self.scheduler.status == Status.running + and self.status != Status.closing + ): await self.scheduler_comm.retire_workers(workers=list(to_close)) tasks = [ asyncio.create_task(self.workers[w].close()) From c60b590316ca4347944ee4fa0b76f674cfd41945 Mon Sep 17 00:00:00 2001 From: Abhisek Behera <123497213+abhisek343@users.noreply.github.com> Date: Mon, 28 Sep 2026 04:08:04 +0530 Subject: [PATCH 2/2] Add regression test for cluster teardown retirement --- distributed/deploy/tests/test_spec_cluster.py | 36 +++++++++++++++++++ 1 file changed, 36 insertions(+) diff --git a/distributed/deploy/tests/test_spec_cluster.py b/distributed/deploy/tests/test_spec_cluster.py index ab1ec27e31e..7aa3d987fbd 100644 --- a/distributed/deploy/tests/test_spec_cluster.py +++ b/distributed/deploy/tests/test_spec_cluster.py @@ -511,6 +511,42 @@ async def test_bad_close(): assert not record +@gen_test() +async def test_correct_state_skips_retirement_while_closing(): + retired = False + + class DummyScheduler: + status = Status.running + + class DummySchedulerComm: + async def retire_workers(self, workers): + nonlocal retired + retired = True + + class DummyWorker: + def __init__(self): + self.closed = False + + async def close(self): + self.closed = True + + cluster = object.__new__(SpecCluster) + cluster._lock = asyncio.Lock() + cluster._correct_state_waiting = None + cluster.status = Status.closing + cluster.scheduler = DummyScheduler() + cluster.scheduler_comm = DummySchedulerComm() + worker = DummyWorker() + cluster.workers = {"worker": worker} + cluster.worker_spec = {} + + await cluster._correct_state_internal() + + assert not retired + assert worker.closed + assert not cluster.workers + + @gen_test() async def test_shutdown_scheduler_disabled(): async with SpecCluster(