Repository navigation
fix: surface prompt stream errors instead of hanging query() forever - #1109
Echolonius wants to merge 1 commit into
Conversation
stream_input() caught every exception from the caller's prompt AsyncIterable (or a failed stdin write) and logged it at debug level. Stdin stayed open, the CLI kept waiting for input that would never come, no result ever arrived, and query() blocked forever with no surfaced error. On failure (outside teardown): log at error level, close stdin so the CLI can wind down, and push an error message into the message stream so receive_messages() raises instead of blocking indefinitely.
avshalomd
left a comment
There was a problem hiding this comment.
Independently verified this branch against the #1108 repros. Environment: SDK from source, CLI 2.1.226, Python 3.12.13, macOS arm64, 45s watchdog. main @ be2d0df, this branch @ e6cab50.
On main — both variants hang until the watchdog kills them (3/3 runs each), with no visible trace at the default log level:
- generator raises before first yield → 0 messages, hang (45s)
- generator yields once, turn completes normally (assistant + result arrive), then raises → hang after the healthy-looking turn (45s)
On this branch — both surface promptly from the message iterator (~1s; the post-turn variant a few seconds more, dominated by its live API turn — consistent across repeated runs), message Error streaming input: <original message>, and the failure is also logged at error level, visible by default. Clean path unaffected: a generator that ends normally still completes with the full message sequence.
ruff check, mypy src, and pytest (1057 passed, 5 skipped, including the regression test added here) all green on the branch. The runs above are end-to-end against the real CLI, complementing that test's mock transport.
One observation, not a blocker: the caller's original exception type is lost — a custom RuntimeError subclass surfaces as a generic Exception with the message embedded, since the failure travels as a wire-format error frame. May well be the intended contract; noting it for the record.
Repro harness (click to expand)
"""Repro harness for issue #1108 — exception in the prompt AsyncIterable.
Variant A: generator raises before yielding anything.
Variant B: generator yields one valid message, waits for the turn to
complete (result message observed), then raises on the next pull.
Each variant runs under a watchdog. On main both are expected to HANG
(watchdog fires); on pr-1109 both are expected to RAISE promptly.
"""
import anyio
import sys
import time
import traceback
from claude_agent_sdk import query, ClaudeAgentOptions
WATCHDOG_S = 45
class DeliberateError(RuntimeError):
pass
async def gen_raise_immediately():
raise DeliberateError("generator failed before first yield")
yield # pragma: no cover
def make_gen_raise_after_turn(turn_done: anyio.Event):
async def gen():
yield {
"type": "user",
"message": {"role": "user", "content": "Reply with the single word: pong"},
}
# Wait until the consumer has seen the result for turn 1,
# then raise on the next pull of the generator.
await turn_done.wait()
raise DeliberateError("generator failed after a completed turn")
return gen
async def run_variant(name, prompt_factory, turn_done=None):
opts = ClaudeAgentOptions(max_turns=1)
t0 = time.monotonic()
messages = []
outcome = None
try:
with anyio.fail_after(WATCHDOG_S):
async for msg in query(prompt=prompt_factory(), options=opts):
messages.append(type(msg).__name__)
if type(msg).__name__ == "ResultMessage" and turn_done is not None:
turn_done.set()
outcome = "COMPLETED"
except TimeoutError:
outcome = "HANG(watchdog)"
except BaseException as e:
outcome = f"RAISED {type(e).__name__}: {e}"
if "--trace" in sys.argv:
traceback.print_exc()
dt = time.monotonic() - t0
print(f"RESULT variant={name} outcome={outcome!r} elapsed={dt:.1f}s messages={messages}")
async def main():
which = sys.argv[1] if len(sys.argv) > 1 else "both"
if which in ("A", "both"):
await run_variant("A-raise-before-first-yield", gen_raise_immediately)
if which in ("B", "both"):
ev = anyio.Event()
await run_variant("B-raise-after-completed-turn", make_gen_raise_after_turn(ev), ev)
anyio.run(main)Observed output — representative run, main @ be2d0df (note: no error line on stderr):
RESULT variant=A-raise-before-first-yield outcome='HANG(watchdog)' elapsed=45.3s messages=[]
RESULT variant=B-raise-after-completed-turn outcome='HANG(watchdog)' elapsed=45.4s messages=['SystemMessage', 'RateLimitEvent', 'AssistantMessage', 'ResultMessage']
Observed output — representative run, this branch @ e6cab50 (error line now on stderr):
Error streaming input: generator failed before first yield
RESULT variant=A-raise-before-first-yield outcome='RAISED Exception: Error streaming input: generator failed before first yield' elapsed=1.2s messages=[]
Error streaming input: generator failed after a completed turn
RESULT variant=B-raise-after-completed-turn outcome='RAISED Exception: Error streaming input: generator failed after a completed turn' elapsed=4.8s messages=['SystemMessage', 'RateLimitEvent', 'AssistantMessage', 'ResultMessage']
Verification runs were agent-assisted; I reviewed the results myself.
tonydzi
left a comment
There was a problem hiding this comment.
Hi, Mycroft here, TonyDzi's synthetic AI co-founder (disclosure: AI-written, human-accountable). I re-tested this PR against current main. Full table is in #1108.
- Merge conflict.
git merge origin/main(763922b) conflicts in thestream_inputexceptblock. #1204 rewrote that block after this PR was opened, andmainno longer hangs: it logs at error level and always callsend_input(). - Case A is still silent here. When the generator raises before its first yield,
end_input()runs before the{"type": "error"}send._read_messagessees EOF, sendsendand closes_message_send, so the error send hits a suppressedClosedResourceError. With an instant-EOF fake transport this fails 3/3. Case B (raise after one message) passes 3/3. - Race-free alternative. Store the exception before
end_input()and raise it fromreceive_messageson theendframe. That's 4 lines and passes A, B and a control 3/3. It flipsmain's twotest_prompt_iterable_that_raises_*tests, which currently lock in the silent behavior. So this is a contract call for maintainers, see #1108.
Repro + diff: https://gist.github.com/tonydzi/1c7c0e0b286e495b3015ae4b685fb79e
— TonyDzi · more of this (agent fleet, consensus tooling): github.com/tonydzi
Fixes #1108.
Problem
stream_input()caught every exception from the caller's promptAsyncIterable(or a failed stdin write /json.dumps) and logged it at debug level. Stdin stayed open, the CLI kept waiting for input that would never come, no result ever arrived, andquery()blocked forever with no surfaced error. Verified against the real CLI: a generator that raises hangsquery()indefinitely whether it raises before the first message or between messages.Fix
In the exception handler, when not already in teardown (
self._closed):end_input(), best-effort) so the CLI can wind down instead of waiting on input,{"type": "error"}message into the message stream — the same mechanism_read_messagesalready uses for fatal read errors — soreceive_messages()raises promptly.During teardown the old debug-and-return behavior is preserved (writes interrupted by
close()are expected noise, not actionable errors). If the read task already closed the stream, its error wins (ClosedResourceErrorsuppressed).Verification
Error streaming input: boom-...in 3-7s.test_streaming_prompt_input_stream_error_surfaces: mock transport whose stdout produces nothing on its own (like the real CLI mid-wait); prompt generator yields one message then raises. Fails on main (bounded 5s timeout — the hang); passes with the fix, and assertsend_input()was called.pytest tests/test_query.py tests/test_tool_callbacks.py tests/test_streaming_client.py: 137 passed.ruff checkandruff format --checkclean.Note: independent of #1106 (different hunks in the same file; no overlap).
🤖 Generated with Claude Code
https://claude.ai/code/session_01YcMi8ny9DUBDf6y6m76HWk