690b903f85
* chore: add pytest-asyncio in auto mode * test: migrate channel and stream tests to native async Convert run_async() wrapper tests to plain 'async def test_*' under pytest-asyncio auto mode. collect_events() in stream_v3_fakes becomes a coroutine awaited at every call site. * test: migrate command and model/middleware tests to native async Convert run_async() wrappers (import, alias, and fixture forms) to plain 'async def test_*'. Multi-call tests merge onto one loop as sequential awaits; none asserted on loop identity. * test: migrate TUI, notifier, gateway, and session tests to native async TUI/notifier/gateway files convert run_async wrappers to plain async tests. test_sessions.py's unittest.TestCase classes move to unittest.IsolatedAsyncioTestCase (pytest-asyncio does not await async methods on plain TestCase; converting blindly would have made ~70 tests silently vacuous). Its setUpClass keeps a one-shot asyncio.run() since IsolatedAsyncioTestCase has no async class-level hook. TestLoadingWidget in test_tui_widgets.py drops its TestCase base for the same reason. * test: replace direct asyncio.run() calls with native async tests Convert tests that called asyncio.run() (directly or via a local _run helper) to plain 'async def test_*'; delete the local helpers. * test: drop undeclared anyio markers and delete run_async helper The @pytest.mark.anyio tests relied on anyio being a transitive dep of httpx; auto-mode pytest-asyncio collects them natively. run_async() and its fixture are unreferenced after the migration, so remove them — pytest-asyncio's per-test loop teardown covers the pending-task cancellation the helper existed for (verified: full suite runs with no 'Event loop is closed' errors or destroyed-task warnings).
128 lines
4.2 KiB
Python
128 lines
4.2 KiB
Python
"""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
|
|
|
|
|
|
async def test_writes_each_event_as_one_jsonl_line():
|
|
"""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()
|
|
|
|
await 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"}
|
|
|
|
|
|
async def test_returns_final_response_from_done_event():
|
|
"""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 = await write_events_as_json(_agen(events), out)
|
|
|
|
assert result == "the answer"
|
|
|
|
|
|
async def test_non_serializable_arg_does_not_crash_the_stream():
|
|
"""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()
|
|
|
|
await 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"
|
|
|
|
|
|
async def test_stream_json_sources_events_from_gateway():
|
|
"""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 = await 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"]
|
|
|
|
|
|
async def test_stream_json_propagates_gateway_errors():
|
|
"""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"):
|
|
await stream_json(_FakeGateway(), object(), out=out)
|