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
18 changes: 17 additions & 1 deletion examples/tutorials/00_sync/060_claude_code/project/acp.py
Original file line number Diff line number Diff line change
Expand Up @@ -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", ""),
Expand All @@ -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",
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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", ""),
Expand All @@ -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",
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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"


Expand Down Expand Up @@ -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",
Expand All @@ -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
Expand Down Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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", ""),
Expand All @@ -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",
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down
18 changes: 17 additions & 1 deletion src/agentex/lib/cli/templates/sync-claude-code/project/acp.py.j2
Original file line number Diff line number Diff line change
Expand Up @@ -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", ""),
Expand All @@ -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",
Expand All @@ -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
Expand Down Expand Up @@ -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:
Comment thread
greptile-apps[bot] marked this conversation as resolved.
proc.kill()
except ProcessLookupError:
pass
await proc.wait()


@acp.on_message_send
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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"


Expand Down Expand Up @@ -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",
Expand All @@ -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
Expand Down Expand Up @@ -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)
Expand Down
88 changes: 88 additions & 0 deletions tests/lib/cli/test_claude_code_template_subprocess.py
Original file line number Diff line number Diff line change
@@ -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