[object Object]

← back to Exo

fix a race condition

f25689d9c240be0b7f691546278830c1809e7095 · 2025-10-15 10:49:53 +0100 · Evan Quiney

Files touched

Diff

commit f25689d9c240be0b7f691546278830c1809e7095
Author: Evan Quiney <evanev7@gmail.com>
Date:   Wed Oct 15 10:49:53 2025 +0100

    fix a race condition
---
 src/exo/master/main.py        | 2 ++
 src/exo/utils/event_buffer.py | 3 +++
 src/exo/worker/main.py        | 7 +++++--
 3 files changed, 10 insertions(+), 2 deletions(-)

diff --git a/src/exo/master/main.py b/src/exo/master/main.py
index ce3643c2..b60b263a 100644
--- a/src/exo/master/main.py
+++ b/src/exo/master/main.py
@@ -193,6 +193,7 @@ class Master:
                     logger.debug(f"Master indexing event: {str(event)[:100]}")
                     indexed = IndexedEvent(event=event, idx=len(self._event_log))
                     self.state = apply(self.state, indexed)
+
                     # TODO: SQL
                     self._event_log.append(event)
                     await self._send_event(indexed)
@@ -225,6 +226,7 @@ class Master:
                 )
                 local_index += 1
 
+    # This function is re-entrant, take care!
     async def _send_event(self, event: IndexedEvent):
         # Convenience method since this line is ugly
         await self.global_event_sender.send(
diff --git a/src/exo/utils/event_buffer.py b/src/exo/utils/event_buffer.py
index eb1b4cf0..8fcf5fa2 100644
--- a/src/exo/utils/event_buffer.py
+++ b/src/exo/utils/event_buffer.py
@@ -19,6 +19,9 @@ class OrderedBuffer[T]:
         if idx < self.next_idx_to_release:
             return
         if idx in self.store:
+            assert self.store[idx] == t, (
+                "Received different messages with identical indices, probable race condition"
+            )
             return
         self.store[idx] = t
 
diff --git a/src/exo/worker/main.py b/src/exo/worker/main.py
index 2dc09559..0c3699fd 100644
--- a/src/exo/worker/main.py
+++ b/src/exo/worker/main.py
@@ -630,18 +630,21 @@ class Worker:
             async for event in self.fail_runner(e, runner_id):
                 yield event
 
+
+
+    # This function is re-entrant, take care!
     async def event_publisher(self, event: Event) -> None:
         fe = ForwarderEvent(
             origin_idx=self.local_event_index,
             origin=self.node_id,
             event=event,
         )
-        await self.local_event_sender.send(fe)
-        self.out_for_delivery[event.event_id] = fe
         logger.debug(
             f"Worker published event {self.local_event_index}: {str(event)[:100]}"
         )
         self.local_event_index += 1
+        await self.local_event_sender.send(fe)
+        self.out_for_delivery[event.event_id] = fe
 
 
 def event_relevant_to_worker(event: Event, worker: Worker):

← 1c6b5ce9 new tagged union  ·  back to Exo  ·  leaf placement 363c98a8 →