From bd54eaa0a430699ab96d8225cf2353834d401617 Mon Sep 17 00:00:00 2001 From: kalisgd0h Date: Sun, 5 Jul 2026 06:53:52 +0200 Subject: [PATCH] feat(cli): add --output-format stream-json for headless clients (#309) * feat(cli): add --output-format stream-json for headless clients Emit EvoScientist's native event stream as line-delimited JSON on stdout in single-shot (-p) mode, with all human output redirected to stderr so stdout stays pure JSONL. Intended as the integration surface for programmatic clients (e.g. an agent runtime) that drive EvoSci headlessly. - stream/json_sink.py: write_events_as_json + stream_json sink, plus redirect_console_to_stderr helper for stdout purity - cli/interactive.py: cmd_run gains output_format; stream-json branch runs the sink instead of the Rich renderer - cli/commands.py: --output-format option + validation (stream-json requires -p; value must be text|stream-json) - docs/stream-json.md: event-schema contract + example transcript - tests: json sink serialization, CLI dispatch, console redirect, validation Co-Authored-By: Claude Opus 4.8 * fix(cli): honor explicit --no-auto-mode over config in stream-json Address CodeRabbit review (discussion_r3514041123): the auto-mode override block only wrote to cli_overrides when the resolved value was True, so an explicit --no-auto-mode silently fell back to a config that enables auto-mode -- breaking "explicit flags always win" and leaving stream-json running unattended despite the warning. Write auto_mode=False when the flag is explicitly False. Add regression tests that capture the overrides. Co-Authored-By: Claude Opus 4.8 --------- Co-authored-by: Claude Opus 4.8 --- EvoScientist/cli/commands.py | 122 ++++++++++++++++++++++---- EvoScientist/stream/json_sink.py | 77 +++++++++++++++++ README.md | 6 ++ docs/stream-json.md | 81 +++++++++++++++++ tests/test_cli_output_format.py | 144 +++++++++++++++++++++++++++++++ tests/test_json_sink.py | 127 +++++++++++++++++++++++++++ 6 files changed, 539 insertions(+), 18 deletions(-) create mode 100644 EvoScientist/stream/json_sink.py create mode 100644 docs/stream-json.md create mode 100644 tests/test_cli_output_format.py create mode 100644 tests/test_json_sink.py diff --git a/EvoScientist/cli/commands.py b/EvoScientist/cli/commands.py index 2c1e0ba..2b95042 100644 --- a/EvoScientist/cli/commands.py +++ b/EvoScientist/cli/commands.py @@ -20,6 +20,7 @@ from ..commands.base import ChannelRuntime, Command, CommandContext from ..gateway import ( GraphGateway, GraphTarget, + RunRequest, RuntimeGateways, create_runtime_gateways, ) @@ -2009,6 +2010,18 @@ def _is_fresh_interactive_session(prompt: str | None, thread_id: str | None) -> return not prompt and not thread_id +def _resolve_stream_json_auto_mode(auto_mode: bool | None, output_format: str) -> bool: + """Resolve the effective ``--auto-mode`` for a single-shot run. + + stream-json is headless, so auto-mode defaults on there when the caller did + not pass the flag (``None``); an explicit ``--auto-mode`` / ``--no-auto-mode`` + always wins, and non-stream-json runs keep the historical off-by-default. + """ + if auto_mode is not None: + return auto_mode + return output_format == "stream-json" + + @app.callback(invoke_without_command=True) def _main_callback( ctx: typer.Context, @@ -2055,10 +2068,11 @@ def _main_callback( "--auto-approve", help="Skip tool approval prompts for HITL actions", ), - auto_mode: bool = typer.Option( - False, - "--auto-mode", - help="Run unattended: skip ask_user and tool approval prompts", + auto_mode: bool | None = typer.Option( + None, + "--auto-mode/--no-auto-mode", + help="Run unattended: skip ask_user and tool approval prompts " + "(default: on when --output-format stream-json)", ), ask_user: bool = typer.Option( False, @@ -2080,6 +2094,14 @@ def _main_callback( "--ui", help="UI backend: tui (default), cli, or webui.", ), + output_format: str | None = typer.Option( + None, + "--output-format", + help=( + "Output format for single-shot (-p) mode: 'text' (default) or " + "'stream-json' (line-delimited JSON events to stdout)." + ), + ), ): """EvoScientist Agent - AI-powered research & code execution CLI""" # If a subcommand was invoked, don't run the default behavior @@ -2089,6 +2111,37 @@ def _main_callback( # Load and apply configuration from ..config import apply_config_to_env, get_effective_config + # Resolve the output format first. In stream-json mode stdout must carry + # only JSONL, so establish the mode and redirect the console to stderr + # BEFORE anything below (ccproxy startup, validation) can print a + # human-readable line to stdout. + effective_output_format = (output_format or "text").lower() + if effective_output_format not in ("text", "stream-json"): + raise typer.BadParameter("--output-format must be 'text' or 'stream-json'") + if effective_output_format == "stream-json": + if not prompt: + raise typer.BadParameter( + "--output-format stream-json requires -p/--prompt (single-shot mode)" + ) + from ..stream.json_sink import redirect_console_to_stderr + + redirect_console_to_stderr() + + # stream-json is headless, so auto-mode defaults on (auto-handle approval and + # ask_user gates) unless the caller explicitly passed --no-auto-mode. Without + # it the run would stall at the first gate and end without doing the work; + # --no-auto-mode is the (experimental) opt-in to receiving interrupt/ask_user + # events and driving resume yourself. + effective_auto_mode = _resolve_stream_json_auto_mode( + auto_mode, effective_output_format + ) + if effective_output_format == "stream-json" and auto_mode is False: + console.print( + "[yellow]--no-auto-mode with stream-json is experimental: an " + "interrupt/ask_user event ends the run early and is not yet " + "resumable.[/yellow]" + ) + # Build CLI overrides dict cli_overrides = {} if mode: @@ -2101,12 +2154,18 @@ def _main_callback( cli_overrides["ui_backend"] = ui if auto_approve: cli_overrides["auto_approve"] = True - if auto_mode: + if effective_auto_mode: cli_overrides["auto_mode"] = True cli_overrides["auto_approve"] = True cli_overrides["enable_ask_user"] = False - elif ask_user: - cli_overrides["enable_ask_user"] = True + else: + # An explicit --no-auto-mode (auto_mode is False, not None) must win over + # a config file that enables auto-mode; without this the resolved-off + # value writes nothing and silently falls back to the config default. + if auto_mode is False: + cli_overrides["auto_mode"] = False + if ask_user: + cli_overrides["enable_ask_user"] = True if dangerous: cli_overrides["dangerous_mode"] = True if auth_mode: @@ -2261,7 +2320,8 @@ def _main_callback( import asyncio from ..sessions import get_checkpointer - from .interactive import cmd_run + from ..stream.json_sink import stream_json + from .interactive import _wait_for_memory_workers_before_exit, cmd_run from .resume_hint import print_resume_hint runtime_gateways = create_runtime_gateways() @@ -2296,16 +2356,42 @@ def _main_callback( config=config, ) try: - cmd_run( - agent, - prompt, - thread_id=tid, - show_thinking=show_thinking, - workspace_dir=workspace_dir, - model=config.model, - ui_backend=config.ui_backend, - runtime_gateways=runtime_gateways, - ) + if effective_output_format == "stream-json": + # Headless JSONL path: drive the sink through the gateway + # directly. We are already inside the async single-shot + # loop, so this is a plain await — no nested-loop juggling, + # and the gateway seam keeps it execution-backend agnostic. + request = RunRequest( + message=prompt, + thread_id=tid, + metadata=build_metadata(workspace_dir, config.model), + target=GraphTarget( + local_graph=agent, workspace_dir=workspace_dir + ), + ) + try: + await stream_json(graph_gateway, request) + except Exception as exc: + # stream_events already emitted a terminal `error` + # event onto the JSON stream before re-raising; exit + # cleanly so the stream ends with that event instead + # of a raw traceback. + raise typer.Exit(1) from exc + finally: + # Let post-run memory workers persist before exit, + # matching the text path (cmd_run does this itself). + _wait_for_memory_workers_before_exit() + else: + cmd_run( + agent, + prompt, + thread_id=tid, + show_thinking=show_thinking, + workspace_dir=workspace_dir, + model=config.model, + ui_backend=config.ui_backend, + runtime_gateways=runtime_gateways, + ) finally: try: print_resume_hint(tid, console=console) diff --git a/EvoScientist/stream/json_sink.py b/EvoScientist/stream/json_sink.py new file mode 100644 index 0000000..0f2eba9 --- /dev/null +++ b/EvoScientist/stream/json_sink.py @@ -0,0 +1,77 @@ +"""Line-delimited JSON (JSONL) output sink for headless / programmatic use. + +Selected by ``--output-format stream-json``. Serializes EvoScientist's native +event stream verbatim — one JSON object per line — onto a writable stream +(stdout in production). External clients (e.g. an agent daemon) read the stream +with a line scanner and ``json.loads`` per line, dispatching on each object's +``type`` field. + +This module deliberately contains no rendering logic: it writes the normalized +event dicts produced by :meth:`EvoScientist.gateway.GraphGateway.stream_events` +verbatim — the same events the Rich renderer consumes — so it stays agnostic to +whether the run executes locally or against a langgraph server. +""" + +from __future__ import annotations + +import json +import sys +from collections.abc import AsyncIterable +from typing import TYPE_CHECKING, Any, TextIO + +if TYPE_CHECKING: + from ..gateway import GraphGateway, RunRequest + + +async def write_events_as_json( + events: AsyncIterable[dict[str, Any]], + out: TextIO, +) -> str: + """Serialize an async stream of event dicts to ``out`` as JSONL. + + Each event is written as a single line and flushed immediately so consumers + receive events in real time. ``default=str`` ensures an unexpected + non-serializable value (e.g. a tool argument carrying a rich object) degrades + to its string form instead of raising and tearing down the stream. + + Returns the final response text carried by the terminal ``done`` event + (empty string if the stream ends without one). + """ + final = "" + async for event in events: + out.write(json.dumps(event, default=str)) + out.write("\n") + out.flush() + if event.get("type") == "done": + final = event.get("response", "") or "" + return final + + +def redirect_console_to_stderr() -> None: + """Route all human-facing Rich output to stderr. + + In stream-json mode stdout must carry only JSONL event lines, but the shared + ``console`` singleton (used everywhere for status lines, separators, resume + hints, error panels) writes to stdout by default. Reassigning its ``file`` + once moves every existing ``console.print`` call to stderr without touching + individual call sites. + """ + from .console import console + + console.file = sys.stderr + + +async def stream_json( + gateway: GraphGateway, + request: RunRequest, + *, + out: TextIO | None = None, +) -> str: + """Run a graph turn through the gateway and emit its event stream as JSONL. + + Sources events from :meth:`GraphGateway.stream_events` — the same normalized + event dicts the Rich renderer consumes — so it works for both local and + langgraph-server execution. ``out`` defaults to stdout. + """ + sink = out if out is not None else sys.stdout + return await write_events_as_json(gateway.stream_events(request), sink) diff --git a/README.md b/README.md index 3006e12..99cee43 100644 --- a/README.md +++ b/README.md @@ -428,8 +428,14 @@ EvoSci --ui cli # classic CLI (lightweight) EvoSci --ui webui # browser workspace UI (needs Node/npx) EvoSci serve # headless mode — channels only, no interactive prompt EvoSci deploy # standalone LangGraph server for external UIs / SDK clients +EvoSci -p "query" --output-format stream-json --auto-mode # JSONL event stream on stdout (for programmatic clients) ``` +`--output-format stream-json` makes a single-shot (`-p`) run emit its native +events as line-delimited JSON on stdout (one object per line), with all human +output on stderr — the integration surface for headless clients (e.g. an agent +runtime). See [docs/stream-json.md](docs/stream-json.md) for the event schema. +
diff --git a/docs/stream-json.md b/docs/stream-json.md new file mode 100644 index 0000000..36f503a --- /dev/null +++ b/docs/stream-json.md @@ -0,0 +1,81 @@ +# `stream-json` output protocol + +`EvoSci --output-format stream-json` runs a single-shot (`-p`) session and emits +EvoScientist's **native event stream** as line-delimited JSON (JSONL) on stdout — +one self-describing JSON object per line. This is the integration surface for +programmatic clients that drive EvoScientist headlessly (for example, an agent +runtime that assigns work and renders live progress). + +```bash +EvoSci -p "Summarize the attention mechanism into notes.md" \ + --output-format stream-json \ + --auto-mode \ + --workdir /path/to/task +``` + +## Stream contract + +- **stdout is pure JSONL.** Every line is one complete JSON object. Parse it with + a line reader + `json.loads` per line. Nothing else is written to stdout. +- **stderr carries everything human** — status lines ("Loading agent…"), + separators, the resume hint, and error panels. A consumer should treat stderr + as logs, not protocol. +- **Each object has a `type` field** used for dispatch. A consumer that does not + recognize a `type` should **ignore that line** rather than fail — new event + types may be added over time, and the protocol is forward-compatible by design. +- **The stream ends with a `done` event** carrying the final response text. On an + unhandled failure an `error` event is emitted instead/also. +- **Unattended by default.** `stream-json` is headless, so `--auto-mode` is + enabled automatically: approval and `ask_user` gates are auto-handled and the + run proceeds straight to its `done` event. Pass `--no-auto-mode` to opt out — + a human-in-the-loop `interrupt` / `ask_user` is then emitted as a normal event + and the single-shot run ends right after it (`… → interrupt → done → EOF`). It + does **not** block waiting for input, but it also stops before finishing the + task; answering the event requires re-invoking with `--resume` (experimental). + +## Event types + +Each line is a self-contained JSON object with a `type` field. Most fields are +scalars, but some events carry nested payloads (e.g. `args`, `action_requests`, +`questions`) — parse each line as a full object, not a flat key/value map. The +fields beyond `type` are listed below. + +| `type` | Fields | Meaning | +|------------------------|------------------------------------------------------------|---------| +| `thinking` | `content`, `id` | Model reasoning text | +| `text` | `content` | Assistant output text | +| `tool_call` | `name`, `args`, `id` | Tool invocation | +| `tool_result` | `name`, `content`, `success`, `id` | Tool result (`id` matches the `tool_call`) | +| `subagent_start` | `name`, `description` | Sub-agent delegation begins | +| `subagent_tool_call` | `subagent`, `name`, `args`, `id` | Tool call inside a sub-agent | +| `subagent_tool_result` | `subagent`, `name`, `content`, `success`, `id` | Tool result inside a sub-agent | +| `subagent_text` | `subagent`, `content`, `instance_id` | Text from a sub-agent | +| `subagent_end` | `name` | Sub-agent delegation completes | +| `tool_selection` | `tools` | Tool-selector middleware picked tools | +| `summarization_start` | — | Context summarization begins | +| `summarization` | `content` | Context summarization output | +| `usage_stats` | `input_tokens`, `output_tokens` | Token usage | +| `interrupt` | `interrupt_id`, `action_requests`, `review_configs` | HITL approval interrupt | +| `ask_user` | `interrupt_id`, `questions`, `tool_call_id` | Agent-initiated clarifying question | +| `error` | `message` | Error during the run | +| `done` | `content`, `response` | Final response; end of stream | + +## Example transcript + +```jsonl +{"type": "thinking", "content": "I should write the notes file.", "id": 0} +{"type": "tool_call", "name": "write_file", "args": {"path": "notes.md", "content": "..."}, "id": "call_1"} +{"type": "tool_result", "name": "write_file", "content": "wrote 412 bytes", "success": true, "id": "call_1"} +{"type": "text", "content": "Done — notes.md now summarizes attention."} +{"type": "usage_stats", "input_tokens": 5123, "output_tokens": 388} +{"type": "done", "content": "Done — notes.md now summarizes attention.", "response": "Done — notes.md now summarizes attention."} +``` + +## Notes for client implementers + +- Read stdout line by line; do not assume a single JSON document. +- Accumulate `text` events for the running assistant message; the terminal + `done.response` is the authoritative final text. +- `tool_call` / `tool_result` correlate by `id`. +- Token usage may arrive across multiple `usage_stats` events; sum them. +- Treat unknown `type` values and unknown fields as non-fatal. diff --git a/tests/test_cli_output_format.py b/tests/test_cli_output_format.py new file mode 100644 index 0000000..ceb4844 --- /dev/null +++ b/tests/test_cli_output_format.py @@ -0,0 +1,144 @@ +"""Tests for the --output-format stream-json CLI wiring. + +Covers: + 1. The console singleton is redirected to stderr so stdout stays pure JSONL. + 2. --output-format is validated (stream-json requires -p; bad values rejected). + +The dispatch itself runs the sink through ``GraphGateway.stream_events`` from the +single-shot path; that sourcing is exercised by the json_sink unit tests. +""" + +import sys + +import pytest + +from EvoScientist.cli.commands import _resolve_stream_json_auto_mode + + +@pytest.mark.parametrize( + ("auto_mode", "output_format", "expected"), + [ + (None, "stream-json", True), # default on for headless stream-json + (None, "text", False), # text keeps historical off-by-default + (False, "stream-json", False), # explicit --no-auto-mode wins + (True, "text", True), # explicit --auto-mode wins + ], +) +def test_resolve_stream_json_auto_mode(auto_mode, output_format, expected): + """auto-mode defaults on for stream-json when unset; explicit flags always win.""" + assert _resolve_stream_json_auto_mode(auto_mode, output_format) is expected + + +def _overrides_for(monkeypatch, argv): + """Invoke the CLI and capture the cli_overrides handed to get_effective_config, + aborting before the run actually starts.""" + from typer.testing import CliRunner + + import EvoScientist.config as cfg_mod + from EvoScientist.cli._app import app + from EvoScientist.stream.console import console + + captured: dict[str, object] = {} + + class _Stop(Exception): + """Sentinel to short-circuit the callback once overrides are captured.""" + + def _grab(overrides): + """Record the overrides dict, then abort before the heavy run path.""" + captured["overrides"] = dict(overrides) + raise _Stop + + monkeypatch.setattr(cfg_mod, "get_effective_config", _grab) + original_console_file = console.file + try: + CliRunner().invoke(app, argv, catch_exceptions=True) + finally: + console.file = original_console_file + return captured.get("overrides", {}) + + +def test_stream_json_defaults_auto_mode_on_in_overrides(monkeypatch): + """stream-json without the flag defaults auto-mode (and auto-approve) on.""" + overrides = _overrides_for( + monkeypatch, ["-p", "hi", "--output-format", "stream-json"] + ) + assert overrides.get("auto_mode") is True + assert overrides.get("auto_approve") is True + + +def test_no_auto_mode_writes_explicit_false_override(monkeypatch): + """Explicit --no-auto-mode must write auto_mode=False so it wins over a config + that enables auto-mode (not silently fall back to the config default).""" + overrides = _overrides_for( + monkeypatch, + ["-p", "hi", "--output-format", "stream-json", "--no-auto-mode"], + ) + assert overrides.get("auto_mode") is False + + +def test_redirect_console_to_stderr(): + """redirect_console_to_stderr moves the shared console's output to stderr.""" + from EvoScientist.stream.console import console + from EvoScientist.stream.json_sink import redirect_console_to_stderr + + original = console.file + try: + redirect_console_to_stderr() + assert console.file is sys.stderr + finally: + console.file = original + + +# ── CLI validation (via typer CliRunner) ── + + +def _invoke(monkeypatch, argv): + """Run the CLI app with argv under stubbed config/dispatch, returning + (dispatch-calls, CliRunner result).""" + from typer.testing import CliRunner + + import EvoScientist.cli.commands as cmds + import EvoScientist.cli.interactive as interactive_mod + import EvoScientist.config as cfg_mod + from EvoScientist.cli._app import app + from EvoScientist.config.settings import EvoScientistConfig + + calls: dict[str, object] = {} + + def _fake_config(overrides): + """Return a minimal config honoring only the ui_backend override.""" + cfg = EvoScientistConfig() + cfg.ui_backend = overrides.get("ui_backend") or "cli" + return cfg + + monkeypatch.setattr(cfg_mod, "get_effective_config", _fake_config) + monkeypatch.setattr(cfg_mod, "apply_config_to_env", lambda cfg: None) + monkeypatch.setattr(cmds, "ensure_dirs", lambda: None) + monkeypatch.setattr(cmds, "_ensure_async_subagent_server", lambda *a, **k: None) + monkeypatch.setattr( + interactive_mod, + "cmd_interactive", + lambda **kw: calls.__setitem__("dispatch", "interactive"), + ) + monkeypatch.setattr( + interactive_mod, + "cmd_run", + lambda *a, **kw: calls.__setitem__("dispatch", "run"), + ) + + result = CliRunner().invoke(app, argv, catch_exceptions=False) + return calls, result + + +def test_stream_json_without_prompt_is_rejected(monkeypatch): + """stream-json requires -p (single-shot); bare --output-format is rejected.""" + calls, result = _invoke(monkeypatch, ["--output-format", "stream-json"]) + assert result.exit_code != 0 + assert "dispatch" not in calls + + +def test_invalid_output_format_is_rejected(monkeypatch): + """An unknown --output-format value is rejected before any dispatch.""" + calls, result = _invoke(monkeypatch, ["-p", "hi", "--output-format", "bogus"]) + assert result.exit_code != 0 + assert "dispatch" not in calls diff --git a/tests/test_json_sink.py b/tests/test_json_sink.py new file mode 100644 index 0000000..22fa7ab --- /dev/null +++ b/tests/test_json_sink.py @@ -0,0 +1,127 @@ +"""Tests for the stream-json output sink (EvoScientist/stream/json_sink.py). + +The sink serializes the agent's native event stream as line-delimited JSON +(JSONL) to a writable stream, one JSON object per line. It is the headless +output path selected by ``--output-format stream-json``. +""" + +import io +import json + +import pytest + +from EvoScientist.stream.json_sink import stream_json, write_events_as_json + + +async def _agen(items): + """Wrap a list as an async generator (a real event source, not a mock).""" + for item in items: + yield item + + +def test_writes_each_event_as_one_jsonl_line(run_async): + """Each event dict is serialized to exactly one JSON line, in order.""" + events = [ + {"type": "thinking", "content": "hmm", "id": 0}, + {"type": "text", "content": "hello"}, + { + "type": "tool_call", + "name": "write_file", + "args": {"path": "a.md"}, + "id": "t1", + }, + {"type": "done", "content": "hello", "response": "hello"}, + ] + out = io.StringIO() + + run_async(write_events_as_json(_agen(events), out)) + + lines = out.getvalue().splitlines() + assert len(lines) == len(events) + parsed = [json.loads(line) for line in lines] + assert [e["type"] for e in parsed] == ["thinking", "text", "tool_call", "done"] + assert parsed[2]["args"] == {"path": "a.md"} + + +def test_returns_final_response_from_done_event(run_async): + """The sink returns the response text carried by the terminal `done` event.""" + events = [ + {"type": "text", "content": "partial"}, + {"type": "done", "content": "the answer", "response": "the answer"}, + ] + out = io.StringIO() + + result = run_async(write_events_as_json(_agen(events), out)) + + assert result == "the answer" + + +def test_non_serializable_arg_does_not_crash_the_stream(run_async): + """A non-JSON-serializable value degrades to its str form instead of raising.""" + + class Weird: + """A value json.dumps cannot serialize, used to exercise the str fallback.""" + + def __str__(self): + """Return a sentinel so the fallback is observable in the output.""" + return "WEIRD" + + events = [ + {"type": "tool_call", "name": "x", "args": {"obj": Weird()}, "id": "t1"}, + {"type": "done", "content": "", "response": ""}, + ] + out = io.StringIO() + + run_async(write_events_as_json(_agen(events), out)) + + lines = out.getvalue().splitlines() + # Both lines must be valid JSON; the non-serializable value falls back to str. + first = json.loads(lines[0]) + assert first["args"]["obj"] == "WEIRD" + + +def test_stream_json_sources_events_from_gateway(run_async): + """stream_json pulls events from gateway.stream_events(request) and serializes + them — it does not reach past the gateway abstraction.""" + seen: dict[str, object] = {} + events = [ + {"type": "text", "content": "hi"}, + {"type": "done", "content": "hi", "response": "hi"}, + ] + + class _FakeGateway: + """A gateway stub whose stream_events yields a fixed event sequence.""" + + def stream_events(self, request): + """Record the request and return the canned event stream.""" + seen["request"] = request + return _agen(events) + + out = io.StringIO() + result = run_async(stream_json(_FakeGateway(), object(), out=out)) + + assert result == "hi" + assert "request" in seen # the request was forwarded to the gateway + types = [json.loads(line)["type"] for line in out.getvalue().splitlines()] + assert types == ["text", "done"] + + +def test_stream_json_propagates_gateway_errors(run_async): + """An error from the gateway stream propagates out of stream_json so the CLI + dispatch can turn it into a clean exit.""" + + async def _boom(): + """Yield one event, then fail like a mid-run graph error.""" + yield {"type": "text", "content": "partial"} + raise RuntimeError("boom") + + class _FakeGateway: + """A gateway stub whose stream raises partway through.""" + + def stream_events(self, request): + """Return a stream that fails after the first event.""" + return _boom() + + out = io.StringIO() + with pytest.raises(RuntimeError, match="boom"): + run_async(stream_json(_FakeGateway(), object(), out=out))