diff --git a/src/agentex/lib/cli/handlers/run_handlers.py b/src/agentex/lib/cli/handlers/run_handlers.py index 18ee84e93..40d509462 100644 --- a/src/agentex/lib/cli/handlers/run_handlers.py +++ b/src/agentex/lib/cli/handlers/run_handlers.py @@ -5,6 +5,7 @@ import asyncio from pathlib import Path +from rich.text import Text from rich.panel import Panel from rich.console import Console @@ -291,7 +292,7 @@ async def stream_process_output(process: asyncio.subprocess.Process, prefix: str continue if decoded_line: # Only print non-empty lines - console.print(f"[dim]{prefix}:[/dim] {decoded_line}") + console.print(Text.assemble((f"{prefix}:", "dim"), " ", decoded_line)) except Exception as e: # The escalation path, including for the re-raise above. Anything reaching # here ends the loop, so the child is now at risk of blocking on a full pipe. diff --git a/tests/lib/cli/test_run_handlers_streaming.py b/tests/lib/cli/test_run_handlers_streaming.py index 8f0ab13b5..5346a2481 100644 --- a/tests/lib/cli/test_run_handlers_streaming.py +++ b/tests/lib/cli/test_run_handlers_streaming.py @@ -8,6 +8,7 @@ from __future__ import annotations +import re import sys import asyncio from typing import Any @@ -40,6 +41,22 @@ print("done") """ +MARKUP_LINES = ( + "[INST] hi [/INST]", + "loaded [/etc/hosts]", + "[info] starting", + "deps: [pkg==1.0]", + "[bold]x[/bold]", +) + +MARKUP_CHILD_SCRIPT = """ +for line in {lines!r}: + print(line) +for i in range(10000): + print("after", i, "y" * 60) +print("done") +""" + async def _drain(limit: int, oversized: int) -> int | None: """Run the child under stream_process_output. None means it never exited.""" @@ -61,6 +78,26 @@ async def _drain(limit: int, oversized: int) -> int | None: return process.returncode +async def _drain_markup_lines(limit: int) -> int | None: + """Run a child whose output looks like rich markup, then floods the pipe.""" + process = await asyncio.create_subprocess_exec( + sys.executable, + "-c", + MARKUP_CHILD_SCRIPT.format(lines=MARKUP_LINES), + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.STDOUT, + limit=limit, + ) + streamer = asyncio.create_task(stream_process_output(process, "TEST")) + try: + await asyncio.wait_for(asyncio.gather(streamer, process.wait()), timeout=60) + except TimeoutError: + process.kill() + await process.wait() + return None + return process.returncode + + async def test_oversized_line_is_skipped_without_stalling_the_child( capsys: pytest.CaptureFixture[str], ) -> None: @@ -98,6 +135,35 @@ async def test_large_line_within_the_limit_is_streamed_in_full( assert out.count(MARKER) == oversized, "the large line was dropped rather than streamed" +async def test_markup_like_output_does_not_stop_the_reader( + capsys: pytest.CaptureFixture[str], +) -> None: + """An unmatched closing tag in the child's output must not end the loop. + + Agents dump prompts and paths, so "[INST] hi [/INST]" is ordinary output. + Parsed as markup it raises MarkupError, which reaches the outer handler and + leaves nothing draining the pipe. The child reaching exit is the assertion, + so the child writes far more than the buffers behind a 64 KiB limit hold. + """ + returncode = await _drain_markup_lines(limit=64 * 1024) + out = capsys.readouterr().out + + assert returncode == 0, "child did not exit: the reader stopped draining its pipe" + assert "done" in out + + +async def test_markup_like_output_is_printed_verbatim( + capsys: pytest.CaptureFixture[str], +) -> None: + """Bracketed words are the child's content, not styling, so none are eaten.""" + returncode = await _drain_markup_lines(limit=SUBPROCESS_STREAM_LIMIT) + out = re.sub(r"\x1b\[[0-9;]*m", "", capsys.readouterr().out) + + assert returncode == 0 + for line in MARKUP_LINES: + assert f"TEST: {line}" in out, f"rich consumed part of {line!r}" + + class _AlwaysFailingReader: """A reader whose readline() raises without consuming anything.