Skip to content
Open
Show file tree
Hide file tree
Changes from 3 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import logging
import os
import time
from typing import Collection
Expand All @@ -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 = {
Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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,
Expand Down
69 changes: 69 additions & 0 deletions packages/opentelemetry-instrumentation-crewai/tests/conftest.py
Original file line number Diff line number Diff line change
@@ -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()
Original file line number Diff line number Diff line change
Expand Up @@ -63,10 +63,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):
Expand All @@ -80,12 +77,73 @@ 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="sk-test")
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Outdated
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"
Comment thread
coderabbitai[bot] marked this conversation as resolved.