mirror of
https://github.com/langgenius/dify.git
synced 2026-09-24 23:22:26 +08:00
feat(trace-aliyun): GenAI spans for LLM TTFT, agent ReAct, tools, and failed LLM nodes (#39339)
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -1,6 +1,6 @@
|
||||
import logging
|
||||
from collections.abc import Sequence
|
||||
from typing import override
|
||||
from collections.abc import Mapping, Sequence
|
||||
from typing import Any, override
|
||||
|
||||
from opentelemetry.trace import SpanKind
|
||||
from sqlalchemy.orm import sessionmaker
|
||||
@@ -29,33 +29,50 @@ from dify_trace_aliyun.data_exporter.traceclient import (
|
||||
from dify_trace_aliyun.entities.aliyun_trace_entity import SpanData, TraceMetadata
|
||||
from dify_trace_aliyun.entities.semconv import (
|
||||
DIFY_APP_ID,
|
||||
GEN_AI_AGENT_NAME,
|
||||
GEN_AI_COMPLETION,
|
||||
GEN_AI_INPUT_MESSAGE,
|
||||
GEN_AI_OPERATION_NAME,
|
||||
GEN_AI_OUTPUT_MESSAGE,
|
||||
GEN_AI_PROMPT,
|
||||
GEN_AI_PROVIDER_NAME,
|
||||
GEN_AI_REACT_FINISH_REASON,
|
||||
GEN_AI_REACT_ROUND,
|
||||
GEN_AI_REQUEST_MODEL,
|
||||
GEN_AI_RESPONSE_FINISH_REASON,
|
||||
GEN_AI_RESPONSE_TIME_TO_FIRST_TOKEN,
|
||||
GEN_AI_USAGE_INPUT_TOKENS,
|
||||
GEN_AI_USAGE_OUTPUT_TOKENS,
|
||||
GEN_AI_USAGE_TOTAL_TOKENS,
|
||||
OPERATION_NAME_CHAT,
|
||||
OPERATION_NAME_INVOKE_AGENT,
|
||||
OPERATION_NAME_REACT,
|
||||
RETRIEVAL_DOCUMENT,
|
||||
RETRIEVAL_QUERY,
|
||||
TOOL_DESCRIPTION,
|
||||
TOOL_NAME,
|
||||
TOOL_PARAMETERS,
|
||||
GenAISpanKind,
|
||||
)
|
||||
from dify_trace_aliyun.utils import (
|
||||
AgentLogEntry,
|
||||
convert_seconds_to_nanoseconds,
|
||||
create_common_span_attributes,
|
||||
create_gen_ai_tool_attributes,
|
||||
create_links_from_trace_id,
|
||||
create_status_from_agent_log_entry,
|
||||
create_status_from_error,
|
||||
extract_model_name_from_thought_label,
|
||||
extract_react_round_number,
|
||||
extract_retrieval_documents,
|
||||
extract_tool_description,
|
||||
extract_tool_name_from_call_label,
|
||||
format_input_messages,
|
||||
format_output_messages,
|
||||
format_retrieval_documents,
|
||||
get_user_id_from_message_data,
|
||||
get_workflow_node_status,
|
||||
is_llm_thought_entry,
|
||||
is_tool_call_entry,
|
||||
map_gen_ai_tool_type,
|
||||
parse_agent_log_entries,
|
||||
serialize_json_data,
|
||||
)
|
||||
from extensions.ext_database import db
|
||||
@@ -130,6 +147,9 @@ class AliyunDataTrace(BaseTraceInstance):
|
||||
for node_execution in workflow_node_executions:
|
||||
node_span = self.build_workflow_node_span(node_execution, trace_info, trace_metadata)
|
||||
self.trace_client.add_span(node_span)
|
||||
if node_span is not None and node_execution.node_type == BuiltinNodeTypes.AGENT:
|
||||
for react_span in self.build_agent_react_spans(node_execution, trace_metadata):
|
||||
self.trace_client.add_span(react_span)
|
||||
|
||||
def message_trace(self, trace_info: MessageTraceInfo):
|
||||
message_data = trace_info.message_data
|
||||
@@ -175,6 +195,28 @@ class AliyunDataTrace(BaseTraceInstance):
|
||||
)
|
||||
self.trace_client.add_span(message_span)
|
||||
|
||||
llm_attributes: dict[str, Any] = {
|
||||
**create_common_span_attributes(
|
||||
session_id=trace_metadata.session_id,
|
||||
user_id=trace_metadata.user_id,
|
||||
span_kind=GenAISpanKind.LLM,
|
||||
inputs=inputs_json,
|
||||
outputs=outputs_str,
|
||||
),
|
||||
GEN_AI_OPERATION_NAME: OPERATION_NAME_CHAT,
|
||||
GEN_AI_REQUEST_MODEL: trace_info.metadata.get("ls_model_name") or "",
|
||||
GEN_AI_PROVIDER_NAME: trace_info.metadata.get("ls_provider") or "",
|
||||
GEN_AI_USAGE_INPUT_TOKENS: str(trace_info.message_tokens),
|
||||
GEN_AI_USAGE_OUTPUT_TOKENS: str(trace_info.answer_tokens),
|
||||
GEN_AI_USAGE_TOTAL_TOKENS: str(trace_info.total_tokens),
|
||||
GEN_AI_PROMPT: inputs_json,
|
||||
GEN_AI_COMPLETION: outputs_str,
|
||||
}
|
||||
if trace_info.gen_ai_server_time_to_first_token is not None:
|
||||
llm_attributes[GEN_AI_RESPONSE_TIME_TO_FIRST_TOKEN] = convert_seconds_to_nanoseconds(
|
||||
trace_info.gen_ai_server_time_to_first_token
|
||||
)
|
||||
|
||||
llm_span = SpanData(
|
||||
trace_id=trace_metadata.trace_id,
|
||||
parent_span_id=message_span_id,
|
||||
@@ -182,22 +224,7 @@ class AliyunDataTrace(BaseTraceInstance):
|
||||
name="llm",
|
||||
start_time=convert_datetime_to_nanoseconds(trace_info.start_time),
|
||||
end_time=convert_datetime_to_nanoseconds(trace_info.end_time),
|
||||
attributes={
|
||||
**create_common_span_attributes(
|
||||
session_id=trace_metadata.session_id,
|
||||
user_id=trace_metadata.user_id,
|
||||
span_kind=GenAISpanKind.LLM,
|
||||
inputs=inputs_json,
|
||||
outputs=outputs_str,
|
||||
),
|
||||
GEN_AI_REQUEST_MODEL: trace_info.metadata.get("ls_model_name") or "",
|
||||
GEN_AI_PROVIDER_NAME: trace_info.metadata.get("ls_provider") or "",
|
||||
GEN_AI_USAGE_INPUT_TOKENS: str(trace_info.message_tokens),
|
||||
GEN_AI_USAGE_OUTPUT_TOKENS: str(trace_info.answer_tokens),
|
||||
GEN_AI_USAGE_TOTAL_TOKENS: str(trace_info.total_tokens),
|
||||
GEN_AI_PROMPT: inputs_json,
|
||||
GEN_AI_COMPLETION: outputs_str,
|
||||
},
|
||||
attributes=llm_attributes,
|
||||
status=status,
|
||||
links=trace_metadata.links,
|
||||
)
|
||||
@@ -258,9 +285,11 @@ class AliyunDataTrace(BaseTraceInstance):
|
||||
links=create_links_from_trace_id(trace_info.trace_id),
|
||||
)
|
||||
|
||||
tool_config_json = serialize_json_data(trace_info.tool_config)
|
||||
tool_config = trace_info.tool_config if isinstance(trace_info.tool_config, Mapping) else {}
|
||||
tool_inputs_json = serialize_json_data(trace_info.tool_inputs)
|
||||
tool_result = str(trace_info.tool_outputs)
|
||||
inputs_json = serialize_json_data(trace_info.inputs)
|
||||
provider_type = tool_config.get("tool_provider_type") or tool_config.get("provider_type")
|
||||
|
||||
tool_span = SpanData(
|
||||
trace_id=trace_metadata.trace_id,
|
||||
@@ -275,11 +304,16 @@ class AliyunDataTrace(BaseTraceInstance):
|
||||
user_id=trace_metadata.user_id,
|
||||
span_kind=GenAISpanKind.TOOL,
|
||||
inputs=inputs_json,
|
||||
outputs=str(trace_info.tool_outputs),
|
||||
outputs=tool_result,
|
||||
),
|
||||
**create_gen_ai_tool_attributes(
|
||||
tool_name=trace_info.tool_name,
|
||||
tool_type=map_gen_ai_tool_type(str(provider_type) if provider_type else None),
|
||||
tool_description=extract_tool_description(tool_config),
|
||||
tool_call_id=str(trace_info.metadata.get("node_execution_id") or ""),
|
||||
tool_call_arguments=tool_inputs_json,
|
||||
tool_call_result=tool_result,
|
||||
),
|
||||
TOOL_NAME: trace_info.tool_name,
|
||||
TOOL_DESCRIPTION: tool_config_json,
|
||||
TOOL_PARAMETERS: tool_inputs_json,
|
||||
},
|
||||
status=status,
|
||||
links=trace_metadata.links,
|
||||
@@ -316,6 +350,8 @@ class AliyunDataTrace(BaseTraceInstance):
|
||||
node_span = self.build_workflow_retrieval_span(trace_info, node_execution, trace_metadata)
|
||||
elif node_execution.node_type == BuiltinNodeTypes.TOOL:
|
||||
node_span = self.build_workflow_tool_span(trace_info, node_execution, trace_metadata)
|
||||
elif node_execution.node_type == BuiltinNodeTypes.AGENT:
|
||||
node_span = self.build_workflow_agent_span(trace_info, node_execution, trace_metadata)
|
||||
else:
|
||||
node_span = self.build_workflow_task_span(trace_info, node_execution, trace_metadata)
|
||||
return node_span
|
||||
@@ -349,12 +385,15 @@ class AliyunDataTrace(BaseTraceInstance):
|
||||
def build_workflow_tool_span(
|
||||
self, trace_info: WorkflowTraceInfo, node_execution: WorkflowNodeExecution, trace_metadata: TraceMetadata
|
||||
) -> SpanData:
|
||||
tool_des = {}
|
||||
if node_execution.metadata:
|
||||
tool_des = node_execution.metadata.get(WorkflowNodeExecutionMetadataKey.TOOL_INFO, {})
|
||||
tool_info: Mapping[str, Any] = {}
|
||||
if isinstance(node_execution.metadata, Mapping):
|
||||
raw_tool_info = node_execution.metadata.get(WorkflowNodeExecutionMetadataKey.TOOL_INFO, {})
|
||||
if isinstance(raw_tool_info, Mapping):
|
||||
tool_info = raw_tool_info
|
||||
|
||||
inputs_json = serialize_json_data(node_execution.inputs or {})
|
||||
outputs_json = serialize_json_data(node_execution.outputs)
|
||||
provider_type = tool_info.get("provider_type")
|
||||
|
||||
return SpanData(
|
||||
trace_id=trace_metadata.trace_id,
|
||||
@@ -371,9 +410,14 @@ class AliyunDataTrace(BaseTraceInstance):
|
||||
inputs=inputs_json,
|
||||
outputs=outputs_json,
|
||||
),
|
||||
TOOL_NAME: node_execution.title,
|
||||
TOOL_DESCRIPTION: serialize_json_data(tool_des),
|
||||
TOOL_PARAMETERS: inputs_json,
|
||||
**create_gen_ai_tool_attributes(
|
||||
tool_name=node_execution.title,
|
||||
tool_type=map_gen_ai_tool_type(str(provider_type) if provider_type else None),
|
||||
tool_description=extract_tool_description(tool_info),
|
||||
tool_call_id=str(node_execution.id or ""),
|
||||
tool_call_arguments=inputs_json,
|
||||
tool_call_result=outputs_json,
|
||||
),
|
||||
},
|
||||
status=get_workflow_node_status(node_execution),
|
||||
links=trace_metadata.links,
|
||||
@@ -414,16 +458,58 @@ class AliyunDataTrace(BaseTraceInstance):
|
||||
def build_workflow_llm_span(
|
||||
self, trace_info: WorkflowTraceInfo, node_execution: WorkflowNodeExecution, trace_metadata: TraceMetadata
|
||||
) -> SpanData:
|
||||
process_data = node_execution.process_data or {}
|
||||
outputs = node_execution.outputs or {}
|
||||
process_data = node_execution.process_data if isinstance(node_execution.process_data, Mapping) else {}
|
||||
inputs = node_execution.inputs if isinstance(node_execution.inputs, Mapping) else {}
|
||||
outputs = node_execution.outputs if isinstance(node_execution.outputs, Mapping) else {}
|
||||
usage_data = process_data.get("usage", {}) if "usage" in process_data else outputs.get("usage", {})
|
||||
if not isinstance(usage_data, Mapping):
|
||||
usage_data = {}
|
||||
|
||||
prompts_json = serialize_json_data(process_data.get("prompts", []))
|
||||
text_output = str(outputs.get("text", ""))
|
||||
# On invoke failure graphon leaves process_data empty, but prep already wrote
|
||||
# model identity (and template variables / context) into node inputs.
|
||||
prompts = process_data.get("prompts") or []
|
||||
if prompts:
|
||||
prompts_json = serialize_json_data(prompts)
|
||||
elif inputs:
|
||||
prompts_json = serialize_json_data(inputs)
|
||||
else:
|
||||
prompts_json = serialize_json_data([])
|
||||
|
||||
text_output = str(outputs.get("text") or "")
|
||||
if not text_output:
|
||||
text_output = str(outputs.get("error_message") or node_execution.error or "")
|
||||
|
||||
finish_reason = outputs.get("finish_reason") or outputs.get("error_type") or ""
|
||||
model_name = process_data.get("model_name") or inputs.get("model_name") or ""
|
||||
model_provider = process_data.get("model_provider") or inputs.get("model_provider") or ""
|
||||
|
||||
gen_ai_input_message = format_input_messages(process_data)
|
||||
gen_ai_output_message = format_output_messages(outputs)
|
||||
|
||||
attributes: dict[str, Any] = {
|
||||
**create_common_span_attributes(
|
||||
session_id=trace_metadata.session_id,
|
||||
user_id=trace_metadata.user_id,
|
||||
span_kind=GenAISpanKind.LLM,
|
||||
inputs=prompts_json,
|
||||
outputs=text_output,
|
||||
),
|
||||
GEN_AI_OPERATION_NAME: OPERATION_NAME_CHAT,
|
||||
GEN_AI_REQUEST_MODEL: str(model_name),
|
||||
GEN_AI_PROVIDER_NAME: str(model_provider),
|
||||
GEN_AI_USAGE_INPUT_TOKENS: str(usage_data.get("prompt_tokens", 0)),
|
||||
GEN_AI_USAGE_OUTPUT_TOKENS: str(usage_data.get("completion_tokens", 0)),
|
||||
GEN_AI_USAGE_TOTAL_TOKENS: str(usage_data.get("total_tokens", 0)),
|
||||
GEN_AI_PROMPT: prompts_json,
|
||||
GEN_AI_COMPLETION: text_output,
|
||||
GEN_AI_RESPONSE_FINISH_REASON: str(finish_reason),
|
||||
GEN_AI_INPUT_MESSAGE: gen_ai_input_message,
|
||||
GEN_AI_OUTPUT_MESSAGE: gen_ai_output_message,
|
||||
}
|
||||
time_to_first_token = usage_data.get("time_to_first_token")
|
||||
if isinstance(time_to_first_token, (int, float)):
|
||||
attributes[GEN_AI_RESPONSE_TIME_TO_FIRST_TOKEN] = convert_seconds_to_nanoseconds(float(time_to_first_token))
|
||||
|
||||
return SpanData(
|
||||
trace_id=trace_metadata.trace_id,
|
||||
parent_span_id=trace_metadata.workflow_span_id,
|
||||
@@ -431,26 +517,225 @@ class AliyunDataTrace(BaseTraceInstance):
|
||||
name=node_execution.title,
|
||||
start_time=convert_datetime_to_nanoseconds(node_execution.created_at),
|
||||
end_time=convert_datetime_to_nanoseconds(node_execution.finished_at),
|
||||
attributes=attributes,
|
||||
status=get_workflow_node_status(node_execution),
|
||||
links=trace_metadata.links,
|
||||
)
|
||||
|
||||
def build_workflow_agent_span(
|
||||
self, trace_info: WorkflowTraceInfo, node_execution: WorkflowNodeExecution, trace_metadata: TraceMetadata
|
||||
) -> SpanData:
|
||||
"""Build an AGENT-kind span for an agent-strategy node (instead of a generic TASK span)."""
|
||||
inputs_json = serialize_json_data(node_execution.inputs)
|
||||
outputs = node_execution.outputs if isinstance(node_execution.outputs, Mapping) else {}
|
||||
usage_data = outputs.get("usage", {})
|
||||
if not isinstance(usage_data, Mapping):
|
||||
usage_data = {}
|
||||
text_output = str(outputs.get("text", ""))
|
||||
|
||||
attributes: dict[str, Any] = {
|
||||
**create_common_span_attributes(
|
||||
session_id=trace_metadata.session_id,
|
||||
user_id=trace_metadata.user_id,
|
||||
span_kind=GenAISpanKind.AGENT,
|
||||
inputs=inputs_json,
|
||||
outputs=text_output,
|
||||
),
|
||||
GEN_AI_OPERATION_NAME: OPERATION_NAME_INVOKE_AGENT,
|
||||
GEN_AI_AGENT_NAME: node_execution.title,
|
||||
GEN_AI_USAGE_INPUT_TOKENS: str(usage_data.get("prompt_tokens", 0)),
|
||||
GEN_AI_USAGE_OUTPUT_TOKENS: str(usage_data.get("completion_tokens", 0)),
|
||||
GEN_AI_USAGE_TOTAL_TOKENS: str(usage_data.get("total_tokens", 0)),
|
||||
}
|
||||
time_to_first_token = usage_data.get("time_to_first_token")
|
||||
if isinstance(time_to_first_token, (int, float)):
|
||||
attributes[GEN_AI_RESPONSE_TIME_TO_FIRST_TOKEN] = convert_seconds_to_nanoseconds(float(time_to_first_token))
|
||||
|
||||
return SpanData(
|
||||
trace_id=trace_metadata.trace_id,
|
||||
parent_span_id=trace_metadata.workflow_span_id,
|
||||
span_id=convert_to_span_id(node_execution.id, "node"),
|
||||
name=node_execution.title,
|
||||
start_time=convert_datetime_to_nanoseconds(node_execution.created_at),
|
||||
end_time=convert_datetime_to_nanoseconds(node_execution.finished_at),
|
||||
attributes=attributes,
|
||||
status=get_workflow_node_status(node_execution),
|
||||
links=trace_metadata.links,
|
||||
)
|
||||
|
||||
def build_agent_react_spans(
|
||||
self, node_execution: WorkflowNodeExecution, trace_metadata: TraceMetadata
|
||||
) -> list[SpanData]:
|
||||
"""Build ReAct STEP spans and child LLM / TOOL spans from the agent execution log.
|
||||
|
||||
The agent log lives in ``outputs["json"]``; ``started_at``/``finished_at`` there are
|
||||
monotonic-clock seconds, so they are mapped onto wall-clock time by anchoring the
|
||||
earliest ``started_at`` to the node's start time. Entries without timing fall back
|
||||
to the node's start/end times. Returns an empty list when no log is available.
|
||||
"""
|
||||
try:
|
||||
outputs = node_execution.outputs or {}
|
||||
round_entries = parse_agent_log_entries(outputs)
|
||||
if not round_entries:
|
||||
return []
|
||||
|
||||
agent_span_id = convert_to_span_id(node_execution.id, "node")
|
||||
node_start_ns = convert_datetime_to_nanoseconds(node_execution.created_at)
|
||||
node_end_ns = convert_datetime_to_nanoseconds(node_execution.finished_at)
|
||||
|
||||
monotonic_starts = [
|
||||
entry.metadata["started_at"]
|
||||
for round_entry in round_entries
|
||||
for entry in [round_entry, *round_entry.children]
|
||||
if isinstance(entry.metadata.get("started_at"), (int, float))
|
||||
]
|
||||
base_monotonic = min(monotonic_starts) if monotonic_starts else None
|
||||
|
||||
def to_wall_clock_ns(monotonic_seconds: Any, fallback: int | None) -> int | None:
|
||||
if (
|
||||
isinstance(monotonic_seconds, (int, float))
|
||||
and base_monotonic is not None
|
||||
and node_start_ns is not None
|
||||
):
|
||||
return node_start_ns + convert_seconds_to_nanoseconds(float(monotonic_seconds) - base_monotonic)
|
||||
return fallback
|
||||
|
||||
spans: list[SpanData] = []
|
||||
for index, round_entry in enumerate(round_entries, start=1):
|
||||
round_number = extract_react_round_number(round_entry.label, index)
|
||||
step_span_id = generate_span_id()
|
||||
step_attributes: dict[str, Any] = {
|
||||
**create_common_span_attributes(
|
||||
session_id=trace_metadata.session_id,
|
||||
user_id=trace_metadata.user_id,
|
||||
span_kind=GenAISpanKind.STEP,
|
||||
inputs="",
|
||||
outputs=serialize_json_data(round_entry.data),
|
||||
),
|
||||
GEN_AI_OPERATION_NAME: OPERATION_NAME_REACT,
|
||||
GEN_AI_REACT_ROUND: round_number,
|
||||
}
|
||||
if round_entry.error:
|
||||
step_attributes[GEN_AI_REACT_FINISH_REASON] = "error"
|
||||
spans.append(
|
||||
SpanData(
|
||||
trace_id=trace_metadata.trace_id,
|
||||
parent_span_id=agent_span_id,
|
||||
span_id=step_span_id,
|
||||
name=round_entry.label or f"react step {round_number}",
|
||||
start_time=to_wall_clock_ns(round_entry.metadata.get("started_at"), node_start_ns),
|
||||
end_time=to_wall_clock_ns(round_entry.metadata.get("finished_at"), node_end_ns),
|
||||
attributes=step_attributes,
|
||||
status=create_status_from_agent_log_entry(round_entry),
|
||||
links=trace_metadata.links,
|
||||
)
|
||||
)
|
||||
|
||||
for child in round_entry.children:
|
||||
child_start = to_wall_clock_ns(child.metadata.get("started_at"), node_start_ns)
|
||||
child_end = to_wall_clock_ns(child.metadata.get("finished_at"), node_end_ns)
|
||||
if is_tool_call_entry(child):
|
||||
spans.append(
|
||||
self._build_agent_tool_call_span(
|
||||
entry=child,
|
||||
step_span_id=step_span_id,
|
||||
trace_metadata=trace_metadata,
|
||||
start_time=child_start,
|
||||
end_time=child_end,
|
||||
)
|
||||
)
|
||||
elif is_llm_thought_entry(child):
|
||||
spans.append(
|
||||
self._build_agent_llm_call_span(
|
||||
entry=child,
|
||||
step_span_id=step_span_id,
|
||||
trace_metadata=trace_metadata,
|
||||
start_time=child_start,
|
||||
end_time=child_end,
|
||||
)
|
||||
)
|
||||
return spans
|
||||
except Exception as e:
|
||||
logger.warning("Error occurred in build_agent_react_spans: %s", e, exc_info=True)
|
||||
return []
|
||||
|
||||
def _build_agent_llm_call_span(
|
||||
self,
|
||||
entry: AgentLogEntry,
|
||||
step_span_id: int,
|
||||
trace_metadata: TraceMetadata,
|
||||
start_time: int | None,
|
||||
end_time: int | None,
|
||||
) -> SpanData:
|
||||
completion = str(entry.data.get("thought") or entry.data.get("action") or "")
|
||||
return SpanData(
|
||||
trace_id=trace_metadata.trace_id,
|
||||
parent_span_id=step_span_id,
|
||||
span_id=generate_span_id(),
|
||||
name=entry.label or "llm",
|
||||
start_time=start_time,
|
||||
end_time=end_time,
|
||||
attributes={
|
||||
**create_common_span_attributes(
|
||||
session_id=trace_metadata.session_id,
|
||||
user_id=trace_metadata.user_id,
|
||||
span_kind=GenAISpanKind.LLM,
|
||||
inputs=prompts_json,
|
||||
outputs=text_output,
|
||||
inputs="",
|
||||
outputs=serialize_json_data(entry.data),
|
||||
),
|
||||
GEN_AI_REQUEST_MODEL: process_data.get("model_name") or "",
|
||||
GEN_AI_PROVIDER_NAME: process_data.get("model_provider") or "",
|
||||
GEN_AI_USAGE_INPUT_TOKENS: str(usage_data.get("prompt_tokens", 0)),
|
||||
GEN_AI_USAGE_OUTPUT_TOKENS: str(usage_data.get("completion_tokens", 0)),
|
||||
GEN_AI_USAGE_TOTAL_TOKENS: str(usage_data.get("total_tokens", 0)),
|
||||
GEN_AI_PROMPT: prompts_json,
|
||||
GEN_AI_COMPLETION: text_output,
|
||||
GEN_AI_RESPONSE_FINISH_REASON: outputs.get("finish_reason") or "",
|
||||
GEN_AI_INPUT_MESSAGE: gen_ai_input_message,
|
||||
GEN_AI_OUTPUT_MESSAGE: gen_ai_output_message,
|
||||
GEN_AI_OPERATION_NAME: OPERATION_NAME_CHAT,
|
||||
GEN_AI_REQUEST_MODEL: extract_model_name_from_thought_label(entry.label),
|
||||
GEN_AI_PROVIDER_NAME: str(entry.metadata.get("provider") or ""),
|
||||
GEN_AI_USAGE_TOTAL_TOKENS: str(entry.metadata.get("total_tokens", 0)),
|
||||
GEN_AI_COMPLETION: completion,
|
||||
},
|
||||
status=get_workflow_node_status(node_execution),
|
||||
status=create_status_from_agent_log_entry(entry),
|
||||
links=trace_metadata.links,
|
||||
)
|
||||
|
||||
def _build_agent_tool_call_span(
|
||||
self,
|
||||
entry: AgentLogEntry,
|
||||
step_span_id: int,
|
||||
trace_metadata: TraceMetadata,
|
||||
start_time: int | None,
|
||||
end_time: int | None,
|
||||
) -> SpanData:
|
||||
tool_name = str(entry.data.get("tool_name") or extract_tool_name_from_call_label(entry.label) or "tool")
|
||||
tool_parameters = entry.data.get("tool_call_args")
|
||||
if tool_parameters is None:
|
||||
tool_parameters = entry.data.get("tool_call_input")
|
||||
if tool_parameters is None:
|
||||
tool_parameters = entry.data
|
||||
tool_arguments_json = serialize_json_data(tool_parameters)
|
||||
tool_result = entry.data.get("output", entry.data)
|
||||
tool_result_json = tool_result if isinstance(tool_result, str) else serialize_json_data(tool_result)
|
||||
provider_type = entry.metadata.get("provider_type") or entry.data.get("provider_type")
|
||||
return SpanData(
|
||||
trace_id=trace_metadata.trace_id,
|
||||
parent_span_id=step_span_id,
|
||||
span_id=generate_span_id(),
|
||||
name=entry.label or f"CALL {tool_name}",
|
||||
start_time=start_time,
|
||||
end_time=end_time,
|
||||
attributes={
|
||||
**create_common_span_attributes(
|
||||
session_id=trace_metadata.session_id,
|
||||
user_id=trace_metadata.user_id,
|
||||
span_kind=GenAISpanKind.TOOL,
|
||||
inputs=tool_arguments_json,
|
||||
outputs=tool_result_json,
|
||||
),
|
||||
**create_gen_ai_tool_attributes(
|
||||
tool_name=tool_name,
|
||||
tool_type=map_gen_ai_tool_type(str(provider_type) if provider_type else None),
|
||||
tool_description=extract_tool_description(entry.data) or extract_tool_description(entry.metadata),
|
||||
tool_call_id=entry.id,
|
||||
tool_call_arguments=tool_arguments_json,
|
||||
tool_call_result=tool_result_json,
|
||||
),
|
||||
},
|
||||
status=create_status_from_agent_log_entry(entry),
|
||||
links=trace_metadata.links,
|
||||
)
|
||||
|
||||
|
||||
@@ -11,6 +11,7 @@ GEN_AI_SESSION_ID: Final[str] = "gen_ai.session.id"
|
||||
GEN_AI_USER_ID: Final[str] = "gen_ai.user.id"
|
||||
GEN_AI_USER_NAME: Final[str] = "gen_ai.user.name"
|
||||
GEN_AI_SPAN_KIND: Final[str] = "gen_ai.span.kind"
|
||||
GEN_AI_OPERATION_NAME: Final[str] = "gen_ai.operation.name"
|
||||
GEN_AI_FRAMEWORK: Final[str] = "gen_ai.framework"
|
||||
|
||||
# Chain attributes
|
||||
@@ -30,14 +31,43 @@ GEN_AI_USAGE_TOTAL_TOKENS: Final[str] = "gen_ai.usage.total_tokens"
|
||||
GEN_AI_PROMPT: Final[str] = "gen_ai.prompt"
|
||||
GEN_AI_COMPLETION: Final[str] = "gen_ai.completion"
|
||||
GEN_AI_RESPONSE_FINISH_REASON: Final[str] = "gen_ai.response.finish_reason"
|
||||
# Time to first token of the model response in a streaming scenario, in nanoseconds.
|
||||
GEN_AI_RESPONSE_TIME_TO_FIRST_TOKEN: Final[str] = "gen_ai.response.time_to_first_token"
|
||||
|
||||
GEN_AI_INPUT_MESSAGE: Final[str] = "gen_ai.input.messages"
|
||||
GEN_AI_OUTPUT_MESSAGE: Final[str] = "gen_ai.output.messages"
|
||||
|
||||
# Tool attributes
|
||||
TOOL_NAME: Final[str] = "tool.name"
|
||||
TOOL_DESCRIPTION: Final[str] = "tool.description"
|
||||
TOOL_PARAMETERS: Final[str] = "tool.parameters"
|
||||
# Tool attributes (GenAI semantic conventions)
|
||||
GEN_AI_TOOL_CALL_ID: Final[str] = "gen_ai.tool.call.id"
|
||||
GEN_AI_TOOL_DESCRIPTION: Final[str] = "gen_ai.tool.description"
|
||||
GEN_AI_TOOL_NAME: Final[str] = "gen_ai.tool.name"
|
||||
GEN_AI_TOOL_TYPE: Final[str] = "gen_ai.tool.type"
|
||||
GEN_AI_TOOL_CALL_ARGUMENTS: Final[str] = "gen_ai.tool.call.arguments"
|
||||
GEN_AI_TOOL_CALL_RESULT: Final[str] = "gen_ai.tool.call.result"
|
||||
|
||||
# Skill attributes (conditionally required when loading a Skill)
|
||||
GEN_AI_SKILL_ID: Final[str] = "gen_ai.skill.id"
|
||||
GEN_AI_SKILL_NAME: Final[str] = "gen_ai.skill.name"
|
||||
GEN_AI_SKILL_DESCRIPTION: Final[str] = "gen_ai.skill.description"
|
||||
GEN_AI_SKILL_VERSION: Final[str] = "gen_ai.skill.version"
|
||||
|
||||
# Agent attributes
|
||||
GEN_AI_AGENT_NAME: Final[str] = "gen_ai.agent.name"
|
||||
|
||||
# ReAct step attributes
|
||||
GEN_AI_REACT_ROUND: Final[str] = "gen_ai.react.round"
|
||||
GEN_AI_REACT_FINISH_REASON: Final[str] = "gen_ai.react.finish_reason"
|
||||
|
||||
# gen_ai.operation.name values (see Aliyun LLM Trace field definitions)
|
||||
OPERATION_NAME_CHAT: Final[str] = "chat"
|
||||
OPERATION_NAME_EXECUTE_TOOL: Final[str] = "execute_tool"
|
||||
OPERATION_NAME_INVOKE_AGENT: Final[str] = "invoke_agent"
|
||||
OPERATION_NAME_REACT: Final[str] = "react"
|
||||
|
||||
# gen_ai.tool.type values
|
||||
TOOL_TYPE_FUNCTION: Final[str] = "function"
|
||||
TOOL_TYPE_EXTENSION: Final[str] = "extension"
|
||||
TOOL_TYPE_DATASTORE: Final[str] = "datastore"
|
||||
|
||||
|
||||
class GenAISpanKind(StrEnum):
|
||||
@@ -49,3 +79,5 @@ class GenAISpanKind(StrEnum):
|
||||
TOOL = "TOOL"
|
||||
AGENT = "AGENT"
|
||||
TASK = "TASK"
|
||||
# Marks one Reasoning-Acting iteration of an agent (ReAct step).
|
||||
STEP = "STEP"
|
||||
|
||||
@@ -1,5 +1,7 @@
|
||||
import json
|
||||
import re
|
||||
from collections.abc import Mapping
|
||||
from dataclasses import dataclass, field
|
||||
from typing import Any, TypedDict
|
||||
|
||||
from opentelemetry.trace import Link, Status, StatusCode
|
||||
@@ -7,11 +9,22 @@ from opentelemetry.trace import Link, Status, StatusCode
|
||||
from core.rag.models.document import Document
|
||||
from dify_trace_aliyun.entities.semconv import (
|
||||
GEN_AI_FRAMEWORK,
|
||||
GEN_AI_OPERATION_NAME,
|
||||
GEN_AI_SESSION_ID,
|
||||
GEN_AI_SPAN_KIND,
|
||||
GEN_AI_TOOL_CALL_ARGUMENTS,
|
||||
GEN_AI_TOOL_CALL_ID,
|
||||
GEN_AI_TOOL_CALL_RESULT,
|
||||
GEN_AI_TOOL_DESCRIPTION,
|
||||
GEN_AI_TOOL_NAME,
|
||||
GEN_AI_TOOL_TYPE,
|
||||
GEN_AI_USER_ID,
|
||||
INPUT_VALUE,
|
||||
OPERATION_NAME_EXECUTE_TOOL,
|
||||
OUTPUT_VALUE,
|
||||
TOOL_TYPE_DATASTORE,
|
||||
TOOL_TYPE_EXTENSION,
|
||||
TOOL_TYPE_FUNCTION,
|
||||
GenAISpanKind,
|
||||
)
|
||||
from extensions.ext_database import db
|
||||
@@ -106,6 +119,47 @@ def create_common_span_attributes(
|
||||
}
|
||||
|
||||
|
||||
def map_gen_ai_tool_type(provider_type: str | None) -> str:
|
||||
"""Map Dify tool provider type to GenAI ``gen_ai.tool.type`` values."""
|
||||
normalized = (provider_type or "").strip().lower()
|
||||
if normalized in {"dataset-retrieval", "datastore"}:
|
||||
return TOOL_TYPE_DATASTORE
|
||||
if normalized == "extension":
|
||||
return TOOL_TYPE_EXTENSION
|
||||
return TOOL_TYPE_FUNCTION
|
||||
|
||||
|
||||
def extract_tool_description(tool_meta: Mapping[str, Any] | None) -> str:
|
||||
if not isinstance(tool_meta, Mapping):
|
||||
return ""
|
||||
for key in ("description", "tool_description"):
|
||||
value = tool_meta.get(key)
|
||||
if value:
|
||||
return str(value)
|
||||
return ""
|
||||
|
||||
|
||||
def create_gen_ai_tool_attributes(
|
||||
*,
|
||||
tool_name: str,
|
||||
tool_type: str = TOOL_TYPE_FUNCTION,
|
||||
tool_description: str = "",
|
||||
tool_call_id: str = "",
|
||||
tool_call_arguments: str = "",
|
||||
tool_call_result: str = "",
|
||||
) -> dict[str, str]:
|
||||
"""Build GenAI-compliant attributes for a TOOL span."""
|
||||
return {
|
||||
GEN_AI_OPERATION_NAME: OPERATION_NAME_EXECUTE_TOOL,
|
||||
GEN_AI_TOOL_NAME: tool_name,
|
||||
GEN_AI_TOOL_TYPE: tool_type or TOOL_TYPE_FUNCTION,
|
||||
GEN_AI_TOOL_DESCRIPTION: tool_description,
|
||||
GEN_AI_TOOL_CALL_ID: tool_call_id,
|
||||
GEN_AI_TOOL_CALL_ARGUMENTS: tool_call_arguments,
|
||||
GEN_AI_TOOL_CALL_RESULT: tool_call_result,
|
||||
}
|
||||
|
||||
|
||||
def format_retrieval_documents(retrieval_documents: list) -> list:
|
||||
try:
|
||||
if not isinstance(retrieval_documents, list):
|
||||
@@ -174,6 +228,121 @@ def format_input_messages(process_data: Mapping[str, Any]) -> str:
|
||||
return serialize_json_data([])
|
||||
|
||||
|
||||
def convert_seconds_to_nanoseconds(seconds: float) -> int:
|
||||
return int(seconds * 1e9)
|
||||
|
||||
|
||||
_REACT_ROUND_LABEL_PATTERN = re.compile(r"ROUND\s+(\d+)", re.IGNORECASE)
|
||||
_LLM_THOUGHT_LABEL_SUFFIX = " Thought"
|
||||
_TOOL_CALL_LABEL_PREFIX = "CALL "
|
||||
|
||||
|
||||
@dataclass
|
||||
class AgentLogEntry:
|
||||
"""One entry of the agent-strategy execution log (``outputs["json"]`` of an agent node).
|
||||
|
||||
Entries form a tree via ``parent_id``: top-level entries are ReAct rounds
|
||||
(label like ``ROUND 1``) and their children are LLM thoughts / tool calls.
|
||||
``started_at``/``finished_at`` in ``metadata`` are monotonic-clock seconds
|
||||
(``time.perf_counter``), not epoch timestamps.
|
||||
"""
|
||||
|
||||
id: str
|
||||
parent_id: str | None
|
||||
label: str
|
||||
status: str
|
||||
error: str | None
|
||||
data: dict[str, Any]
|
||||
metadata: dict[str, Any]
|
||||
children: list["AgentLogEntry"] = field(default_factory=list)
|
||||
|
||||
|
||||
def parse_agent_log_entries(outputs: Mapping[str, Any]) -> list[AgentLogEntry]:
|
||||
"""Parse agent node outputs into a tree of log entries, returning top-level rounds in order.
|
||||
|
||||
Entries without an ``id`` (e.g. the trailing ``{"data": []}`` element) are skipped.
|
||||
Children whose parent is missing are dropped.
|
||||
"""
|
||||
raw_entries = outputs.get("json")
|
||||
if not isinstance(raw_entries, list):
|
||||
return []
|
||||
|
||||
entries: list[AgentLogEntry] = []
|
||||
entries_by_id: dict[str, AgentLogEntry] = {}
|
||||
for raw_entry in raw_entries:
|
||||
if not isinstance(raw_entry, dict):
|
||||
continue
|
||||
entry_id = raw_entry.get("id")
|
||||
if not entry_id:
|
||||
continue
|
||||
data = raw_entry.get("data")
|
||||
metadata = raw_entry.get("metadata")
|
||||
entry = AgentLogEntry(
|
||||
id=str(entry_id),
|
||||
parent_id=raw_entry.get("parent_id"),
|
||||
label=str(raw_entry.get("label") or ""),
|
||||
status=str(raw_entry.get("status") or ""),
|
||||
error=raw_entry.get("error"),
|
||||
data=data if isinstance(data, dict) else {},
|
||||
metadata=metadata if isinstance(metadata, dict) else {},
|
||||
)
|
||||
entries.append(entry)
|
||||
entries_by_id[entry.id] = entry
|
||||
|
||||
roots: list[AgentLogEntry] = []
|
||||
for entry in entries:
|
||||
if entry.parent_id is None:
|
||||
roots.append(entry)
|
||||
else:
|
||||
parent = entries_by_id.get(entry.parent_id)
|
||||
if parent is not None:
|
||||
parent.children.append(entry)
|
||||
return roots
|
||||
|
||||
|
||||
def extract_react_round_number(label: str, fallback: int) -> int:
|
||||
match = _REACT_ROUND_LABEL_PATTERN.search(label)
|
||||
if match:
|
||||
return int(match.group(1))
|
||||
return fallback
|
||||
|
||||
|
||||
def is_tool_call_entry(entry: AgentLogEntry) -> bool:
|
||||
"""Tool invocations use labels like ``CALL {tool_name}`` (also carry a provider in metadata)."""
|
||||
return entry.label.startswith(_TOOL_CALL_LABEL_PREFIX)
|
||||
|
||||
|
||||
def is_llm_thought_entry(entry: AgentLogEntry) -> bool:
|
||||
"""LLM thought entries use labels like ``{model} Thought``.
|
||||
|
||||
Do not key off ``metadata.provider`` alone: tool CALL entries also set provider
|
||||
(to the tool provider), which previously misclassified them as LLM spans.
|
||||
"""
|
||||
if is_tool_call_entry(entry):
|
||||
return False
|
||||
return entry.label.endswith(_LLM_THOUGHT_LABEL_SUFFIX) or bool(entry.metadata.get("provider"))
|
||||
|
||||
|
||||
def extract_model_name_from_thought_label(label: str) -> str:
|
||||
if label.endswith(_LLM_THOUGHT_LABEL_SUFFIX):
|
||||
return label.removesuffix(_LLM_THOUGHT_LABEL_SUFFIX)
|
||||
return ""
|
||||
|
||||
|
||||
def extract_tool_name_from_call_label(label: str) -> str:
|
||||
if label.startswith(_TOOL_CALL_LABEL_PREFIX):
|
||||
return label.removeprefix(_TOOL_CALL_LABEL_PREFIX).strip()
|
||||
return ""
|
||||
|
||||
|
||||
def create_status_from_agent_log_entry(entry: AgentLogEntry) -> Status:
|
||||
if entry.error:
|
||||
return Status(StatusCode.ERROR, str(entry.error))
|
||||
if entry.status == "success":
|
||||
return Status(StatusCode.OK)
|
||||
return Status(StatusCode.UNSET)
|
||||
|
||||
|
||||
def format_output_messages(outputs: Mapping[str, Any]) -> str:
|
||||
try:
|
||||
if not isinstance(outputs, dict):
|
||||
|
||||
+40
-7
@@ -1,27 +1,43 @@
|
||||
from dify_trace_aliyun.entities.semconv import (
|
||||
ACS_ARMS_SERVICE_FEATURE,
|
||||
GEN_AI_AGENT_NAME,
|
||||
GEN_AI_COMPLETION,
|
||||
GEN_AI_FRAMEWORK,
|
||||
GEN_AI_INPUT_MESSAGE,
|
||||
GEN_AI_OPERATION_NAME,
|
||||
GEN_AI_OUTPUT_MESSAGE,
|
||||
GEN_AI_PROMPT,
|
||||
GEN_AI_PROVIDER_NAME,
|
||||
GEN_AI_REACT_FINISH_REASON,
|
||||
GEN_AI_REACT_ROUND,
|
||||
GEN_AI_REQUEST_MODEL,
|
||||
GEN_AI_RESPONSE_FINISH_REASON,
|
||||
GEN_AI_RESPONSE_TIME_TO_FIRST_TOKEN,
|
||||
GEN_AI_SESSION_ID,
|
||||
GEN_AI_SKILL_DESCRIPTION,
|
||||
GEN_AI_SKILL_ID,
|
||||
GEN_AI_SKILL_NAME,
|
||||
GEN_AI_SKILL_VERSION,
|
||||
GEN_AI_SPAN_KIND,
|
||||
GEN_AI_TOOL_CALL_ARGUMENTS,
|
||||
GEN_AI_TOOL_CALL_ID,
|
||||
GEN_AI_TOOL_CALL_RESULT,
|
||||
GEN_AI_TOOL_DESCRIPTION,
|
||||
GEN_AI_TOOL_NAME,
|
||||
GEN_AI_TOOL_TYPE,
|
||||
GEN_AI_USAGE_INPUT_TOKENS,
|
||||
GEN_AI_USAGE_OUTPUT_TOKENS,
|
||||
GEN_AI_USAGE_TOTAL_TOKENS,
|
||||
GEN_AI_USER_ID,
|
||||
GEN_AI_USER_NAME,
|
||||
INPUT_VALUE,
|
||||
OPERATION_NAME_EXECUTE_TOOL,
|
||||
OUTPUT_VALUE,
|
||||
RETRIEVAL_DOCUMENT,
|
||||
RETRIEVAL_QUERY,
|
||||
TOOL_DESCRIPTION,
|
||||
TOOL_NAME,
|
||||
TOOL_PARAMETERS,
|
||||
TOOL_TYPE_DATASTORE,
|
||||
TOOL_TYPE_EXTENSION,
|
||||
TOOL_TYPE_FUNCTION,
|
||||
GenAISpanKind,
|
||||
)
|
||||
|
||||
@@ -47,9 +63,25 @@ def test_constants():
|
||||
assert GEN_AI_RESPONSE_FINISH_REASON == "gen_ai.response.finish_reason"
|
||||
assert GEN_AI_INPUT_MESSAGE == "gen_ai.input.messages"
|
||||
assert GEN_AI_OUTPUT_MESSAGE == "gen_ai.output.messages"
|
||||
assert TOOL_NAME == "tool.name"
|
||||
assert TOOL_DESCRIPTION == "tool.description"
|
||||
assert TOOL_PARAMETERS == "tool.parameters"
|
||||
assert GEN_AI_TOOL_CALL_ID == "gen_ai.tool.call.id"
|
||||
assert GEN_AI_TOOL_DESCRIPTION == "gen_ai.tool.description"
|
||||
assert GEN_AI_TOOL_NAME == "gen_ai.tool.name"
|
||||
assert GEN_AI_TOOL_TYPE == "gen_ai.tool.type"
|
||||
assert GEN_AI_TOOL_CALL_ARGUMENTS == "gen_ai.tool.call.arguments"
|
||||
assert GEN_AI_TOOL_CALL_RESULT == "gen_ai.tool.call.result"
|
||||
assert GEN_AI_SKILL_ID == "gen_ai.skill.id"
|
||||
assert GEN_AI_SKILL_NAME == "gen_ai.skill.name"
|
||||
assert GEN_AI_SKILL_DESCRIPTION == "gen_ai.skill.description"
|
||||
assert GEN_AI_SKILL_VERSION == "gen_ai.skill.version"
|
||||
assert GEN_AI_OPERATION_NAME == "gen_ai.operation.name"
|
||||
assert OPERATION_NAME_EXECUTE_TOOL == "execute_tool"
|
||||
assert TOOL_TYPE_FUNCTION == "function"
|
||||
assert TOOL_TYPE_EXTENSION == "extension"
|
||||
assert TOOL_TYPE_DATASTORE == "datastore"
|
||||
assert GEN_AI_RESPONSE_TIME_TO_FIRST_TOKEN == "gen_ai.response.time_to_first_token"
|
||||
assert GEN_AI_AGENT_NAME == "gen_ai.agent.name"
|
||||
assert GEN_AI_REACT_ROUND == "gen_ai.react.round"
|
||||
assert GEN_AI_REACT_FINISH_REASON == "gen_ai.react.finish_reason"
|
||||
|
||||
|
||||
def test_gen_ai_span_kind_enum():
|
||||
@@ -61,8 +93,9 @@ def test_gen_ai_span_kind_enum():
|
||||
assert GenAISpanKind.TOOL == "TOOL"
|
||||
assert GenAISpanKind.AGENT == "AGENT"
|
||||
assert GenAISpanKind.TASK == "TASK"
|
||||
assert GenAISpanKind.STEP == "STEP"
|
||||
|
||||
# Verify iteration works (covers the class definition)
|
||||
kinds = list(GenAISpanKind)
|
||||
assert len(kinds) == 8
|
||||
assert len(kinds) == 9
|
||||
assert "LLM" in kinds
|
||||
|
||||
+330
-12
@@ -1,5 +1,6 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
from datetime import UTC, datetime
|
||||
from types import SimpleNamespace
|
||||
@@ -12,18 +13,34 @@ from dify_trace_aliyun.aliyun_trace import AliyunDataTrace
|
||||
from dify_trace_aliyun.config import AliyunConfig
|
||||
from dify_trace_aliyun.entities.aliyun_trace_entity import SpanData, TraceMetadata
|
||||
from dify_trace_aliyun.entities.semconv import (
|
||||
GEN_AI_AGENT_NAME,
|
||||
GEN_AI_COMPLETION,
|
||||
GEN_AI_INPUT_MESSAGE,
|
||||
GEN_AI_OPERATION_NAME,
|
||||
GEN_AI_OUTPUT_MESSAGE,
|
||||
GEN_AI_PROMPT,
|
||||
GEN_AI_PROVIDER_NAME,
|
||||
GEN_AI_REACT_ROUND,
|
||||
GEN_AI_REQUEST_MODEL,
|
||||
GEN_AI_RESPONSE_FINISH_REASON,
|
||||
GEN_AI_RESPONSE_TIME_TO_FIRST_TOKEN,
|
||||
GEN_AI_TOOL_CALL_ARGUMENTS,
|
||||
GEN_AI_TOOL_CALL_ID,
|
||||
GEN_AI_TOOL_CALL_RESULT,
|
||||
GEN_AI_TOOL_DESCRIPTION,
|
||||
GEN_AI_TOOL_NAME,
|
||||
GEN_AI_TOOL_TYPE,
|
||||
GEN_AI_USAGE_TOTAL_TOKENS,
|
||||
INPUT_VALUE,
|
||||
OPERATION_NAME_CHAT,
|
||||
OPERATION_NAME_EXECUTE_TOOL,
|
||||
OPERATION_NAME_INVOKE_AGENT,
|
||||
OPERATION_NAME_REACT,
|
||||
OUTPUT_VALUE,
|
||||
RETRIEVAL_DOCUMENT,
|
||||
RETRIEVAL_QUERY,
|
||||
TOOL_DESCRIPTION,
|
||||
TOOL_NAME,
|
||||
TOOL_PARAMETERS,
|
||||
TOOL_TYPE_DATASTORE,
|
||||
TOOL_TYPE_FUNCTION,
|
||||
GenAISpanKind,
|
||||
)
|
||||
from opentelemetry.trace import Link, SpanContext, SpanKind, Status, StatusCode, TraceFlags
|
||||
@@ -397,8 +414,10 @@ def test_tool_trace_creates_span(trace_instance: AliyunDataTrace, monkeypatch: p
|
||||
_make_tool_trace_info(
|
||||
tool_name="my-tool",
|
||||
tool_inputs={"a": 1},
|
||||
tool_config={"description": "x"},
|
||||
tool_outputs="tool-out",
|
||||
tool_config={"description": "x", "tool_provider_type": "builtin"},
|
||||
inputs={"i": 1},
|
||||
metadata={"conversation_id": "conv", "user_id": "u", "node_execution_id": "exec-1"},
|
||||
)
|
||||
)
|
||||
|
||||
@@ -407,8 +426,13 @@ def test_tool_trace_creates_span(trace_instance: AliyunDataTrace, monkeypatch: p
|
||||
span = spans[0]
|
||||
assert span.name == "my-tool"
|
||||
assert span.status == status
|
||||
assert span.attributes[TOOL_NAME] == "my-tool"
|
||||
assert span.attributes[TOOL_DESCRIPTION] == '{"description": "x"}'
|
||||
assert span.attributes[GEN_AI_OPERATION_NAME] == OPERATION_NAME_EXECUTE_TOOL
|
||||
assert span.attributes[GEN_AI_TOOL_NAME] == "my-tool"
|
||||
assert span.attributes[GEN_AI_TOOL_DESCRIPTION] == "x"
|
||||
assert span.attributes[GEN_AI_TOOL_TYPE] == TOOL_TYPE_FUNCTION
|
||||
assert span.attributes[GEN_AI_TOOL_CALL_ID] == "exec-1"
|
||||
assert span.attributes[GEN_AI_TOOL_CALL_ARGUMENTS] == '{"a": 1}'
|
||||
assert span.attributes[GEN_AI_TOOL_CALL_RESULT] == "tool-out"
|
||||
|
||||
|
||||
def test_get_workflow_node_executions_requires_app_id(trace_instance: AliyunDataTrace):
|
||||
@@ -534,20 +558,30 @@ def test_build_workflow_tool_span(trace_instance: AliyunDataTrace, monkeypatch:
|
||||
node_execution.outputs = {"b": 2}
|
||||
node_execution.created_at = _dt()
|
||||
node_execution.finished_at = _dt()
|
||||
node_execution.metadata = {WorkflowNodeExecutionMetadataKey.TOOL_INFO: {"k": "v"}}
|
||||
node_execution.metadata = {
|
||||
WorkflowNodeExecutionMetadataKey.TOOL_INFO: {
|
||||
"provider_type": "dataset-retrieval",
|
||||
"description": "search docs",
|
||||
}
|
||||
}
|
||||
|
||||
span = trace_instance.build_workflow_tool_span(_make_workflow_trace_info(), node_execution, trace_metadata)
|
||||
assert span.attributes[TOOL_NAME] == "my-tool"
|
||||
assert span.attributes[TOOL_DESCRIPTION] == '{"k": "v"}'
|
||||
assert span.attributes[TOOL_PARAMETERS] == '{"a": 1}'
|
||||
assert span.attributes[GEN_AI_OPERATION_NAME] == OPERATION_NAME_EXECUTE_TOOL
|
||||
assert span.attributes[GEN_AI_TOOL_NAME] == "my-tool"
|
||||
assert span.attributes[GEN_AI_TOOL_DESCRIPTION] == "search docs"
|
||||
assert span.attributes[GEN_AI_TOOL_TYPE] == TOOL_TYPE_DATASTORE
|
||||
assert span.attributes[GEN_AI_TOOL_CALL_ID] == "node-id"
|
||||
assert span.attributes[GEN_AI_TOOL_CALL_ARGUMENTS] == '{"a": 1}'
|
||||
assert span.attributes[GEN_AI_TOOL_CALL_RESULT] == '{"b": 2}'
|
||||
assert span.status.status_code == StatusCode.OK
|
||||
|
||||
# Cover metadata is None and inputs is None
|
||||
node_execution.metadata = None
|
||||
node_execution.inputs = None
|
||||
span2 = trace_instance.build_workflow_tool_span(_make_workflow_trace_info(), node_execution, trace_metadata)
|
||||
assert span2.attributes[TOOL_DESCRIPTION] == "{}"
|
||||
assert span2.attributes[TOOL_PARAMETERS] == "{}"
|
||||
assert span2.attributes[GEN_AI_TOOL_DESCRIPTION] == ""
|
||||
assert span2.attributes[GEN_AI_TOOL_TYPE] == TOOL_TYPE_FUNCTION
|
||||
assert span2.attributes[GEN_AI_TOOL_CALL_ARGUMENTS] == "{}"
|
||||
|
||||
|
||||
def test_build_workflow_retrieval_span(trace_instance: AliyunDataTrace, monkeypatch: pytest.MonkeyPatch):
|
||||
@@ -592,6 +626,8 @@ def test_build_workflow_llm_span(trace_instance: AliyunDataTrace, monkeypatch: p
|
||||
node_execution = MagicMock(spec=WorkflowNodeExecution)
|
||||
node_execution.id = "node-id"
|
||||
node_execution.title = "llm"
|
||||
node_execution.inputs = {}
|
||||
node_execution.error = None
|
||||
node_execution.process_data = {
|
||||
"usage": {"prompt_tokens": 1, "completion_tokens": 2, "total_tokens": 3},
|
||||
"prompts": ["p"],
|
||||
@@ -618,6 +654,288 @@ def test_build_workflow_llm_span(trace_instance: AliyunDataTrace, monkeypatch: p
|
||||
assert span2.attributes[GEN_AI_USAGE_TOTAL_TOKENS] == "10"
|
||||
|
||||
|
||||
def test_build_workflow_llm_span_falls_back_to_inputs_on_invoke_failure(
|
||||
trace_instance: AliyunDataTrace, monkeypatch: pytest.MonkeyPatch
|
||||
):
|
||||
"""Model invoke failures leave process_data empty; prep still populates inputs."""
|
||||
monkeypatch.setattr(aliyun_trace_module, "convert_to_span_id", lambda _, __: 9)
|
||||
monkeypatch.setattr(aliyun_trace_module, "convert_datetime_to_nanoseconds", lambda _: 123)
|
||||
monkeypatch.setattr(aliyun_trace_module, "get_workflow_node_status", lambda _: Status(StatusCode.ERROR, "boom"))
|
||||
|
||||
trace_metadata = _make_trace_metadata()
|
||||
node_execution = MagicMock(spec=WorkflowNodeExecution)
|
||||
node_execution.id = "node-id"
|
||||
node_execution.title = "llm"
|
||||
node_execution.process_data = {}
|
||||
node_execution.inputs = {
|
||||
"model_name": "qwen-plus",
|
||||
"model_provider": "langgenius/tongyi/tongyi",
|
||||
"#context#": "ctx",
|
||||
"query": "hello",
|
||||
}
|
||||
node_execution.outputs = {"error_message": "API rate limit", "error_type": "InvokeError"}
|
||||
node_execution.error = "API rate limit"
|
||||
node_execution.created_at = _dt()
|
||||
node_execution.finished_at = _dt()
|
||||
|
||||
span = trace_instance.build_workflow_llm_span(_make_workflow_trace_info(), node_execution, trace_metadata)
|
||||
|
||||
expected_prompt = json.dumps(node_execution.inputs, ensure_ascii=False)
|
||||
assert span.attributes[GEN_AI_REQUEST_MODEL] == "qwen-plus"
|
||||
assert span.attributes[GEN_AI_PROVIDER_NAME] == "langgenius/tongyi/tongyi"
|
||||
assert span.attributes[GEN_AI_PROMPT] == expected_prompt
|
||||
assert span.attributes[INPUT_VALUE] == expected_prompt
|
||||
assert span.attributes[GEN_AI_COMPLETION] == "API rate limit"
|
||||
assert span.attributes[OUTPUT_VALUE] == "API rate limit"
|
||||
assert span.attributes[GEN_AI_RESPONSE_FINISH_REASON] == "InvokeError"
|
||||
# prompts were never persisted, so structured input.messages stays empty
|
||||
assert span.attributes[GEN_AI_INPUT_MESSAGE] == "[]"
|
||||
|
||||
|
||||
def _make_agent_outputs() -> dict:
|
||||
return {
|
||||
"text": " pong!",
|
||||
"usage": {
|
||||
"prompt_tokens": 100,
|
||||
"completion_tokens": 225,
|
||||
"total_tokens": 325,
|
||||
"time_to_first_token": 0.5,
|
||||
},
|
||||
"files": [],
|
||||
"json": [
|
||||
{
|
||||
"id": "round-1",
|
||||
"parent_id": None,
|
||||
"error": None,
|
||||
"status": "success",
|
||||
"data": {"action_input": "", "action_name": "", "observation": "", "thought": ""},
|
||||
"label": "ROUND 1",
|
||||
"metadata": {"started_at": 6055.211092814, "finished_at": 6056.990671049, "total_tokens": 325},
|
||||
"node_id": "n1",
|
||||
},
|
||||
{
|
||||
"id": "thought-1",
|
||||
"parent_id": "round-1",
|
||||
"error": None,
|
||||
"status": "success",
|
||||
"data": {"action": " pong!", "thought": ""},
|
||||
"label": "deepseek-v4-flash Thought",
|
||||
"metadata": {
|
||||
"started_at": 6055.211450345,
|
||||
"finished_at": 6056.990473571,
|
||||
"provider": "langgenius/deepseek/deepseek",
|
||||
"total_tokens": 325,
|
||||
},
|
||||
"node_id": "n1",
|
||||
},
|
||||
{
|
||||
"id": "tool-1",
|
||||
"parent_id": "round-1",
|
||||
"error": None,
|
||||
"status": "success",
|
||||
"data": {
|
||||
"tool_name": "current_time",
|
||||
"tool_call_args": {"timezone": "Asia/Shanghai"},
|
||||
"output": "2026-07-28 18:00:00",
|
||||
},
|
||||
"label": "CALL current_time",
|
||||
# Tool CALL logs also set provider (tool provider); must not become LLM spans.
|
||||
"metadata": {
|
||||
"started_at": 6056.990500000,
|
||||
"finished_at": 6056.990600000,
|
||||
"provider": "langgenius/time/time",
|
||||
},
|
||||
"node_id": "n1",
|
||||
},
|
||||
{"data": []},
|
||||
],
|
||||
}
|
||||
|
||||
|
||||
def test_build_workflow_node_span_routes_agent_type(trace_instance: AliyunDataTrace, monkeypatch: pytest.MonkeyPatch):
|
||||
node_execution = MagicMock(spec=WorkflowNodeExecution)
|
||||
trace_info = _make_workflow_trace_info()
|
||||
trace_metadata = _make_trace_metadata()
|
||||
|
||||
monkeypatch.setattr(trace_instance, "build_workflow_agent_span", MagicMock(return_value="agent"))
|
||||
|
||||
node_execution.node_type = BuiltinNodeTypes.AGENT
|
||||
assert trace_instance.build_workflow_node_span(node_execution, trace_info, trace_metadata) == "agent"
|
||||
|
||||
|
||||
def test_build_workflow_agent_span(trace_instance: AliyunDataTrace, monkeypatch: pytest.MonkeyPatch):
|
||||
monkeypatch.setattr(aliyun_trace_module, "convert_to_span_id", lambda _, __: 9)
|
||||
monkeypatch.setattr(aliyun_trace_module, "convert_datetime_to_nanoseconds", lambda _: 123)
|
||||
status = Status(StatusCode.OK)
|
||||
monkeypatch.setattr(aliyun_trace_module, "get_workflow_node_status", lambda _: status)
|
||||
|
||||
trace_metadata = _make_trace_metadata()
|
||||
node_execution = MagicMock(spec=WorkflowNodeExecution)
|
||||
node_execution.id = "node-id"
|
||||
node_execution.title = "my-agent"
|
||||
node_execution.inputs = {"query": "ping"}
|
||||
node_execution.outputs = _make_agent_outputs()
|
||||
node_execution.created_at = _dt()
|
||||
node_execution.finished_at = _dt()
|
||||
|
||||
span = trace_instance.build_workflow_agent_span(_make_workflow_trace_info(), node_execution, trace_metadata)
|
||||
assert span.attributes["gen_ai.span.kind"] == GenAISpanKind.AGENT
|
||||
assert span.attributes[GEN_AI_OPERATION_NAME] == OPERATION_NAME_INVOKE_AGENT
|
||||
assert span.attributes[GEN_AI_AGENT_NAME] == "my-agent"
|
||||
assert span.attributes[GEN_AI_USAGE_TOTAL_TOKENS] == "325"
|
||||
assert span.attributes[GEN_AI_RESPONSE_TIME_TO_FIRST_TOKEN] == 500_000_000
|
||||
assert span.attributes["output.value"] == " pong!"
|
||||
|
||||
# TTFT attribute must be absent when usage does not carry it (e.g. blocking mode)
|
||||
node_execution.outputs = {"text": "t", "usage": {"total_tokens": 1, "time_to_first_token": None}}
|
||||
span2 = trace_instance.build_workflow_agent_span(_make_workflow_trace_info(), node_execution, trace_metadata)
|
||||
assert GEN_AI_RESPONSE_TIME_TO_FIRST_TOKEN not in span2.attributes
|
||||
|
||||
# Malformed usage payloads must not break agent span building
|
||||
node_execution.outputs = {"text": "t", "usage": "not-a-mapping"}
|
||||
span3 = trace_instance.build_workflow_agent_span(_make_workflow_trace_info(), node_execution, trace_metadata)
|
||||
assert span3.attributes[GEN_AI_USAGE_TOTAL_TOKENS] == "0"
|
||||
assert GEN_AI_RESPONSE_TIME_TO_FIRST_TOKEN not in span3.attributes
|
||||
|
||||
|
||||
def test_build_agent_react_spans(trace_instance: AliyunDataTrace, monkeypatch: pytest.MonkeyPatch):
|
||||
node_start_ns = 1_000_000_000
|
||||
monkeypatch.setattr(aliyun_trace_module, "convert_to_span_id", lambda _, __: 9)
|
||||
monkeypatch.setattr(aliyun_trace_module, "convert_datetime_to_nanoseconds", lambda _: node_start_ns)
|
||||
span_ids = iter([100, 200, 300])
|
||||
monkeypatch.setattr(aliyun_trace_module, "generate_span_id", lambda: next(span_ids))
|
||||
|
||||
trace_metadata = _make_trace_metadata()
|
||||
node_execution = MagicMock(spec=WorkflowNodeExecution)
|
||||
node_execution.id = "node-id"
|
||||
node_execution.outputs = _make_agent_outputs()
|
||||
node_execution.created_at = _dt()
|
||||
node_execution.finished_at = _dt()
|
||||
|
||||
spans = trace_instance.build_agent_react_spans(node_execution, trace_metadata)
|
||||
assert len(spans) == 3
|
||||
step_span, llm_span, tool_span = spans
|
||||
|
||||
assert step_span.parent_span_id == 9
|
||||
assert step_span.span_id == 100
|
||||
assert step_span.name == "ROUND 1"
|
||||
assert step_span.attributes["gen_ai.span.kind"] == GenAISpanKind.STEP
|
||||
assert step_span.attributes[GEN_AI_OPERATION_NAME] == OPERATION_NAME_REACT
|
||||
assert step_span.attributes[GEN_AI_REACT_ROUND] == 1
|
||||
# Round starts at the earliest monotonic timestamp, anchored to node start time
|
||||
assert step_span.start_time == node_start_ns
|
||||
assert step_span.end_time == node_start_ns + int((6056.990671049 - 6055.211092814) * 1e9)
|
||||
assert step_span.status.status_code == StatusCode.OK
|
||||
|
||||
assert llm_span.parent_span_id == 100
|
||||
assert llm_span.span_id == 200
|
||||
assert llm_span.name == "deepseek-v4-flash Thought"
|
||||
assert llm_span.attributes["gen_ai.span.kind"] == GenAISpanKind.LLM
|
||||
assert llm_span.attributes[GEN_AI_OPERATION_NAME] == OPERATION_NAME_CHAT
|
||||
assert llm_span.attributes[GEN_AI_REQUEST_MODEL] == "deepseek-v4-flash"
|
||||
assert llm_span.attributes[GEN_AI_PROVIDER_NAME] == "langgenius/deepseek/deepseek"
|
||||
assert llm_span.attributes[GEN_AI_USAGE_TOTAL_TOKENS] == "325"
|
||||
assert llm_span.attributes[GEN_AI_COMPLETION] == " pong!"
|
||||
assert llm_span.start_time == node_start_ns + int((6055.211450345 - 6055.211092814) * 1e9)
|
||||
|
||||
assert tool_span.parent_span_id == 100
|
||||
assert tool_span.span_id == 300
|
||||
assert tool_span.name == "CALL current_time"
|
||||
assert tool_span.attributes["gen_ai.span.kind"] == GenAISpanKind.TOOL
|
||||
assert tool_span.attributes[GEN_AI_OPERATION_NAME] == OPERATION_NAME_EXECUTE_TOOL
|
||||
assert tool_span.attributes[GEN_AI_TOOL_NAME] == "current_time"
|
||||
assert tool_span.attributes[GEN_AI_TOOL_TYPE] == TOOL_TYPE_FUNCTION
|
||||
assert tool_span.attributes[GEN_AI_TOOL_CALL_ID] == "tool-1"
|
||||
assert tool_span.attributes[GEN_AI_TOOL_CALL_ARGUMENTS] == '{"timezone": "Asia/Shanghai"}'
|
||||
assert tool_span.attributes[GEN_AI_TOOL_CALL_RESULT] == "2026-07-28 18:00:00"
|
||||
assert tool_span.start_time == node_start_ns + int((6056.990500000 - 6055.211092814) * 1e9)
|
||||
|
||||
|
||||
def test_build_agent_react_spans_returns_empty_without_log(trace_instance: AliyunDataTrace):
|
||||
node_execution = MagicMock(spec=WorkflowNodeExecution)
|
||||
node_execution.id = "node-id"
|
||||
node_execution.outputs = {"text": "t"}
|
||||
node_execution.created_at = _dt()
|
||||
node_execution.finished_at = _dt()
|
||||
|
||||
assert trace_instance.build_agent_react_spans(node_execution, _make_trace_metadata()) == []
|
||||
|
||||
|
||||
def test_workflow_trace_adds_react_spans_for_agent_nodes(
|
||||
trace_instance: AliyunDataTrace, monkeypatch: pytest.MonkeyPatch
|
||||
):
|
||||
monkeypatch.setattr(aliyun_trace_module, "convert_to_trace_id", lambda _: 111)
|
||||
monkeypatch.setattr(aliyun_trace_module, "convert_to_span_id", lambda _, __: 222)
|
||||
monkeypatch.setattr(aliyun_trace_module, "create_links_from_trace_id", lambda _: [])
|
||||
|
||||
agent_node = MagicMock(spec=WorkflowNodeExecution)
|
||||
agent_node.node_type = BuiltinNodeTypes.AGENT
|
||||
code_node = MagicMock(spec=WorkflowNodeExecution)
|
||||
code_node.node_type = BuiltinNodeTypes.CODE
|
||||
|
||||
monkeypatch.setattr(trace_instance, "add_workflow_span", MagicMock())
|
||||
monkeypatch.setattr(trace_instance, "get_workflow_node_executions", MagicMock(return_value=[agent_node, code_node]))
|
||||
monkeypatch.setattr(trace_instance, "build_workflow_node_span", MagicMock(side_effect=["agent-span", "task-span"]))
|
||||
build_agent_react_spans = MagicMock(return_value=["step-span", "llm-span"])
|
||||
monkeypatch.setattr(trace_instance, "build_agent_react_spans", build_agent_react_spans)
|
||||
|
||||
trace_instance.workflow_trace(_make_workflow_trace_info())
|
||||
|
||||
build_agent_react_spans.assert_called_once()
|
||||
assert _recording_trace_client(trace_instance).added_spans == [
|
||||
"agent-span",
|
||||
"step-span",
|
||||
"llm-span",
|
||||
"task-span",
|
||||
]
|
||||
|
||||
|
||||
def test_build_workflow_llm_span_records_time_to_first_token(
|
||||
trace_instance: AliyunDataTrace, monkeypatch: pytest.MonkeyPatch
|
||||
):
|
||||
monkeypatch.setattr(aliyun_trace_module, "convert_to_span_id", lambda _, __: 9)
|
||||
monkeypatch.setattr(aliyun_trace_module, "convert_datetime_to_nanoseconds", lambda _: 123)
|
||||
monkeypatch.setattr(aliyun_trace_module, "get_workflow_node_status", lambda _: Status(StatusCode.OK))
|
||||
|
||||
trace_metadata = _make_trace_metadata()
|
||||
node_execution = MagicMock(spec=WorkflowNodeExecution)
|
||||
node_execution.id = "node-id"
|
||||
node_execution.title = "llm"
|
||||
node_execution.inputs = {}
|
||||
node_execution.error = None
|
||||
node_execution.process_data = {"prompts": []}
|
||||
node_execution.outputs = {"text": "t", "usage": {"total_tokens": 1, "time_to_first_token": 0.123}}
|
||||
node_execution.created_at = _dt()
|
||||
node_execution.finished_at = _dt()
|
||||
|
||||
span = trace_instance.build_workflow_llm_span(_make_workflow_trace_info(), node_execution, trace_metadata)
|
||||
assert span.attributes[GEN_AI_RESPONSE_TIME_TO_FIRST_TOKEN] == 123_000_000
|
||||
assert span.attributes[GEN_AI_OPERATION_NAME] == OPERATION_NAME_CHAT
|
||||
|
||||
# Absent when usage does not carry TTFT (blocking mode)
|
||||
node_execution.outputs = {"text": "t", "usage": {"total_tokens": 1, "time_to_first_token": None}}
|
||||
span2 = trace_instance.build_workflow_llm_span(_make_workflow_trace_info(), node_execution, trace_metadata)
|
||||
assert GEN_AI_RESPONSE_TIME_TO_FIRST_TOKEN not in span2.attributes
|
||||
|
||||
# Non-numeric TTFT must not raise and must be omitted, so the span is still built
|
||||
node_execution.outputs = {"text": "t", "usage": {"total_tokens": 1, "time_to_first_token": "n/a"}}
|
||||
span3 = trace_instance.build_workflow_llm_span(_make_workflow_trace_info(), node_execution, trace_metadata)
|
||||
assert GEN_AI_RESPONSE_TIME_TO_FIRST_TOKEN not in span3.attributes
|
||||
|
||||
|
||||
def test_message_trace_records_time_to_first_token(trace_instance: AliyunDataTrace, monkeypatch: pytest.MonkeyPatch):
|
||||
monkeypatch.setattr(aliyun_trace_module, "convert_to_trace_id", lambda _: 10)
|
||||
monkeypatch.setattr(aliyun_trace_module, "convert_to_span_id", lambda _, span_type: 0)
|
||||
monkeypatch.setattr(aliyun_trace_module, "convert_datetime_to_nanoseconds", lambda _: 123)
|
||||
monkeypatch.setattr(aliyun_trace_module, "get_user_id_from_message_data", lambda _: "user")
|
||||
monkeypatch.setattr(aliyun_trace_module, "create_links_from_trace_id", lambda _: [])
|
||||
|
||||
trace_instance.message_trace(_make_message_trace_info(gen_ai_server_time_to_first_token=0.25))
|
||||
|
||||
llm_span = _recorded_span_data(trace_instance)[1]
|
||||
assert llm_span.attributes[GEN_AI_RESPONSE_TIME_TO_FIRST_TOKEN] == 250_000_000
|
||||
|
||||
|
||||
def test_add_workflow_span(trace_instance: AliyunDataTrace, monkeypatch: pytest.MonkeyPatch):
|
||||
monkeypatch.setattr(
|
||||
aliyun_trace_module, "convert_to_span_id", lambda _, span_type: {"message": 20}.get(span_type, 0)
|
||||
|
||||
+37
@@ -8,22 +8,33 @@ from unittest.mock import MagicMock
|
||||
import pytest
|
||||
from dify_trace_aliyun.entities.semconv import (
|
||||
GEN_AI_FRAMEWORK,
|
||||
GEN_AI_OPERATION_NAME,
|
||||
GEN_AI_SESSION_ID,
|
||||
GEN_AI_SPAN_KIND,
|
||||
GEN_AI_TOOL_CALL_ARGUMENTS,
|
||||
GEN_AI_TOOL_NAME,
|
||||
GEN_AI_TOOL_TYPE,
|
||||
GEN_AI_USER_ID,
|
||||
INPUT_VALUE,
|
||||
OPERATION_NAME_EXECUTE_TOOL,
|
||||
OUTPUT_VALUE,
|
||||
TOOL_TYPE_DATASTORE,
|
||||
TOOL_TYPE_EXTENSION,
|
||||
TOOL_TYPE_FUNCTION,
|
||||
)
|
||||
from dify_trace_aliyun.utils import (
|
||||
create_common_span_attributes,
|
||||
create_gen_ai_tool_attributes,
|
||||
create_links_from_trace_id,
|
||||
create_status_from_error,
|
||||
extract_retrieval_documents,
|
||||
extract_tool_description,
|
||||
format_input_messages,
|
||||
format_output_messages,
|
||||
format_retrieval_documents,
|
||||
get_user_id_from_message_data,
|
||||
get_workflow_node_status,
|
||||
map_gen_ai_tool_type,
|
||||
serialize_json_data,
|
||||
)
|
||||
from opentelemetry.trace import Link, StatusCode
|
||||
@@ -279,3 +290,29 @@ def test_format_output_messages():
|
||||
# Exception path
|
||||
# Trigger exception in serialize_json_data by passing non-serializable
|
||||
assert format_output_messages({"text": MagicMock()}) == serialize_json_data([])
|
||||
|
||||
|
||||
def test_map_gen_ai_tool_type():
|
||||
assert map_gen_ai_tool_type("dataset-retrieval") == TOOL_TYPE_DATASTORE
|
||||
assert map_gen_ai_tool_type("extension") == TOOL_TYPE_EXTENSION
|
||||
assert map_gen_ai_tool_type("builtin") == TOOL_TYPE_FUNCTION
|
||||
assert map_gen_ai_tool_type(None) == TOOL_TYPE_FUNCTION
|
||||
|
||||
|
||||
def test_extract_tool_description():
|
||||
assert extract_tool_description({"description": "d"}) == "d"
|
||||
assert extract_tool_description({"tool_description": "td"}) == "td"
|
||||
assert extract_tool_description({"other": 1}) == ""
|
||||
assert extract_tool_description(None) == ""
|
||||
|
||||
|
||||
def test_create_gen_ai_tool_attributes():
|
||||
attrs = create_gen_ai_tool_attributes(
|
||||
tool_name="search",
|
||||
tool_type=TOOL_TYPE_DATASTORE,
|
||||
tool_call_arguments='{"q": 1}',
|
||||
)
|
||||
assert attrs[GEN_AI_OPERATION_NAME] == OPERATION_NAME_EXECUTE_TOOL
|
||||
assert attrs[GEN_AI_TOOL_NAME] == "search"
|
||||
assert attrs[GEN_AI_TOOL_TYPE] == TOOL_TYPE_DATASTORE
|
||||
assert attrs[GEN_AI_TOOL_CALL_ARGUMENTS] == '{"q": 1}'
|
||||
|
||||
Reference in New Issue
Block a user