fix(streaming): stop suppressing context errors and publish DONE once - #539
Open
michaelxu2288 wants to merge 2 commits into
Open
michaelxu2288 wants to merge 2 commits into
michaelxu2288 wants to merge 2 commits into
Conversation
StreamingTaskMessageContext.__aexit__ returned close()'s TaskMessage, and a truthy __aexit__ return value suppresses the exception. An error raised inside `async with streaming_task_message_context(...)` was therefore swallowed: the caller carried on past the block, the half-written message was persisted as DONE, and a Temporal activity saw a success and never retried. __aexit__ now closes the context and returns False, matching the convention the inference-call context manager already follows. Separately, stream_update() published an explicit StreamTaskMessageDone itself and then called close(), which published a second DONE because _is_closed was still False. Two DONE frames reached the stream, and since close() drains the coalescing buffer before publishing, buffered deltas could land after the first one. The Done branch now delegates to close(), which reaps the buffer, publishes exactly one DONE and persists. Full and delta handling are unchanged. Verified with `uv run pytest -q -n 0` on tests/lib/core/services/adk/test_streaming.py: the three added regression tests fail before the fix (two "DID NOT RAISE RuntimeError", one "expected exactly one DONE publish, got 2") and the file is 37 passed after. tests/lib/adk, tests/lib/core/harness, tests/lib/test_claude_agents_hooks.py and the openai_agents test_streaming_model.py suite stay green.
Routing an explicit StreamTaskMessageDone through close() dropped its index, and consumers use that index to stop the matching progress indicator. The Done branch now finishes with the caller's index, so the single Done published after the buffer drains carries it. close() is unchanged for callers.
This branch has not been deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Problem
Two bugs in
StreamingTaskMessageContext(src/agentex/lib/core/services/adk/streaming.py).1. Errors raised inside the context are swallowed.
__aexit__wasreturn await self.close(), andclose()returns theTaskMessage. A truthy return from__aexit__suppresses the exception, so inthe exception disappears. Execution continues after the block, the half-written message is persisted as
DONE, and a Temporal activity reports success, so its retry policy never runs.2. An explicit
StreamTaskMessageDoneis published twice.stream_update()published the Done and then calledclose()._is_closedwas stillFalseat that point, soclose()published a second Done. The message was also updated twice. And becauseclose()drains the coalescing buffer before publishing, buffered deltas could land after the first Done.Fix
__aexit__closes the context and returnsFalse. On the error path the context is still closed andDONEis still persisted, but the exception is no longer suppressed. The inference-call context manager already follows this convention.StreamTaskMessageDonepassed tostream_update()now goes straight toclose(), which reaps the buffer, publishes exactly one Done and persists once. Delta and Full handling are unchanged.Verification
uv run pytest -n 0 tests/lib/core/services/adk/test_streaming.pytest_exception_inside_context_propagates: DID NOT RAISE.test_context_is_still_closed_when_body_raises: DID NOT RAISE.test_explicit_done_publishes_exactly_one_done: expected exactly one DONE publish, got 2.tests/lib/adk,tests/lib/core/harness,tests/lib/test_claude_agents_hooks.py, and the openai_agentstest_streaming_model.py.ruff checkis clean on both files.The PR appears safe to merge; the earlier Done-index issue is fixed, and this review found no new issue.
What we checked:
StreamTaskMessageDone.indexalready defaults toNone, so plain close still sends the same index value.Summary
The streaming context now lets errors from its body reach the caller while still closing the message. Explicit DONE updates now finish after queued deltas and publish and save the final message once.
Reviews (2) · Last reviewed commit: "fix(streaming): keep the caller's index ..."