1010
1111import sys
1212import asyncio
13+ from typing import Any
1314
15+ import pytest
16+
17+ from agentex .lib .cli .handlers import run_handlers
1418from agentex .lib .cli .handlers .run_handlers import (
1519 SUBPROCESS_STREAM_LIMIT ,
20+ start_acp_server ,
21+ start_temporal_worker ,
1622 stream_process_output ,
1723)
1824
19- # Emits a line over the reader's limit, then enough further output to more than
20- # fill a 64 KiB pipe. If the reader stops draining, the child cannot finish its
21- # writes and never exits.
25+ # Emits a line of MARKER over the reader's limit, then enough further output to
26+ # more than fill a 64 KiB pipe. If the reader stops draining, the child cannot
27+ # finish its writes and never exits.
28+ MARKER = "X"
29+
2230CHILD_SCRIPT = """
23- import sys
2431print("before")
25- print("X " * {oversized})
32+ print("{marker} " * {oversized})
2633for i in range(2000):
2734 print("after", i, "y" * 60)
2835print("done")
2936"""
3037
3138
3239async def _drain (limit : int , oversized : int ) -> int | None :
40+ """Run the child under stream_process_output. None means it never exited."""
3341 process = await asyncio .create_subprocess_exec (
3442 sys .executable ,
3543 "-c" ,
36- CHILD_SCRIPT .format (oversized = oversized ),
44+ CHILD_SCRIPT .format (marker = MARKER , oversized = oversized ),
3745 stdout = asyncio .subprocess .PIPE ,
3846 stderr = asyncio .subprocess .STDOUT ,
3947 limit = limit ,
@@ -48,25 +56,57 @@ async def _drain(limit: int, oversized: int) -> int | None:
4856 return process .returncode
4957
5058
51- async def test_oversized_line_is_skipped_without_stalling_the_child () -> None :
52- """A line past the reader's limit is dropped and streaming continues.
59+ async def test_oversized_line_is_skipped_without_stalling_the_child (
60+ capsys : pytest .CaptureFixture [str ],
61+ ) -> None :
62+ """A line past the reader's limit is dropped, and streaming continues.
5363
5464 Before this was handled per line, readline() raised, the loop exited, and the
5565 child deadlocked on a full pipe. The child reaching exit is the assertion.
5666 """
5767 limit = 64 * 1024
58- returncode = await _drain (limit = limit , oversized = limit + 16_000 )
68+ oversized = limit + 16_000
69+
70+ returncode = await _drain (limit = limit , oversized = oversized )
71+ out = capsys .readouterr ().out
5972
6073 assert returncode == 0 , "child did not exit: the reader stopped draining its pipe"
74+ # The offending line is gone, but everything after it still streamed.
75+ assert out .count (MARKER ) == 0
76+ assert "done" in out
77+
6178
79+ async def test_large_line_within_the_limit_is_streamed_in_full (
80+ capsys : pytest .CaptureFixture [str ],
81+ ) -> None :
82+ """A line over asyncio's 64 KiB default still reaches the console under our limit.
6283
63- async def test_large_line_within_limit_is_streamed () -> None :
64- """A line larger than asyncio's 64 KiB default still streams under our limit."""
65- returncode = await _drain (limit = SUBPROCESS_STREAM_LIMIT , oversized = 82_000 )
84+ Counts marker characters rather than matching the line, because rich wraps
85+ long output across terminal-width lines.
86+ """
87+ oversized = 82_000
88+
89+ returncode = await _drain (limit = SUBPROCESS_STREAM_LIMIT , oversized = oversized )
90+ out = capsys .readouterr ().out
6691
6792 assert returncode == 0
93+ assert out .count (MARKER ) == oversized , "the large line was dropped rather than streamed"
94+
95+
96+ async def test_agent_subprocesses_are_spawned_with_the_larger_limit (
97+ monkeypatch : pytest .MonkeyPatch , tmp_path : Any
98+ ) -> None :
99+ """The helpers must pass limit=, or large lines are dropped in production."""
100+ seen : list [int | None ] = []
101+
102+ async def fake_exec (* _args : Any , ** kwargs : Any ) -> None :
103+ seen .append (kwargs .get ("limit" ))
104+
105+ monkeypatch .setattr (asyncio , "create_subprocess_exec" , fake_exec )
106+ monkeypatch .setattr (run_handlers , "calculate_uvicorn_target_for_local" , lambda * _ : "project.acp" )
68107
108+ await start_acp_server (tmp_path / "acp.py" , 8000 , {}, tmp_path )
109+ await start_temporal_worker (tmp_path / "run_worker.py" , {}, tmp_path )
69110
70- async def test_subprocess_stream_limit_exceeds_asyncio_default () -> None :
71- """The whole point of the constant: asyncio's default is what breaks readline()."""
72- assert SUBPROCESS_STREAM_LIMIT > 64 * 1024
111+ assert seen == [SUBPROCESS_STREAM_LIMIT , SUBPROCESS_STREAM_LIMIT ]
112+ assert SUBPROCESS_STREAM_LIMIT > 64 * 1024 , "asyncio's default is what breaks readline()"
0 commit comments