Skip to content

Commit 88bf5b9

Browse files
committed
fix(harness): close the turn's event source when delivery stops early
The delivery adapters (auto_send for the async push path, yield_events and UnifiedEmitter.yield_turn for the sync yield path) closed their streaming contexts and flushed span derivation in `finally`, but never closed the event source they were iterating. Every tap relies on that close: the Claude Code and Codex taps close their raw line iterator in a `finally`, which is what terminates the CLI subprocess in the scaffolds. When a turn is cancelled while delivery is awaiting the backend (stream_update or a context close) rather than the source, or a sync client disconnects mid-stream, the source stays suspended at a yield. The turn object keeps a reference to its event generator, so garbage collection does not finalize it while the caller still holds the turn, and the CLI subprocess and its stdout handle leak. Both adapters now close the source in `finally`, after the existing cleanup and guarded so a failing aclose cannot mask the original exception. yield_turn closes its delivery generator explicitly instead of leaving it to async-generator finalization. Plain async iterators without aclose and already-exhausted generators are unaffected. Three new tests (cancel during a blocked stream_update; early close of yield_events; early close of yield_turn with a turn that pins its generator) fail on main with `assert [] == [True]` and pass here. tests/lib/core/harness and tests/lib/adk: 571 passed, 1 skipped.
1 parent 6d68f3a commit 88bf5b9

6 files changed

Lines changed: 241 additions & 11 deletions

File tree

‎src/agentex/lib/core/harness/auto_send.py‎

Lines changed: 18 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22

33
from __future__ import annotations
44

5+
import contextlib
56
from typing import Any, AsyncIterator
67
from datetime import datetime
78

@@ -46,7 +47,13 @@ async def auto_send(
4647
Index-keyed routing: each Start(index=i) opens a context stored in
4748
ctx_map[i]; Delta(index=i) routes to ctx_map.get(i); Done(index=i) closes
4849
and removes ctx_map[i]. Events with index is None are skipped. The finally
49-
block closes all remaining open contexts.
50+
block closes all remaining open contexts, then closes `events` itself (when
51+
it exposes `aclose`) so delivery stopping early — cancelled while awaiting
52+
the backend rather than the source — tears the tap down instead of leaving it
53+
suspended at a yield: the turn object pins its event generator, so GC does
54+
not rescue it and the harness subprocess leaks. Closing an exhausted
55+
generator is a no-op, and a failure there is suppressed so it cannot mask the
56+
original exception.
5057
5158
final_text last-segment semantics: a new Start(TextContent) resets
5259
final_text_parts so that multi-step turns return the LAST text segment.
@@ -148,9 +155,15 @@ async def _close_all() -> None:
148155
pass
149156

150157
finally:
151-
await _close_all()
152-
if deriver is not None and tracer is not None:
153-
for signal in deriver.flush():
154-
await tracer.handle(signal)
158+
try:
159+
await _close_all()
160+
if deriver is not None and tracer is not None:
161+
for signal in deriver.flush():
162+
await tracer.handle(signal)
163+
finally:
164+
aclose = getattr(events, "aclose", None)
165+
if aclose is not None:
166+
with contextlib.suppress(Exception):
167+
await aclose()
155168

156169
return TurnResult(final_text="".join(final_text_parts), usage=usage or TurnUsage())

‎src/agentex/lib/core/harness/emitter.py‎

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -53,9 +53,19 @@ def __init__(
5353
self.tracer = None
5454

5555
async def yield_turn(self, turn: HarnessTurn) -> AsyncGenerator[StreamTaskMessage, None]:
56-
"""Sync HTTP ACP delivery: forward events, trace as side effect."""
57-
async for event in yield_events(turn.events, tracer=self.tracer):
58-
yield event
56+
"""Sync HTTP ACP delivery: forward events, trace as side effect.
57+
58+
The finally closes the delivery generator, which closes the turn's event
59+
source in turn, so a consumer that stops early (client disconnect) tears
60+
the tap down here rather than leaving it to async-generator finalization,
61+
which never runs promptly while the turn still pins its generator.
62+
"""
63+
delivery = yield_events(turn.events, tracer=self.tracer)
64+
try:
65+
async for event in delivery:
66+
yield event
67+
finally:
68+
await delivery.aclose()
5969

6070
async def auto_send_turn(self, turn: HarnessTurn, created_at: datetime | None = None) -> TurnResult:
6171
"""Async/temporal delivery: push to the task stream, return TurnResult.

‎src/agentex/lib/core/harness/yield_delivery.py‎

Lines changed: 17 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22

33
from __future__ import annotations
44

5+
import contextlib
56
from typing import AsyncIterator, AsyncGenerator
67

78
from agentex.lib.core.harness.types import StreamTaskMessage
@@ -17,6 +18,13 @@ async def yield_events(
1718
1819
For sync HTTP ACP agents that yield events back over the response. When
1920
`tracer` is None, this is a pure passthrough.
21+
22+
The finally also closes `events` (when it exposes `aclose`), so a consumer
23+
that stops early — a client disconnect closes this generator — tears the tap
24+
down instead of leaving it suspended at a yield: the turn object pins its
25+
event generator, so GC does not rescue it and the harness subprocess leaks.
26+
Closing an exhausted generator is a no-op, and a failure there is suppressed
27+
so it cannot mask the original exception.
2028
"""
2129
deriver = SpanDeriver() if tracer is not None else None
2230
try:
@@ -26,6 +34,12 @@ async def yield_events(
2634
await tracer.handle(signal)
2735
yield event
2836
finally:
29-
if deriver is not None and tracer is not None:
30-
for signal in deriver.flush():
31-
await tracer.handle(signal)
37+
try:
38+
if deriver is not None and tracer is not None:
39+
for signal in deriver.flush():
40+
await tracer.handle(signal)
41+
finally:
42+
aclose = getattr(events, "aclose", None)
43+
if aclose is not None:
44+
with contextlib.suppress(Exception):
45+
await aclose()

‎tests/lib/core/harness/test_auto_send.py‎

Lines changed: 115 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,8 @@
99
This mirrors _langgraph_async.py lines 62-78 and 100-127.
1010
"""
1111

12+
import asyncio
13+
from typing import override
1214
from datetime import datetime
1315

1416
import pytest
@@ -478,3 +480,116 @@ async def test_auto_send_created_at_forwarded():
478480
await auto_send(_gen(events), task_id="task1", tracer=None, streaming=streaming, created_at=dt)
479481

480482
assert all(ts == dt for ts in streaming.recorded_created_at)
483+
484+
485+
class _BlockingCtx(_FakeCtx):
486+
"""A context whose stream_update never returns (a stalled backend).
487+
488+
Sets `blocked` once stream_update is awaited so the test can cancel exactly
489+
while delivery is suspended on the backend rather than on the source.
490+
"""
491+
492+
def __init__(self, sink, content_type, initial_content, blocked):
493+
super().__init__(sink, content_type, initial_content)
494+
self.blocked = blocked
495+
496+
@override
497+
async def stream_update(self, update):
498+
self.sink.append(("update", update))
499+
self.blocked.set()
500+
await asyncio.Event().wait()
501+
502+
503+
class _BlockingStreaming(_FakeStreaming):
504+
"""_FakeStreaming whose contexts block forever inside stream_update."""
505+
506+
def __init__(self):
507+
super().__init__()
508+
self.blocked = asyncio.Event()
509+
510+
@override
511+
def streaming_task_message_context(self, task_id, initial_content, streaming_mode="coalesced", created_at=None):
512+
ctype = getattr(initial_content, "type", None)
513+
self.sink.append(("ctx", ctype))
514+
self.recorded_created_at.append(created_at)
515+
return _BlockingCtx(self.sink, ctype, initial_content, self.blocked)
516+
517+
518+
class _PlainAsyncIterator:
519+
"""An async iterator with no aclose (the AsyncIterator contract minimum)."""
520+
521+
def __init__(self, events):
522+
self._events = iter(events)
523+
524+
def __aiter__(self):
525+
return self
526+
527+
async def __anext__(self):
528+
try:
529+
return next(self._events)
530+
except StopIteration:
531+
raise StopAsyncIteration from None
532+
533+
534+
@pytest.mark.asyncio
535+
async def test_auto_send_closes_source_when_cancelled_mid_delivery():
536+
"""Cancelling delivery while the backend blocks must close the event source.
537+
538+
The turn object pins its event generator, so GC cannot rescue it: when
539+
auto_send returns with the source still suspended at a yield, the tap's
540+
finally never runs and the harness CLI subprocess leaks. The source is held
541+
by a local here for the whole test, and nothing calls gc.collect().
542+
"""
543+
streaming = _BlockingStreaming()
544+
closed: list[bool] = []
545+
546+
async def _recording_source():
547+
try:
548+
yield StreamTaskMessageStart(
549+
type="start",
550+
index=0,
551+
content=TextContent(type="text", author="agent", content=""),
552+
)
553+
yield StreamTaskMessageDelta(
554+
type="delta",
555+
index=0,
556+
delta=TextDelta(type="text", text_delta="hi"),
557+
)
558+
yield StreamTaskMessageDone(type="done", index=0)
559+
finally:
560+
closed.append(True)
561+
562+
source = _recording_source()
563+
task = asyncio.create_task(auto_send(source, task_id="task1", tracer=None, streaming=streaming))
564+
await asyncio.wait_for(streaming.blocked.wait(), timeout=5)
565+
task.cancel()
566+
with pytest.raises(asyncio.CancelledError):
567+
await task
568+
569+
assert closed == [True]
570+
assert ("close", "text") in [(s[0], s[1]) for s in streaming.sink]
571+
572+
573+
@pytest.mark.asyncio
574+
async def test_auto_send_accepts_source_without_aclose():
575+
"""A plain async iterator (no aclose) must deliver exactly as before."""
576+
streaming = _FakeStreaming()
577+
events = [
578+
StreamTaskMessageStart(
579+
type="start",
580+
index=0,
581+
content=TextContent(type="text", author="agent", content=""),
582+
),
583+
StreamTaskMessageDelta(
584+
type="delta",
585+
index=0,
586+
delta=TextDelta(type="text", text_delta="Hi"),
587+
),
588+
StreamTaskMessageDone(type="done", index=0),
589+
]
590+
result = await auto_send(_PlainAsyncIterator(events), task_id="task1", tracer=None, streaming=streaming)
591+
592+
assert result.final_text == "Hi"
593+
kinds = [s[0] for s in streaming.sink]
594+
assert kinds.count("open") == 1
595+
assert kinds.count("close") == 1

‎tests/lib/core/harness/test_emitter.py‎

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -140,3 +140,54 @@ async def test_emitter_auto_send_turn_reads_usage_after_exhaustion():
140140
result = await emitter.auto_send_turn(turn)
141141
assert result.usage == real_usage
142142
assert result.usage.input_tokens == 11 and result.usage.total_tokens == 33
143+
144+
145+
class _PinnedTurn:
146+
"""A turn that pins its event generator, as the real CLI taps do.
147+
148+
`closed` records that the generator's finally ran (what terminates the CLI
149+
subprocess in the scaffolds).
150+
"""
151+
152+
def __init__(self, events_list):
153+
self._events_list = events_list
154+
self.closed: list[bool] = []
155+
self._gen = None
156+
157+
@property
158+
def events(self):
159+
if self._gen is None:
160+
self._gen = self._stream()
161+
return self._gen
162+
163+
async def _stream(self):
164+
try:
165+
for e in self._events_list:
166+
yield e
167+
finally:
168+
self.closed.append(True)
169+
170+
def usage(self):
171+
return TurnUsage()
172+
173+
174+
@pytest.mark.asyncio
175+
async def test_emitter_yield_turn_closes_source_on_early_close():
176+
"""Closing the delivery generator must close the turn's event source.
177+
178+
No gc.collect() and no sleep: the turn keeps the generator referenced, and
179+
async-generator finalization is too late for a disconnected client anyway.
180+
"""
181+
events = [
182+
StreamTaskMessageStart(type="start", index=0, content=TextContent(type="text", author="agent", content="")),
183+
StreamTaskMessageDelta(type="delta", index=0, delta=TextDelta(type="text", text_delta="hi")),
184+
StreamTaskMessageDone(type="done", index=0),
185+
]
186+
turn = _PinnedTurn(events)
187+
emitter = UnifiedEmitter(task_id="t", trace_id=None, parent_span_id=None)
188+
gen = emitter.yield_turn(turn)
189+
first = await gen.__anext__()
190+
await gen.aclose()
191+
192+
assert first.index == 0
193+
assert turn.closed == [True]

‎tests/lib/core/harness/test_yield_delivery.py‎

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -76,3 +76,30 @@ async def test_flush_runs_on_early_close():
7676
await gen.aclose() # triggers the finally -> flush()
7777
assert fake.started_names == ["Bash"]
7878
assert fake.ended_outputs == [None] # flush closed the unpaired span (incomplete, no output)
79+
80+
81+
@pytest.mark.asyncio
82+
async def test_source_closed_when_consumer_closes_early():
83+
"""Closing the delivery generator must close the upstream event source.
84+
85+
A client disconnect (or any early break) closes the generator handed to the
86+
caller. The source is pinned by the turn object, so if it is left suspended
87+
at a yield the tap's finally never runs and the CLI subprocess leaks. No
88+
gc.collect() here: the source stays referenced for the whole test.
89+
"""
90+
closed: list[bool] = []
91+
92+
async def _recording_source():
93+
try:
94+
yield StreamTaskMessageDone(type="done", index=0)
95+
yield StreamTaskMessageDone(type="done", index=1)
96+
finally:
97+
closed.append(True)
98+
99+
source = _recording_source()
100+
gen = yield_events(source, tracer=None)
101+
first = await gen.__anext__()
102+
await gen.aclose()
103+
104+
assert first.index == 0
105+
assert closed == [True]

0 commit comments

Comments
 (0)