Skip to content

Closing a live event stream early can hang forever in _merge_live_event_streams聽#7357

Description

@harshal-96

馃敶 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:

  1. Install ADK from main (checked at 4d06641).
  2. 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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

No labels
No labels

Type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions