From 5d7a005a13e7616964bafac2dcdca1016cbfc876 Mon Sep 17 00:00:00 2001 From: Ryuichi Leo Takashige Date: Thu, 12 Mar 2026 15:33:32 +0000 Subject: [PATCH] Destroy process group on Keyboard Interrupt --- .../worker/runner/vllm_inference/runner.py | 22 +++++++++++-------- 1 file changed, 13 insertions(+), 9 deletions(-) diff --git a/src/exo/worker/runner/vllm_inference/runner.py b/src/exo/worker/runner/vllm_inference/runner.py index e443bff6..69660981 100644 --- a/src/exo/worker/runner/vllm_inference/runner.py +++ b/src/exo/worker/runner/vllm_inference/runner.py @@ -297,15 +297,19 @@ class Runner: self.event_sender.send(TaskAcknowledged(task_id=task.task_id)) def main(self): - with self.task_receiver: - for task in self.task_receiver: - if task.task_id in self.seen: - logger.warning("repeat task - potential error") - continue - self.seen.add(task.task_id) - self.handle_first_task(task) - if isinstance(self.current_status, RunnerShutdown): - break + try: + with self.task_receiver: + for task in self.task_receiver: + if task.task_id in self.seen: + logger.warning("repeat task - potential error") + continue + self.seen.add(task.task_id) + self.handle_first_task(task) + if isinstance(self.current_status, RunnerShutdown): + break + finally: + if torch.distributed.is_initialized(): + torch.distributed.destroy_process_group() def handle_first_task(self, task: Task): self.send_task_status(task.task_id, TaskStatus.Running)