馃敶 Required Information
Describe the Bug:
_merge_live_event_streams in src/google/adk/live/_runner_utils.py merges the live agent's events and the events queued on ic._event_queue into a one-slot queue, merged.
If the caller stops reading and closes the stream while merged holds an event, the cleanup in the merge's finally cancels _pump_queued_events. That pump's own finally then runs await merged.put(done_sentinel). Nothing reads merged anymore and the queue is full, so the put never finishes. The cleanup waits on the pump forever, and aclose() never returns.
It depends on timing: it only hangs when merged is full at the moment of closing.
Steps to Reproduce:
- Install ADK from
main (checked at 4d06641).
- Run this script:
import asyncio
from google.adk.agents.base_agent import BaseAgent
from google.adk.events.event import Event
from google.adk.live import _runner_utils
from google.adk.live import LiveRequestQueue
from google.adk.runners import Runner
from google.adk.sessions.in_memory_session_service import InMemorySessionService
class Agent(BaseAgent):
async def _run_impl(self, ctx):
yield Event(author=self.name)
async def main():
runner = Runner(
app_name="repro",
agent=Agent(name="root"),
session_service=InMemorySessionService(),
)
session = await runner.session_service.create_session(
app_name="repro", user_id="u", session_id="s"
)
ic = runner._new_invocation_context_for_live(
session, live_request_queue=LiveRequestQueue()
)
ic._event_queue = asyncio.Queue()
for i in range(3):
ic._event_queue.put_nowait((Event(author=f"queued_{i}", partial=True), None))
async def agent_events():
await asyncio.Event().wait() # a live agent that is still running
yield Event(author="never")
merged = _runner_utils._merge_live_event_streams(runner, ic, agent_events())
print("got", (await anext(merged)).author)
await asyncio.sleep(0.1) # the queue pump refills the one-slot queue
print("closing...")
await merged.aclose() # never returns
print("closed")
asyncio.run(main())
Expected Behavior:
aclose() returns, both pumps are cleaned up, and the script prints closed.
Observed Behavior:
The script prints got queued_0 and closing..., then hangs. The pending task is stuck here:
_merge_live_event_streams.<locals>._pump_queued_events
google/adk/live/_runner_utils.py:302 await merged.put(done_sentinel)
It also makes the LiveKit tests that drive a real run_live hang on my machine. On every run, one or both of these tests time out:
tests/unittests/integrations/livekit/test_livekit_runner.py::test_room_drives_a_real_run_live
tests/unittests/integrations/livekit/test_livekit_call.py::test_real_tools_reach_the_call_during_a_live_session
They pass on GitHub Actions, most likely because of different timing there.
Environment Details:
- ADK Library Version:
main at 4d06641. The released google-adk 2.10.0 on PyPI has the same one-slot queue and the same sentinel put in finally.
- Desktop OS: Linux (Ubuntu on WSL2) and Windows 11
- Python Version: 3.10, 3.11, 3.12, 3.13, 3.14
Model Information:
- Are you using LiteLLM: No
- Which model is being used: N/A (no model call is needed to reproduce)
馃煛 Optional Information
Regression:
Not sure. This code is unchanged since at least the earliest commit in this repository's history (0961ced), and 2.10.0 has it too.
Logs:
I ran a stress script for 500 trials. Each trial used a random number of queued and agent events and closed the stream at a random point. On main with Python 3.12, 131 of 500 trials hung on Linux and 157 of 500 on Windows.
Additional Context:
I have a small fix with a regression test ready and will open a PR that links here.
馃敶 Required Information
Describe the Bug:
_merge_live_event_streamsinsrc/google/adk/live/_runner_utils.pymerges the live agent's events and the events queued onic._event_queueinto a one-slot queue,merged.If the caller stops reading and closes the stream while
mergedholds an event, the cleanup in the merge'sfinallycancels_pump_queued_events. That pump's ownfinallythen runsawait merged.put(done_sentinel). Nothing readsmergedanymore and the queue is full, so the put never finishes. The cleanup waits on the pump forever, andaclose()never returns.It depends on timing: it only hangs when
mergedis full at the moment of closing.Steps to Reproduce:
main(checked at 4d06641).Expected Behavior:
aclose()returns, both pumps are cleaned up, and the script printsclosed.Observed Behavior:
The script prints
got queued_0andclosing..., then hangs. The pending task is stuck here:It also makes the LiveKit tests that drive a real
run_livehang on my machine. On every run, one or both of these tests time out:tests/unittests/integrations/livekit/test_livekit_runner.py::test_room_drives_a_real_run_livetests/unittests/integrations/livekit/test_livekit_call.py::test_real_tools_reach_the_call_during_a_live_sessionThey pass on GitHub Actions, most likely because of different timing there.
Environment Details:
mainat 4d06641. The released google-adk 2.10.0 on PyPI has the same one-slot queue and the same sentinel put infinally.Model Information:
馃煛 Optional Information
Regression:
Not sure. This code is unchanged since at least the earliest commit in this repository's history (0961ced), and 2.10.0 has it too.
Logs:
I ran a stress script for 500 trials. Each trial used a random number of queued and agent events and closed the stream at a random point. On
mainwith Python 3.12, 131 of 500 trials hung on Linux and 157 of 500 on Windows.Additional Context:
I have a small fix with a regression test ready and will open a PR that links here.