Conversation
…completion
Interim run events (tool progress, message deltas) are emitted from the
executor thread running agent.run_conversation() via
loop.call_soon_threadsafe(q.put_nowait, event). That only schedules the
put; it doesn't wait for the event loop to apply it. Once
run_in_executor() returns, the run coroutine resumes on the loop thread
and enqueues run.completed/run.failed and the closing None sentinel
directly on that same thread. A still-pending scheduled put for an
interim event can lose that race and land after run.completed, or after
the SSE consumer has already seen the sentinel and stopped reading --
silently dropping the event from the /v1/runs/{run_id}/events stream.
Observed on Python 3.13.
Add _enqueue_run_event(), a small ack-wait wrapper around
call_soon_threadsafe: when called off the loop thread it blocks (bounded
by a 2s timeout so loop shutdown can't strand the executor thread) until
the loop has actually run the given put, so the caller cannot proceed
toward run completion until its event is durably queued ahead of
anything the loop thread does afterward. Wire it into both interim event
paths that cross the executor thread/event loop boundary: the
tool_progress callback (_push) and the stream_delta_callback (_text_cb).
_text_cb routes through the existing _put_event_if_active staleness
guard rather than a bare queue put, so _enqueue_run_event takes a zero-arg
callable instead of a fixed (queue, event) pair to accommodate both call
sites without bypassing that guard.
Add regression coverage: unit tests for _enqueue_run_event itself
(blocks until applied off-loop, synchronous on-loop, FIFO against a
racing direct put, and respects the _put_event_if_active staleness
guard) plus an SSE-level test that an interim tool event fired just
before run_conversation() returns still shows up ahead of run.completed
in the stream. Tests: tests/gateway/test_api_server_runs.py.
cbd7065 to
0cb2b7e
Compare
Reviewed by reviewer-e (AI automated review). Correct race analysis and fix:
|
Interim run events (tool progress, reasoning previews, message deltas) are
emitted from the executor thread running
agent.run_conversation()vialoop.call_soon_threadsafe(q.put_nowait, event). That only schedules theput; it doesn't wait for the event loop to apply it. Once
run_in_executor()returns, the run coroutine resumes on the loop threadand enqueues
run.completed/run.failedand the closingNonesentineldirectly via
q.put_nowait()on that same thread. A still-pendingscheduled put for an interim event can lose that race and land after
run.completed, or after the SSE consumer has already seen the sentineland stopped reading -- silently dropping the event from the
/v1/runs/{run_id}/eventsstream. Observed on Python 3.13.Add
_enqueue_run_event(), a small ack-wait wrapper aroundcall_soon_threadsafe: when called off the loop thread it blocks (boundedby a 2s timeout so loop shutdown can't strand the executor thread) until
the loop has actually applied the put, so the caller cannot proceed toward
run completion until its event is durably queued ahead of anything the loop
thread does afterward. Wired into both interim event paths that cross the
executor thread/event loop boundary: the
tool_progress_callbackand thestream_delta_callback.Testing
applied when called off the loop thread, synchronous passthrough when
already on the loop thread, and FIFO ordering preserved against a
directly-enqueued completion event. Confirmed they fail against the
pre-fix code and pass after.
run_conversation()returns still appears ahead ofrun.completedinthe stream.
tests/gateway/test_api_server_runs.pysuite passes (26/26), plusthe broader
test_api_server*suite (267/267), no regressions.