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) 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