Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions .stats.yml
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
configured_endpoints: 75
openapi_spec_url: https://storage.googleapis.com/stainless-sdk-openapi-specs/sgp/agentex-sdk-f639b00012e6e4329ead09cceb032a38f83db74efb88a0890505b01011942606.yml
openapi_spec_hash: 1fef57cdbefab1119cdb5b81a261f06d
openapi_spec_url: https://storage.googleapis.com/stainless-sdk-openapi-specs/sgp/agentex-sdk-7acaeb315af90255109ae17afc71e32a8e5851bb8a956a2a284cb4d344dfab51.yml
openapi_spec_hash: 3044e94b48d60311b6048e8df88e7552
config_hash: 593e89b291976a5e84e4c3c3f8324354
8 changes: 8 additions & 0 deletions src/agentex/lib/adk/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,10 @@
from agentex.lib.adk._modules.tasks import TasksModule
from agentex.lib.adk._modules.tracing import TracingModule, TurnSpan

# Data-source refs for lineage (SGP-6513); implementation lives in core.tracing
from agentex.lib.core.tracing import lineage
from agentex.lib.core.tracing.lineage import DataSourceRef, data_sources

# Unified harness surface (AGX1-375)
from agentex.lib.core.harness import (
UnifiedEmitter,
Expand Down Expand Up @@ -67,6 +71,10 @@
"events",
"agent_task_tracker",
"TurnSpan",
# Lineage data-source refs (SGP-6513)
"lineage",
"DataSourceRef",
"data_sources",
# Checkpointing / LangGraph
"create_checkpointer",
"stream_langgraph_events",
Expand Down
7 changes: 7 additions & 0 deletions src/agentex/lib/adk/providers/_modules/sync_provider.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
from agentex import AsyncAgentex
from agentex.lib.utils.logging import make_logger
from agentex.lib.core.tracing.tracer import AsyncTracer
from agentex.lib.core.tracing.lineage import merge_refs_into_data, resolve_refs_from_items

logger = make_logger(__name__)

Expand Down Expand Up @@ -185,6 +186,9 @@ async def get_response(
"new_items": new_items,
"final_output": final_output,
}
lineage_refs = resolve_refs_from_items(new_items)
if lineage_refs:
span.data = merge_refs_into_data(span.data, lineage_refs)

return response
else:
Expand Down Expand Up @@ -303,6 +307,9 @@ async def stream_response(
"new_items": new_items,
"final_output": final_response_text if final_response_text else None,
}
lineage_refs = resolve_refs_from_items(new_items)
if lineage_refs:
span.data = merge_refs_into_data(span.data, lineage_refs)
finally:
# End the span after all events have been yielded
await trace.end_span(span)
Expand Down
16 changes: 16 additions & 0 deletions src/agentex/lib/core/harness/tracer.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,17 @@

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

try:
from agentex.lib.core.tracing.lineage import resolve_refs, merge_refs_into_data
except Exception: # keep the harness importable without optional tracing deps

def resolve_refs(tool_name: str, arguments: dict[str, Any] | None) -> list[dict[str, Any]]: # noqa: ARG001
return []

def merge_refs_into_data(data: dict[str, Any] | None, refs: list[dict[str, Any]]) -> dict[str, Any]: # noqa: ARG001
return dict(data or {})


try:
from agentex.lib.utils.logging import make_logger

Expand Down Expand Up @@ -80,6 +91,11 @@ async def handle(self, signal: SpanSignal) -> None:
task_id=self.task_id,
)
if span is not None:
if signal.kind == "tool":
refs = resolve_refs(signal.name, signal.input if isinstance(signal.input, dict) else {})
if refs:
data = span.data if isinstance(span.data, dict) else {}
span.data = merge_refs_into_data(data, refs)
self._open[signal.key] = span
elif isinstance(signal, CloseSpan):
span = self._open.pop(signal.key, None)
Expand Down
49 changes: 33 additions & 16 deletions src/agentex/lib/core/services/adk/providers/openai.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@
from agentex.lib.utils.temporal import heartbeat_if_in_workflow
from agentex.lib.core.tracing.tracer import AsyncTracer
from agentex.lib.core.harness.emitter import UnifiedEmitter
from agentex.lib.core.tracing.lineage import merge_refs_into_data, resolve_refs_from_items
from agentex.types.task_message_update import StreamTaskMessageFull
from agentex.types.task_message_content import (
TextContent,
Expand Down Expand Up @@ -286,13 +287,17 @@ async def run_agent(
result = await Runner.run(starting_agent=agent, input=input_list)

if span:
serialized_items = [
item.raw_item.model_dump() if isinstance(item.raw_item, BaseModel) else item.raw_item
for item in result.new_items
]
span.output = {
"new_items": [
item.raw_item.model_dump() if isinstance(item.raw_item, BaseModel) else item.raw_item
for item in result.new_items
],
"new_items": serialized_items,
"final_output": result.final_output,
}
lineage_refs = resolve_refs_from_items(serialized_items)
if lineage_refs:
span.data = merge_refs_into_data(span.data, lineage_refs)

return result

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

if span:
serialized_items = [
item.raw_item.model_dump() if isinstance(item.raw_item, BaseModel) else item.raw_item
for item in result.new_items
]
span.output = {
"new_items": [
item.raw_item.model_dump() if isinstance(item.raw_item, BaseModel) else item.raw_item
for item in result.new_items
],
"new_items": serialized_items,
"final_output": result.final_output,
}
lineage_refs = resolve_refs_from_items(serialized_items)
if lineage_refs:
span.data = merge_refs_into_data(span.data, lineage_refs)

tool_call_map: dict[str, Any] = {}

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

if span:
serialized_items = [
item.raw_item.model_dump() if isinstance(item.raw_item, BaseModel) else item.raw_item
for item in result.new_items
]
span.output = {
"new_items": [
item.raw_item.model_dump() if isinstance(item.raw_item, BaseModel) else item.raw_item
for item in result.new_items
],
"new_items": serialized_items,
"final_output": result.final_output,
}
lineage_refs = resolve_refs_from_items(serialized_items)
if lineage_refs:
span.data = merge_refs_into_data(span.data, lineage_refs)

return result

Expand Down Expand Up @@ -906,12 +919,16 @@ async def run_agent_streamed_auto_send(
raise

if span:
serialized_items = [
item.raw_item.model_dump() if isinstance(item.raw_item, BaseModel) else item.raw_item
for item in result.new_items
]
span.output = {
"new_items": [
item.raw_item.model_dump() if isinstance(item.raw_item, BaseModel) else item.raw_item
for item in result.new_items
],
"new_items": serialized_items,
"final_output": result.final_output,
}
lineage_refs = resolve_refs_from_items(serialized_items)
if lineage_refs:
span.data = merge_refs_into_data(span.data, lineage_refs)

return result
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,7 @@
from agentex.lib import adk
from agentex.lib.utils.logging import make_logger
from agentex.lib.core.tracing.tracer import AsyncTracer
from agentex.lib.core.tracing.lineage import merge_refs_into_data, resolve_refs_from_items
from agentex.types.task_message_delta import TextDelta, ToolRequestDelta, ReasoningContentDelta, ReasoningSummaryDelta
from agentex.types.task_message_update import StreamTaskMessageFull, StreamTaskMessageDelta
from agentex.types.task_message_content import TextContent, ReasoningContent, ToolRequestContent, ToolResponseContent
Expand Down Expand Up @@ -1257,6 +1258,9 @@ async def get_response(
output_data["tool_outputs"] = tool_outputs

span.output = output_data
lineage_refs = resolve_refs_from_items(new_items)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Input tool lineage is dropped

When a Temporal model invocation receives a registered function_call in its input and the response does not repeat it in new_items, lineage resolution examines only new_items, causing the span to record the call in tool_calls while omitting its registered data-source references.

Suggested change
lineage_refs = resolve_refs_from_items(new_items)
lineage_refs = resolve_refs_from_items([*new_items, *tool_calls])

Knowledge Base Used:

Prompt To Fix With AI
This is a comment left during a code review.
Path: src/agentex/lib/core/temporal/plugins/openai_agents/models/temporal_streaming_model.py
Line: 1261

Comment:
**Input tool lineage is dropped**

When a Temporal model invocation receives a registered `function_call` in its input and the response does not repeat it in `new_items`, lineage resolution examines only `new_items`, causing the span to record the call in `tool_calls` while omitting its registered data-source references.

```suggestion
                    lineage_refs = resolve_refs_from_items([*new_items, *tool_calls])
```

**Knowledge Base Used:**
- [Agent framework integrations](https://app.greptile.com/scale-ai/-/custom-context/knowledge-base/scaleapi/scale-agentex-python/-/docs/agent-framework-integrations.md)
- [Tracing pipeline](https://app.greptile.com/scale-ai/-/custom-context/knowledge-base/scaleapi/scale-agentex-python/-/docs/tracing-pipeline.md)

---

For each issue above, determine whether it is valid and should be fixed. If so, fix it directly.

Fix in Cursor Fix in Claude Code Fix in Codex

if lineage_refs:
span.data = merge_refs_into_data(span.data, lineage_refs)

# Streaming-only metrics. Token counters and the success request
# counter are emitted by LLMMetricsHooks.on_llm_end so they fire
Expand Down
9 changes: 8 additions & 1 deletion src/agentex/lib/core/temporal/workers/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,7 @@ def _validate_interceptors(interceptors: list) -> None:
async def get_temporal_client(
temporal_address: str,
metrics_url: str | None = None,
metrics_headers: dict[str, str] | None = None,
plugins: list = [],
payload_codec: PayloadCodec | None = None,
data_converter: DataConverter | None = None,
Expand Down Expand Up @@ -143,7 +144,10 @@ async def get_temporal_client(
if not metrics_url:
client = await Client.connect(**connect_kwargs)
else:
runtime = Runtime(telemetry=TelemetryConfig(metrics=OpenTelemetryConfig(url=metrics_url)))
runtime = Runtime(telemetry=TelemetryConfig(metrics=OpenTelemetryConfig(
url=metrics_url,
headers=metrics_headers or {},
)))
connect_kwargs["runtime"] = runtime
client = await Client.connect(**connect_kwargs)
return client
Expand All @@ -159,6 +163,7 @@ def __init__(
plugins: list = [],
interceptors: list = [],
metrics_url: str | None = None,
metrics_headers: dict[str, str] | None = None,
payload_codec: PayloadCodec | None = None,
data_converter: DataConverter | None = None,
):
Expand All @@ -174,6 +179,7 @@ def __init__(
self.plugins = plugins
self.interceptors = interceptors
self.metrics_url = metrics_url
self.metrics_headers = metrics_headers
self.payload_codec = payload_codec
self.data_converter = data_converter

Expand Down Expand Up @@ -211,6 +217,7 @@ async def run(
temporal_address=os.environ.get("TEMPORAL_ADDRESS", "localhost:7233"),
plugins=self.plugins,
metrics_url=self.metrics_url,
metrics_headers=self.metrics_headers,
payload_codec=self.payload_codec,
data_converter=self.data_converter,
)
Expand Down
Loading
Loading