diff --git a/packages/opentelemetry-instrumentation-crewai/opentelemetry/instrumentation/crewai/instrumentation.py b/packages/opentelemetry-instrumentation-crewai/opentelemetry/instrumentation/crewai/instrumentation.py index 6d57888fd4..3d169adc4c 100644 --- a/packages/opentelemetry-instrumentation-crewai/opentelemetry/instrumentation/crewai/instrumentation.py +++ b/packages/opentelemetry-instrumentation-crewai/opentelemetry/instrumentation/crewai/instrumentation.py @@ -1,3 +1,4 @@ +import logging import os import time from typing import Collection @@ -22,6 +23,26 @@ _instruments = ("crewai >= 1.0.0",) +# crewai >= 1.x routes `LLM(...)` through `LLM.__new__`, which returns native +# provider classes (e.g. `OpenAICompletion`) instead of an `LLM` instance — +# those classes inherit from `BaseLLM`, not `LLM`, so wrapping only +# `crewai.llm.LLM.call` (the LiteLLM fallback path) misses every native call. +# Each provider class overrides `call`, so the base class cannot be wrapped +# instead; every native class needs its own wrap. Provider modules import +# their SDK lazily and may be absent, hence the try/except around each. +# The third element is the authoritative OTel provider identity for that +# class: crewai constructs native instances with an explicit `provider=` +# kwarg (e.g. `AzureCompletion` with `provider="azure"` even for `gpt-4` +# model names), so model-name inference would misattribute those spans. +_NATIVE_LLM_WRAPPED_METHODS = [ + # (module, class qualname, method, otel provider value) — import failures are non-fatal. + ("crewai.llms.providers.openai.completion", "OpenAICompletion", "call", GenAISystem.OPENAI.value), + ("crewai.llms.providers.azure.completion", "AzureCompletion", "call", GenAiSystemValues.AZURE_AI_OPENAI.value), + ("crewai.llms.providers.anthropic.completion", "AnthropicCompletion", "call", GenAISystem.ANTHROPIC.value), + ("crewai.llms.providers.gemini.completion", "GeminiCompletion", "call", GenAiSystemValues.GCP_GEMINI.value), + ("crewai.llms.providers.bedrock.completion", "BedrockCompletion", "call", GenAISystem.AWS.value), +] + # Maps LiteLLM vendor prefixes (e.g. "openai" in "openai/gpt-4") to OTel provider name values. # Uses GenAISystem (semconv-ai) and GenAiSystemValues (OTel upstream) — no raw strings. _LITELLM_PREFIX_TO_OTEL_PROVIDER = { @@ -92,20 +113,34 @@ def _instrument(self, **kwargs): wrap_task_execute(tracer, duration_histogram, token_histogram)) wrap_function_wrapper("crewai.llm", "LLM.call", wrap_llm_call(tracer, duration_histogram, token_histogram)) + for module, class_name, method, otel_provider in _NATIVE_LLM_WRAPPED_METHODS: + try: + wrap_function_wrapper( + module, f"{class_name}.{method}", + wrap_llm_call(tracer, duration_histogram, token_histogram, fixed_provider=otel_provider), + ) + except (ImportError, AttributeError): + # Provider SDK not installed — the class is never used. + logging.debug("crewai native LLM provider %s.%s not importable; skipping wrap", class_name, method) def _uninstrument(self, **kwargs): unwrap("crewai.crew.Crew", "kickoff") unwrap("crewai.agent.Agent", "execute_task") unwrap("crewai.task.Task", "execute_sync") unwrap("crewai.llm.LLM", "call") + for module, class_name, method, _otel_provider in _NATIVE_LLM_WRAPPED_METHODS: + try: + unwrap(f"{module}.{class_name}", method) + except (ImportError, AttributeError): + pass def with_tracer_wrapper(func): """Helper for providing tracer for wrapper functions.""" - def _with_tracer(tracer, duration_histogram, token_histogram): + def _with_tracer(tracer, duration_histogram, token_histogram, **wrapper_kwargs): def wrapper(wrapped, instance, args, kwargs): - return func(tracer, duration_histogram, token_histogram, wrapped, instance, args, kwargs) + return func(tracer, duration_histogram, token_histogram, wrapped, instance, args, kwargs, **wrapper_kwargs) return wrapper return _with_tracer @@ -214,9 +249,15 @@ def wrap_task_execute(tracer, duration_histogram, token_histogram, wrapped, inst @with_tracer_wrapper -def wrap_llm_call(tracer, duration_histogram, token_histogram, wrapped, instance, args, kwargs): +def wrap_llm_call(tracer, duration_histogram, token_histogram, wrapped, instance, args, kwargs, fixed_provider=None): model = str(instance.model) if hasattr(instance, "model") else "llm" - provider = _infer_llm_provider_from_model(getattr(instance, "model", None)) + if fixed_provider: + # Native provider classes get their authoritative OTel identity at + # wrap time; model-name inference would misattribute them (e.g. + # AzureCompletion serving a "gpt-4" model name is azure, not openai). + provider = fixed_provider + else: + provider = _infer_llm_provider_from_model(getattr(instance, "model", None)) span_attrs = { GenAIAttributes.GEN_AI_OPERATION_NAME: GenAiOperationNameValues.CHAT.value, diff --git a/packages/opentelemetry-instrumentation-crewai/tests/conftest.py b/packages/opentelemetry-instrumentation-crewai/tests/conftest.py new file mode 100644 index 0000000000..b63feecae1 --- /dev/null +++ b/packages/opentelemetry-instrumentation-crewai/tests/conftest.py @@ -0,0 +1,69 @@ +"""Unit tests configuration module.""" + +import pytest +from opentelemetry.sdk.metrics import Counter, Histogram, MeterProvider +from opentelemetry.sdk.metrics.export import ( + AggregationTemporality, + InMemoryMetricReader, +) +from opentelemetry.sdk.resources import Resource +from opentelemetry.sdk.trace import TracerProvider +from opentelemetry.sdk.trace.export import SimpleSpanProcessor +from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter + +pytest_plugins = [] + + +@pytest.fixture(scope="function", name="span_exporter") +def fixture_span_exporter(): + exporter = InMemorySpanExporter() + yield exporter + + +@pytest.fixture(scope="function", name="tracer_provider") +def fixture_tracer_provider(span_exporter): + provider = TracerProvider() + provider.add_span_processor(SimpleSpanProcessor(span_exporter)) + return provider + + +@pytest.fixture(scope="function", name="reader") +def fixture_reader(): + reader = InMemoryMetricReader( + {Counter: AggregationTemporality.DELTA, Histogram: AggregationTemporality.DELTA} + ) + return reader + + +@pytest.fixture(scope="function", name="meter_provider") +def fixture_meter_provider(reader): + resource = Resource.create() + meter_provider = MeterProvider(metric_readers=[reader], resource=resource) + return meter_provider + + +@pytest.fixture(scope="function") +def instrument(reader, tracer_provider, meter_provider): + """Real instrumentation against an in-memory exporter. + + BaseInstrumentor is a singleton (its __new__ returns a shared instance), + so any instance-attribute mocks left behind by other tests must be + removed before the real methods can run again. + """ + from opentelemetry.instrumentation.crewai import CrewAIInstrumentor + + instrumentor = CrewAIInstrumentor() + # Restore real methods in case a previous test mocked them on the singleton. + instrumentor.__dict__.pop("instrument", None) + instrumentor.__dict__.pop("uninstrument", None) + if instrumentor._is_instrumented_by_opentelemetry: + instrumentor.uninstrument() + + instrumentor.instrument( + tracer_provider=tracer_provider, + meter_provider=meter_provider, + ) + + yield instrumentor + + instrumentor.uninstrument() diff --git a/packages/opentelemetry-instrumentation-crewai/tests/test_crewai_instrumentation.py b/packages/opentelemetry-instrumentation-crewai/tests/test_crewai_instrumentation.py index dccae7a27a..2dc1ae7eeb 100644 --- a/packages/opentelemetry-instrumentation-crewai/tests/test_crewai_instrumentation.py +++ b/packages/opentelemetry-instrumentation-crewai/tests/test_crewai_instrumentation.py @@ -1,3 +1,5 @@ +import os + import pytest from unittest.mock import MagicMock from pydantic import BaseModel @@ -63,10 +65,7 @@ def test_crewai_instrumentation(mock_crew, mock_instrumentor): assert len(mock_crew.agents) == 1 assert mock_crew.agents[0].role == "Data Collector" assert len(mock_crew.tasks) == 1 - assert ( - mock_crew.tasks[0].description - == "Collect stock data for AAPL for the past month" - ) + assert mock_crew.tasks[0].description == "Collect stock data for AAPL for the past month" def test_trace_status(mock_crew, mock_instrumentor): @@ -80,12 +79,74 @@ def test_trace_status(mock_crew, mock_instrumentor): mock_span.set_status.assert_called_with(StatusCode.ERROR) memory_exporter = MagicMock() - memory_exporter.get_finished_spans.return_value = [ - MagicMock(status=MagicMock(status_code=StatusCode.ERROR)) - ] + memory_exporter.get_finished_spans.return_value = [MagicMock(status=MagicMock(status_code=StatusCode.ERROR))] spans = memory_exporter.get_finished_spans() assert spans[-1].status.status_code == StatusCode.ERROR mock_instrumentor.uninstrument() mock_instrumentor.uninstrument.assert_called_once() + + +def test_native_provider_call_is_wrapped(instrument): + """crewai's native provider classes (BaseLLM subclasses, not LLM) must have + their call wrapped — previously only crewai.llm.LLM.call was patched, which + the LLM.__new__ factory never returns on the native path.""" + from crewai.llms.providers.openai.completion import OpenAICompletion + + assert getattr(OpenAICompletion.call, "__wrapped__", None) is not None + + +def test_native_provider_uninstrument_restores_call(instrument): + from crewai.llms.providers.openai.completion import OpenAICompletion + + instrument.uninstrument() + try: + assert getattr(OpenAICompletion.call, "__wrapped__", None) is None + finally: + # re-instrument so the autouse-style fixture teardown stays balanced + instrument.instrument() + + +def test_litellm_fallback_still_infers_provider_from_model(mock_crew, mock_instrumentor): + """The LLM.call wrap has no fixed provider and keeps model-name inference + (used for the LiteLLM fallback path).""" + from opentelemetry.instrumentation.crewai.instrumentation import _infer_llm_provider_from_model + + assert _infer_llm_provider_from_model("openai/gpt-4") == "openai" + assert _infer_llm_provider_from_model("claude-3") == "anthropic" + assert _infer_llm_provider_from_model("totally-unknown-xyz") is None + + +def test_native_wrapper_emits_span_and_duration_metric(span_exporter, reader, instrument): + """Behavioral check: a native provider class (BaseLLM subclass, not LLM) + going through the wrapper must emit an llm span with the authoritative + provider attribute and record a duration histogram sample.""" + from unittest.mock import MagicMock, patch + + from crewai.llms.providers.openai.completion import OpenAICompletion + + llm = OpenAICompletion(model="gpt-4o-mini", api_key=os.environ.get("OPENAI_API_KEY", "test-key")) + fake_response = MagicMock() + fake_response.choices = [MagicMock()] + fake_response.choices[0].message.content = "stubbed" + fake_response.model = "gpt-4o-mini" + + with patch.object(llm._client.chat.completions, "create", return_value=fake_response): + llm.call([{"role": "user", "content": "hello"}]) + + llm_spans = [s for s in span_exporter.get_finished_spans() if s.name == "gpt-4o-mini.llm"] + assert len(llm_spans) == 1 + attrs = dict(llm_spans[0].attributes) + assert attrs.get("gen_ai.provider.name") == "openai" + assert attrs.get("gen_ai.request.model") == "gpt-4o-mini" + + metrics = reader.get_metrics_data() + duration_points = [] + for resource_metrics in metrics.resource_metrics: + for scope_metrics in resource_metrics.scope_metrics: + for metric in scope_metrics.metrics: + if metric.name == "gen_ai.client.operation.duration": + duration_points.extend(metric.data.data_points) + assert duration_points, "duration histogram recorded nothing" + assert all(dp.attributes.get("gen_ai.provider.name") == "openai" for dp in duration_points)