← back to Exo
fix skipping logic in worker plan (#1342)
cd946742f7b4eb9377cef64e15cc6ecd825c5cb7 · 2026-01-30 14:31:40 +0000 · Evan Quiney
the worker plan function had some skipping logic missing, leading to
double-submitting tasks.
Files touched
M src/exo/worker/runner/runner.pyM src/exo/worker/runner/runner_supervisor.py
Diff
commit cd946742f7b4eb9377cef64e15cc6ecd825c5cb7
Author: Evan Quiney <evanev7@gmail.com>
Date: Fri Jan 30 14:31:40 2026 +0000
fix skipping logic in worker plan (#1342)
the worker plan function had some skipping logic missing, leading to
double-submitting tasks.
---
src/exo/worker/runner/runner.py | 5 +++++
src/exo/worker/runner/runner_supervisor.py | 11 ++++++++---
2 files changed, 13 insertions(+), 3 deletions(-)
diff --git a/src/exo/worker/runner/runner.py b/src/exo/worker/runner/runner.py
index bf272aaa..61205f38 100644
--- a/src/exo/worker/runner/runner.py
+++ b/src/exo/worker/runner/runner.py
@@ -37,6 +37,7 @@ from exo.shared.types.tasks import (
Shutdown,
StartWarmup,
Task,
+ TaskId,
TaskStatus,
)
from exo.shared.types.worker.instances import BoundInstance
@@ -111,8 +112,12 @@ def main(
event_sender.send(
RunnerStatusUpdated(runner_id=runner_id, runner_status=current_status)
)
+ seen = set[TaskId]()
with task_receiver as tasks:
for task in tasks:
+ if task.task_id in seen:
+ logger.warning("repeat task - potential error")
+ seen.add(task.task_id)
event_sender.send(
TaskStatusUpdated(task_id=task.task_id, task_status=TaskStatus.Running)
)
diff --git a/src/exo/worker/runner/runner_supervisor.py b/src/exo/worker/runner/runner_supervisor.py
index fc17cddc..a1951d80 100644
--- a/src/exo/worker/runner/runner_supervisor.py
+++ b/src/exo/worker/runner/runner_supervisor.py
@@ -127,20 +127,25 @@ class RunnerSupervisor:
self._tg.cancel_scope.cancel()
async def start_task(self, task: Task):
+ if task.task_id in self.pending:
+ logger.warning(
+ f"Skipping invalid task {task} as it has already been submitted"
+ )
+ return
if task.task_id in self.completed:
- logger.info(
+ logger.warning(
f"Skipping invalid task {task} as it has already been completed"
)
+ return
logger.info(f"Starting task {task}")
event = anyio.Event()
self.pending[task.task_id] = event
try:
- self._task_sender.send(task)
+ await self._task_sender.send_async(task)
except ClosedResourceError:
logger.warning(f"Task {task} dropped, runner closed communication.")
return
await event.wait()
- logger.info(f"Finished task {task}")
async def _forward_events(self):
with self._ev_recv as events:
← a5bc38ad Check all nodes to evict (#1341)
·
back to Exo
·
nix: add macmon to PATH in wrapper scripts on Darwin b2579c78 →