diff --git a/instrumentation/opentelemetry-instrumentation-genai-crewai/.changelog/363.added b/instrumentation/opentelemetry-instrumentation-genai-crewai/.changelog/363.added new file mode 100644 index 000000000..baf37d209 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-crewai/.changelog/363.added @@ -0,0 +1 @@ +Add CrewAI LLM inference spans and metrics with content capture and completion-hook support. diff --git a/instrumentation/opentelemetry-instrumentation-genai-crewai/README.rst b/instrumentation/opentelemetry-instrumentation-genai-crewai/README.rst index 828fe4be5..7e65e36f0 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-crewai/README.rst +++ b/instrumentation/opentelemetry-instrumentation-genai-crewai/README.rst @@ -1,9 +1,9 @@ OpenTelemetry CrewAI Instrumentation ==================================== -This package provides the setup for instrumenting CrewAI with OpenTelemetry -Generative AI semantic conventions. CrewAI operation instrumentation will be -added in follow-up changes. +This package instruments CrewAI LLM calls using the OpenTelemetry Generative AI +semantic conventions. It emits inference spans and client duration and token +usage metrics from CrewAI's public LLM lifecycle events. Installation ------------ @@ -24,6 +24,16 @@ Usage Configuration ------------- +CrewAI's native telemetry is disabled while this instrumentation is active to +avoid emitting two independent sets of spans. To retain CrewAI's native +telemetry, explicitly enable it before instrumenting:: + + export CREWAI_DISABLE_TELEMETRY=false + +An existing ``CREWAI_DISABLE_TELEMETRY`` value is always preserved. When the +instrumentation supplies the default value, it removes that value again during +``uninstrument()``. + By default, prompts and completions are not captured. To capture message content, set the environment variable ``OTEL_INSTRUMENTATION_GENAI_CAPTURE_MESSAGE_CONTENT`` to one of ``NO_CONTENT``, ``SPAN_ONLY``, ``EVENT_ONLY``, or ``SPAN_AND_EVENT``: @@ -31,3 +41,22 @@ environment variable ``OTEL_INSTRUMENTATION_GENAI_CAPTURE_MESSAGE_CONTENT`` to o :: export OTEL_INSTRUMENTATION_GENAI_CAPTURE_MESSAGE_CONTENT=SPAN_AND_EVENT + +Completion hooks +---------------- + +To forward captured prompts and completions to the built-in upload hook, set +``OTEL_INSTRUMENTATION_GENAI_COMPLETION_HOOK=upload`` and configure an +``fsspec``-compatible destination with +``OTEL_INSTRUMENTATION_GENAI_UPLOAD_BASE_PATH``:: + + export OTEL_INSTRUMENTATION_GENAI_COMPLETION_HOOK=upload + export OTEL_INSTRUMENTATION_GENAI_UPLOAD_BASE_PATH=/path/to/prompts + +Install ``opentelemetry-util-genai[upload]`` to use the upload hook. A hook can +also be supplied programmatically; it takes precedence over the environment +variable:: + + CrewAIInstrumentor().instrument(completion_hook=my_hook) + +See ``examples/custom_hook.py`` for a minimal custom hook. diff --git a/instrumentation/opentelemetry-instrumentation-genai-crewai/examples/custom_hook.py b/instrumentation/opentelemetry-instrumentation-genai-crewai/examples/custom_hook.py new file mode 100644 index 000000000..c5abacbe3 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-crewai/examples/custom_hook.py @@ -0,0 +1,42 @@ +# Copyright The OpenTelemetry Authors +# SPDX-License-Identifier: Apache-2.0 + +"""Run a CrewAI agent with a custom completion hook.""" + +from crewai import Agent + +from opentelemetry.instrumentation.genai.crewai import CrewAIInstrumentor +from opentelemetry.util.genai.completion_hook import CompletionHook +from opentelemetry.util.genai.types import ( + InputMessage, + MessagePart, + OutputMessage, + ToolDefinition, +) + + +class PrintCompletionHook(CompletionHook): + """Print content after each CrewAI LLM call.""" + + def on_completion( + self, + *, + inputs: list[InputMessage], + outputs: list[OutputMessage], + system_instruction: list[MessagePart], + tool_definitions: list[ToolDefinition] | None = None, + span=None, + log_record=None, + ) -> None: + print(f"inputs: {inputs}") + print(f"outputs: {outputs}") + + +CrewAIInstrumentor().instrument(completion_hook=PrintCompletionHook()) + +agent = Agent( + role="Assistant", + goal="Answer questions concisely", + backstory="You are a helpful assistant.", +) +print(agent.kickoff("What is OpenTelemetry?")) diff --git a/instrumentation/opentelemetry-instrumentation-genai-crewai/pyproject.toml b/instrumentation/opentelemetry-instrumentation-genai-crewai/pyproject.toml index 447561a84..e18e30a9c 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-crewai/pyproject.toml +++ b/instrumentation/opentelemetry-instrumentation-genai-crewai/pyproject.toml @@ -22,7 +22,6 @@ classifiers = [ "Programming Language :: Python :: 3.11", "Programming Language :: Python :: 3.12", "Programming Language :: Python :: 3.13", - "Programming Language :: Python :: 3.14", ] dependencies = [ "opentelemetry-api ~= 1.43", diff --git a/instrumentation/opentelemetry-instrumentation-genai-crewai/src/opentelemetry/instrumentation/genai/crewai/__init__.py b/instrumentation/opentelemetry-instrumentation-genai-crewai/src/opentelemetry/instrumentation/genai/crewai/__init__.py index d0496762a..8deecfd12 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-crewai/src/opentelemetry/instrumentation/genai/crewai/__init__.py +++ b/instrumentation/opentelemetry-instrumentation-genai-crewai/src/opentelemetry/instrumentation/genai/crewai/__init__.py @@ -5,6 +5,7 @@ from __future__ import annotations +import os from collections.abc import Collection from typing import Any @@ -15,6 +16,8 @@ __all__ = ["CrewAIInstrumentor"] +_CREWAI_DISABLE_TELEMETRY = "CREWAI_DISABLE_TELEMETRY" + class CrewAIInstrumentor(BaseInstrumentor): """An instrumentor for CrewAI.""" @@ -24,17 +27,42 @@ def instrumentation_dependencies(self) -> Collection[str]: def _instrument(self, **kwargs: Any) -> None: """Enable CrewAI instrumentation.""" - completion_hook = ( - kwargs.get("completion_hook") or load_completion_hook() - ) - TelemetryHandler( - tracer_provider=kwargs.get("tracer_provider"), - meter_provider=kwargs.get("meter_provider"), - logger_provider=kwargs.get("logger_provider"), - completion_hook=completion_hook, + self._disabled_crewai_telemetry = ( + _CREWAI_DISABLE_TELEMETRY not in os.environ ) - # CrewAI patching will be added in a follow-up change. + if self._disabled_crewai_telemetry: + os.environ[_CREWAI_DISABLE_TELEMETRY] = "true" + + try: + completion_hook = ( + kwargs.get("completion_hook") or load_completion_hook() + ) + telemetry_handler = TelemetryHandler( + tracer_provider=kwargs.get("tracer_provider"), + meter_provider=kwargs.get("meter_provider"), + logger_provider=kwargs.get("logger_provider"), + completion_hook=completion_hook, + ) + from opentelemetry.instrumentation.genai.crewai.event_listener import ( + CrewAIInferenceEventListener, + ) + + self._event_listener = CrewAIInferenceEventListener( + telemetry_handler + ) + except BaseException: + self._restore_crewai_telemetry() + raise def _uninstrument(self, **kwargs: Any) -> None: """Disable CrewAI instrumentation.""" - # CrewAI unpatching will be added in a follow-up change. + listener = getattr(self, "_event_listener", None) + if listener is not None: + listener.shutdown() + self._event_listener = None + self._restore_crewai_telemetry() + + def _restore_crewai_telemetry(self) -> None: + if getattr(self, "_disabled_crewai_telemetry", False): + os.environ.pop(_CREWAI_DISABLE_TELEMETRY, None) + self._disabled_crewai_telemetry = False diff --git a/instrumentation/opentelemetry-instrumentation-genai-crewai/src/opentelemetry/instrumentation/genai/crewai/event_listener.py b/instrumentation/opentelemetry-instrumentation-genai-crewai/src/opentelemetry/instrumentation/genai/crewai/event_listener.py new file mode 100644 index 000000000..cfef61f50 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-crewai/src/opentelemetry/instrumentation/genai/crewai/event_listener.py @@ -0,0 +1,343 @@ +# Copyright The OpenTelemetry Authors +# SPDX-License-Identifier: Apache-2.0 + +"""Translate CrewAI LLM events into GenAI inference telemetry.""" + +from __future__ import annotations + +import json +import threading +from collections.abc import Mapping, Sequence +from typing import Any, cast +from urllib.parse import urlparse + +from crewai.events.base_event_listener import BaseEventListener +from crewai.events.event_bus import CrewAIEventsBus +from crewai.events.types.llm_events import ( + LLMCallCompletedEvent, + LLMCallFailedEvent, + LLMCallStartedEvent, +) + +from opentelemetry.instrumentation.utils import is_instrumentation_enabled +from opentelemetry.util.genai.handler import TelemetryHandler +from opentelemetry.util.genai.invocation import InferenceInvocation +from opentelemetry.util.genai.types import ( + FunctionToolDefinition, + InputMessage, + MessagePart, + OutputMessage, + Text, + ToolCallRequest, + ToolDefinition, +) + + +def _event_key(event: Any) -> str | None: + event_id = getattr(event, "started_event_id", None) or getattr( + event, "event_id", None + ) + return str(event_id) if event_id else None + + +def _tool_call_part(value: Mapping[str, Any]) -> ToolCallRequest | None: + function = value.get("function") + if isinstance(function, Mapping): + function = cast(Mapping[str, Any], function) + name = function.get("name") + arguments = function.get("arguments") + else: + name = value.get("name") + arguments = value.get("arguments", value.get("args")) + if not name: + return None + if isinstance(arguments, str): + try: + arguments = json.loads(arguments) + except json.JSONDecodeError: + pass + call_id = value.get("id") + return ToolCallRequest( + name=str(name), + id=str(call_id) if call_id else None, + arguments=arguments, + ) + + +def _message_parts(message: Mapping[str, Any]) -> list[MessagePart]: + parts: list[MessagePart] = [] + content = message.get("content") + if isinstance(content, str): + parts.append(Text(content=content)) + elif isinstance(content, Sequence): + for block in cast(Sequence[Any], content): + if not isinstance(block, Mapping): + continue + block = cast(Mapping[str, Any], block) + text = block.get("text") + if isinstance(text, str): + parts.append(Text(content=text)) + + tool_calls = message.get("tool_calls") + if isinstance(tool_calls, Sequence) and not isinstance(tool_calls, str): + for tool_call in cast(Sequence[Any], tool_calls): + if isinstance(tool_call, Mapping) and ( + part := _tool_call_part(cast(Mapping[str, Any], tool_call)) + ): + parts.append(part) + return parts + + +def _input_messages(value: Any) -> list[InputMessage]: + if isinstance(value, str): + return [InputMessage(role="user", parts=[Text(content=value)])] + if isinstance(value, Mapping): + value = [value] + if not isinstance(value, Sequence): + return [] + messages: list[InputMessage] = [] + for message in cast(Sequence[Any], value): + if not isinstance(message, Mapping): + continue + message = cast(Mapping[str, Any], message) + parts = _message_parts(message) + if parts: + messages.append( + InputMessage( + role=str(message.get("role") or "user"), parts=parts + ) + ) + return messages + + +def _output_messages( + value: Any, finish_reason: str | None +) -> list[OutputMessage]: + if isinstance(value, str): + value = {"role": "assistant", "content": value} + elif not isinstance(value, (Mapping, Sequence)): + content = getattr(value, "content", None) + if isinstance(content, str): + value = {"role": "assistant", "content": content} + + inputs = _input_messages(value) + return [ + OutputMessage( + role=message.role, + parts=message.parts, + finish_reason=finish_reason or "", + ) + for message in inputs + ] + + +def _tool_definitions(value: Any) -> list[ToolDefinition] | None: + if not isinstance(value, Sequence) or isinstance(value, str): + return None + definitions: list[ToolDefinition] = [] + for tool in cast(Sequence[Any], value): + if not isinstance(tool, Mapping): + continue + tool = cast(Mapping[str, Any], tool) + function = tool.get("function") + definition = ( + cast(Mapping[str, Any], function) + if isinstance(function, Mapping) + else tool + ) + name = definition.get("name") + if not name: + continue + definitions.append( + FunctionToolDefinition( + name=str(name), + description=( + str(definition["description"]) + if definition.get("description") is not None + else None + ), + parameters=definition.get("parameters", {}), + ) + ) + return definitions or None + + +def _server(source: Any) -> tuple[str | None, int | None]: + endpoint = getattr(source, "base_url", None) or getattr( + source, "api_base", None + ) + if not isinstance(endpoint, str): + return None, None + parsed = urlparse(endpoint if "://" in endpoint else f"//{endpoint}") + try: + return parsed.hostname, parsed.port + except ValueError: + return None, None + + +def _int_or_none(value: Any) -> int | None: + try: + return int(value) if value is not None else None + except (TypeError, ValueError): + return None + + +def _first_not_none(*values: Any) -> Any: + return next((value for value in values if value is not None), None) + + +def _set_request_attributes( + invocation: InferenceInvocation, source: Any, event: LLMCallStartedEvent +) -> None: + invocation.input_messages = _input_messages( + getattr(event, "messages", None) + ) + invocation.tool_definitions = _tool_definitions( + getattr(event, "tools", None) + ) + invocation.temperature = _first_not_none( + getattr(event, "temperature", None), + getattr(source, "temperature", None), + ) + invocation.top_p = _first_not_none( + getattr(event, "top_p", None), getattr(source, "top_p", None) + ) + invocation.frequency_penalty = _first_not_none( + getattr(event, "frequency_penalty", None), + getattr(source, "frequency_penalty", None), + ) + invocation.presence_penalty = _first_not_none( + getattr(event, "presence_penalty", None), + getattr(source, "presence_penalty", None), + ) + invocation.max_tokens = _int_or_none( + _first_not_none( + getattr(event, "max_tokens", None), + getattr(source, "max_tokens", None), + getattr(source, "max_completion_tokens", None), + ) + ) + invocation.stop_sequences = _first_not_none( + getattr(event, "stop_sequences", None), + getattr(source, "stop_sequences", None), + ) + invocation.seed = _first_not_none( + getattr(event, "seed", None), getattr(source, "seed", None) + ) + invocation.request_choice_count = _first_not_none( + getattr(event, "n", None), getattr(source, "n", None) + ) + + +def _set_usage(invocation: InferenceInvocation, usage: Any) -> None: + if not isinstance(usage, Mapping): + return + usage = cast(Mapping[str, Any], usage) + input_tokens = usage.get("prompt_tokens", usage.get("input_tokens")) + output_tokens = usage.get("completion_tokens", usage.get("output_tokens")) + cached_tokens = usage.get( + "cached_tokens", usage.get("cached_prompt_tokens") + ) + if input_tokens is not None: + invocation.input_tokens = _int_or_none(input_tokens) + if output_tokens is not None: + invocation.output_tokens = _int_or_none(output_tokens) + if cached_tokens is not None: + invocation.cache_read_input_tokens = _int_or_none(cached_tokens) + + +class CrewAIInferenceEventListener(BaseEventListener): + """Listen for CrewAI LLM lifecycle events.""" + + def __init__(self, telemetry_handler: TelemetryHandler) -> None: + self._telemetry_handler = telemetry_handler + self._invocations: dict[str, InferenceInvocation] = {} + self._handlers: list[tuple[type[Any], Any]] = [] + self._event_bus: CrewAIEventsBus | None = None + self._lock = threading.RLock() + super().__init__() + + def setup_listeners(self, crewai_event_bus: CrewAIEventsBus) -> None: + self._event_bus = crewai_event_bus + self._register(LLMCallStartedEvent, self._on_started) + self._register(LLMCallCompletedEvent, self._on_completed) + self._register(LLMCallFailedEvent, self._on_failed) + + def _register(self, event_type: type[Any], handler: Any) -> None: + if self._event_bus is None: + return + registered = self._event_bus.on(event_type)(handler) + self._handlers.append((event_type, registered)) + + def _on_started(self, source: Any, event: LLMCallStartedEvent) -> None: + if not is_instrumentation_enabled(): + return + key = _event_key(event) + if key is None: + return + provider = str(getattr(source, "provider", None) or "unknown") + model = getattr(event, "model", None) or getattr(source, "model", None) + server_address, server_port = _server(source) + invocation = self._telemetry_handler.inference( + provider=provider, + request_model=str(model) if model else None, + server_address=server_address, + server_port=server_port, + ) + _set_request_attributes(invocation, source, event) + with self._lock: + previous = self._invocations.setdefault(key, invocation) + if previous is not invocation: + invocation.stop() + + def _pop(self, event: Any) -> InferenceInvocation | None: + key = _event_key(event) + if key is None: + return None + with self._lock: + return self._invocations.pop(key, None) + + def _on_completed(self, source: Any, event: LLMCallCompletedEvent) -> None: + invocation = self._pop(event) + if invocation is None: + return + finish_reason = getattr(event, "finish_reason", None) or getattr( + event.response, "finish_reason", None + ) + invocation.output_messages = _output_messages( + event.response, str(finish_reason) if finish_reason else None + ) + finish_reason = str(finish_reason) if finish_reason else None + invocation.finish_reasons = [finish_reason] if finish_reason else None + response_id = getattr(event, "response_id", None) or getattr( + event.response, "id", None + ) + invocation.response_id = str(response_id) if response_id else None + response_model = getattr(event.response, "model", None) + if response_model: + invocation.response_model_name = str(response_model) + _set_usage(invocation, getattr(event, "usage", None)) + invocation.stop() + + def _on_failed(self, source: Any, event: LLMCallFailedEvent) -> None: + invocation = self._pop(event) + if invocation is None: + return + error = getattr(event, "error", None) + invocation.fail( + error + if isinstance(error, BaseException) + else RuntimeError(str(error or "CrewAI LLM call failed")) + ) + + def shutdown(self) -> None: + if self._event_bus is not None: + for event_type, handler in self._handlers: + self._event_bus.off(event_type, handler) + with self._lock: + invocations = list(self._invocations.values()) + self._invocations.clear() + for invocation in invocations: + invocation.stop() + self._handlers.clear() + self._event_bus = None diff --git a/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/conformance/__init__.py b/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/conformance/__init__.py new file mode 100644 index 000000000..e57cf4aba --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/conformance/__init__.py @@ -0,0 +1,2 @@ +# Copyright The OpenTelemetry Authors +# SPDX-License-Identifier: Apache-2.0 diff --git a/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/conformance/inference.py b/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/conformance/inference.py new file mode 100644 index 000000000..4f03a2212 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/conformance/inference.py @@ -0,0 +1,81 @@ +# Copyright The OpenTelemetry Authors +# SPDX-License-Identifier: Apache-2.0 + +"""Conformance scenario for CrewAI LLM inference events.""" + +from __future__ import annotations + +from types import SimpleNamespace +from typing import Any + +from crewai.events.event_bus import crewai_event_bus +from crewai.events.types.llm_events import ( + LLMCallCompletedEvent, + LLMCallStartedEvent, + LLMCallType, +) + +from opentelemetry.instrumentation.genai.crewai import CrewAIInstrumentor +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 InferenceScenario(Scenario): + expected_spans = {"chat": 1} + expected_metrics = ( + "gen_ai.client.operation.duration", + "gen_ai.client.token.usage", + ) + + def run( + self, + *, + tracer_provider: TracerProvider, + meter_provider: MeterProvider, + logger_provider: LoggerProvider, + vcr: Any, + ) -> None: + source = SimpleNamespace( + provider="openai", + model="gpt-4.1-nano", + base_url="https://api.openai.com/v1", + ) + started = LLMCallStartedEvent( + call_id="conformance-call", + model="gpt-4.1-nano", + messages=[{"role": "user", "content": "What is 2 + 2?"}], + temperature=0.0, + max_tokens=32, + ) + completed = LLMCallCompletedEvent( + call_id="conformance-call", + model="gpt-4.1-nano", + started_event_id=started.event_id, + response=SimpleNamespace( + content="2 + 2 equals 4.", + model="gpt-4.1-nano-2025-04-14", + id="chatcmpl-conformance", + finish_reason="stop", + ), + call_type=LLMCallType.LLM_CALL, + usage={"prompt_tokens": 12, "completion_tokens": 8}, + finish_reason="stop", + response_id="chatcmpl-conformance", + ) + + with instrument( + CrewAIInstrumentor(), + tracer_provider=tracer_provider, + meter_provider=meter_provider, + logger_provider=logger_provider, + content_capture="SPAN_ONLY", + ): + start_future = crewai_event_bus.emit(source, started) + if start_future is not None: + start_future.result(timeout=5) + completed_future = crewai_event_bus.emit(source, completed) + if completed_future is not None: + completed_future.result(timeout=5) diff --git a/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/conftest.py b/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/conftest.py index 3a65a0c00..b799e2a25 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/conftest.py +++ b/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/conftest.py @@ -28,5 +28,6 @@ def instrument_crewai( tracer_provider=tracer_provider, logger_provider=logger_provider, meter_provider=meter_provider, + content_capture="SPAN_ONLY", ) as instrumentor: yield instrumentor diff --git a/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/overrides.txt b/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/overrides.txt index afcdad773..997a9c21f 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/overrides.txt +++ b/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/overrides.txt @@ -1,5 +1,6 @@ -# CrewAI 1.10.1 caps opentelemetry-api and opentelemetry-sdk below the +# CrewAI caps the OpenTelemetry API, SDK, and OTLP HTTP exporter below the # workspace's supported ~=1.43 line. Relax those transitive caps so the # instrumentation package can be tested with the repository's OTel stack. opentelemetry-api ~= 1.43 opentelemetry-sdk ~= 1.43 +opentelemetry-exporter-otlp-proto-http ~= 1.43 diff --git a/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/requirements.latest.txt b/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/requirements.latest.txt index a943bf0d1..d5c80c4a7 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/requirements.latest.txt +++ b/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/requirements.latest.txt @@ -1,8 +1,7 @@ # Copyright The OpenTelemetry Authors # SPDX-License-Identifier: Apache-2.0 -# Exercise the latest supported CrewAI release with workspace instrumentation. -crewai - +# Exercise the latest supported CrewAI release through the package's declared +# instruments extra so the test environment cannot drift below its lower bound. -e util/opentelemetry-util-genai --e instrumentation/opentelemetry-instrumentation-genai-crewai +-e instrumentation/opentelemetry-instrumentation-genai-crewai[instruments] diff --git a/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/test_conformance.py b/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/test_conformance.py new file mode 100644 index 000000000..ea7e54b71 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/test_conformance.py @@ -0,0 +1,25 @@ +# Copyright The OpenTelemetry Authors +# SPDX-License-Identifier: Apache-2.0 + +from __future__ import annotations + +import pytest + +pytest.importorskip("opentelemetry.test.weaver_live_check") +pytest.importorskip("opentelemetry.exporter.otlp.proto.grpc") + +from opentelemetry.test.weaver_live_check import WeaverLiveCheck +from opentelemetry.test_util_genai.conformance import ( + Scenario, + run_conformance, +) + +from .conformance.inference import InferenceScenario + + +@pytest.mark.parametrize("scenario", [InferenceScenario()]) +def test_conformance( + scenario: Scenario, + weaver_live_check: WeaverLiveCheck, +) -> None: + run_conformance(scenario, vcr=None, weaver=weaver_live_check) diff --git a/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/test_inference.py b/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/test_inference.py new file mode 100644 index 000000000..ca54b230e --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/test_inference.py @@ -0,0 +1,237 @@ +# Copyright The OpenTelemetry Authors +# SPDX-License-Identifier: Apache-2.0 + +from __future__ import annotations + +from types import SimpleNamespace +from typing import Any +from unittest.mock import MagicMock + +from crewai.events.event_bus import crewai_event_bus +from crewai.events.types.llm_events import ( + LLMCallCompletedEvent, + LLMCallFailedEvent, + LLMCallStartedEvent, + LLMCallType, +) + +from opentelemetry.instrumentation.genai.crewai import CrewAIInstrumentor +from opentelemetry.sdk._logs import LoggerProvider +from opentelemetry.sdk.metrics import MeterProvider +from opentelemetry.sdk.trace import TracerProvider +from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( + InMemorySpanExporter, +) +from opentelemetry.semconv._incubating.attributes import ( + gen_ai_attributes as GenAI, +) +from opentelemetry.semconv.attributes import ( + error_attributes, + server_attributes, +) +from opentelemetry.test_util_genai.instrumentor import instrument +from opentelemetry.trace import SpanKind, StatusCode + + +def _source() -> Any: + return SimpleNamespace( + provider="openai", + model="gpt-4.1-nano", + base_url="https://api.openai.com/v1", + temperature=0.2, + top_p=0.9, + max_tokens=64, + seed=7, + stop_sequences=["done"], + frequency_penalty=0.1, + presence_penalty=0.3, + n=1, + ) + + +def _start() -> LLMCallStartedEvent: + return LLMCallStartedEvent( + call_id="call-1", + model="gpt-4.1-nano", + messages=[{"role": "user", "content": "What is 2 + 2?"}], + tools=[ + { + "type": "function", + "function": { + "name": "calculator", + "description": "Evaluate an expression", + "parameters": { + "type": "object", + "properties": {"expression": {"type": "string"}}, + }, + }, + } + ], + temperature=0.2, + top_p=0.9, + max_tokens=64, + seed=7, + stop_sequences=["done"], + frequency_penalty=0.1, + presence_penalty=0.3, + n=1, + ) + + +def test_completed_call_emits_semantic_inference_telemetry( + instrument_crewai: CrewAIInstrumentor, + span_exporter: InMemorySpanExporter, +) -> None: + listener = instrument_crewai._event_listener + started = _start() + listener._on_started(_source(), started) + completed = LLMCallCompletedEvent( + call_id="call-1", + model="gpt-4.1-nano", + started_event_id=started.event_id, + response=SimpleNamespace( + content="2 + 2 equals 4.", + model="gpt-4.1-nano-2025-04-14", + id="chatcmpl-test", + finish_reason="stop", + ), + call_type=LLMCallType.LLM_CALL, + usage={"prompt_tokens": 12, "completion_tokens": 8}, + finish_reason="stop", + response_id="chatcmpl-test", + ) + listener._on_completed(_source(), completed) + + (span,) = span_exporter.get_finished_spans() + assert span.name == "chat gpt-4.1-nano" + assert span.kind == SpanKind.CLIENT + assert span.status.status_code == StatusCode.UNSET + attributes = span.attributes + assert attributes is not None + assert attributes[GenAI.GEN_AI_OPERATION_NAME] == "chat" + assert attributes[GenAI.GEN_AI_PROVIDER_NAME] == "openai" + assert attributes[GenAI.GEN_AI_REQUEST_MODEL] == "gpt-4.1-nano" + assert attributes[GenAI.GEN_AI_RESPONSE_MODEL] == "gpt-4.1-nano-2025-04-14" + assert attributes[GenAI.GEN_AI_RESPONSE_ID] == "chatcmpl-test" + assert attributes[server_attributes.SERVER_ADDRESS] == "api.openai.com" + assert attributes[GenAI.GEN_AI_REQUEST_TEMPERATURE] == 0.2 + assert attributes[GenAI.GEN_AI_REQUEST_TOP_P] == 0.9 + assert attributes[GenAI.GEN_AI_REQUEST_MAX_TOKENS] == 64 + assert isinstance(attributes[GenAI.GEN_AI_REQUEST_MAX_TOKENS], int) + assert attributes[GenAI.GEN_AI_REQUEST_SEED] == 7 + assert tuple(attributes[GenAI.GEN_AI_REQUEST_STOP_SEQUENCES]) == ("done",) + if getattr(completed, "usage", None) is not None: + assert attributes[GenAI.GEN_AI_USAGE_INPUT_TOKENS] == 12 + assert isinstance(attributes[GenAI.GEN_AI_USAGE_INPUT_TOKENS], int) + assert attributes[GenAI.GEN_AI_USAGE_OUTPUT_TOKENS] == 8 + assert isinstance(attributes[GenAI.GEN_AI_USAGE_OUTPUT_TOKENS], int) + assert isinstance(attributes[GenAI.GEN_AI_INPUT_MESSAGES], str) + assert isinstance(attributes[GenAI.GEN_AI_OUTPUT_MESSAGES], str) + assert isinstance(attributes[GenAI.GEN_AI_TOOL_DEFINITIONS], str) + + +def test_failed_call_ends_span_with_error( + instrument_crewai: CrewAIInstrumentor, + span_exporter: InMemorySpanExporter, +) -> None: + listener = instrument_crewai._event_listener + started = _start() + listener._on_started(_source(), started) + listener._on_failed( + _source(), + LLMCallFailedEvent( + call_id="call-1", + model="gpt-4.1-nano", + started_event_id=started.event_id, + error="provider unavailable", + ), + ) + + (span,) = span_exporter.get_finished_spans() + assert span.status.status_code == StatusCode.ERROR + assert span.attributes is not None + assert span.attributes[error_attributes.ERROR_TYPE].endswith( + "RuntimeError" + ) + + +def test_unmatched_completion_is_ignored( + instrument_crewai: CrewAIInstrumentor, + span_exporter: InMemorySpanExporter, +) -> None: + instrument_crewai._event_listener._on_completed( + _source(), + LLMCallCompletedEvent( + call_id="missing", + model="gpt-4.1-nano", + started_event_id="missing", + response="ignored", + call_type=LLMCallType.LLM_CALL, + ), + ) + assert not span_exporter.get_finished_spans() + + +def test_instrumentor_registers_with_crewai_event_bus( + instrument_crewai: CrewAIInstrumentor, + span_exporter: InMemorySpanExporter, +) -> None: + started = _start() + start_future = crewai_event_bus.emit(_source(), started) + assert start_future is not None + start_future.result(timeout=5) + + completed_future = crewai_event_bus.emit( + _source(), + LLMCallCompletedEvent( + call_id="call-1", + model="gpt-4.1-nano", + started_event_id=started.event_id, + response="2 + 2 equals 4.", + call_type=LLMCallType.LLM_CALL, + usage={"prompt_tokens": 12, "completion_tokens": 8}, + finish_reason="stop", + ), + ) + assert completed_future is not None + completed_future.result(timeout=5) + + (span,) = span_exporter.get_finished_spans() + assert span.attributes is not None + assert span.attributes[GenAI.GEN_AI_PROVIDER_NAME] == "openai" + + +def test_explicit_completion_hook_receives_inference_content( + tracer_provider: TracerProvider, + logger_provider: LoggerProvider, + meter_provider: MeterProvider, +) -> None: + hook = MagicMock() + with instrument( + CrewAIInstrumentor(), + tracer_provider=tracer_provider, + logger_provider=logger_provider, + meter_provider=meter_provider, + completion_hook=hook, + ) as instrumentor: + listener = instrumentor._event_listener + started = _start() + listener._on_started(_source(), started) + listener._on_completed( + _source(), + LLMCallCompletedEvent( + call_id="call-1", + model="gpt-4.1-nano", + started_event_id=started.event_id, + response="2 + 2 equals 4.", + call_type=LLMCallType.LLM_CALL, + finish_reason="stop", + ), + ) + + hook.on_completion.assert_called_once() + call = hook.on_completion.call_args.kwargs + assert call["inputs"] + assert call["outputs"] + assert call["tool_definitions"] + assert call["span"] is not None diff --git a/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/test_instrumentor.py b/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/test_instrumentor.py index db0c6080a..6e98b59b1 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/test_instrumentor.py +++ b/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/test_instrumentor.py @@ -3,6 +3,8 @@ """Tests for the CrewAI instrumentor lifecycle.""" +import os + from opentelemetry.instrumentation.genai.crewai import CrewAIInstrumentor from opentelemetry.sdk._logs import LoggerProvider from opentelemetry.sdk.metrics import MeterProvider @@ -35,3 +37,43 @@ def test_instrument_with_global_providers() -> None: instrumentor = CrewAIInstrumentor() instrumentor.instrument() instrumentor.uninstrument() + + +def test_native_telemetry_is_disabled_while_instrumented( + monkeypatch, + tracer_provider: TracerProvider, + logger_provider: LoggerProvider, + meter_provider: MeterProvider, +) -> None: + monkeypatch.delenv("CREWAI_DISABLE_TELEMETRY", raising=False) + instrumentor = CrewAIInstrumentor() + + instrumentor.instrument( + tracer_provider=tracer_provider, + logger_provider=logger_provider, + meter_provider=meter_provider, + ) + + assert os.environ["CREWAI_DISABLE_TELEMETRY"] == "true" + instrumentor.uninstrument() + assert "CREWAI_DISABLE_TELEMETRY" not in os.environ + + +def test_explicit_native_telemetry_configuration_is_preserved( + monkeypatch, + tracer_provider: TracerProvider, + logger_provider: LoggerProvider, + meter_provider: MeterProvider, +) -> None: + monkeypatch.setenv("CREWAI_DISABLE_TELEMETRY", "false") + instrumentor = CrewAIInstrumentor() + + instrumentor.instrument( + tracer_provider=tracer_provider, + logger_provider=logger_provider, + meter_provider=meter_provider, + ) + + assert os.environ["CREWAI_DISABLE_TELEMETRY"] == "false" + instrumentor.uninstrument() + assert os.environ["CREWAI_DISABLE_TELEMETRY"] == "false" diff --git a/tox.ini b/tox.ini index ce3694aa0..d454a58f9 100644 --- a/tox.ini +++ b/tox.ini @@ -54,8 +54,9 @@ envlist = # No pypy3: jiter (an anthropic dependency) ships no PyPy wheels, so it would need a Rust source build lint-instrumentation-genai-claude-agent-sdk ; instrumentation-genai-crewai - py3{10,11,12,13,14}-test-instrumentation-genai-crewai-latest + py3{10,11,12,13}-test-instrumentation-genai-crewai-latest py310-test-instrumentation-genai-crewai-oldest + py313-test-instrumentation-genai-crewai-conformance lint-instrumentation-genai-crewai ; instrumentation-genai-langchain @@ -180,6 +181,8 @@ deps = crewai-latest: {[testenv]test_deps} crewai-latest: {[testenv]pytest_deps} crewai-latest: -r {toxinidir}/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/requirements.latest.txt + crewai-conformance: {[testenv]pytest_deps} + crewai-conformance: -r {toxinidir}/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/requirements.latest.txt qwen-agent-oldest: {[testenv]pytest_deps} qwen-agent-oldest: -e {toxinidir}/instrumentation/opentelemetry-instrumentation-genai-qwen-agent[instruments] @@ -265,7 +268,8 @@ commands = test-instrumentation-genai-claude-agent-sdk: pytest {toxinidir}/instrumentation/opentelemetry-instrumentation-genai-claude-agent-sdk/tests --vcr-record=none {posargs} lint-instrumentation-genai-claude-agent-sdk: sh -c "cd instrumentation && ruff check opentelemetry-instrumentation-genai-claude-agent-sdk" - test-instrumentation-genai-crewai-{oldest,latest}: pytest {toxinidir}/instrumentation/opentelemetry-instrumentation-genai-crewai/tests --vcr-record=none {posargs} + test-instrumentation-genai-crewai-{oldest,latest}: pytest --ignore={toxinidir}/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/test_conformance.py {toxinidir}/instrumentation/opentelemetry-instrumentation-genai-crewai/tests --vcr-record=none {posargs} + test-instrumentation-genai-crewai-conformance: pytest {toxinidir}/instrumentation/opentelemetry-instrumentation-genai-crewai/tests/test_conformance.py --vcr-record=none {posargs} lint-instrumentation-genai-crewai: sh -c "cd instrumentation && ruff check opentelemetry-instrumentation-genai-crewai" test-instrumentation-genai-langchain-{oldest,latest}: pytest --ignore={toxinidir}/instrumentation/opentelemetry-instrumentation-genai-langchain/tests/test_conformance.py {toxinidir}/instrumentation/opentelemetry-instrumentation-genai-langchain/tests --vcr-record=none {posargs}