Skip to content

Commit 2740bac

Browse files
NiteshDhanpalclaude
andcommitted
refactor(tracing): low-cardinality ACP dispatch span names
acp.task_create:{task.id} / acp.event_send:{task.id} put the task id in the span NAME, which is high-cardinality and breaks span-name aggregation in Tempo. Use static names (acp.task_create / acp.event_send) and carry the id as the agentex.task_id span attribute instead. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
1 parent 044b3f4 commit 2740bac

1 file changed

Lines changed: 7 additions & 4 deletions

File tree

src/agentex/lib/core/temporal/services/temporal_task_service.py

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

1616

1717
@contextmanager
18-
def _acp_dispatch_span(name: str) -> Iterator[None]:
18+
def _acp_dispatch_span(name: str, task_id: str | None = None) -> Iterator[None]:
1919
"""Wrap an ACP -> Temporal dispatch (start_workflow / signal) in an OTel span.
2020
2121
The Temporal OpenTelemetry interceptor propagates trace context by injecting
@@ -38,7 +38,10 @@ def _acp_dispatch_span(name: str) -> Iterator[None]:
3838
yield
3939
return
4040
tracer = _otel_trace.get_tracer("agentex.acp")
41-
with tracer.start_as_current_span(name, kind=_otel_trace.SpanKind.PRODUCER):
41+
# task_id goes on an attribute, NOT in the span name: a per-task span name is
42+
# high-cardinality and breaks span-name aggregation in Tempo.
43+
attributes = {"agentex.task_id": task_id} if task_id else None
44+
with tracer.start_as_current_span(name, kind=_otel_trace.SpanKind.PRODUCER, attributes=attributes):
4245
yield
4346

4447

@@ -66,7 +69,7 @@ async def submit_task(self, agent: Agent, task: Task, params: dict[str, Any] | N
6669
# value bounds the whole continue-as-new chain's wall-clock lifetime.
6770
timeout_seconds = self._env_vars.WORKFLOW_EXECUTION_TIMEOUT_SECONDS
6871
execution_timeout = timedelta(seconds=timeout_seconds) if timeout_seconds and timeout_seconds > 0 else None
69-
with _acp_dispatch_span(f"acp.task_create:{task.id}"):
72+
with _acp_dispatch_span("acp.task_create", task_id=task.id):
7073
return await self._temporal_client.start_workflow(
7174
workflow=self._env_vars.WORKFLOW_NAME,
7275
arg=CreateTaskParams(
@@ -88,7 +91,7 @@ async def get_state(self, task_id: str) -> WorkflowState:
8891
)
8992

9093
async def send_event(self, agent: Agent, task: Task, event: Event, request: dict | None = None) -> None:
91-
with _acp_dispatch_span(f"acp.event_send:{task.id}"):
94+
with _acp_dispatch_span("acp.event_send", task_id=task.id):
9295
return await self._temporal_client.send_signal(
9396
workflow_id=task.id,
9497
signal=SignalName.RECEIVE_EVENT.value,

0 commit comments

Comments
 (0)