From a35652aceea2a5876f458b5d177e06c046bc686a Mon Sep 17 00:00:00 2001 From: Michael Xu Date: Thu, 1 Oct 2026 12:51:17 -0500 Subject: [PATCH 1/2] fix(templates): read long stream-json lines and bound Claude Code shutdown The Claude Code scaffolds (sync, async and Temporal) and their tutorial copies read `claude -p --output-format stream-json` stdout through asyncio's StreamReader with the default 64 KiB line limit. Claude Code writes one JSON line per event, so a tool_result that echoes a large file read is a single line well past 64 KiB; readline() then raises "Separator is found, but chunk is longer than limit" and the turn aborts. Read with an 8 MiB limit, the same value `agentex agents run` uses for its child processes. Their cleanup also sent SIGTERM and then awaited proc.wait() with no bound, so a CLI that ignores or delays SIGTERM hangs request or activity cancellation. Wait at most 5 s, then SIGKILL. Verified with a fake `claude` on PATH, against the rendered sync and Temporal templates: - a 100,000-character line: main raises the ValueError, this branch reads it whole; - a CLI that ignores SIGTERM: main's aclose() is still hanging after 20 s, this branch returns after 5 s. tests/lib/cli/test_init_templates.py passes, and the three tutorials' test_agent_offline.py suites pass (7, 5 and 5 tests). --- .../00_sync/060_claude_code/project/acp.py | 18 +++++++++++++++++- .../00_base/130_claude_code/project/acp.py | 18 +++++++++++++++++- .../140_claude_code/project/activities.py | 18 +++++++++++++++++- .../default-claude-code/project/acp.py.j2 | 18 +++++++++++++++++- .../sync-claude-code/project/acp.py.j2 | 18 +++++++++++++++++- .../project/activities.py.j2 | 18 +++++++++++++++++- 6 files changed, 102 insertions(+), 6 deletions(-) diff --git a/examples/tutorials/00_sync/060_claude_code/project/acp.py b/examples/tutorials/00_sync/060_claude_code/project/acp.py index aad53801a..184abd677 100644 --- a/examples/tutorials/00_sync/060_claude_code/project/acp.py +++ b/examples/tutorials/00_sync/060_claude_code/project/acp.py @@ -35,6 +35,9 @@ logger = make_logger(__name__) +STDOUT_LINE_LIMIT = 8 * 1024 * 1024 +TERMINATE_TIMEOUT_SECONDS = 5.0 + add_tracing_processor_config( SGPTracingProcessorConfig( sgp_api_key=os.environ.get("SGP_API_KEY", ""), @@ -51,6 +54,11 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]: This is a seam: tests replace it with a fake async iterator of pre-recorded lines so no real CLI invocation is needed offline. + + Lines up to ``STDOUT_LINE_LIMIT`` are read whole: a ``tool_result`` that + echoes a file is a single stream-json line, often past asyncio's 64 KiB + default. If the consumer stops early, the CLI gets SIGTERM and, after + ``TERMINATE_TIMEOUT_SECONDS``, SIGKILL, so cancellation cannot hang. """ proc = await asyncio.create_subprocess_exec( "claude", @@ -61,6 +69,7 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]: stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, + limit=STDOUT_LINE_LIMIT, ) assert proc.stdout is not None assert proc.stdin is not None @@ -108,7 +117,14 @@ async def _drain_stderr() -> None: proc.terminate() except ProcessLookupError: pass - await proc.wait() + try: + await asyncio.wait_for(proc.wait(), timeout=TERMINATE_TIMEOUT_SECONDS) + except TimeoutError: + try: + proc.kill() + except ProcessLookupError: + pass + await proc.wait() @acp.on_message_send diff --git a/examples/tutorials/10_async/00_base/130_claude_code/project/acp.py b/examples/tutorials/10_async/00_base/130_claude_code/project/acp.py index b6681f6a8..eaed174e6 100644 --- a/examples/tutorials/10_async/00_base/130_claude_code/project/acp.py +++ b/examples/tutorials/10_async/00_base/130_claude_code/project/acp.py @@ -33,6 +33,9 @@ logger = make_logger(__name__) +STDOUT_LINE_LIMIT = 8 * 1024 * 1024 +TERMINATE_TIMEOUT_SECONDS = 5.0 + add_tracing_processor_config( SGPTracingProcessorConfig( sgp_api_key=os.environ.get("SGP_API_KEY", ""), @@ -52,6 +55,11 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]: Injectable seam: tests monkeypatch this with a fake async iterator of pre-recorded lines so no real CLI invocation is needed offline. + + Lines up to ``STDOUT_LINE_LIMIT`` are read whole: a ``tool_result`` that + echoes a file is a single stream-json line, often past asyncio's 64 KiB + default. If the consumer stops early, the CLI gets SIGTERM and, after + ``TERMINATE_TIMEOUT_SECONDS``, SIGKILL, so cancellation cannot hang. """ proc = await asyncio.create_subprocess_exec( "claude", @@ -62,6 +70,7 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]: stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, + limit=STDOUT_LINE_LIMIT, ) assert proc.stdout is not None assert proc.stdin is not None @@ -109,7 +118,14 @@ async def _drain_stderr() -> None: proc.terminate() except ProcessLookupError: pass - await proc.wait() + try: + await asyncio.wait_for(proc.wait(), timeout=TERMINATE_TIMEOUT_SECONDS) + except TimeoutError: + try: + proc.kill() + except ProcessLookupError: + pass + await proc.wait() @acp.on_task_create diff --git a/examples/tutorials/10_async/10_temporal/140_claude_code/project/activities.py b/examples/tutorials/10_async/10_temporal/140_claude_code/project/activities.py index dcba0f9a7..8092ca977 100644 --- a/examples/tutorials/10_async/10_temporal/140_claude_code/project/activities.py +++ b/examples/tutorials/10_async/10_temporal/140_claude_code/project/activities.py @@ -27,6 +27,9 @@ logger = make_logger(__name__) +STDOUT_LINE_LIMIT = 8 * 1024 * 1024 +TERMINATE_TIMEOUT_SECONDS = 5.0 + RUN_CLAUDE_CODE_TURN_ACTIVITY = "run_claude_code_turn" @@ -56,6 +59,11 @@ async def _spawn_claude(prompt: str, session_id: str | None = None) -> AsyncIter Injectable seam: tests monkeypatch this with a fake async iterator so no real CLI invocation is needed offline. + + Lines up to ``STDOUT_LINE_LIMIT`` are read whole: a ``tool_result`` that + echoes a file is a single stream-json line, often past asyncio's 64 KiB + default. If the consumer stops early, the CLI gets SIGTERM and, after + ``TERMINATE_TIMEOUT_SECONDS``, SIGKILL, so cancellation cannot hang. """ cmd = [ "claude", @@ -72,6 +80,7 @@ async def _spawn_claude(prompt: str, session_id: str | None = None) -> AsyncIter stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, + limit=STDOUT_LINE_LIMIT, ) assert proc.stdout is not None assert proc.stdin is not None @@ -119,7 +128,14 @@ async def _drain_stderr() -> None: proc.terminate() except ProcessLookupError: pass - await proc.wait() + try: + await asyncio.wait_for(proc.wait(), timeout=TERMINATE_TIMEOUT_SECONDS) + except TimeoutError: + try: + proc.kill() + except ProcessLookupError: + pass + await proc.wait() @activity.defn(name=RUN_CLAUDE_CODE_TURN_ACTIVITY) diff --git a/src/agentex/lib/cli/templates/default-claude-code/project/acp.py.j2 b/src/agentex/lib/cli/templates/default-claude-code/project/acp.py.j2 index 42512c601..854e3fcfd 100644 --- a/src/agentex/lib/cli/templates/default-claude-code/project/acp.py.j2 +++ b/src/agentex/lib/cli/templates/default-claude-code/project/acp.py.j2 @@ -33,6 +33,9 @@ from agentex.lib.core.tracing.tracing_processor_manager import add_tracing_proce logger = make_logger(__name__) +STDOUT_LINE_LIMIT = 8 * 1024 * 1024 +TERMINATE_TIMEOUT_SECONDS = 5.0 + add_tracing_processor_config( SGPTracingProcessorConfig( sgp_api_key=os.environ.get("SGP_API_KEY", ""), @@ -52,6 +55,11 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]: Injectable seam: tests can monkeypatch this with a fake async iterator of pre-recorded lines so no real CLI invocation is needed offline. + + Lines up to ``STDOUT_LINE_LIMIT`` are read whole: a ``tool_result`` that + echoes a file is a single stream-json line, often past asyncio's 64 KiB + default. If the consumer stops early, the CLI gets SIGTERM and, after + ``TERMINATE_TIMEOUT_SECONDS``, SIGKILL, so cancellation cannot hang. """ proc = await asyncio.create_subprocess_exec( "claude", @@ -62,6 +70,7 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]: stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, + limit=STDOUT_LINE_LIMIT, ) assert proc.stdout is not None assert proc.stdin is not None @@ -123,7 +132,14 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]: proc.terminate() except ProcessLookupError: pass - await proc.wait() + try: + await asyncio.wait_for(proc.wait(), timeout=TERMINATE_TIMEOUT_SECONDS) + except TimeoutError: + try: + proc.kill() + except ProcessLookupError: + pass + await proc.wait() @acp.on_task_create diff --git a/src/agentex/lib/cli/templates/sync-claude-code/project/acp.py.j2 b/src/agentex/lib/cli/templates/sync-claude-code/project/acp.py.j2 index 33a89a51e..d0b512f0c 100644 --- a/src/agentex/lib/cli/templates/sync-claude-code/project/acp.py.j2 +++ b/src/agentex/lib/cli/templates/sync-claude-code/project/acp.py.j2 @@ -35,6 +35,9 @@ from agentex.lib.core.tracing.tracing_processor_manager import add_tracing_proce logger = make_logger(__name__) +STDOUT_LINE_LIMIT = 8 * 1024 * 1024 +TERMINATE_TIMEOUT_SECONDS = 5.0 + add_tracing_processor_config( SGPTracingProcessorConfig( sgp_api_key=os.environ.get("SGP_API_KEY", ""), @@ -51,6 +54,11 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]: This is a seam: tests can replace it with a fake async iterator of pre-recorded lines so no real CLI invocation is needed offline. + + Lines up to ``STDOUT_LINE_LIMIT`` are read whole: a ``tool_result`` that + echoes a file is a single stream-json line, often past asyncio's 64 KiB + default. If the consumer stops early, the CLI gets SIGTERM and, after + ``TERMINATE_TIMEOUT_SECONDS``, SIGKILL, so cancellation cannot hang. """ proc = await asyncio.create_subprocess_exec( "claude", @@ -61,6 +69,7 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]: stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, + limit=STDOUT_LINE_LIMIT, ) assert proc.stdout is not None assert proc.stdin is not None @@ -122,7 +131,14 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]: proc.terminate() except ProcessLookupError: pass - await proc.wait() + try: + await asyncio.wait_for(proc.wait(), timeout=TERMINATE_TIMEOUT_SECONDS) + except TimeoutError: + try: + proc.kill() + except ProcessLookupError: + pass + await proc.wait() @acp.on_message_send diff --git a/src/agentex/lib/cli/templates/temporal-claude-code/project/activities.py.j2 b/src/agentex/lib/cli/templates/temporal-claude-code/project/activities.py.j2 index 94055c7df..2b667f143 100644 --- a/src/agentex/lib/cli/templates/temporal-claude-code/project/activities.py.j2 +++ b/src/agentex/lib/cli/templates/temporal-claude-code/project/activities.py.j2 @@ -28,6 +28,9 @@ from agentex.lib.utils.model_utils import BaseModel logger = make_logger(__name__) +STDOUT_LINE_LIMIT = 8 * 1024 * 1024 +TERMINATE_TIMEOUT_SECONDS = 5.0 + RUN_CLAUDE_CODE_TURN_ACTIVITY = "run_claude_code_turn" @@ -57,6 +60,11 @@ async def _spawn_claude(prompt: str, session_id: str | None = None) -> AsyncIter Injectable seam: tests can monkeypatch this with a fake async iterator so no real CLI invocation is needed offline. + + Lines up to ``STDOUT_LINE_LIMIT`` are read whole: a ``tool_result`` that + echoes a file is a single stream-json line, often past asyncio's 64 KiB + default. If the consumer stops early, the CLI gets SIGTERM and, after + ``TERMINATE_TIMEOUT_SECONDS``, SIGKILL, so cancellation cannot hang. """ cmd = [ "claude", @@ -73,6 +81,7 @@ async def _spawn_claude(prompt: str, session_id: str | None = None) -> AsyncIter stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, + limit=STDOUT_LINE_LIMIT, ) assert proc.stdout is not None assert proc.stdin is not None @@ -135,7 +144,14 @@ async def _spawn_claude(prompt: str, session_id: str | None = None) -> AsyncIter proc.terminate() except ProcessLookupError: pass - await proc.wait() + try: + await asyncio.wait_for(proc.wait(), timeout=TERMINATE_TIMEOUT_SECONDS) + except TimeoutError: + try: + proc.kill() + except ProcessLookupError: + pass + await proc.wait() @activity.defn(name=RUN_CLAUDE_CODE_TURN_ACTIVITY) From 1c487a260f708afdae59b0dfb0fd03380aeb6472 Mon Sep 17 00:00:00 2001 From: Michael Xu Date: Fri, 2 Oct 2026 15:52:48 -0500 Subject: [PATCH 2/2] test(templates): run the Claude Code scaffold against a real child process Render the sync Claude Code template, put a fake `claude` on PATH that prints a 100,000-character stream-json line and ignores SIGTERM, and drive _spawn_claude(): the long line must arrive whole and aclose() must return once the (shortened) grace period ends. Against main's template the test fails with "Separator is found, but chunk is longer than limit" after the unbounded wait sits out the child's 60 s sleep; here it passes in under 4 s. --- .../test_claude_code_template_subprocess.py | 88 +++++++++++++++++++ 1 file changed, 88 insertions(+) create mode 100644 tests/lib/cli/test_claude_code_template_subprocess.py diff --git a/tests/lib/cli/test_claude_code_template_subprocess.py b/tests/lib/cli/test_claude_code_template_subprocess.py new file mode 100644 index 000000000..234e15dc9 --- /dev/null +++ b/tests/lib/cli/test_claude_code_template_subprocess.py @@ -0,0 +1,88 @@ +"""The rendered Claude Code scaffold against a real child process. + +A fake ``claude`` on PATH stands in for the CLI, so these run offline but +exercise the template's actual subprocess handling: the stdout line limit +and the SIGTERM-then-SIGKILL shutdown. +""" + +from __future__ import annotations + +import sys +import time +import types +import asyncio +import importlib.util +from pathlib import Path + +import pytest + +from agentex.lib.core.tracing import tracing_processor_manager +from agentex.lib.cli.commands.init import TemplateType, get_project_context, create_project_structure + +pytestmark = pytest.mark.skipif(sys.platform == "win32", reason="needs POSIX signals") + +LONG_LINE_CHARS = 100_000 +FAST_TERMINATE_SECONDS = 0.5 +SHUTDOWN_BUDGET_SECONDS = 10.0 + +FAKE_CLAUDE = """#!{python} +import json, signal, sys, time +signal.signal(signal.SIGTERM, signal.SIG_IGN) +sys.stdin.read() +print(json.dumps({{"type": "system", "subtype": "init", "session_id": "s"}}), flush=True) +print(json.dumps({{"type": "user", "message": {{"content": [{{"type": "tool_result", "content": "x" * {chars}}}]}}}}), flush=True) +time.sleep(60) +""" + + +def _render_sync_template(tmp_path: Path) -> Path: + template = TemplateType.SYNC_CLAUDE_CODE + answers = { + "template_type": template, + "project_path": str(tmp_path), + "agent_name": "cc-bot", + "agent_directory_name": "cc-bot", + "description": "test", + "use_uv": True, + } + context = get_project_context(answers, tmp_path, Path("../../")) + context["template_type"] = template.value + context["use_uv"] = True + create_project_structure(tmp_path, context, template, True) + return tmp_path / "cc_bot" / "project" / "acp.py" + + +@pytest.fixture +def scaffold(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> types.ModuleType: + bin_dir = tmp_path / "bin" + bin_dir.mkdir() + fake = bin_dir / "claude" + fake.write_text(FAKE_CLAUDE.format(python=sys.executable, chars=LONG_LINE_CHARS)) + fake.chmod(0o755) + monkeypatch.setenv("PATH", f"{bin_dir}:{Path(sys.executable).parent}") + monkeypatch.setenv("AGENT_NAME", "cc-bot") + monkeypatch.setenv("ACP_URL", "http://localhost:8000") + monkeypatch.setattr(tracing_processor_manager, "add_tracing_processor_config", lambda _config: None) + + spec = importlib.util.spec_from_file_location("cc_bot_acp", _render_sync_template(tmp_path)) + assert spec is not None and spec.loader is not None + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + monkeypatch.setattr(module, "TERMINATE_TIMEOUT_SECONDS", FAST_TERMINATE_SECONDS, raising=False) + return module + + +async def test_long_stream_json_line_is_read_whole_and_shutdown_is_bounded(scaffold: types.ModuleType) -> None: + """A tool_result past asyncio's 64 KiB default arrives intact, and a CLI that + ignores SIGTERM is killed after the grace period instead of hanging cleanup.""" + lines = scaffold._spawn_claude("hi") + first = await lines.__anext__() + second = await lines.__anext__() + + started = time.monotonic() + await asyncio.wait_for(lines.aclose(), SHUTDOWN_BUDGET_SECONDS) + elapsed = time.monotonic() - started + + assert '"init"' in first + assert second.count("x") == LONG_LINE_CHARS + assert elapsed >= FAST_TERMINATE_SECONDS