diff --git a/src/exo/main.py b/src/exo/main.py index 8a56ab19..85e27b29 100644 --- a/src/exo/main.py +++ b/src/exo/main.py @@ -252,7 +252,7 @@ def main(): target = min(max(soft, 65535), hard) resource.setrlimit(resource.RLIMIT_NOFILE, (target, hard)) - mp.set_start_method("spawn") + mp.set_start_method("spawn", force=True) # TODO: Refactor the current verbosity system logger_setup(EXO_LOG, args.verbosity) logger.info("Starting EXO") diff --git a/src/exo/worker/runner/runner_supervisor.py b/src/exo/worker/runner/runner_supervisor.py index 4aa4e6f7..565ca315 100644 --- a/src/exo/worker/runner/runner_supervisor.py +++ b/src/exo/worker/runner/runner_supervisor.py @@ -106,13 +106,18 @@ class RunnerSupervisor: def shutdown(self): logger.info("Runner supervisor shutting down") self._tg.cancel_tasks() - self._ev_recv.close() - self._task_sender.close() if not self._cancel_watch_runner.cancel_called: self._cancel_watch_runner.cancel() + with contextlib.suppress(ClosedResourceError): + self._ev_recv.close() + with contextlib.suppress(ClosedResourceError): + self._task_sender.close() + with contextlib.suppress(ClosedResourceError): + self._event_sender.close() with contextlib.suppress(ClosedResourceError): self._cancel_sender.send(TaskId("CANCEL_CURRENT_TASK")) - self._cancel_sender.close() + with contextlib.suppress(ClosedResourceError): + self._cancel_sender.close() self.runner_process.join(5) if not self.runner_process.is_alive(): logger.info("Runner process succesfully terminated")