Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
76 changes: 72 additions & 4 deletions bubus/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -295,6 +295,33 @@ async def wait_for_handlers_to_complete_then_return_event():
max_iterations = 1000 # Prevent infinite loops
iterations = 0

# Snapshot the buses once for the whole drain: the ancestry
# walk and the per-candidate descendant walk both look up
# events across this set, and rebuilding it per hop would make
# the drain O(candidates × bus_count × chain_depth) per
# iteration (review feedback on #5509).
buses_snapshot = list(EventBus.all_instances)

@cubic-dev-ai cubic-dev-ai Bot Aug 29, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: The lookup snapshot is frozen once, but the drain loop below still re-enumerates list(EventBus.all_instances) on every iteration (for bus in list(EventBus.all_instances)). A bus created mid-drain is iterated by the drain but is invisible to _find_event_by_id, which only scans buses_snapshot. Because the snapshot was taken before those buses existed, any of their queued events that belong to this waiting chain get parent_candidate = None during the descendant walk and are never drained, while unrelated events on them are scanned in the loop — so the drain's behavior silently depends on when the snapshot happened and contradicts the stated "snapshot once per drain" intent.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At bubus/models.py, line 303:

<comment>The lookup snapshot is frozen once, but the drain loop below still re-enumerates `list(EventBus.all_instances)` on every iteration (`for bus in list(EventBus.all_instances)`). A bus created mid-drain is iterated by the drain but is invisible to `_find_event_by_id`, which only scans `buses_snapshot`. Because the snapshot was taken before those buses existed, any of their queued events that belong to this waiting chain get `parent_candidate = None` during the descendant walk and are never drained, while unrelated events on them are scanned in the loop — so the drain's behavior silently depends on when the snapshot happened and contradicts the stated "snapshot once per drain" intent.</comment>

<file context>
@@ -295,6 +295,19 @@ async def wait_for_handlers_to_complete_then_return_event():
+                # events across this set, and rebuilding it per hop would make
+                # the drain O(candidates × bus_count × chain_depth) per
+                # iteration (review feedback on #5509).
+                buses_snapshot = list(EventBus.all_instances)
+
+                def _find_event_by_id(event_id: str) -> BaseEvent[Any] | None:
</file context>
Fix with cubic


def _find_event_by_id(event_id: str) -> BaseEvent[Any] | None:
for bus_ in buses_snapshot:
if bus_ and event_id in bus_.event_history:
return bus_.event_history[event_id]
return None

# Compute this event's ancestry chain once, so the drain loop
# below only processes events that belong to this waiting
# chain. Draining unrelated events from other buses (cross-loop
# contamination) consumes them with the wrong handler table and
# can stall or drop the owning bus's work — see issue #5509.
ancestor_ids: set[str] = {self.event_id}
cursor: BaseEvent[Any] | None = self
while cursor is not None and cursor.event_parent_id:
parent_event = _find_event_by_id(cursor.event_parent_id)
if parent_event is None or parent_event.event_id in ancestor_ids:
break
ancestor_ids.add(parent_event.event_id)
cursor = parent_event

try:
while not self.event_completed_signal.is_set() and iterations < max_iterations:
iterations += 1
Expand All @@ -309,10 +336,51 @@ async def wait_for_handlers_to_complete_then_return_event():
# Process one event from this bus if available
try:
if bus.event_queue.qsize() > 0:
event = bus.event_queue.get_nowait()
await bus.process_event(event)
bus.event_queue.task_done()
processed_any = True
# Only process events that belong to this
# waiting chain. Locate the first such event
# in the queue and process just that one;
# unrelated events (from other buses or
# independent tasks) keep their FIFO position
# untouched and are left for their own bus's
# run loop — draining them here is cross-loop
# contamination that steals the owning bus's
# scheduling (issue #5509).
# The underlying deque is indexed directly
# (same approach as the memory-usage check);
# task_done() pairs with the put() that
# enqueued the event.
for idx, candidate in enumerate(bus.event_queue._queue):
# The candidate belongs to this waiting
# chain iff it is this event or a
# descendant of it (its ancestor chain
# contains self.event_id). Matching on
# parent-id alone would also drain
# siblings that merely share a parent
# with the awaited event, stealing
# their normal scheduling.
if candidate.event_id in ancestor_ids:
matches = True
else:
matches = False
cursor_candidate = candidate
seen_ids: set[str] = set()
while cursor_candidate.event_parent_id and cursor_candidate.event_parent_id not in seen_ids:
seen_ids.add(cursor_candidate.event_parent_id)
if cursor_candidate.event_parent_id == self.event_id:
matches = True
break
parent_candidate = _find_event_by_id(cursor_candidate.event_parent_id)
if parent_candidate is None:
break
cursor_candidate = parent_candidate
if matches:
del bus.event_queue._queue[idx]
try:
await bus.process_event(candidate)
finally:
bus.event_queue.task_done()
processed_any = True
break
# Check if the event we're waiting for is now complete
if self.event_completed_signal.is_set():
break
Expand Down
144 changes: 144 additions & 0 deletions tests/test_issue5509_cross_loop.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,144 @@
"""
Reproduction for issue #5509: cross-loop contamination in the EventBus drain loop.

Scenario: two EventBus instances running in parallel (like two concurrent
agent sessions in browser-use). A handler chain on bus A (parent -> child ->
grandchild) enters the inner `__await__` drain loop, which iterates over
`EventBus.all_instances` and processes queued events from ALL buses.

Before the fix, a bus B event queued during that window is processed by the
drain loop — i.e. INSIDE bus A's handler context (while bus A holds the global
lock and bus B's own run loop is starved). This cross-loop execution steals
bus B's scheduling, and with many parallel sessions it is what drives the
EventBus capacity errors reported in the issue.

After the fix, the drain loop only processes events belonging to its own
waiting chain; bus B's independent event stays queued and is handled by bus
B's own run loop once the lock is released.
"""

import asyncio

import pytest

from bubus import BaseEvent, EventBus


class ParentAEvent(BaseEvent[str]):
message: str


class ChildAEvent(BaseEvent[str]):
data: str


class GrandchildAEvent(BaseEvent[str]):
value: int


class BusBEvent(BaseEvent[str]):
"""Independent event that should only be handled on bus B."""

payload: str


class StartEvent(BaseEvent[str]):
"""Handler-less event used only to auto-start a bus's run loop."""

data: str


@pytest.fixture
async def buses():
"""Two isolated buses simulating parallel sessions."""
# Create bus_b first so it appears before bus_a in the
# `EventBus.all_instances` iteration order (WeakSet keeps insertion
# order); the drain loop scans buses in that order, so bus B's queued
# event is hit before the drain completes its own chain and breaks.
bus_b = EventBus(name="bus_b")
bus_a = EventBus(name="bus_a")
# First dispatch auto-starts each bus's run loop (and creates event_queue).
bus_a.dispatch(StartEvent(data="__start__"))
bus_b.dispatch(StartEvent(data="__start__"))
await asyncio.gather(
bus_a.wait_until_idle(timeout=5),
bus_b.wait_until_idle(timeout=5),
)
yield bus_a, bus_b
await bus_a.stop(clear=True)
await bus_b.stop(clear=True)


@pytest.mark.asyncio
async def test_bus_b_event_not_processed_inside_bus_a_drain(buses):
bus_a, bus_b = buses

ready = asyncio.Event()
proceed = asyncio.Event()
order: list[str] = []

async def grandchild_handler(event: GrandchildAEvent) -> int:
await asyncio.sleep(0.5)
return event.value

async def child_handler(event: ChildAEvent) -> str:
# Tell the test the chain is armed; wait until bus B's independent
# event is queued BEFORE forwarding the grandchild, so bus B's queue
# holds [independent event, grandchild] when the inner drain loop
# scans it — the drain must skip the unrelated event and still
# process the grandchild.
ready.set()
await proceed.wait()
order.append("a_drain_enter")
grandchild = GrandchildAEvent(value=1)
bus_b.dispatch(grandchild)
await grandchild
order.append("a_drain_exit")
return "child_handled"

async def handler_a(event: ParentAEvent) -> str:
child = ChildAEvent(data=f"child_of_{event.message}")
bus_a.dispatch(child)
await child
return "a_handled"

async def handler_b(event: BusBEvent) -> str:
order.append("b_handled")
return "b_handled"

bus_a.on("GrandchildAEvent", grandchild_handler)
bus_a.on("ChildAEvent", child_handler)
bus_a.on("ParentAEvent", handler_a)
bus_b.on("BusBEvent", handler_b)

# Start bus A's event chain; wait until its handler holds the global lock
# (child_handler is parked on `proceed` inside bus A's run loop).
bus_a.dispatch(ParentAEvent(message="trigger"))
await asyncio.wait_for(ready.wait(), timeout=5)

# Dispatch an event on bus B: its run loop `get()`s it (outside the lock)
# and then blocks acquiring the global lock held by bus A — so bus B's
# consumer is parked and will not `get()` again until the lock is free.
bus_b.dispatch(BusBEvent(payload="first"))
await asyncio.sleep(0.05)

# Queue a SECOND independent event on bus B (directly into its queue):
# with bus B's run loop parked on the lock, this event is only reachable
# by bus A's drain loop — the cross-loop window (issue #5509).
bus_b.event_queue.put_nowait(BusBEvent(payload="second"))
proceed.set()

await asyncio.gather(
bus_a.wait_until_idle(timeout=5),
bus_b.wait_until_idle(timeout=5),
)

# Bus B's events must NOT be processed while bus A's drain loop is running
# (inside bus A's handler context, holding the global lock): that is the
# cross-loop contamination. Both events must be handled afterwards by bus
# B's own run loop, once the lock is released.
assert order[:2] == ["a_drain_enter", "a_drain_exit"], f"drain markers missing: {order!r}"
assert order[2:] == ["b_handled", "b_handled"], (
f"cross-loop contamination: bus B events were processed inside bus A's "
f"drain loop (order was {order!r}, expected both bus B events after the drain)"
)