[object Object]

← 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

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 →