Skip to content

Commit a61f8c0

Browse files
declan-scaleclaude
andcommitted
feat(harness): UnifiedEmitter facade tying delivery + tracing + usage
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
1 parent df1a0eb commit a61f8c0

3 files changed

Lines changed: 139 additions & 0 deletions

File tree

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

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,3 +4,27 @@
44
package derives spans from it and delivers it (yield or auto-send), so every
55
harness tap gets streaming + tracing + turn usage uniformly.
66
"""
7+
8+
from agentex.lib.core.harness.emitter import UnifiedEmitter
9+
from agentex.lib.core.harness.tracer import SpanTracer
10+
from agentex.lib.core.harness.types import (
11+
CloseSpan,
12+
HarnessTurn,
13+
OpenSpan,
14+
SpanSignal,
15+
StreamTaskMessage,
16+
TurnResult,
17+
TurnUsage,
18+
)
19+
20+
__all__ = [
21+
"UnifiedEmitter",
22+
"SpanTracer",
23+
"OpenSpan",
24+
"CloseSpan",
25+
"SpanSignal",
26+
"StreamTaskMessage",
27+
"TurnUsage",
28+
"TurnResult",
29+
"HarnessTurn",
30+
]
Lines changed: 59 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,59 @@
1+
"""UnifiedEmitter: the single facade agent authors use for either delivery mode."""
2+
3+
from __future__ import annotations
4+
5+
from typing import AsyncIterator
6+
7+
from agentex.lib.core.harness.auto_send import auto_send
8+
from agentex.lib.core.harness.tracer import SpanTracer
9+
from agentex.lib.core.harness.types import HarnessTurn, StreamTaskMessage, TurnResult
10+
from agentex.lib.core.harness.yield_delivery import yield_events
11+
12+
13+
class UnifiedEmitter:
14+
"""Ties trace context + chosen delivery together.
15+
16+
Tracing is default-on whenever `trace_id` is truthy; pass `tracer=False` to
17+
disable, or a custom `SpanTracer` to override.
18+
"""
19+
20+
tracer: SpanTracer | None
21+
22+
def __init__(
23+
self,
24+
task_id: str,
25+
trace_id: str | None,
26+
parent_span_id: str | None,
27+
tracer: SpanTracer | bool | None = None,
28+
tracing: object | None = None,
29+
):
30+
self.task_id = task_id
31+
self.trace_id = trace_id
32+
self.parent_span_id = parent_span_id
33+
if tracer is False:
34+
self.tracer = None
35+
elif isinstance(tracer, SpanTracer):
36+
self.tracer = tracer
37+
elif trace_id:
38+
self.tracer = SpanTracer(
39+
trace_id=trace_id,
40+
parent_span_id=parent_span_id,
41+
task_id=task_id,
42+
tracing=tracing,
43+
)
44+
else:
45+
self.tracer = None
46+
47+
async def yield_turn(self, turn: HarnessTurn) -> AsyncIterator[StreamTaskMessage]:
48+
"""Sync HTTP ACP delivery: forward events, trace as side effect."""
49+
async for event in yield_events(turn.events, tracer=self.tracer):
50+
yield event
51+
52+
async def auto_send_turn(self, turn: HarnessTurn) -> TurnResult:
53+
"""Async/temporal delivery: push to the task stream, return TurnResult."""
54+
return await auto_send(
55+
turn.events,
56+
task_id=self.task_id,
57+
tracer=self.tracer,
58+
usage=turn.usage(),
59+
)
Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,56 @@
1+
import pytest
2+
3+
from agentex.lib.core.harness.emitter import UnifiedEmitter
4+
from agentex.lib.core.harness.types import TurnUsage
5+
from agentex.types.task_message_update import StreamTaskMessageStart, StreamTaskMessageDone
6+
from agentex.types.text_content import TextContent
7+
8+
9+
class _FakeTracing:
10+
async def start_span(self, **kw):
11+
return None
12+
13+
async def end_span(self, **kw):
14+
pass
15+
16+
17+
class _Turn:
18+
def __init__(self, events_list, usage):
19+
self._events_list = events_list
20+
self._usage = usage
21+
22+
@property
23+
async def events(self):
24+
for e in self._events_list:
25+
yield e
26+
27+
def usage(self):
28+
return self._usage
29+
30+
31+
@pytest.mark.asyncio
32+
async def test_emitter_yield_mode_passes_through():
33+
events = [
34+
StreamTaskMessageStart(type="start", index=0,
35+
content=TextContent(type="text", author="agent", content="hi")),
36+
StreamTaskMessageDone(type="done", index=0),
37+
]
38+
turn = _Turn(events, TurnUsage(model="m"))
39+
emitter = UnifiedEmitter(task_id="t", trace_id=None, parent_span_id=None)
40+
out = [e async for e in emitter.yield_turn(turn)]
41+
assert out == events
42+
43+
44+
@pytest.mark.asyncio
45+
async def test_emitter_tracing_default_on_when_trace_id_present():
46+
# Inject a fake tracing backend so the test env doesn't need temporalio.
47+
# This exercises the default-on path (tracer=None) when trace_id is truthy.
48+
emitter = UnifiedEmitter(task_id="t", trace_id="trace1", parent_span_id="p",
49+
tracing=_FakeTracing())
50+
assert emitter.tracer is not None
51+
52+
53+
@pytest.mark.asyncio
54+
async def test_emitter_tracing_overridable_off():
55+
emitter = UnifiedEmitter(task_id="t", trace_id="trace1", parent_span_id="p", tracer=False)
56+
assert emitter.tracer is None

0 commit comments

Comments
 (0)