Skip to content

fix: surface prompt stream errors instead of hanging query() forever - #1109

Open
Echolonius wants to merge 1 commit into
anthropics:mainfrom
Echolonius:fix/stream-input-error-surfacing
Open

Echolonius wants to merge 1 commit into
anthropics:mainfrom
Echolonius:fix/stream-input-error-surfacing

Conversation

@Echolonius

Copy link
Copy Markdown

Fixes #1108.

Problem

stream_input() caught every exception from the caller's prompt AsyncIterable (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, and query() blocked forever with no surfaced error. Verified against the real CLI: a generator that raises hangs query() indefinitely whether it raises before the first message or between messages.

Fix

In the exception handler, when not already in teardown (self._closed):

  1. log at error level instead of debug,
  2. close stdin (end_input(), best-effort) so the CLI can wind down instead of waiting on input,
  3. push an {"type": "error"} message into the message stream — the same mechanism _read_messages already uses for fatal read errors — so receive_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 (ClosedResourceError suppressed).

Verification

  • Real CLI (SDK 0.2.116, CLI 2.1.207): both repro variants from Exception in the prompt AsyncIterable is swallowed at debug level and query() hangs forever #1108 hung for the full 45s watchdog on main; with this fix they raise Error streaming input: boom-... in 3-7s.
  • New regression test 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 asserts end_input() was called.
  • pytest tests/test_query.py tests/test_tool_callbacks.py tests/test_streaming_client.py: 137 passed. ruff check and ruff format --check clean.

Note: independent of #1106 (different hunks in the same file; no overlap).

🤖 Generated with Claude Code

https://claude.ai/code/session_01YcMi8ny9DUBDf6y6m76HWk

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 avshalomd left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 tonydzi left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

  1. Merge conflict. git merge origin/main (763922b) conflicts in the stream_input except block. #1204 rewrote that block after this PR was opened, and main no longer hangs: it logs at error level and always calls end_input().
  2. Case A is still silent here. When the generator raises before its first yield, end_input() runs before the {"type": "error"} send. _read_messages sees EOF, sends end and closes _message_send, so the error send hits a suppressed ClosedResourceError. With an instant-EOF fake transport this fails 3/3. Case B (raise after one message) passes 3/3.
  3. Race-free alternative. Store the exception before end_input() and raise it from receive_messages on the end frame. That's 4 lines and passes A, B and a control 3/3. It flips main's two test_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

This branch has not been deployed

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Exception in the prompt AsyncIterable is swallowed at debug level and query() hangs forever

3 participants