Skip to content

Commit 62b97be

Browse files
fix(tracing): keep extra span keys, make the processor lock reentrant, drop the tracer's dead client
Review follow-ups on the span-processor removal: - Span keeps unknown keys like the generated model it replaces, so a custom processor's own attributes and older worker payloads survive. - set_processor_configs re-enters add_processor_config under the manager lock, which deadlocked on a plain Lock. Reentrant now, with tests for the SGP happy path and the batch registration. - Trace and Tracer no longer use the client they receive, so it is optional and the ADK tracing module stops building an httpx client per event loop. - Changelog entry for the breaking removals, stale prose updated. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
1 parent d7d8720 commit 62b97be

10 files changed

Lines changed: 109 additions & 63 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@
44

55
### ⚠ BREAKING CHANGES
66

7+
* **tracing:** removed the Agentex-native span processor and its `AgentexTracingProcessorConfig` (the Agentex server is retiring its Postgres spans API), along with `Trace.get_span` / `Trace.list_spans` and their async twins. `SGPTracingProcessorConfig` is the only processor config and registering any other type raises `ValueError`. The in-memory `Span` is now `agentex.lib.types.tracing.Span` (also exported as `agentex.lib.core.tracing.Span`); the generated `agentex.types.span.Span` disappears with the next client generation. Its `to_dict()` / `to_json()` return the full JSON-mode dump rather than only the fields that were set.
78
* **harness:** removed the deprecated bespoke LangGraph tracing handler `create_langgraph_tracing_handler` (and its `AgentexLangGraphTracingHandler` class) from the public `agentex.lib.adk` surface. Span tracing is now derived from the canonical `StreamTaskMessage*` stream by `UnifiedEmitter` — wrap your run in the harness `*Turn` and drive `UnifiedEmitter.yield_turn` / `auto_send_turn`. The `agentex init` templates were migrated accordingly.
89
* **harness:** removed the deprecated bespoke Pydantic-AI tracing handler `create_pydantic_ai_tracing_handler` (and its `AgentexPydanticAITracingHandler` class) from the public `agentex.lib.adk` surface. Span tracing is now derived from the canonical `StreamTaskMessage*` stream by `UnifiedEmitter` — wrap your run in `PydanticAITurn` and drive `UnifiedEmitter.yield_turn` / `auto_send_turn`. The `agentex init` templates were migrated accordingly.
910
* **harness:** each harness now exposes exactly `_<harness>_sync.py` + `_<harness>_turn.py` under `agentex.lib.adk._modules`. The OpenAI harness `OpenAITurn` and `convert_openai_to_agentex_events` moved to `agentex.lib.adk._modules._openai_turn` / `_openai_sync`; back-compat shims remain at `agentex.lib.adk.providers._modules.{openai_turn,sync_provider}` for one release. Public facade names (`stream_pydantic_ai_events`, `stream_langgraph_events`, `emit_langgraph_messages`, etc.) are unchanged.

‎src/agentex/lib/adk/_modules/tracing.py‎

Lines changed: 3 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,6 @@
1111
from temporalio.exceptions import ActivityError, TimeoutError as TemporalTimeoutError, is_cancelled_exception
1212

1313
from agentex import AsyncAgentex # noqa: F401
14-
from agentex.lib.adk.utils._modules.client import create_async_agentex_client
1514
from agentex.lib.core.services.adk.tracing import TracingService
1615
from agentex.lib.core.temporal.activities.activity_helpers import ActivityHelpers
1716
from agentex.lib.core.temporal.activities.adk.tracing_activities import (
@@ -145,46 +144,17 @@ def __init__(self, tracing_service: TracingService | None = None):
145144
146145
Args:
147146
tracing_service (Optional[TracingService]): Optional pre-configured tracing service.
148-
If None, will be lazily created on first use so the httpx client is
149-
bound to the correct running event loop.
147+
If None, one is created on first use.
150148
"""
151149
self._tracing_service_explicit = tracing_service
152150
self._tracing_service_lazy: TracingService | None = None
153-
self._bound_loop_id: int | None = None
154151

155152
@property
156153
def _tracing_service(self) -> TracingService:
157154
if self._tracing_service_explicit is not None:
158155
return self._tracing_service_explicit
159-
160-
import asyncio
161-
162-
# Determine the current event loop (if any).
163-
try:
164-
loop = asyncio.get_running_loop()
165-
loop_id = id(loop)
166-
except RuntimeError:
167-
loop_id = None
168-
169-
# Re-create the underlying httpx client when the event loop changes
170-
# (e.g. between HTTP requests in a sync ASGI server) to avoid
171-
# "Event loop is closed" / "bound to a different event loop" errors.
172-
if self._tracing_service_lazy is None or (loop_id is not None and loop_id != self._bound_loop_id):
173-
import httpx
174-
175-
# Keepalive ON: connections are reused within a single event
176-
# loop, eliminating the TLS-handshake-per-span penalty under
177-
# load. Cross-loop safety is preserved by rebuilding the
178-
# client whenever loop_id changes (the conditional above).
179-
agentex_client = create_async_agentex_client(
180-
http_client=httpx.AsyncClient(
181-
limits=httpx.Limits(max_keepalive_connections=20),
182-
),
183-
)
184-
tracer = AsyncTracer(agentex_client)
185-
self._tracing_service_lazy = TracingService(tracer=tracer)
186-
self._bound_loop_id = loop_id
187-
156+
if self._tracing_service_lazy is None:
157+
self._tracing_service_lazy = TracingService(tracer=AsyncTracer())
188158
return self._tracing_service_lazy
189159

190160
@asynccontextmanager

‎src/agentex/lib/core/tracing/processors/sgp_tracing_processor.py‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -74,8 +74,8 @@ def _sgp_metadata(span: Span) -> Any:
7474
7575
Returns a COPY rather than mutating ``span``. ``trace.py`` hands the same
7676
Span instance to every registered processor, so anything written onto
77-
``span.data`` here would also be serialized by the Agentex processor and
78-
show up in caller-visible span data. ``__commit_sha__`` is opt-in and
77+
``span.data`` here would also reach every other processor and show up in
78+
caller-visible span data. ``__commit_sha__`` is opt-in and
7979
SGP-scoped, so it must not leak that way.
8080
8181
(The ``__source__`` / ``__agent_*`` keys set by ``_add_source_to_span`` do

‎src/agentex/lib/core/tracing/span_error.py‎

Lines changed: 3 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -12,12 +12,9 @@
1212
from agentex.lib.types.tracing import Span
1313

1414
# Reserved key under ``Span.data`` carrying failure info for a span whose
15-
# context-manager body raised. Mirrors the existing ``__span_type__`` /
16-
# ``__source__`` reserved-key convention already read/written by the SGP
17-
# processor. Stored in ``data`` because the Span model is generated from the
18-
# OpenAPI spec and has no first-class status/error field; ``data`` is a real
19-
# field, so it survives ``model_copy(deep=True)`` and round-trips to both the
20-
# SGP and agentex-native span stores.
15+
# context-manager body raised, alongside the ``__span_type__`` / ``__source__``
16+
# keys the SGP processor already reads. Kept in ``data`` so it survives
17+
# ``model_copy(deep=True)`` and reaches the processors with the span.
2118
SPAN_ERROR_KEY = "__error__"
2219

2320
ERROR_CATEGORY_UNKNOWN: ErrorCategory = "unknown"

‎src/agentex/lib/core/tracing/trace.py‎

Lines changed: 12 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -218,22 +218,24 @@ def _begin_obs(
218218

219219
class Trace:
220220
"""
221-
Trace is a wrapper around the Agentex API for tracing.
222-
It provides a context manager for spans and a way to start and end spans.
221+
Trace groups the spans of one trace id and hands each span to the
222+
registered processors. It provides a context manager for spans and a way
223+
to start and end spans.
223224
"""
224225

225226
def __init__(
226227
self,
227228
processors: list[SyncTracingProcessor],
228-
client: Agentex,
229+
client: Agentex | None = None,
229230
trace_id: str | None = None,
230231
):
231232
"""
232233
Initialize a new trace with the specified trace ID.
233234
234235
Args:
235-
trace_id: Required trace ID to use for this trace.
236-
processors: Optional list of tracing processors to use for this trace.
236+
processors: Tracing processors every span is handed to.
237+
client: Kept for backward compatibility, no longer used.
238+
trace_id: Trace ID to use for this trace.
237239
"""
238240
self.processors = processors
239241
self.client = client
@@ -251,7 +253,7 @@ def start_span(
251253
task_id: str | None = None,
252254
) -> Span:
253255
"""
254-
Start a new span and register it with the API.
256+
Start a new span and hand it to the registered processors.
255257
256258
Args:
257259
name: Name of the span.
@@ -325,8 +327,6 @@ def end_span(
325327

326328
return span
327329

328-
329-
330330
@contextmanager
331331
def span(
332332
self,
@@ -355,14 +355,14 @@ def span(
355355

356356
class AsyncTrace:
357357
"""
358-
AsyncTrace is a wrapper around the Agentex API for tracing.
359-
It provides a context manager for spans and a way to start and end spans.
358+
AsyncTrace is the async version of Trace. It provides a context manager
359+
for spans and a way to start and end spans.
360360
"""
361361

362362
def __init__(
363363
self,
364364
processors: list[AsyncTracingProcessor],
365-
client: AsyncAgentex,
365+
client: AsyncAgentex | None = None,
366366
trace_id: str | None = None,
367367
span_queue: AsyncSpanQueue | None = None,
368368
):
@@ -391,7 +391,7 @@ async def start_span(
391391
task_id: str | None = None,
392392
) -> Span:
393393
"""
394-
Start a new span and register it with the API.
394+
Start a new span and hand it to the registered processors.
395395
396396
Args:
397397
name: Name of the span.
@@ -480,8 +480,6 @@ async def end_span(
480480

481481
return span
482482

483-
484-
485483
@asynccontextmanager
486484
async def span(
487485
self,

‎src/agentex/lib/core/tracing/tracer.py‎

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -15,12 +15,12 @@ class Tracer:
1515
It manages the client connection and creates traces.
1616
"""
1717

18-
def __init__(self, client: Agentex):
18+
def __init__(self, client: Agentex | None = None):
1919
"""
20-
Initialize a new sync tracer with the provided client.
20+
Initialize a new sync tracer.
2121
2222
Args:
23-
client: Agentex client instance used for API communication.
23+
client: Kept for backward compatibility, no longer used.
2424
"""
2525
self.client = client
2626

@@ -47,12 +47,12 @@ class AsyncTracer:
4747
It manages the async client connection and creates async traces.
4848
"""
4949

50-
def __init__(self, client: AsyncAgentex):
50+
def __init__(self, client: AsyncAgentex | None = None):
5151
"""
52-
Initialize a new async tracer with the provided client.
52+
Initialize a new async tracer.
5353
5454
Args:
55-
client: AsyncAgentex client instance used for API communication.
55+
client: Kept for backward compatibility, no longer used.
5656
"""
5757
self.client = client
5858

‎src/agentex/lib/core/tracing/tracing_processor_manager.py‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
from __future__ import annotations
22

3-
from threading import Lock
3+
from threading import RLock
44

55
from agentex.lib.types.tracing import TracingProcessorConfig
66
from agentex.lib.core.tracing.processors.sgp_tracing_processor import (
@@ -24,7 +24,8 @@ def __init__(self):
2424
# Cache for processors
2525
self.sync_processors: list[SyncTracingProcessor] = []
2626
self.async_processors: list[AsyncTracingProcessor] = []
27-
self.lock = Lock()
27+
# Reentrant: set_processor_configs holds it while calling add_processor_config.
28+
self.lock = RLock()
2829

2930
def add_processor_config(self, processor_config: TracingProcessorConfig) -> None:
3031
with self.lock:

‎src/agentex/lib/types/tracing.py‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,8 @@
33
from typing import Any, Literal
44
from datetime import datetime
55

6+
from pydantic import ConfigDict
7+
68
from agentex.lib.utils.model_utils import BaseModel
79

810

@@ -22,6 +24,9 @@ class BaseModelWithTraceParams(BaseModel):
2224
class Span(BaseModel):
2325
"""In-memory span handed to tracing processors. Owned here, not by the generated client."""
2426

27+
# The generated model kept unknown keys, and custom processors may stash their own.
28+
model_config = ConfigDict(extra="allow")
29+
2530
id: str
2631
name: str
2732
start_time: datetime
Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,21 @@
1+
from datetime import UTC, datetime
2+
3+
from agentex.lib.types.tracing import Span
4+
5+
6+
def _span(**extra) -> Span:
7+
return Span(id="s1", name="n", trace_id="t1", start_time=datetime(2026, 1, 1, tzinfo=UTC), **extra)
8+
9+
10+
def test_unknown_keys_survive_validation_and_a_json_round_trip():
11+
span = Span.model_validate({**_span().model_dump(), "extension": {"sampled": True}})
12+
13+
assert span.extension == {"sampled": True} # type: ignore[attr-defined]
14+
assert Span.model_validate_json(span.model_dump_json()).model_dump()["extension"] == {"sampled": True}
15+
16+
17+
def test_processors_can_attach_their_own_attributes():
18+
span = _span()
19+
span.annotation = "custom" # type: ignore[attr-defined]
20+
21+
assert span.model_copy(deep=True).model_dump()["annotation"] == "custom"

‎tests/lib/core/tracing/test_tracing_processor_manager.py‎

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,32 @@
1+
import threading
12
from types import SimpleNamespace
3+
from unittest.mock import MagicMock, patch
24

35
import pytest
46

7+
from agentex.lib.types.tracing import SGPTracingProcessorConfig
58
from agentex.lib.core.tracing.tracing_processor_manager import TracingProcessorManager
9+
from agentex.lib.core.tracing.processors.sgp_tracing_processor import (
10+
SGPSyncTracingProcessor,
11+
SGPAsyncTracingProcessor,
12+
)
13+
14+
SGP_MODULE = "agentex.lib.core.tracing.processors.sgp_tracing_processor"
15+
16+
17+
def _sgp_config() -> SGPTracingProcessorConfig:
18+
return SGPTracingProcessorConfig(sgp_api_key="k", sgp_account_id="a", sgp_base_url="http://sgp.test")
19+
20+
21+
def _patched_sgp():
22+
env = MagicMock()
23+
env.refresh.return_value = MagicMock(ACP_TYPE=None, AGENT_NAME=None, AGENT_ID=None, AGENT_VERSION=None)
24+
return (
25+
patch(f"{SGP_MODULE}.SGPClient"),
26+
patch(f"{SGP_MODULE}.AsyncSGPClient"),
27+
patch(f"{SGP_MODULE}.tracing.init"),
28+
patch(f"{SGP_MODULE}.EnvironmentVariables", env),
29+
)
630

731

832
def test_unknown_processor_type_is_rejected_by_name():
@@ -13,3 +37,32 @@ def test_unknown_processor_type_is_rejected_by_name():
1337

1438
assert manager.get_sync_processors() == []
1539
assert manager.get_async_processors() == []
40+
41+
42+
def test_sgp_config_registers_one_sync_and_one_async_processor():
43+
p1, p2, p3, p4 = _patched_sgp()
44+
with p1, p2, p3, p4:
45+
manager = TracingProcessorManager()
46+
manager.add_processor_config(_sgp_config())
47+
48+
(sync_processor,) = manager.get_sync_processors()
49+
(async_processor,) = manager.get_async_processors()
50+
assert isinstance(sync_processor, SGPSyncTracingProcessor)
51+
assert isinstance(async_processor, SGPAsyncTracingProcessor)
52+
53+
54+
def test_set_processor_configs_registers_every_config_without_deadlocking():
55+
p1, p2, p3, p4 = _patched_sgp()
56+
manager = TracingProcessorManager()
57+
done = threading.Event()
58+
59+
def register():
60+
with p1, p2, p3, p4:
61+
manager.set_processor_configs([_sgp_config(), _sgp_config()])
62+
done.set()
63+
64+
threading.Thread(target=register, daemon=True).start()
65+
66+
assert done.wait(timeout=5), "set_processor_configs hung: the manager lock must be reentrant"
67+
assert len(manager.get_sync_processors()) == 2
68+
assert len(manager.get_async_processors()) == 2

0 commit comments

Comments
 (0)