Skip to content

Commit a35652a

Browse files
committed
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).
1 parent 6d68f3a commit a35652a

6 files changed

Lines changed: 102 additions & 6 deletions

File tree

‎examples/tutorials/00_sync/060_claude_code/project/acp.py‎

Lines changed: 17 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,9 @@
3535

3636
logger = make_logger(__name__)
3737

38+
STDOUT_LINE_LIMIT = 8 * 1024 * 1024
39+
TERMINATE_TIMEOUT_SECONDS = 5.0
40+
3841
add_tracing_processor_config(
3942
SGPTracingProcessorConfig(
4043
sgp_api_key=os.environ.get("SGP_API_KEY", ""),
@@ -51,6 +54,11 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]:
5154
5255
This is a seam: tests replace it with a fake async iterator of
5356
pre-recorded lines so no real CLI invocation is needed offline.
57+
58+
Lines up to ``STDOUT_LINE_LIMIT`` are read whole: a ``tool_result`` that
59+
echoes a file is a single stream-json line, often past asyncio's 64 KiB
60+
default. If the consumer stops early, the CLI gets SIGTERM and, after
61+
``TERMINATE_TIMEOUT_SECONDS``, SIGKILL, so cancellation cannot hang.
5462
"""
5563
proc = await asyncio.create_subprocess_exec(
5664
"claude",
@@ -61,6 +69,7 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]:
6169
stdin=asyncio.subprocess.PIPE,
6270
stdout=asyncio.subprocess.PIPE,
6371
stderr=asyncio.subprocess.PIPE,
72+
limit=STDOUT_LINE_LIMIT,
6473
)
6574
assert proc.stdout is not None
6675
assert proc.stdin is not None
@@ -108,7 +117,14 @@ async def _drain_stderr() -> None:
108117
proc.terminate()
109118
except ProcessLookupError:
110119
pass
111-
await proc.wait()
120+
try:
121+
await asyncio.wait_for(proc.wait(), timeout=TERMINATE_TIMEOUT_SECONDS)
122+
except TimeoutError:
123+
try:
124+
proc.kill()
125+
except ProcessLookupError:
126+
pass
127+
await proc.wait()
112128

113129

114130
@acp.on_message_send

‎examples/tutorials/10_async/00_base/130_claude_code/project/acp.py‎

Lines changed: 17 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,9 @@
3333

3434
logger = make_logger(__name__)
3535

36+
STDOUT_LINE_LIMIT = 8 * 1024 * 1024
37+
TERMINATE_TIMEOUT_SECONDS = 5.0
38+
3639
add_tracing_processor_config(
3740
SGPTracingProcessorConfig(
3841
sgp_api_key=os.environ.get("SGP_API_KEY", ""),
@@ -52,6 +55,11 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]:
5255
5356
Injectable seam: tests monkeypatch this with a fake async iterator of
5457
pre-recorded lines so no real CLI invocation is needed offline.
58+
59+
Lines up to ``STDOUT_LINE_LIMIT`` are read whole: a ``tool_result`` that
60+
echoes a file is a single stream-json line, often past asyncio's 64 KiB
61+
default. If the consumer stops early, the CLI gets SIGTERM and, after
62+
``TERMINATE_TIMEOUT_SECONDS``, SIGKILL, so cancellation cannot hang.
5563
"""
5664
proc = await asyncio.create_subprocess_exec(
5765
"claude",
@@ -62,6 +70,7 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]:
6270
stdin=asyncio.subprocess.PIPE,
6371
stdout=asyncio.subprocess.PIPE,
6472
stderr=asyncio.subprocess.PIPE,
73+
limit=STDOUT_LINE_LIMIT,
6574
)
6675
assert proc.stdout is not None
6776
assert proc.stdin is not None
@@ -109,7 +118,14 @@ async def _drain_stderr() -> None:
109118
proc.terminate()
110119
except ProcessLookupError:
111120
pass
112-
await proc.wait()
121+
try:
122+
await asyncio.wait_for(proc.wait(), timeout=TERMINATE_TIMEOUT_SECONDS)
123+
except TimeoutError:
124+
try:
125+
proc.kill()
126+
except ProcessLookupError:
127+
pass
128+
await proc.wait()
113129

114130

115131
@acp.on_task_create

‎examples/tutorials/10_async/10_temporal/140_claude_code/project/activities.py‎

Lines changed: 17 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,9 @@
2727

2828
logger = make_logger(__name__)
2929

30+
STDOUT_LINE_LIMIT = 8 * 1024 * 1024
31+
TERMINATE_TIMEOUT_SECONDS = 5.0
32+
3033
RUN_CLAUDE_CODE_TURN_ACTIVITY = "run_claude_code_turn"
3134

3235

@@ -56,6 +59,11 @@ async def _spawn_claude(prompt: str, session_id: str | None = None) -> AsyncIter
5659
5760
Injectable seam: tests monkeypatch this with a fake async iterator so no
5861
real CLI invocation is needed offline.
62+
63+
Lines up to ``STDOUT_LINE_LIMIT`` are read whole: a ``tool_result`` that
64+
echoes a file is a single stream-json line, often past asyncio's 64 KiB
65+
default. If the consumer stops early, the CLI gets SIGTERM and, after
66+
``TERMINATE_TIMEOUT_SECONDS``, SIGKILL, so cancellation cannot hang.
5967
"""
6068
cmd = [
6169
"claude",
@@ -72,6 +80,7 @@ async def _spawn_claude(prompt: str, session_id: str | None = None) -> AsyncIter
7280
stdin=asyncio.subprocess.PIPE,
7381
stdout=asyncio.subprocess.PIPE,
7482
stderr=asyncio.subprocess.PIPE,
83+
limit=STDOUT_LINE_LIMIT,
7584
)
7685
assert proc.stdout is not None
7786
assert proc.stdin is not None
@@ -119,7 +128,14 @@ async def _drain_stderr() -> None:
119128
proc.terminate()
120129
except ProcessLookupError:
121130
pass
122-
await proc.wait()
131+
try:
132+
await asyncio.wait_for(proc.wait(), timeout=TERMINATE_TIMEOUT_SECONDS)
133+
except TimeoutError:
134+
try:
135+
proc.kill()
136+
except ProcessLookupError:
137+
pass
138+
await proc.wait()
123139

124140

125141
@activity.defn(name=RUN_CLAUDE_CODE_TURN_ACTIVITY)

‎src/agentex/lib/cli/templates/default-claude-code/project/acp.py.j2‎

Lines changed: 17 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,9 @@ from agentex.lib.core.tracing.tracing_processor_manager import add_tracing_proce
3333

3434
logger = make_logger(__name__)
3535

36+
STDOUT_LINE_LIMIT = 8 * 1024 * 1024
37+
TERMINATE_TIMEOUT_SECONDS = 5.0
38+
3639
add_tracing_processor_config(
3740
SGPTracingProcessorConfig(
3841
sgp_api_key=os.environ.get("SGP_API_KEY", ""),
@@ -52,6 +55,11 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]:
5255

5356
Injectable seam: tests can monkeypatch this with a fake async iterator of
5457
pre-recorded lines so no real CLI invocation is needed offline.
58+
59+
Lines up to ``STDOUT_LINE_LIMIT`` are read whole: a ``tool_result`` that
60+
echoes a file is a single stream-json line, often past asyncio's 64 KiB
61+
default. If the consumer stops early, the CLI gets SIGTERM and, after
62+
``TERMINATE_TIMEOUT_SECONDS``, SIGKILL, so cancellation cannot hang.
5563
"""
5664
proc = await asyncio.create_subprocess_exec(
5765
"claude",
@@ -62,6 +70,7 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]:
6270
stdin=asyncio.subprocess.PIPE,
6371
stdout=asyncio.subprocess.PIPE,
6472
stderr=asyncio.subprocess.PIPE,
73+
limit=STDOUT_LINE_LIMIT,
6574
)
6675
assert proc.stdout is not None
6776
assert proc.stdin is not None
@@ -123,7 +132,14 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]:
123132
proc.terminate()
124133
except ProcessLookupError:
125134
pass
126-
await proc.wait()
135+
try:
136+
await asyncio.wait_for(proc.wait(), timeout=TERMINATE_TIMEOUT_SECONDS)
137+
except TimeoutError:
138+
try:
139+
proc.kill()
140+
except ProcessLookupError:
141+
pass
142+
await proc.wait()
127143

128144

129145
@acp.on_task_create

‎src/agentex/lib/cli/templates/sync-claude-code/project/acp.py.j2‎

Lines changed: 17 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,9 @@ from agentex.lib.core.tracing.tracing_processor_manager import add_tracing_proce
3535

3636
logger = make_logger(__name__)
3737

38+
STDOUT_LINE_LIMIT = 8 * 1024 * 1024
39+
TERMINATE_TIMEOUT_SECONDS = 5.0
40+
3841
add_tracing_processor_config(
3942
SGPTracingProcessorConfig(
4043
sgp_api_key=os.environ.get("SGP_API_KEY", ""),
@@ -51,6 +54,11 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]:
5154

5255
This is a seam: tests can replace it with a fake async iterator of
5356
pre-recorded lines so no real CLI invocation is needed offline.
57+
58+
Lines up to ``STDOUT_LINE_LIMIT`` are read whole: a ``tool_result`` that
59+
echoes a file is a single stream-json line, often past asyncio's 64 KiB
60+
default. If the consumer stops early, the CLI gets SIGTERM and, after
61+
``TERMINATE_TIMEOUT_SECONDS``, SIGKILL, so cancellation cannot hang.
5462
"""
5563
proc = await asyncio.create_subprocess_exec(
5664
"claude",
@@ -61,6 +69,7 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]:
6169
stdin=asyncio.subprocess.PIPE,
6270
stdout=asyncio.subprocess.PIPE,
6371
stderr=asyncio.subprocess.PIPE,
72+
limit=STDOUT_LINE_LIMIT,
6473
)
6574
assert proc.stdout is not None
6675
assert proc.stdin is not None
@@ -122,7 +131,14 @@ async def _spawn_claude(prompt: str) -> AsyncIterator[str]:
122131
proc.terminate()
123132
except ProcessLookupError:
124133
pass
125-
await proc.wait()
134+
try:
135+
await asyncio.wait_for(proc.wait(), timeout=TERMINATE_TIMEOUT_SECONDS)
136+
except TimeoutError:
137+
try:
138+
proc.kill()
139+
except ProcessLookupError:
140+
pass
141+
await proc.wait()
126142

127143

128144
@acp.on_message_send

‎src/agentex/lib/cli/templates/temporal-claude-code/project/activities.py.j2‎

Lines changed: 17 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,9 @@ from agentex.lib.utils.model_utils import BaseModel
2828

2929
logger = make_logger(__name__)
3030

31+
STDOUT_LINE_LIMIT = 8 * 1024 * 1024
32+
TERMINATE_TIMEOUT_SECONDS = 5.0
33+
3134
RUN_CLAUDE_CODE_TURN_ACTIVITY = "run_claude_code_turn"
3235

3336

@@ -57,6 +60,11 @@ async def _spawn_claude(prompt: str, session_id: str | None = None) -> AsyncIter
5760

5861
Injectable seam: tests can monkeypatch this with a fake async iterator so no
5962
real CLI invocation is needed offline.
63+
64+
Lines up to ``STDOUT_LINE_LIMIT`` are read whole: a ``tool_result`` that
65+
echoes a file is a single stream-json line, often past asyncio's 64 KiB
66+
default. If the consumer stops early, the CLI gets SIGTERM and, after
67+
``TERMINATE_TIMEOUT_SECONDS``, SIGKILL, so cancellation cannot hang.
6068
"""
6169
cmd = [
6270
"claude",
@@ -73,6 +81,7 @@ async def _spawn_claude(prompt: str, session_id: str | None = None) -> AsyncIter
7381
stdin=asyncio.subprocess.PIPE,
7482
stdout=asyncio.subprocess.PIPE,
7583
stderr=asyncio.subprocess.PIPE,
84+
limit=STDOUT_LINE_LIMIT,
7685
)
7786
assert proc.stdout is not None
7887
assert proc.stdin is not None
@@ -135,7 +144,14 @@ async def _spawn_claude(prompt: str, session_id: str | None = None) -> AsyncIter
135144
proc.terminate()
136145
except ProcessLookupError:
137146
pass
138-
await proc.wait()
147+
try:
148+
await asyncio.wait_for(proc.wait(), timeout=TERMINATE_TIMEOUT_SECONDS)
149+
except TimeoutError:
150+
try:
151+
proc.kill()
152+
except ProcessLookupError:
153+
pass
154+
await proc.wait()
139155

140156

141157
@activity.defn(name=RUN_CLAUDE_CODE_TURN_ACTIVITY)

0 commit comments

Comments
 (0)