Visitar URL original
[Agent Engines] `_wrap_async_stream_query_operation` blocks asyncio event loop due to synchronous gRPC iteration · Issue #7136 · googleapis/python-aiplatform · GitHub
Skip to content

[Agent Engines] _wrap_async_stream_query_operation blocks asyncio event loop due to synchronous gRPC iteration #7136

Description

@steffanianigro

agent_engine.async_stream_query(...) blocks the calling asyncio event loop thread during stream generation.

In vertexai/agent_engines/_agent_engines.py, _wrap_async_stream_query_operation wraps a synchronous client call in an async def and uses a blocking for loop:

# CURRENT IMPLEMENTATION (vertexai/agent_engines/_agent_engines.py)
def _wrap_async_stream_query_operation(*, method_name: str):
    async def _method(self, **kwargs):
        # 1. Uses sync client instead of self.execution_async_client
        response = self.execution_api_client.stream_query_reasoning_engine(...)
        # 2. Synchronous iteration blocks the asyncio event loop on socket reads
        for chunk in response:
            for parsed_json in _utils.yield_parsed_json(chunk):
                if parsed_json is not None:
                    yield parsed_json
    return _method

Reproduction:

import asyncio
from vertexai import agent_engines

agent = agent_engines.get("projects/<P>/locations/<L>/reasoningEngines/<ID>")

async def main():
    async for chunk in agent.async_stream_query(user_id="u", message="Long prompt"):
        pass

asyncio.run(main(), debug=True)
# Result: Emits "Executing <Task ...> took X.XX seconds" because the thread is blocked on socket reads without yielding.

Activity

  1. ArushKhasru commented on Sep 13, 2026

    @ArushKhasru

    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!

  2. self-assigned this
    on Sep 18, 2026
  3. spectakural commented on Sep 22, 2026

    @spectakural

    I 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_operation contains no suspension point - execution_api_client.stream_query_reasoning_engine() is the synchronous client and for 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_engine is already annotated -> Awaitable[AsyncIterable[HttpBody]] and its docstring already documents the stream = await client.stream_query_reasoning_engine(...) / async for response in stream: pattern. It works today even though the method is a plain def that returns rpc(...) without awaiting, because api_core's grpc_helpers_async._wrap_stream_errors wraps the RPC in a coroutine function - so calling it returns a coroutine, and awaiting that yields the async iterable. The fix is contained to vertexai/agent_engines/_agent_engines.py.

    2. The existing async tests cannot observe this bug.

    The three async_stream_query tests in tests/unit/vertex_langchain/test_agent_engines.py patch the synchronous ReasoningEngineExecutionServiceClient, 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 patching ReasoningEngineExecutionServiceAsyncClient.

    Note these tests live under tests/unit/vertex_langchain, which nox -s unit explicitly ignores; the session that covers them is nox -s unit_langchain.

    I have a branch with the fix, the test-fixture change, and a regression test that asserts asyncio.wait_for can 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.

  4. maxgasztych commented on Sep 24, 2026

    @maxgasztych

    Sorry the fix is already merged internally, there was a delay due to broken tests

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

Metadata

Metadata

Assignees

Labels

api: vertex-aiIssues related to the googleapis/python-aiplatform API.

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions