Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion src/agentex/lib/cli/handlers/run_handlers.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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.
Expand Down
66 changes: 66 additions & 0 deletions tests/lib/cli/test_run_handlers_streaming.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@

from __future__ import annotations

import re
import sys
import asyncio
from typing import Any
Expand Down Expand Up @@ -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."""
Expand All @@ -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:
Expand Down Expand Up @@ -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.

Expand Down