Skip to content

Commit dab044f

Browse files
declan-scaleclaude
andcommitted
refactor(harness): simplify yield_events guard + cover finally-flush on early close
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
1 parent 803191b commit dab044f

2 files changed

Lines changed: 23 additions & 4 deletions

File tree

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

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

33
from __future__ import annotations
44

5-
from typing import AsyncIterator
5+
from typing import AsyncGenerator, AsyncIterator
66

77
from agentex.lib.core.harness.span_derivation import SpanDeriver
88
from agentex.lib.core.harness.tracer import SpanTracer
@@ -12,7 +12,7 @@
1212
async def yield_events(
1313
events: AsyncIterator[StreamTaskMessage],
1414
tracer: SpanTracer | None = None,
15-
) -> AsyncIterator[StreamTaskMessage]:
15+
) -> AsyncGenerator[StreamTaskMessage, None]:
1616
"""Forward each event to the caller; derive + trace spans as a side effect.
1717
1818
For sync HTTP ACP agents that yield events back over the response. When
@@ -21,11 +21,11 @@ async def yield_events(
2121
deriver = SpanDeriver() if tracer is not None else None
2222
try:
2323
async for event in events:
24-
if deriver is not None and tracer is not None:
24+
if deriver is not None: # tracer is non-None whenever deriver is set
2525
for signal in deriver.observe(event):
2626
await tracer.handle(signal)
2727
yield event
2828
finally:
29-
if deriver is not None and tracer is not None:
29+
if deriver is not None: # tracer is non-None whenever deriver is set
3030
for signal in deriver.flush():
3131
await tracer.handle(signal)

tests/lib/core/harness/test_yield_delivery.py

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -56,3 +56,22 @@ async def test_yield_without_tracer_is_pure_passthrough():
5656
]
5757
out = [e async for e in yield_events(_gen(events), tracer=None)]
5858
assert out == events
59+
60+
61+
@pytest.mark.asyncio
62+
async def test_flush_runs_on_early_close():
63+
fake = _RecordTracing()
64+
tracer = SpanTracer(trace_id="t", parent_span_id="p", tracing=fake)
65+
events = [
66+
StreamTaskMessageStart(type="start", index=0,
67+
content=ToolRequestContent(type="tool_request", author="agent",
68+
tool_call_id="c", name="Bash", arguments={})),
69+
StreamTaskMessageDone(type="done", index=0),
70+
# response intentionally never arrives
71+
]
72+
gen = yield_events(_gen(events), tracer=tracer)
73+
first = await gen.__anext__() # Start
74+
second = await gen.__anext__() # Done -> tool span opens here
75+
await gen.aclose() # triggers the finally -> flush()
76+
assert fake.started == ["Bash"]
77+
assert fake.ended == [None] # flush closed the unpaired span (incomplete, no output)

0 commit comments

Comments
 (0)