diff --git a/instrumentation/opentelemetry-instrumentation-genai-llama-index/.changelog/495.added b/instrumentation/opentelemetry-instrumentation-genai-llama-index/.changelog/495.added new file mode 100644 index 000000000..2637cab07 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-llama-index/.changelog/495.added @@ -0,0 +1 @@ +Add tracing for LlamaIndex AgentWorkflow runs and member agent executions. diff --git a/instrumentation/opentelemetry-instrumentation-genai-llama-index/README.rst b/instrumentation/opentelemetry-instrumentation-genai-llama-index/README.rst index 9f89244c0..73b769c42 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-llama-index/README.rst +++ b/instrumentation/opentelemetry-instrumentation-genai-llama-index/README.rst @@ -9,9 +9,10 @@ OpenTelemetry LlamaIndex Instrumentation This package contains OpenTelemetry instrumentation for `LlamaIndex `_. -It emits ``invoke_agent`` spans for LlamaIndex ``FunctionAgent`` and -``ReActAgent`` runs, and ``execute_tool`` spans when LlamaIndex executes -function tools. Model calls +It emits ``invoke_workflow`` spans for ``AgentWorkflow`` runs, +``invoke_agent`` spans for standalone and workflow-member ``FunctionAgent`` +and ``ReActAgent`` executions, and ``execute_tool`` spans when LlamaIndex +executes tools. Model calls delegated to provider SDKs are intentionally left to those SDKs' OpenTelemetry instrumentations. diff --git a/instrumentation/opentelemetry-instrumentation-genai-llama-index/src/opentelemetry/instrumentation/genai/llama_index/_handler.py b/instrumentation/opentelemetry-instrumentation-genai-llama-index/src/opentelemetry/instrumentation/genai/llama_index/_handler.py index 87801c306..80f3599d7 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-llama-index/src/opentelemetry/instrumentation/genai/llama_index/_handler.py +++ b/instrumentation/opentelemetry-instrumentation-genai-llama-index/src/opentelemetry/instrumentation/genai/llama_index/_handler.py @@ -12,7 +12,10 @@ from typing import Any, cast from llama_index.core.agent.workflow.base_agent import BaseWorkflowAgent +from llama_index.core.agent.workflow.multi_agent_workflow import AgentWorkflow from llama_index.core.agent.workflow.workflow_events import ( + AgentOutput, + AgentSetup, ToolCall, ToolCallResult, ) @@ -35,6 +38,7 @@ AgentInvocation, GenAIInvocation, ToolInvocation, + WorkflowInvocation, ) from opentelemetry.util.genai.types import ( BlobPart, @@ -216,6 +220,25 @@ def _agent_input(bound_args: inspect.BoundArguments) -> list[InputMessage]: return messages +def _agent_step_input( + event: AgentSetup, system_prompt: str | None +) -> list[InputMessage]: + """Recover the member agent input from an AgentWorkflow step. + + AgentWorkflow prepends the member's system prompt to ``AgentSetup.input``; + it is captured separately as the agent's system instruction. + """ + messages = list(event.input) + if ( + system_prompt + and messages + and messages[0].role.value == "system" + and messages[0].content == system_prompt + ): + messages.pop(0) + return [_input_message(message) for message in messages] + + def _request_model(agent: BaseWorkflowAgent) -> str | None: """Best-effort extraction of the model name across LLM integrations.""" try: @@ -324,6 +347,20 @@ def _set_agent_output(invocation: AgentInvocation, result: Any) -> None: invocation.output_messages = [_output_message(response)] +def _set_agent_step_output(invocation: AgentInvocation, result: Any) -> None: + """Copy a member agent's response out of an AgentWorkflow step.""" + if isinstance(result, AgentOutput): + invocation.output_messages = [_output_message(result.response)] + + +def _set_workflow_output(invocation: WorkflowInvocation, result: Any) -> None: + """Copy the final response out of an AgentWorkflow stop event.""" + output = getattr(result, "result", None) + response = getattr(output, "response", None) + if isinstance(response, ChatMessage): + invocation.output_messages = [_output_message(response)] + + def _tool_arguments( tool: FunctionTool, bound_args: inspect.BoundArguments ) -> dict[str, Any]: @@ -362,6 +399,8 @@ class _LlamaIndexInvocation(BaseSpan): _tool_attributes_token: ( Token[dict[str, _ToolExecutionAttributes] | None] | None ) = PrivateAttr() + _workflow_agents: dict[str, BaseWorkflowAgent] = PrivateAttr() + _workflow_agents_by_run_id: dict[str, BaseWorkflowAgent] = PrivateAttr() def __init__( self, @@ -373,11 +412,36 @@ def __init__( dict[str, _ToolExecutionAttributes] | None ] | None = None, + workflow_agents: Mapping[str, BaseWorkflowAgent] | None = None, + workflow_run_id: str | None = None, + workflow_agent: BaseWorkflowAgent | None = None, ) -> None: """Create the adapter used by LlamaIndex's span-handler lifecycle.""" super().__init__(id_=id_, parent_id=parent_id) self._invocation = invocation self._tool_attributes_token = tool_attributes_token + self._workflow_agents = dict(workflow_agents or {}) + self._workflow_agents_by_run_id = {} + if workflow_run_id is not None and workflow_agent is not None: + self.register_workflow_agent(workflow_run_id, workflow_agent) + + def workflow_agent(self, name: str) -> BaseWorkflowAgent | None: + """Return a member agent owned by this workflow invocation.""" + return self._workflow_agents.get(name) + + def register_workflow_agent( + self, run_id: str, agent: BaseWorkflowAgent + ) -> None: + """Associate a workflow run with its currently executing agent.""" + self._workflow_agents_by_run_id[run_id] = agent + + def workflow_agent_for_run_id( + self, run_id: str | None + ) -> BaseWorkflowAgent | None: + """Return the agent executing the current step for a workflow run.""" + if run_id is None: + return None + return self._workflow_agents_by_run_id.get(run_id) def reset_tool_attributes(self) -> None: """Restore task-local tool metadata after an agent run finishes.""" @@ -418,8 +482,30 @@ def new_span( tool_attributes_token: ( Token[dict[str, _ToolExecutionAttributes] | None] | None ) = None + workflow_agents: Mapping[str, BaseWorkflowAgent] | None = None + workflow_run_id: str | None = None + workflow_agent: BaseWorkflowAgent | None = None - if isinstance(instance, BaseWorkflowAgent) and method_name == "run": + if isinstance(instance, AgentWorkflow) and method_name == "run": + capture_content = self._handler.should_capture_content() + input_messages = ( + _agent_input(bound_args) if capture_content else [] + ) + workflow_agents = instance.agents + workflow_name = getattr(instance, "workflow_name", None) + default_workflow_name = ( + f"{type(instance).__module__}.{type(instance).__qualname__}" + ) + if ( + not isinstance(workflow_name, str) + or not workflow_name + or workflow_name == default_workflow_name + ): + workflow_name = type(instance).__name__ + workflow_invocation = self._handler.workflow(name=workflow_name) + workflow_invocation.input_messages = input_messages + invocation = workflow_invocation + elif isinstance(instance, BaseWorkflowAgent) and method_name == "run": capture_content = self._handler.should_capture_content() agent_name = instance.name or type(instance).__name__ request_model = _request_model(instance) @@ -429,7 +515,7 @@ def new_span( ) tool_definitions = _tool_definitions(instance) system_prompt = instance.system_prompt - system_instruction: list[SystemInstructionPart] = ( + agent_system_instruction: list[SystemInstructionPart] = ( [TextPart(content=system_prompt)] if capture_content and system_prompt else [] @@ -441,18 +527,82 @@ def new_span( agent_invocation.agent_description = agent_description agent_invocation.input_messages = input_messages agent_invocation.tool_definitions = tool_definitions - agent_invocation.system_instruction = system_instruction + agent_invocation.system_instruction = agent_system_instruction invocation = agent_invocation tool_attributes_token = _AGENT_TOOL_ATTRIBUTES.set( _agent_tool_attribute_map(instance) ) + elif method_name == "run_agent_step" and isinstance( + (agent_setup := bound_args.arguments.get("ev")), AgentSetup + ): + parent = self.open_spans.get(parent_span_id or "") + agent = ( + parent.workflow_agent(agent_setup.current_agent_name) + if parent is not None + else None + ) + if agent is None: + return None + capture_content = self._handler.should_capture_content() + agent_name = agent.name or type(agent).__name__ + request_model = _request_model(agent) + agent_description = agent.description + input_messages = ( + _agent_step_input(agent_setup, agent.system_prompt) + if capture_content + else [] + ) + tool_definitions = _tool_definitions(agent) + system_instruction: list[MessagePart] = ( + [TextPart(content=agent.system_prompt)] + if capture_content and agent.system_prompt + else [] + ) + agent_invocation = self._handler.invoke_local_agent( + request_model=request_model, + agent_name=agent_name, + ) + agent_invocation.agent_description = agent_description + agent_invocation.input_messages = input_messages + agent_invocation.tool_definitions = tool_definitions + agent_invocation.system_instruction = system_instruction + invocation = agent_invocation + if parent is not None: + workflow_run_id = ( + tags.get("llamaindex.run_id") if tags is not None else None + ) + if workflow_run_id is not None: + parent.register_workflow_agent(workflow_run_id, agent) + else: + workflow_run_id = ( + tags.get("llamaindex.run_id") if tags is not None else None + ) + workflow_agent = agent elif method_name == "call_tool" and isinstance( (tool_call := bound_args.arguments.get("ev")), ToolCall ): - tool_type, tool_description = _agent_tool_attributes( - instance or bound_args.arguments.get("self"), - tool_call.tool_name, + parent = self.open_spans.get(parent_span_id or "") + active_agent = ( + parent.workflow_agent_for_run_id( + tags.get("llamaindex.run_id") if tags is not None else None + ) + if parent is not None + else None + ) + tool_type, tool_description = ( + _agent_tool_attributes(active_agent, tool_call.tool_name) + if active_agent is not None + else (None, None) ) + if tool_type is None: + tool_type, tool_description = _agent_tool_attributes( + instance or bound_args.arguments.get("self"), + tool_call.tool_name, + ) + if tool_type is None and tool_call.tool_name == "handoff": + # AgentWorkflow's built-in handoff is emitted as a ToolCall, + # although its generated tool metadata is not available here. + tool_type = "function" tool_invocation = self._handler.tool( tool_call.tool_name, tool_type=tool_type, @@ -474,6 +624,11 @@ def new_span( if parent is not None and isinstance( parent._invocation, ToolInvocation ): + # The workflow callback identifies the tool by name only; the + # nested FunctionTool call is the authoritative executing tool. + parent._invocation.tool_description = ( + instance.metadata.description or None + ) return None metadata = instance.metadata tool_invocation = self._handler.tool( @@ -494,6 +649,9 @@ def new_span( parent_id=parent_span_id, invocation=invocation, tool_attributes_token=tool_attributes_token, + workflow_agents=workflow_agents, + workflow_run_id=workflow_run_id, + workflow_agent=workflow_agent, ) def prepare_to_exit_span( @@ -512,10 +670,16 @@ def prepare_to_exit_span( span = self.open_spans.get(id_) if span is None: return None - if isinstance(span._invocation, AgentInvocation): + if isinstance(span._invocation, WorkflowInvocation): + if self._handler.should_capture_content(): + _set_workflow_output(span._invocation, result) + elif isinstance(span._invocation, AgentInvocation): span.reset_tool_attributes() if self._handler.should_capture_content(): - _set_agent_output(span._invocation, result) + if isinstance(result, AgentOutput): + _set_agent_step_output(span._invocation, result) + else: + _set_agent_output(span._invocation, result) elif isinstance(span._invocation, ToolInvocation): tool_output: ToolOutput | None = None if isinstance(result, ToolCallResult): diff --git a/instrumentation/opentelemetry-instrumentation-genai-llama-index/tests/conformance/workflow.py b/instrumentation/opentelemetry-instrumentation-genai-llama-index/tests/conformance/workflow.py new file mode 100644 index 000000000..e443085ab --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-llama-index/tests/conformance/workflow.py @@ -0,0 +1,102 @@ +# Copyright The OpenTelemetry Authors +# SPDX-License-Identifier: Apache-2.0 + +from __future__ import annotations + +import asyncio +from typing import Any + +from llama_index.core.agent.workflow import ( + AgentWorkflow, + FunctionAgent, + ReActAgent, +) +from llama_index.core.base.llms.types import ToolCallBlock +from llama_index.core.llms import ChatMessage, MockFunctionCallingLLM + +from opentelemetry.instrumentation.genai.llama_index import ( + LlamaIndexInstrumentor, +) +from opentelemetry.sdk._logs import LoggerProvider +from opentelemetry.sdk.metrics import MeterProvider +from opentelemetry.sdk.trace import TracerProvider +from opentelemetry.test_util_genai.conformance import Scenario +from opentelemetry.test_util_genai.instrumentor import instrument + + +class WorkflowScenario(Scenario): + expected_spans = { + "invoke_workflow": 1, + "invoke_agent": 2, + "execute_tool": 1, + } + expected_metrics = ("gen_ai.client.operation.duration",) + + def run( + self, + *, + tracer_provider: TracerProvider, + meter_provider: MeterProvider, + logger_provider: LoggerProvider, + vcr: Any, + ) -> None: + def function_response( + messages: list[ChatMessage], **kwargs: Any + ) -> ChatMessage: + return ChatMessage( + role="assistant", + blocks=[ + ToolCallBlock( + tool_call_id="handoff-call", + tool_name="handoff", + tool_kwargs={ + "to_agent": "react-member", + "reason": "The ReAct agent should answer.", + }, + ) + ], + ) + + def react_response( + messages: list[ChatMessage], **kwargs: Any + ) -> ChatMessage: + return ChatMessage( + role="assistant", + content="Thought: I can answer.\nAnswer: complete", + ) + + function_agent = FunctionAgent( + name="function-member", + description="Routes the request.", + llm=MockFunctionCallingLLM( + is_chat_model=True, + response_generator=function_response, + ), + streaming=False, + ) + react_agent = ReActAgent( + name="react-member", + description="Answers the request.", + llm=MockFunctionCallingLLM( + is_chat_model=True, + response_generator=react_response, + ), + streaming=False, + ) + workflow = AgentWorkflow( + agents=[function_agent, react_agent], + root_agent="function-member", + ) + + with instrument( + LlamaIndexInstrumentor(), + tracer_provider=tracer_provider, + logger_provider=logger_provider, + meter_provider=meter_provider, + content_capture="SPAN_ONLY", + ): + + async def run_workflow() -> None: + await workflow.run(user_msg="Complete the request") + + asyncio.run(run_workflow()) diff --git a/instrumentation/opentelemetry-instrumentation-genai-llama-index/tests/test_agent.py b/instrumentation/opentelemetry-instrumentation-genai-llama-index/tests/test_agent.py index 42519c394..16a5cbba6 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-llama-index/tests/test_agent.py +++ b/instrumentation/opentelemetry-instrumentation-genai-llama-index/tests/test_agent.py @@ -837,7 +837,7 @@ def response_generator(messages, **kwargs): @pytest.mark.asyncio -async def test_agent_workflow_emits_tool_span( +async def test_agent_workflow_emits_span_hierarchy( span_exporter, instrument_llama_index ) -> None: def echo(value: str) -> str: @@ -871,12 +871,308 @@ def response_generator(messages, **kwargs): result = await workflow.run(user_msg="Call echo") assert result.response.content == "done" + + workflow_span = _spans_named( + span_exporter, "invoke_workflow AgentWorkflow" + )[0] + workflow_attrs = dict(workflow_span.attributes or {}) + assert workflow_span.kind == SpanKind.INTERNAL + assert workflow_span.parent is None + assert ( + workflow_attrs[GenAIAttributes.GEN_AI_OPERATION_NAME] + == "invoke_workflow" + ) + assert ( + workflow_attrs[GenAIAttributes.GEN_AI_WORKFLOW_NAME] == "AgentWorkflow" + ) + + agent_spans = _spans_named(span_exporter, "invoke_agent workflow-agent") + assert len(agent_spans) == 2 + assert all( + span.parent is not None + and span.parent.span_id == workflow_span.context.span_id + and span.context.trace_id == workflow_span.context.trace_id + for span in agent_spans + ) + tool_spans = _spans_named(span_exporter, "execute_tool echo") assert len(tool_spans) == 1 - assert tool_spans[0].parent is None + assert tool_spans[0].parent is not None + assert tool_spans[0].parent.span_id == workflow_span.context.span_id + assert tool_spans[0].context.trace_id == workflow_span.context.trace_id FunctionTool.from_defaults(echo)(value="after workflow") - assert len(_spans_named(span_exporter, "execute_tool echo")) == 2 + tool_spans = _spans_named(span_exporter, "execute_tool echo") + assert len(tool_spans) == 2 + assert tool_spans[1].parent is None + + +@pytest.mark.asyncio +async def test_agent_workflow_uses_configured_workflow_name( + span_exporter, instrument_llama_index +) -> None: + def response_generator(messages, **kwargs): + return ChatMessage(role="assistant", content="workflow complete") + + agent = FunctionAgent( + name="named-workflow-agent", + llm=MockFunctionCallingLLM( + is_chat_model=True, + response_generator=response_generator, + ), + streaming=False, + ) + workflow = AgentWorkflow( + agents=[agent], + workflow_name="customer-support-workflow", + ) + + await workflow.run(user_msg="Run the workflow") + + workflow_span = _spans_named( + span_exporter, "invoke_workflow customer-support-workflow" + )[0] + assert workflow_span.attributes[GenAIAttributes.GEN_AI_WORKFLOW_NAME] == ( + "customer-support-workflow" + ) + + +@pytest.mark.asyncio +async def test_agent_workflow_uses_executing_agent_tool_metadata( + span_exporter, instrument_llama_index +) -> None: + def response_generator(messages, **kwargs): + if any(message.role.value == "tool" for message in messages): + return ChatMessage(role="assistant", content="workflow complete") + return ChatMessage( + role="assistant", + blocks=[ + ToolCallBlock( + tool_call_id="duplicate-tool-call", + tool_name="lookup", + tool_kwargs={"value": "hello"}, + ) + ], + ) + + def first_lookup(value: str) -> str: + return f"first: {value}" + + class GenericLookupTool(AsyncBaseTool): + @property + def metadata(self) -> ToolMetadata: + return ToolMetadata( + name="lookup", description="Second lookup description." + ) + + def call(self, value: str) -> ToolOutput: + return ToolOutput( + tool_name="lookup", + content=f"second: {value}", + raw_input={"value": value}, + raw_output=value, + ) + + async def acall(self, value: str) -> ToolOutput: + return self.call(value) + + first_agent = FunctionAgent( + name="first-agent", + description="First agent.", + llm=MockFunctionCallingLLM( + is_chat_model=True, + response_generator=lambda messages, **kwargs: ChatMessage( + role="assistant", content="first complete" + ), + ), + tools=[ + FunctionTool.from_defaults( + first_lookup, + name="lookup", + description="First lookup description.", + ) + ], + streaming=False, + ) + second_agent = FunctionAgent( + name="second-agent", + description="Second agent.", + llm=MockFunctionCallingLLM( + is_chat_model=True, + response_generator=response_generator, + ), + tools=[GenericLookupTool()], + streaming=False, + ) + workflow = AgentWorkflow( + agents=[first_agent, second_agent], + root_agent="second-agent", + ) + + await workflow.run(user_msg="Use lookup") + + tool_span = _spans_named(span_exporter, "execute_tool lookup")[0] + assert tool_span.attributes[GenAIAttributes.GEN_AI_TOOL_DESCRIPTION] == ( + "Second lookup description." + ) + assert tool_span.attributes[GenAIAttributes.GEN_AI_TOOL_TYPE] == ( + "GenericLookupTool" + ) + + +@pytest.mark.asyncio +async def test_agent_workflow_captures_content( + span_exporter, instrument_llama_index_with_content +) -> None: + def response_generator(messages, **kwargs): + return ChatMessage(role="assistant", content="workflow complete") + + agent = FunctionAgent( + name="content-agent", + system_prompt="Answer briefly.", + llm=MockFunctionCallingLLM( + is_chat_model=True, + response_generator=response_generator, + ), + streaming=False, + ) + workflow = AgentWorkflow(agents=[agent]) + + await workflow.run(user_msg="Run the workflow") + + workflow_span = _spans_named( + span_exporter, "invoke_workflow AgentWorkflow" + )[0] + agent_span = _spans_named(span_exporter, "invoke_agent content-agent")[0] + workflow_attrs = dict(workflow_span.attributes or {}) + agent_attrs = dict(agent_span.attributes or {}) + assert json.loads( + workflow_attrs[GenAIAttributes.GEN_AI_INPUT_MESSAGES] + ) == [ + { + "role": "user", + "parts": [{"type": "text", "content": "Run the workflow"}], + "name": None, + } + ] + assert json.loads(workflow_attrs[GenAIAttributes.GEN_AI_OUTPUT_MESSAGES])[ + 0 + ]["parts"] == [{"type": "text", "content": "workflow complete"}] + assert json.loads(agent_attrs[GenAIAttributes.GEN_AI_INPUT_MESSAGES]) == [ + { + "role": "user", + "parts": [{"type": "text", "content": "Run the workflow"}], + "name": None, + } + ] + assert json.loads( + agent_attrs[GenAIAttributes.GEN_AI_SYSTEM_INSTRUCTIONS] + ) == [{"type": "text", "content": "Answer briefly."}] + assert json.loads(agent_attrs[GenAIAttributes.GEN_AI_OUTPUT_MESSAGES])[0][ + "parts" + ] == [{"type": "text", "content": "workflow complete"}] + + +@pytest.mark.asyncio +async def test_agent_workflow_error_marks_workflow_and_agent_spans( + span_exporter, instrument_llama_index +) -> None: + error = RuntimeError("workflow agent failed") + + def response_generator(messages, **kwargs): + raise error + + agent = FunctionAgent( + name="failing-workflow-agent", + llm=MockFunctionCallingLLM( + is_chat_model=True, + response_generator=response_generator, + ), + streaming=False, + ) + workflow = AgentWorkflow(agents=[agent]) + + with pytest.raises(RuntimeError) as caught: + await workflow.run(user_msg="Fail") + assert caught.value is error + + workflow_span = _spans_named( + span_exporter, "invoke_workflow AgentWorkflow" + )[0] + agent_span = _spans_named( + span_exporter, "invoke_agent failing-workflow-agent" + )[0] + for span in (workflow_span, agent_span): + assert span.status.status_code == StatusCode.ERROR + assert span.attributes[ErrorAttributes.ERROR_TYPE] == "RuntimeError" + + +@pytest.mark.asyncio +async def test_agent_workflow_instruments_function_and_react_members( + span_exporter, instrument_llama_index +) -> None: + def function_response(messages, **kwargs): + return ChatMessage( + role="assistant", + blocks=[ + ToolCallBlock( + tool_call_id="handoff-call", + tool_name="handoff", + tool_kwargs={ + "to_agent": "react-member", + "reason": "The ReAct agent should answer.", + }, + ) + ], + ) + + def react_response(messages, **kwargs): + return ChatMessage( + role="assistant", + content="Thought: I can answer.\nAnswer: complete", + ) + + function_agent = FunctionAgent( + name="function-member", + description="Routes the request.", + llm=MockFunctionCallingLLM( + is_chat_model=True, + response_generator=function_response, + ), + streaming=False, + ) + react_agent = ReActAgent( + name="react-member", + description="Answers the request.", + llm=MockFunctionCallingLLM( + is_chat_model=True, + response_generator=react_response, + ), + streaming=False, + ) + workflow = AgentWorkflow( + agents=[function_agent, react_agent], + root_agent="function-member", + ) + + result = await workflow.run(user_msg="Complete the request") + assert result.response.content == "complete" + + workflow_span = _spans_named( + span_exporter, "invoke_workflow AgentWorkflow" + )[0] + function_span = _spans_named( + span_exporter, "invoke_agent function-member" + )[0] + react_span = _spans_named(span_exporter, "invoke_agent react-member")[0] + handoff_span = _spans_named(span_exporter, "execute_tool handoff")[0] + assert ( + handoff_span.attributes[GenAIAttributes.GEN_AI_TOOL_TYPE] == "function" + ) + for span in (function_span, react_span, handoff_span): + assert span.parent is not None + assert span.parent.span_id == workflow_span.context.span_id + assert span.context.trace_id == workflow_span.context.trace_id def test_sync_tool_span( diff --git a/instrumentation/opentelemetry-instrumentation-genai-llama-index/tests/test_composition.py b/instrumentation/opentelemetry-instrumentation-genai-llama-index/tests/test_composition.py index dffb45cf8..d8feeaf37 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-llama-index/tests/test_composition.py +++ b/instrumentation/opentelemetry-instrumentation-genai-llama-index/tests/test_composition.py @@ -4,7 +4,11 @@ from __future__ import annotations import pytest -from llama_index.core.agent.workflow import FunctionAgent, ReActAgent +from llama_index.core.agent.workflow import ( + AgentWorkflow, + FunctionAgent, + ReActAgent, +) from llama_index.core.base.llms.types import ToolCallBlock from llama_index.core.llms import ChatMessage, MockFunctionCallingLLM from llama_index.core.tools import FunctionTool @@ -84,6 +88,7 @@ def react_response(messages, **kwargs): llm=openai_llm, streaming=False, ) + provider_workflow = AgentWorkflow(agents=[provider_agent]) providers = { "tracer_provider": tracer_provider, @@ -95,7 +100,7 @@ def react_response(messages, **kwargs): await function_agent.run(user_msg="What is the weather in Paris?") await react_agent.run(user_msg="What is two plus two?") with vcr.use_cassette("inference.yaml"): - await provider_agent.run(user_msg="Hello!") + await provider_workflow.run(user_msg="Hello!") spans = span_exporter.get_finished_spans() operations = [ @@ -103,6 +108,7 @@ def react_response(messages, **kwargs): for span in spans ] assert operations.count("invoke_agent") == 3 + assert operations.count("invoke_workflow") == 1 assert operations.count("execute_tool") == 1 assert operations.count("chat") == 1 assert all(isinstance(operation, str) for operation in operations) @@ -112,6 +118,7 @@ def react_response(messages, **kwargs): function_span = spans_by_name["invoke_agent weather-agent"] inference_span = spans_by_name["chat gpt-4o-mini"] provider_span = spans_by_name["invoke_agent provider-agent"] + workflow_span = spans_by_name["invoke_workflow AgentWorkflow"] assert tool_span.context.trace_id == function_span.context.trace_id assert tool_span.parent is not None @@ -119,6 +126,9 @@ def react_response(messages, **kwargs): assert inference_span.context.trace_id == provider_span.context.trace_id assert inference_span.parent is not None assert inference_span.parent.span_id == provider_span.context.span_id + assert provider_span.context.trace_id == workflow_span.context.trace_id + assert provider_span.parent is not None + assert provider_span.parent.span_id == workflow_span.context.span_id @pytest.mark.asyncio diff --git a/instrumentation/opentelemetry-instrumentation-genai-llama-index/tests/test_conformance.py b/instrumentation/opentelemetry-instrumentation-genai-llama-index/tests/test_conformance.py index 0d2280819..c8103fb31 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-llama-index/tests/test_conformance.py +++ b/instrumentation/opentelemetry-instrumentation-genai-llama-index/tests/test_conformance.py @@ -14,9 +14,14 @@ from opentelemetry.test_util_genai.conformance import Scenario, run_conformance from .conformance.agent import AgentScenario +from .conformance.workflow import WorkflowScenario -@pytest.mark.parametrize("scenario", [AgentScenario()]) +@pytest.mark.parametrize( + "scenario", + [AgentScenario(), WorkflowScenario()], + ids=lambda scenario: type(scenario).__name__, +) def test_conformance( scenario: Scenario, vcr: Any, weaver_live_check: WeaverLiveCheck ) -> None: