Skip to content

Commit e3f591a

Browse files
feat(lineage): capture tool data-source refs and agent build version in span data
Implements the trace data-source-ref convention (SGP-6513): tools declare which data sources they touch — statically, via an args resolver, or by name-keyed registry for MCP/unowned tools — and every tool-span path merges the resolved refs into span data under sgp.lineage.refs, which the SGP tracing processor already ships as span metadata. Capture is decoupled from lineage derivation so agents instrument from day one and edges backfill later. Also stamps __agent_version__ from a new AGENT_VERSION env var (same mechanism as __agent_name__), completing the trace-side join-key set: span -> agent version snapshot is the runtime half of SGP-6132. Convention spec: scaleapi packages/sgp-lineage/docs/specs/ 2026-07-22-sgp-6513-trace-data-source-ref-convention.md (PR #153026). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
1 parent 199fd6a commit e3f591a

12 files changed

Lines changed: 462 additions & 24 deletions

File tree

src/agentex/lib/adk/__init__.py

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,10 @@
2929
from agentex.lib.adk._modules.tasks import TasksModule
3030
from agentex.lib.adk._modules.tracing import TracingModule, TurnSpan
3131

32+
# Data-source refs for lineage (SGP-6513); implementation lives in core.tracing
33+
from agentex.lib.core.tracing import lineage
34+
from agentex.lib.core.tracing.lineage import DataSourceRef, data_sources
35+
3236
# Unified harness surface (AGX1-375)
3337
from agentex.lib.core.harness import (
3438
UnifiedEmitter,
@@ -67,6 +71,10 @@
6771
"events",
6872
"agent_task_tracker",
6973
"TurnSpan",
74+
# Lineage data-source refs (SGP-6513)
75+
"lineage",
76+
"DataSourceRef",
77+
"data_sources",
7078
# Checkpointing / LangGraph
7179
"create_checkpointer",
7280
"stream_langgraph_events",

src/agentex/lib/adk/providers/_modules/sync_provider.py

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919
from agentex import AsyncAgentex
2020
from agentex.lib.utils.logging import make_logger
2121
from agentex.lib.core.tracing.tracer import AsyncTracer
22+
from agentex.lib.core.tracing.lineage import merge_refs_into_data, resolve_refs_from_items
2223

2324
logger = make_logger(__name__)
2425

@@ -185,6 +186,9 @@ async def get_response(
185186
"new_items": new_items,
186187
"final_output": final_output,
187188
}
189+
lineage_refs = resolve_refs_from_items(new_items)
190+
if lineage_refs:
191+
span.data = merge_refs_into_data(span.data, lineage_refs)
188192

189193
return response
190194
else:
@@ -303,6 +307,9 @@ async def stream_response(
303307
"new_items": new_items,
304308
"final_output": final_response_text if final_response_text else None,
305309
}
310+
lineage_refs = resolve_refs_from_items(new_items)
311+
if lineage_refs:
312+
span.data = merge_refs_into_data(span.data, lineage_refs)
306313
finally:
307314
# End the span after all events have been yielded
308315
await trace.end_span(span)

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

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,12 @@
66

77
from agentex.lib.core.harness.types import OpenSpan, CloseSpan, SpanSignal
88

9+
try:
10+
from agentex.lib.core.tracing.lineage import resolve_refs, merge_refs_into_data
11+
except Exception: # keep the harness importable without optional tracing deps
12+
resolve_refs = None # type: ignore[assignment]
13+
merge_refs_into_data = None # type: ignore[assignment]
14+
915
try:
1016
from agentex.lib.utils.logging import make_logger
1117

@@ -80,6 +86,13 @@ async def handle(self, signal: SpanSignal) -> None:
8086
task_id=self.task_id,
8187
)
8288
if span is not None:
89+
if signal.kind == "tool" and resolve_refs is not None:
90+
refs = resolve_refs(
91+
signal.name, signal.input if isinstance(signal.input, dict) else {}
92+
)
93+
if refs:
94+
data = span.data if isinstance(span.data, dict) else {}
95+
span.data = merge_refs_into_data(data, refs)
8396
self._open[signal.key] = span
8497
elif isinstance(signal, CloseSpan):
8598
span = self._open.pop(signal.key, None)

src/agentex/lib/core/services/adk/providers/openai.py

Lines changed: 33 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@
2525
from agentex.lib.utils.temporal import heartbeat_if_in_workflow
2626
from agentex.lib.core.tracing.tracer import AsyncTracer
2727
from agentex.lib.core.harness.emitter import UnifiedEmitter
28+
from agentex.lib.core.tracing.lineage import merge_refs_into_data, resolve_refs_from_items
2829
from agentex.types.task_message_update import StreamTaskMessageFull
2930
from agentex.types.task_message_content import (
3031
TextContent,
@@ -286,13 +287,17 @@ async def run_agent(
286287
result = await Runner.run(starting_agent=agent, input=input_list)
287288

288289
if span:
290+
serialized_items = [
291+
item.raw_item.model_dump() if isinstance(item.raw_item, BaseModel) else item.raw_item
292+
for item in result.new_items
293+
]
289294
span.output = {
290-
"new_items": [
291-
item.raw_item.model_dump() if isinstance(item.raw_item, BaseModel) else item.raw_item
292-
for item in result.new_items
293-
],
295+
"new_items": serialized_items,
294296
"final_output": result.final_output,
295297
}
298+
lineage_refs = resolve_refs_from_items(serialized_items)
299+
if lineage_refs:
300+
span.data = merge_refs_into_data(span.data, lineage_refs)
296301

297302
return result
298303

@@ -431,13 +436,17 @@ async def run_agent_auto_send(
431436
result = await Runner.run(starting_agent=agent, input=input_list)
432437

433438
if span:
439+
serialized_items = [
440+
item.raw_item.model_dump() if isinstance(item.raw_item, BaseModel) else item.raw_item
441+
for item in result.new_items
442+
]
434443
span.output = {
435-
"new_items": [
436-
item.raw_item.model_dump() if isinstance(item.raw_item, BaseModel) else item.raw_item
437-
for item in result.new_items
438-
],
444+
"new_items": serialized_items,
439445
"final_output": result.final_output,
440446
}
447+
lineage_refs = resolve_refs_from_items(serialized_items)
448+
if lineage_refs:
449+
span.data = merge_refs_into_data(span.data, lineage_refs)
441450

442451
tool_call_map: dict[str, Any] = {}
443452

@@ -646,13 +655,17 @@ async def run_agent_streamed(
646655
result = Runner.run_streamed(starting_agent=agent, input=input_list)
647656

648657
if span:
658+
serialized_items = [
659+
item.raw_item.model_dump() if isinstance(item.raw_item, BaseModel) else item.raw_item
660+
for item in result.new_items
661+
]
649662
span.output = {
650-
"new_items": [
651-
item.raw_item.model_dump() if isinstance(item.raw_item, BaseModel) else item.raw_item
652-
for item in result.new_items
653-
],
663+
"new_items": serialized_items,
654664
"final_output": result.final_output,
655665
}
666+
lineage_refs = resolve_refs_from_items(serialized_items)
667+
if lineage_refs:
668+
span.data = merge_refs_into_data(span.data, lineage_refs)
656669

657670
return result
658671

@@ -906,12 +919,16 @@ async def run_agent_streamed_auto_send(
906919
raise
907920

908921
if span:
922+
serialized_items = [
923+
item.raw_item.model_dump() if isinstance(item.raw_item, BaseModel) else item.raw_item
924+
for item in result.new_items
925+
]
909926
span.output = {
910-
"new_items": [
911-
item.raw_item.model_dump() if isinstance(item.raw_item, BaseModel) else item.raw_item
912-
for item in result.new_items
913-
],
927+
"new_items": serialized_items,
914928
"final_output": result.final_output,
915929
}
930+
lineage_refs = resolve_refs_from_items(serialized_items)
931+
if lineage_refs:
932+
span.data = merge_refs_into_data(span.data, lineage_refs)
916933

917934
return result

src/agentex/lib/core/temporal/plugins/openai_agents/models/temporal_streaming_model.py

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,7 @@
6464
from agentex.lib import adk
6565
from agentex.lib.utils.logging import make_logger
6666
from agentex.lib.core.tracing.tracer import AsyncTracer
67+
from agentex.lib.core.tracing.lineage import merge_refs_into_data, resolve_refs_from_items
6768
from agentex.types.task_message_delta import TextDelta, ToolRequestDelta, ReasoningContentDelta, ReasoningSummaryDelta
6869
from agentex.types.task_message_update import StreamTaskMessageFull, StreamTaskMessageDelta
6970
from agentex.types.task_message_content import TextContent, ReasoningContent, ToolRequestContent, ToolResponseContent
@@ -1257,6 +1258,9 @@ async def get_response(
12571258
output_data["tool_outputs"] = tool_outputs
12581259

12591260
span.output = output_data
1261+
lineage_refs = resolve_refs_from_items(new_items)
1262+
if lineage_refs:
1263+
span.data = merge_refs_into_data(span.data, lineage_refs)
12601264

12611265
# Streaming-only metrics. Token counters and the success request
12621266
# counter are emitted by LLMMetricsHooks.on_llm_end so they fire
Lines changed: 170 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,170 @@
1+
"""Data-source reference capture for lineage: tools declare which sources they
2+
touch and the refs land in span data under the ``sgp.lineage.refs`` key."""
3+
4+
from __future__ import annotations
5+
6+
import re
7+
import json
8+
from typing import Any, Literal, Callable, Iterable
9+
10+
from pydantic import Field, BaseModel, field_validator
11+
12+
try:
13+
from agentex.lib.utils.logging import make_logger
14+
15+
logger = make_logger(__name__)
16+
except Exception: # ddtrace may be absent in some envs; fall back to stdlib
17+
import logging
18+
19+
logger = logging.getLogger(__name__)
20+
21+
LINEAGE_REFS_KEY = "sgp.lineage.refs"
22+
23+
_URI_NAMESPACE_RE = re.compile(r"^[a-z][a-z0-9+.-]*://.+$")
24+
25+
RefResolver = Callable[[dict[str, Any]], "list[DataSourceRef]"]
26+
27+
28+
class DataSourceRef(BaseModel):
29+
"""One data source a tool call touched, as a lineage coordinate."""
30+
31+
namespace: str = Field(max_length=512)
32+
name: str = Field(min_length=1, max_length=512)
33+
version: str | None = Field(default=None, max_length=256)
34+
role: Literal["input", "output"] = "input"
35+
36+
def __init__(self, namespace: str | None = None, name: str | None = None, **kwargs: Any) -> None:
37+
if namespace is not None:
38+
kwargs["namespace"] = namespace
39+
if name is not None:
40+
kwargs["name"] = name
41+
super().__init__(**kwargs)
42+
43+
@field_validator("namespace")
44+
@classmethod
45+
def _namespace_is_uri_form(cls, value: str) -> str:
46+
if not _URI_NAMESPACE_RE.match(value):
47+
raise ValueError(f"namespace must be URI-form (scheme://system), got: {value!r}")
48+
return value
49+
50+
51+
class _ToolSources(BaseModel):
52+
refs: list[DataSourceRef] = Field(default_factory=list)
53+
resolver: RefResolver | None = None
54+
55+
model_config = {"arbitrary_types_allowed": True}
56+
57+
58+
_tool_sources: dict[str, _ToolSources] = {}
59+
60+
61+
def register_tool_sources(
62+
tool_name: str,
63+
refs: Iterable[DataSourceRef] | None = None,
64+
resolver: RefResolver | None = None,
65+
) -> None:
66+
"""Declare the data sources a tool touches, keyed by its tool name.
67+
68+
Use for tools the agent does not own (e.g. MCP proxy tools). Static refs and
69+
a resolver over the tool's parsed arguments may be combined; repeated
70+
registration for the same name replaces the prior entry.
71+
"""
72+
_tool_sources[tool_name] = _ToolSources(refs=list(refs or []), resolver=resolver)
73+
74+
75+
def data_sources(*refs: DataSourceRef, resolver: RefResolver | None = None) -> Callable[[Any], Any]:
76+
"""Decorator form of ``register_tool_sources`` for tools the agent owns.
77+
78+
Works below or above ``@function_tool``: the tool name is taken from the
79+
decorated object's ``name`` attribute when present, else ``__name__``.
80+
"""
81+
82+
def _register(obj: Any) -> Any:
83+
tool_name = getattr(obj, "name", None) or getattr(obj, "__name__", None)
84+
if isinstance(tool_name, str) and tool_name:
85+
register_tool_sources(tool_name, refs=refs, resolver=resolver)
86+
else:
87+
logger.warning("data_sources could not determine a tool name for %r; refs not registered", obj)
88+
return obj
89+
90+
return _register
91+
92+
93+
def clear_tool_sources() -> None:
94+
"""Reset the registry (test isolation)."""
95+
_tool_sources.clear()
96+
97+
98+
def resolve_refs(tool_name: str, arguments: dict[str, Any] | None) -> list[dict[str, Any]]:
99+
"""Resolve registered refs for one tool call to serialized, deduplicated dicts.
100+
101+
Resolver failures are logged and swallowed: ref capture must never break a
102+
tool call or its tracing.
103+
"""
104+
entry = _tool_sources.get(tool_name)
105+
if entry is None:
106+
return []
107+
refs = list(entry.refs)
108+
if entry.resolver is not None:
109+
try:
110+
refs.extend(entry.resolver(arguments or {}))
111+
except Exception:
112+
logger.warning("data-source resolver for tool %s failed; static refs kept", tool_name, exc_info=True)
113+
return _dedupe(refs)
114+
115+
116+
def resolve_refs_from_items(items: Iterable[Any]) -> list[dict[str, Any]]:
117+
"""Resolve refs across serialized run items, matching ``function_call`` entries.
118+
119+
Accepts the item dicts the providers already build for span output; string
120+
``arguments`` are parsed as JSON for resolver-based registrations.
121+
"""
122+
refs: list[dict[str, Any]] = []
123+
for item in items:
124+
if not isinstance(item, dict) or item.get("type") != "function_call":
125+
continue
126+
tool_name = item.get("name")
127+
if not isinstance(tool_name, str) or not tool_name:
128+
continue
129+
arguments = item.get("arguments")
130+
if isinstance(arguments, str):
131+
try:
132+
arguments = json.loads(arguments)
133+
except (ValueError, TypeError):
134+
arguments = {}
135+
refs.extend(resolve_refs(tool_name, arguments if isinstance(arguments, dict) else {}))
136+
return _dedupe_dicts(refs)
137+
138+
139+
def record(span: Any, refs: Iterable[DataSourceRef]) -> None:
140+
"""Attach refs to a manually managed span (no-op when the span is None)."""
141+
if span is None:
142+
return
143+
merged = merge_refs_into_data(getattr(span, "data", None), _dedupe(list(refs)))
144+
span.data = merged
145+
146+
147+
def merge_refs_into_data(data: dict[str, Any] | None, refs: list[dict[str, Any]]) -> dict[str, Any]:
148+
"""Merge serialized refs into a span data dict, deduplicating with any present."""
149+
out = dict(data) if isinstance(data, dict) else {}
150+
if refs:
151+
existing = out.get(LINEAGE_REFS_KEY)
152+
combined = list(existing) if isinstance(existing, list) else []
153+
combined.extend(refs)
154+
out[LINEAGE_REFS_KEY] = _dedupe_dicts(combined)
155+
return out
156+
157+
158+
def _dedupe(refs: list[DataSourceRef]) -> list[dict[str, Any]]:
159+
return _dedupe_dicts([ref.model_dump(exclude_none=True) for ref in refs])
160+
161+
162+
def _dedupe_dicts(refs: list[dict[str, Any]]) -> list[dict[str, Any]]:
163+
seen: set[tuple[Any, ...]] = set()
164+
out: list[dict[str, Any]] = []
165+
for ref in refs:
166+
key = (ref.get("namespace"), ref.get("name"), ref.get("version"), ref.get("role"))
167+
if key not in seen:
168+
seen.add(key)
169+
out.append(ref)
170+
return out

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

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -65,6 +65,8 @@ def _add_source_to_span(span: Span, env_vars: EnvironmentVariables) -> None:
6565
span.data["__agent_name__"] = env_vars.AGENT_NAME
6666
if env_vars.AGENT_ID is not None:
6767
span.data["__agent_id__"] = env_vars.AGENT_ID
68+
if env_vars.AGENT_VERSION is not None:
69+
span.data["__agent_version__"] = env_vars.AGENT_VERSION
6870

6971

7072
def _build_sgp_span(span: Span, env_vars: EnvironmentVariables) -> SGPSpan:

src/agentex/lib/environment_variables.py

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@ class EnvVarKeys(str, Enum):
2424
AGENT_NAME = "AGENT_NAME"
2525
AGENT_DESCRIPTION = "AGENT_DESCRIPTION"
2626
AGENT_ID = "AGENT_ID"
27+
AGENT_VERSION = "AGENT_VERSION"
2728
AGENT_API_KEY = "AGENT_API_KEY"
2829
# ACP Configuration
2930
ACP_URL = "ACP_URL"
@@ -66,6 +67,8 @@ class EnvironmentVariables(BaseModel):
6667
AGENT_NAME: str
6768
AGENT_DESCRIPTION: str | None = None
6869
AGENT_ID: str | None = None
70+
# Build/version discriminator (image tag or git sha), set by the deployment
71+
AGENT_VERSION: str | None = None
6972
AGENT_API_KEY: str | None = None
7073
ACP_TYPE: str | None = "async"
7174
AGENT_INPUT_TYPE: str | None = None

0 commit comments

Comments
 (0)