Skip to content

fix(messaging): complete request futures on owning loop - #885

Merged
thomwebb merged 3 commits into
mpfaffenberger:mainfrom
kvandre12-commits:fix/issue-438-threadsafe-futures
Sep 22, 2026
Merged

thomwebb merged 3 commits into
mpfaffenberger:mainfrom
kvandre12-commits:fix/issue-438-threadsafe-futures

Conversation

@kvandre12-commits

Copy link
Copy Markdown
Contributor

Complements #859 by fixing the remaining request/response correctness issue from #438's first point.

MessageBus previously stored only each asyncio.Future. A UI response arriving from another thread could therefore call Future.set_result() directly on a Future owned by another event-loop thread. That is not thread-safe and can leave the awaiting task suspended.

What changed

  • Store each pending request as (loop, future).
  • Atomically consume the pending correlation under self._lock before scheduling completion.
  • Complete the Future on its owning loop using loop.call_soon_threadsafe().
  • Define duplicate/racing response behavior: first correlated response wins.
  • If the owning loop is closed, safely decline direct Future mutation.
  • Check future.done() on the owning loop before setting the result.

Regression coverage

Tests cover:

  • response arriving from a foreign thread
  • pending correlation consumed before scheduling
  • two racing responses scheduling exactly one completion
  • normal event-loop completion
  • closed-loop cleanup without unsafe Future mutation
  • cancellation between scheduling and callback execution
  • already-completed Futures
  • unknown prompt IDs

Focused MessageBus suite:

44 passed

Additional checks:

  • ruff format --check — passed
  • ruff E/F/I checks — passed
  • git diff --check — passed

The queue backoff introduced by #859 is intentionally unchanged.

The separate MessageBus locking/refactor work in #875 is also left untouched so the remaining #438 concerns stay narrowly separated.

Thanks to @StarsExpress for coordinating the remaining #438 work and reviewing this approach.

@StarsExpress

Copy link
Copy Markdown
Contributor

Awesome work @kvandre12-commits 🚀🚀
This one is really heavily enhancing thread safety behind agents 💪💪

@thomwebb thomwebb left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Dug into this one carefully since it's a concurrency fix — those are exactly the kind of "looks right, is subtly wrong" changes worth slowing down for. This one holds up.

What I checked

  • Full diff on bus.py + test_bus.py
  • Ran the focused suite: 44 passed on tests/messaging/test_bus.py, matching the PR description
  • ruff format --check clean; no new ruff check findings in the actual diff hunks
  • Traced every _pending_requests call site to confirm the (loop, future) tuple shape is handled consistently everywhere, including the finally-block cleanup in request_input/request_confirmation/request_selection

Confirmed this fixes a live bug, not a theoretical one

I checked whether self._event_loop (the field this PR removes) was ever actually populated anywhere in bus.py on main:

$ git show main:code_puppy/messaging/bus.py | grep -n "_event_loop"
83:        self._event_loop: Optional[asyncio.AbstractEventLoop] = None
413:            if self._event_loop is not None:
415:                    self._event_loop.call_soon_threadsafe(

It's initialized to None and never assigned anywhere else (confirmed with a repo-wide grep too). So the "thread-safe" branch in the old _complete_request was dead code — every cross-thread response completion on main was always hitting the direct, unsafe future.set_result() path. This PR isn't hardening against a hypothetical race, it's fixing the only code path that was ever actually running.

Race handling checks out

  • _complete_request pops the (loop, future) tuple atomically under self._lock before scheduling anything, so of two concurrent responses for the same prompt_id, only one wins the pop. test_complete_request_racing_responses_schedule_once verifies this with a genuine threading.Barrier-synchronized two-thread race and asserts call_soon_threadsafe fires exactly once — a real regression test, not a mocked assumption.
  • Cancellation-in-flight is safe by construction: _set_future_result (the scheduled callback) and any future.cancel() both execute on the same owning event-loop thread, so there's no TOCTOU window between the done() check and set_result().
  • Closed-loop cleanup correctly avoids ever touching the Future directly from a foreign thread — it just drops the completion, which is the right call since a closed loop can't resume the waiting coroutine anyway.

One tiny nit, non-blocking

The ASCII diagram in the module docstring (around line 29) still says prompt_id → Future; could use a follow-up touch to prompt_id → (loop, Future) for accuracy, but doesn't affect correctness or block this.

Approving — nice, tightly-scoped fix with real regression coverage for the exact race it claims to close.

@thomwebb
thomwebb merged commit 22f9311 into mpfaffenberger:main Sep 22, 2026
3 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants