From 1973a2e96f2062d5351a7d10ebe74dbed6c01656 Mon Sep 17 00:00:00 2001 From: Liu Ziming Date: Thu, 6 Aug 2026 13:51:28 +0800 Subject: [PATCH] feat(trace-aliyun): GenAI spans for LLM TTFT, agent ReAct, tools, and failed LLM nodes (#39339) Co-authored-by: Cursor --- .../src/dify_trace_aliyun/aliyun_trace.py | 383 +++++++++++++++--- .../src/dify_trace_aliyun/entities/semconv.py | 40 +- .../src/dify_trace_aliyun/utils.py | 169 ++++++++ .../aliyun_trace/entities/test_semconv.py | 47 ++- .../aliyun_trace/test_aliyun_trace.py | 342 +++++++++++++++- .../aliyun_trace/test_aliyun_trace_utils.py | 37 ++ 6 files changed, 946 insertions(+), 72 deletions(-) diff --git a/api/providers/trace/trace-aliyun/src/dify_trace_aliyun/aliyun_trace.py b/api/providers/trace/trace-aliyun/src/dify_trace_aliyun/aliyun_trace.py index b2de91860e0..8a519a04ad5 100644 --- a/api/providers/trace/trace-aliyun/src/dify_trace_aliyun/aliyun_trace.py +++ b/api/providers/trace/trace-aliyun/src/dify_trace_aliyun/aliyun_trace.py @@ -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, ) diff --git a/api/providers/trace/trace-aliyun/src/dify_trace_aliyun/entities/semconv.py b/api/providers/trace/trace-aliyun/src/dify_trace_aliyun/entities/semconv.py index b6e46c5262a..19925410a64 100644 --- a/api/providers/trace/trace-aliyun/src/dify_trace_aliyun/entities/semconv.py +++ b/api/providers/trace/trace-aliyun/src/dify_trace_aliyun/entities/semconv.py @@ -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" diff --git a/api/providers/trace/trace-aliyun/src/dify_trace_aliyun/utils.py b/api/providers/trace/trace-aliyun/src/dify_trace_aliyun/utils.py index 5678c66adbf..c9145a95b54 100644 --- a/api/providers/trace/trace-aliyun/src/dify_trace_aliyun/utils.py +++ b/api/providers/trace/trace-aliyun/src/dify_trace_aliyun/utils.py @@ -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): diff --git a/api/providers/trace/trace-aliyun/tests/unit_tests/aliyun_trace/entities/test_semconv.py b/api/providers/trace/trace-aliyun/tests/unit_tests/aliyun_trace/entities/test_semconv.py index 9cab40748f6..e07aa718f29 100644 --- a/api/providers/trace/trace-aliyun/tests/unit_tests/aliyun_trace/entities/test_semconv.py +++ b/api/providers/trace/trace-aliyun/tests/unit_tests/aliyun_trace/entities/test_semconv.py @@ -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 diff --git a/api/providers/trace/trace-aliyun/tests/unit_tests/aliyun_trace/test_aliyun_trace.py b/api/providers/trace/trace-aliyun/tests/unit_tests/aliyun_trace/test_aliyun_trace.py index 06bdb61b4d1..cf9f5c41978 100644 --- a/api/providers/trace/trace-aliyun/tests/unit_tests/aliyun_trace/test_aliyun_trace.py +++ b/api/providers/trace/trace-aliyun/tests/unit_tests/aliyun_trace/test_aliyun_trace.py @@ -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) diff --git a/api/providers/trace/trace-aliyun/tests/unit_tests/aliyun_trace/test_aliyun_trace_utils.py b/api/providers/trace/trace-aliyun/tests/unit_tests/aliyun_trace/test_aliyun_trace_utils.py index 4cdbb773a90..1895f5a08aa 100644 --- a/api/providers/trace/trace-aliyun/tests/unit_tests/aliyun_trace/test_aliyun_trace_utils.py +++ b/api/providers/trace/trace-aliyun/tests/unit_tests/aliyun_trace/test_aliyun_trace_utils.py @@ -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}'