Repository navigation
[Agent Engines] _wrap_async_stream_query_operation blocks asyncio event loop due to synchronous gRPC iteration #7136
Description
Activity
- addedapi: vertex-aiIssues related to the googleapis/python-aiplatform API.Issues related to the googleapis/python-aiplatform API.
on Sep 8, 2026 Hi maintainers, I'd like to work on this issue. Could you assign it to me if it's available?
I noticed that #5773 and the open PR #6295 cover related async streaming behavior. I'd be happy to help with any remaining work, including regression tests to verify that streaming doesn't block the asyncio event loop. Please let me know whether you'd prefer contributions to the existing fix or a separate PR for any remaining gaps.
Thanks!
Reacted by steffanianigroI hit this too, and I can confirm the diagnosis in this issue. Adding a repro and two findings that I think matter for whoever picks it up.
Repro: the timeout simply never fires while the agent is generating:
import asyncio from vertexai import agent_engines agent = agent_engines.get("projects/.../locations/.../reasoningEngines/...") async def heartbeat(): while True: print("tick") await asyncio.sleep(0.5) async def main(): asyncio.create_task(heartbeat()) try: async with asyncio.timeout(5): async for chunk in agent.async_stream_query(input="..."): print(chunk) except TimeoutError: print("timed out") asyncio.run(main())
The
ticks stop for the whole duration of the stream and the 5s timeout does not fire, because_wrap_async_stream_query_operationcontains no suspension point -execution_api_client.stream_query_reasoning_engine()is the synchronous client andfor chunk in response:blocks the thread on the socket read, so the event loop never gets control back to run its timer.1. The generated GAPIC client does not need to change.
ReasoningEngineExecutionServiceAsyncClient.stream_query_reasoning_engineis already annotated-> Awaitable[AsyncIterable[HttpBody]]and its docstring already documents thestream = await client.stream_query_reasoning_engine(...)/async for response in stream:pattern. It works today even though the method is a plaindefthat returnsrpc(...)without awaiting, becauseapi_core'sgrpc_helpers_async._wrap_stream_errorswraps the RPC in a coroutine function - so calling it returns a coroutine, and awaiting that yields the async iterable. The fix is contained tovertexai/agent_engines/_agent_engines.py.2. The existing async tests cannot observe this bug.
The three
async_stream_querytests intests/unit/vertex_langchain/test_agent_engines.pypatch the synchronousReasoningEngineExecutionServiceClient, so they assert against the blocking path. That is why CI has stayed green. It also means the source fix on its own turns those three tests red - I verified this locally - and any fix needs to repoint them at a fixture patchingReasoningEngineExecutionServiceAsyncClient.Note these tests live under
tests/unit/vertex_langchain, whichnox -s unitexplicitly ignores; the session that covers them isnox -s unit_langchain.I have a branch with the fix, the test-fixture change, and a regression test that asserts
asyncio.wait_forcan actually cancel a slow stream (it fails against the current implementation and passes against the fix). Happy to open it as a PR, or to fold the test changes into #6295 if its author prefers.Sorry the fix is already merged internally, there was a delay due to broken tests
agent_engine.async_stream_query(...)blocks the callingasyncioevent loop thread during stream generation.In
vertexai/agent_engines/_agent_engines.py,_wrap_async_stream_query_operationwraps a synchronous client call in anasync defand uses a blockingforloop:Reproduction: