From fa04ffe8c83abc0b6c25da1a2ce3e57179799f63 Mon Sep 17 00:00:00 2001 From: Emmett Butler Date: Fri, 2 Oct 2026 10:46:49 -0700 Subject: [PATCH 1/9] refactor: move LLM stream handler utilities into contrib base_stream_handler.py has no LLM Observability dependencies, but it lived under ddtrace/llmobs/_integrations, so every contrib that traces streamed responses (and AI Guard) imported from the LLMObs product. Move it to ddtrace/contrib/internal/stream_handler.py and update all importers. The unit tests are renamed to tests/llmobs/test_stream_handler.py, and the new path is added to the llmobs suitespec component so changes to it keep triggering the LLM integration suites. Co-Authored-By: Claude Opus 5.5 --- .claude/skills/llmobs-integrations/SKILL.md | 2 +- .../references/implementation-guide.md | 2 +- ddtrace/aiguard/_streaming.py | 6 +++--- ddtrace/contrib/internal/anthropic/_streaming.py | 6 +++--- .../contrib/internal/botocore/services/bedrock.py | 4 ++-- .../contrib/internal/claude_agent_sdk/_streaming.py | 4 ++-- ddtrace/contrib/internal/google_genai/_utils.py | 4 ++-- ddtrace/contrib/internal/google_genai/patch.py | 2 +- ddtrace/contrib/internal/langchain/utils.py | 6 +++--- ddtrace/contrib/internal/litellm/patch.py | 2 +- ddtrace/contrib/internal/litellm/utils.py | 4 ++-- ddtrace/contrib/internal/llama_index/_streaming.py | 6 +++--- ddtrace/contrib/internal/mistralai/_utils.py | 6 +++--- ddtrace/contrib/internal/mistralai/patch.py | 2 +- ddtrace/contrib/internal/openai/_endpoint_hooks.py | 2 +- ddtrace/contrib/internal/openai/utils.py | 4 ++-- .../internal/stream_handler.py} | 2 +- ddtrace/contrib/internal/vertexai/_utils.py | 4 ++-- ddtrace/contrib/internal/vertexai/patch.py | 2 +- mypy.ini | 12 ++++++------ tests/llmobs/suitespec.yml | 1 + ...base_stream_handler.py => test_stream_handler.py} | 8 ++++---- 22 files changed, 46 insertions(+), 45 deletions(-) rename ddtrace/{llmobs/_integrations/base_stream_handler.py => contrib/internal/stream_handler.py} (99%) rename tests/llmobs/{test_base_stream_handler.py => test_stream_handler.py} (98%) diff --git a/.claude/skills/llmobs-integrations/SKILL.md b/.claude/skills/llmobs-integrations/SKILL.md index ecb1570c5de..28147590877 100644 --- a/.claude/skills/llmobs-integrations/SKILL.md +++ b/.claude/skills/llmobs-integrations/SKILL.md @@ -34,7 +34,7 @@ Both layers must work together. The patch layer identifies the operation and pas | Purpose | File | |---------|------| | Base LLM integration class | `ddtrace/llmobs/_integrations/base.py` (`BaseLLMIntegration`) | -| Stream handler base classes | `ddtrace/llmobs/_integrations/base_stream_handler.py` (`BaseStreamHandler`, `StreamHandler`, `AsyncStreamHandler`) | +| Stream handler base classes | `ddtrace/contrib/internal/stream_handler.py` (`BaseStreamHandler`, `StreamHandler`, `AsyncStreamHandler`) | | Shared utilities | `ddtrace/llmobs/_integrations/utils.py` | | LLMObs annotation helper | `ddtrace/llmobs/_utils.py` (`_annotate_llmobs_span_data`) | | LLMObs constants | `ddtrace/llmobs/_constants.py` | diff --git a/.claude/skills/llmobs-integrations/references/implementation-guide.md b/.claude/skills/llmobs-integrations/references/implementation-guide.md index 2d2ef9f0d58..ff5a836ed48 100644 --- a/.claude/skills/llmobs-integrations/references/implementation-guide.md +++ b/.claude/skills/llmobs-integrations/references/implementation-guide.md @@ -105,7 +105,7 @@ Some older or specialized integrations still call `integration.trace()` and `int ## Streaming -Subclass `StreamHandler`/`AsyncStreamHandler` from `ddtrace/llmobs/_integrations/base_stream_handler.py`: +Subclass `StreamHandler`/`AsyncStreamHandler` from `ddtrace/contrib/internal/stream_handler.py`: - `initialize_chunk_storage()` — set up accumulators for content, usage, role - `process_chunk(chunk)` — accumulate text, tool blocks, usage from each chunk diff --git a/ddtrace/aiguard/_streaming.py b/ddtrace/aiguard/_streaming.py index 45f73e8a53d..d37c47982df 100644 --- a/ddtrace/aiguard/_streaming.py +++ b/ddtrace/aiguard/_streaming.py @@ -22,10 +22,10 @@ import wrapt from ddtrace.aiguard._context import is_aiguard_context_active +from ddtrace.contrib.internal.stream_handler import AsyncStreamHandler +from ddtrace.contrib.internal.stream_handler import BaseStreamHandler import ddtrace.internal.logger as ddlogger from ddtrace.internal.settings.aiguard import aiguard_config -from ddtrace.llmobs._integrations.base_stream_handler import AsyncStreamHandler -from ddtrace.llmobs._integrations.base_stream_handler import BaseStreamHandler logger = ddlogger.get_logger(__name__) @@ -146,7 +146,7 @@ def __next__(self) -> Any: # ------------------------------------------------------------------ # Context-manager protocol # - # TracedStream.__enter__() (base_stream_handler.py) has two branches: + # TracedStream.__enter__() (ddtrace/contrib/internal/stream_handler.py) has two branches: # - non-manager (raw Stream): returns ``self`` (the TracedStream). # - manager (MessageStreamManager): returns a NEW TracedStream # wrapping the inner MessageStream. diff --git a/ddtrace/contrib/internal/anthropic/_streaming.py b/ddtrace/contrib/internal/anthropic/_streaming.py index ed83a93ffb6..be548da49bf 100644 --- a/ddtrace/contrib/internal/anthropic/_streaming.py +++ b/ddtrace/contrib/internal/anthropic/_streaming.py @@ -3,11 +3,11 @@ import anthropic +from ddtrace.contrib.internal.stream_handler import AsyncStreamHandler +from ddtrace.contrib.internal.stream_handler import StreamHandler +from ddtrace.contrib.internal.stream_handler import make_traced_stream from ddtrace.internal.logger import get_logger from ddtrace.internal.span_bus import span_from_context -from ddtrace.llmobs._integrations.base_stream_handler import AsyncStreamHandler -from ddtrace.llmobs._integrations.base_stream_handler import StreamHandler -from ddtrace.llmobs._integrations.base_stream_handler import make_traced_stream from ddtrace.llmobs._utils import _get_attr from ddtrace.llmobs._utils import safe_load_json diff --git a/ddtrace/contrib/internal/botocore/services/bedrock.py b/ddtrace/contrib/internal/botocore/services/bedrock.py index 6c866aa292c..cdb794b32be 100644 --- a/ddtrace/contrib/internal/botocore/services/bedrock.py +++ b/ddtrace/contrib/internal/botocore/services/bedrock.py @@ -5,6 +5,8 @@ from typing import Optional from ddtrace import config +from ddtrace.contrib.internal.stream_handler import StreamHandler +from ddtrace.contrib.internal.stream_handler import make_traced_stream from ddtrace.contrib.internal.trace_utils import ext_service from ddtrace.ext import SpanTypes from ddtrace.internal import core @@ -16,8 +18,6 @@ from ddtrace.llmobs._integrations._bedrock_inference_profiles import lookup_inference_profile from ddtrace.llmobs._integrations._bedrock_inference_profiles import record_inference_profile from ddtrace.llmobs._integrations._bedrock_inference_profiles import record_resolve_failure -from ddtrace.llmobs._integrations.base_stream_handler import StreamHandler -from ddtrace.llmobs._integrations.base_stream_handler import make_traced_stream from ddtrace.llmobs._integrations.bedrock_utils import _AI21 from ddtrace.llmobs._integrations.bedrock_utils import _AMAZON from ddtrace.llmobs._integrations.bedrock_utils import _ANTHROPIC diff --git a/ddtrace/contrib/internal/claude_agent_sdk/_streaming.py b/ddtrace/contrib/internal/claude_agent_sdk/_streaming.py index 8c88df28c5f..253be415eb4 100644 --- a/ddtrace/contrib/internal/claude_agent_sdk/_streaming.py +++ b/ddtrace/contrib/internal/claude_agent_sdk/_streaming.py @@ -10,10 +10,10 @@ from ddtrace.contrib.internal.claude_agent_sdk.utils import _extract_model_from_response from ddtrace.contrib.internal.claude_agent_sdk.utils import _retrieve_context from ddtrace.contrib.internal.claude_agent_sdk.utils import extract_partial_message_usage +from ddtrace.contrib.internal.stream_handler import AsyncStreamHandler +from ddtrace.contrib.internal.stream_handler import make_traced_stream from ddtrace.internal.logger import get_logger from ddtrace.internal.utils.formats import format_trace_id -from ddtrace.llmobs._integrations.base_stream_handler import AsyncStreamHandler -from ddtrace.llmobs._integrations.base_stream_handler import make_traced_stream from ddtrace.llmobs._utils import add_span_link from ddtrace.llmobs._utils import safe_json from ddtrace.llmobs.types import Message diff --git a/ddtrace/contrib/internal/google_genai/_utils.py b/ddtrace/contrib/internal/google_genai/_utils.py index eef190dc09e..8df8405b2cb 100644 --- a/ddtrace/contrib/internal/google_genai/_utils.py +++ b/ddtrace/contrib/internal/google_genai/_utils.py @@ -1,8 +1,8 @@ from typing import Any from typing import Optional -from ddtrace.llmobs._integrations.base_stream_handler import AsyncStreamHandler -from ddtrace.llmobs._integrations.base_stream_handler import StreamHandler +from ddtrace.contrib.internal.stream_handler import AsyncStreamHandler +from ddtrace.contrib.internal.stream_handler import StreamHandler from ddtrace.llmobs._integrations.google_utils import GOOGLE_GENAI_DEFAULT_MODEL_ROLE from ddtrace.llmobs._utils import _get_attr diff --git a/ddtrace/contrib/internal/google_genai/patch.py b/ddtrace/contrib/internal/google_genai/patch.py index c6735a9e6f2..bbe6b87d670 100644 --- a/ddtrace/contrib/internal/google_genai/patch.py +++ b/ddtrace/contrib/internal/google_genai/patch.py @@ -5,10 +5,10 @@ from ddtrace import config from ddtrace.contrib.internal.google_genai._utils import GoogleGenAIAsyncStreamHandler from ddtrace.contrib.internal.google_genai._utils import GoogleGenAIStreamHandler +from ddtrace.contrib.internal.stream_handler import make_traced_stream from ddtrace.contrib.internal.trace_utils import unwrap from ddtrace.contrib.internal.trace_utils import wrap from ddtrace.llmobs._integrations import GoogleGenAIIntegration -from ddtrace.llmobs._integrations.base_stream_handler import make_traced_stream from ddtrace.llmobs._integrations.google_utils import extract_provider_and_model_name diff --git a/ddtrace/contrib/internal/langchain/utils.py b/ddtrace/contrib/internal/langchain/utils.py index b7a758f7606..e66e175c393 100644 --- a/ddtrace/contrib/internal/langchain/utils.py +++ b/ddtrace/contrib/internal/langchain/utils.py @@ -1,11 +1,11 @@ import inspect import sys +from ddtrace.contrib.internal.stream_handler import AsyncStreamHandler +from ddtrace.contrib.internal.stream_handler import StreamHandler +from ddtrace.contrib.internal.stream_handler import make_traced_stream from ddtrace.internal import core from ddtrace.internal._exceptions import DDBlockException -from ddtrace.llmobs._integrations.base_stream_handler import AsyncStreamHandler -from ddtrace.llmobs._integrations.base_stream_handler import StreamHandler -from ddtrace.llmobs._integrations.base_stream_handler import make_traced_stream class BaseLangchainStreamHandler: diff --git a/ddtrace/contrib/internal/litellm/patch.py b/ddtrace/contrib/internal/litellm/patch.py index 01ddb941c5d..af958113022 100644 --- a/ddtrace/contrib/internal/litellm/patch.py +++ b/ddtrace/contrib/internal/litellm/patch.py @@ -6,12 +6,12 @@ from ddtrace.contrib.internal.litellm.utils import LiteLLMAsyncStreamHandler from ddtrace.contrib.internal.litellm.utils import LiteLLMStreamHandler from ddtrace.contrib.internal.litellm.utils import extract_host_tag +from ddtrace.contrib.internal.stream_handler import make_traced_stream from ddtrace.contrib.trace_utils import unwrap from ddtrace.contrib.trace_utils import wrap from ddtrace.internal.utils import get_argument_value from ddtrace.llmobs._constants import LITELLM_ROUTER_INSTANCE_KEY from ddtrace.llmobs._integrations import LiteLLMIntegration -from ddtrace.llmobs._integrations.base_stream_handler import make_traced_stream config._add("litellm", {}) diff --git a/ddtrace/contrib/internal/litellm/utils.py b/ddtrace/contrib/internal/litellm/utils.py index 6497a2bd4d5..f22673e18eb 100644 --- a/ddtrace/contrib/internal/litellm/utils.py +++ b/ddtrace/contrib/internal/litellm/utils.py @@ -1,11 +1,11 @@ from collections import defaultdict from dataclasses import FrozenInstanceError +from ddtrace.contrib.internal.stream_handler import AsyncStreamHandler +from ddtrace.contrib.internal.stream_handler import StreamHandler from ddtrace.internal.logger import get_logger from ddtrace.internal.utils import get_argument_value from ddtrace.llmobs._constants import LITELLM_ROUTER_INSTANCE_KEY -from ddtrace.llmobs._integrations.base_stream_handler import AsyncStreamHandler -from ddtrace.llmobs._integrations.base_stream_handler import StreamHandler from ddtrace.llmobs._integrations.utils import openai_construct_completion_from_streamed_chunks from ddtrace.llmobs._integrations.utils import openai_construct_message_from_streamed_chunks diff --git a/ddtrace/contrib/internal/llama_index/_streaming.py b/ddtrace/contrib/internal/llama_index/_streaming.py index 60a629c54c2..579f22973de 100644 --- a/ddtrace/contrib/internal/llama_index/_streaming.py +++ b/ddtrace/contrib/internal/llama_index/_streaming.py @@ -6,13 +6,13 @@ from typing import Optional from typing import Union +from ddtrace.contrib.internal.stream_handler import AsyncStreamHandler +from ddtrace.contrib.internal.stream_handler import StreamHandler +from ddtrace.contrib.internal.stream_handler import make_traced_stream from ddtrace.internal import core from ddtrace.internal.logger import get_logger from ddtrace.internal.span_bus import span_from_context from ddtrace.llmobs._integrations import LlamaIndexIntegration -from ddtrace.llmobs._integrations.base_stream_handler import AsyncStreamHandler -from ddtrace.llmobs._integrations.base_stream_handler import StreamHandler -from ddtrace.llmobs._integrations.base_stream_handler import make_traced_stream if TYPE_CHECKING: diff --git a/ddtrace/contrib/internal/mistralai/_utils.py b/ddtrace/contrib/internal/mistralai/_utils.py index 5c95a3495c5..a65816342d4 100644 --- a/ddtrace/contrib/internal/mistralai/_utils.py +++ b/ddtrace/contrib/internal/mistralai/_utils.py @@ -1,9 +1,9 @@ from typing import Any from typing import Optional -from ddtrace.llmobs._integrations.base_stream_handler import AsyncStreamHandler -from ddtrace.llmobs._integrations.base_stream_handler import BaseStreamHandler -from ddtrace.llmobs._integrations.base_stream_handler import StreamHandler +from ddtrace.contrib.internal.stream_handler import AsyncStreamHandler +from ddtrace.contrib.internal.stream_handler import BaseStreamHandler +from ddtrace.contrib.internal.stream_handler import StreamHandler from ddtrace.llmobs._utils import _get_attr diff --git a/ddtrace/contrib/internal/mistralai/patch.py b/ddtrace/contrib/internal/mistralai/patch.py index ac3d7bc9079..125e5efb089 100644 --- a/ddtrace/contrib/internal/mistralai/patch.py +++ b/ddtrace/contrib/internal/mistralai/patch.py @@ -12,10 +12,10 @@ from ddtrace import config from ddtrace.contrib.internal.mistralai._utils import MistralAIAsyncStreamHandler from ddtrace.contrib.internal.mistralai._utils import MistralAIStreamHandler +from ddtrace.contrib.internal.stream_handler import make_traced_stream from ddtrace.contrib.internal.trace_utils import unwrap from ddtrace.contrib.internal.trace_utils import wrap from ddtrace.llmobs._integrations import MistralAIIntegration -from ddtrace.llmobs._integrations.base_stream_handler import make_traced_stream from ddtrace.llmobs._integrations.mistralai_utils import extract_provider diff --git a/ddtrace/contrib/internal/openai/_endpoint_hooks.py b/ddtrace/contrib/internal/openai/_endpoint_hooks.py index 0f6e6b06406..2e0f2629ef1 100644 --- a/ddtrace/contrib/internal/openai/_endpoint_hooks.py +++ b/ddtrace/contrib/internal/openai/_endpoint_hooks.py @@ -6,9 +6,9 @@ from ddtrace.contrib.internal.openai.utils import _is_generator from ddtrace.contrib.internal.openai.utils import _loop_handler from ddtrace.contrib.internal.openai.utils import _process_finished_stream +from ddtrace.contrib.internal.stream_handler import make_traced_stream from ddtrace.internal.utils.version import parse_version from ddtrace.llmobs._constants import OAI_HANDOFF_TOOL_ARG -from ddtrace.llmobs._integrations.base_stream_handler import make_traced_stream from ddtrace.llmobs._utils import _get_attr from ddtrace.llmobs._utils import safe_load_json diff --git a/ddtrace/contrib/internal/openai/utils.py b/ddtrace/contrib/internal/openai/utils.py index 33b2ac34127..7f8a94b1bdb 100644 --- a/ddtrace/contrib/internal/openai/utils.py +++ b/ddtrace/contrib/internal/openai/utils.py @@ -3,9 +3,9 @@ from typing import Generator from typing import Optional +from ddtrace.contrib.internal.stream_handler import AsyncStreamHandler +from ddtrace.contrib.internal.stream_handler import StreamHandler from ddtrace.internal.logger import get_logger -from ddtrace.llmobs._integrations.base_stream_handler import AsyncStreamHandler -from ddtrace.llmobs._integrations.base_stream_handler import StreamHandler from ddtrace.llmobs._integrations.utils import openai_construct_completion_from_streamed_chunks from ddtrace.llmobs._integrations.utils import openai_construct_message_from_streamed_chunks diff --git a/ddtrace/llmobs/_integrations/base_stream_handler.py b/ddtrace/contrib/internal/stream_handler.py similarity index 99% rename from ddtrace/llmobs/_integrations/base_stream_handler.py rename to ddtrace/contrib/internal/stream_handler.py index d363a544fda..e88cec8caab 100644 --- a/ddtrace/llmobs/_integrations/base_stream_handler.py +++ b/ddtrace/contrib/internal/stream_handler.py @@ -1,5 +1,5 @@ """ -This file contains shared utilities for tracing streams in LLMobs integrations. Integrations should +This file contains shared utilities for tracing streams in LLM integrations. Integrations should implement a StreamHandler and / or AsyncStreamHandler subclass to be passed into the make_traced_stream factory function along with the stream to wrap. """ diff --git a/ddtrace/contrib/internal/vertexai/_utils.py b/ddtrace/contrib/internal/vertexai/_utils.py index 091f7ba9a3a..e44d70c1776 100644 --- a/ddtrace/contrib/internal/vertexai/_utils.py +++ b/ddtrace/contrib/internal/vertexai/_utils.py @@ -1,5 +1,5 @@ -from ddtrace.llmobs._integrations.base_stream_handler import AsyncStreamHandler -from ddtrace.llmobs._integrations.base_stream_handler import StreamHandler +from ddtrace.contrib.internal.stream_handler import AsyncStreamHandler +from ddtrace.contrib.internal.stream_handler import StreamHandler class BaseVertexAIStreamHandler: diff --git a/ddtrace/contrib/internal/vertexai/patch.py b/ddtrace/contrib/internal/vertexai/patch.py index be991b3d695..eb426790821 100644 --- a/ddtrace/contrib/internal/vertexai/patch.py +++ b/ddtrace/contrib/internal/vertexai/patch.py @@ -6,12 +6,12 @@ from vertexai.generative_models import GenerativeModel # noqa:F401 from ddtrace import config +from ddtrace.contrib.internal.stream_handler import make_traced_stream from ddtrace.contrib.internal.trace_utils import unwrap from ddtrace.contrib.internal.trace_utils import wrap from ddtrace.contrib.internal.vertexai._utils import VertexAIAsyncStreamHandler from ddtrace.contrib.internal.vertexai._utils import VertexAIStreamHandler from ddtrace.llmobs._integrations import VertexAIIntegration -from ddtrace.llmobs._integrations.base_stream_handler import make_traced_stream from ddtrace.llmobs._integrations.google_utils import extract_provider_and_model_name diff --git a/mypy.ini b/mypy.ini index 9e96b75fd9f..7b461e85879 100644 --- a/mypy.ini +++ b/mypy.ini @@ -733,12 +733,6 @@ disallow_any_generics = false disallow_untyped_defs = false disallow_incomplete_defs = false -[mypy-ddtrace.llmobs._integrations.base_stream_handler] -disallow_subclassing_any = false -disallow_untyped_calls = false -disallow_untyped_defs = false -disallow_incomplete_defs = false - [mypy-ddtrace.llmobs._integrations.bedrock] disallow_any_generics = false disallow_untyped_calls = false @@ -3027,6 +3021,12 @@ disallow_untyped_defs = false disallow_incomplete_defs = false check_untyped_defs = false +[mypy-ddtrace.contrib.internal.stream_handler] +disallow_subclassing_any = false +disallow_untyped_calls = false +disallow_untyped_defs = false +disallow_incomplete_defs = false + [mypy-ddtrace.contrib.internal.structlog.patch] disallow_untyped_calls = false disallow_untyped_defs = false diff --git a/tests/llmobs/suitespec.yml b/tests/llmobs/suitespec.yml index de8a0a0e3e8..a038817c8a8 100644 --- a/tests/llmobs/suitespec.yml +++ b/tests/llmobs/suitespec.yml @@ -18,6 +18,7 @@ components: - ddtrace/contrib/internal/litellm/* llmobs: - ddtrace/llmobs/* + - ddtrace/contrib/internal/stream_handler.py mcp: - ddtrace/contrib/internal/mcp/* mistralai: diff --git a/tests/llmobs/test_base_stream_handler.py b/tests/llmobs/test_stream_handler.py similarity index 98% rename from tests/llmobs/test_base_stream_handler.py rename to tests/llmobs/test_stream_handler.py index 9050cf2bfee..dc0b221bf79 100644 --- a/tests/llmobs/test_base_stream_handler.py +++ b/tests/llmobs/test_stream_handler.py @@ -10,11 +10,11 @@ import pytest +from ddtrace.contrib.internal.stream_handler import AsyncStreamHandler +from ddtrace.contrib.internal.stream_handler import BaseStreamHandler +from ddtrace.contrib.internal.stream_handler import StreamHandler +from ddtrace.contrib.internal.stream_handler import make_traced_stream from ddtrace.internal._exceptions import DDBlockException -from ddtrace.llmobs._integrations.base_stream_handler import AsyncStreamHandler -from ddtrace.llmobs._integrations.base_stream_handler import BaseStreamHandler -from ddtrace.llmobs._integrations.base_stream_handler import StreamHandler -from ddtrace.llmobs._integrations.base_stream_handler import make_traced_stream class _RecordingMixin: From 72cb2aeaf6f21eacafd7a275d297b99b0fa5a700 Mon Sep 17 00:00:00 2001 From: Emmett Butler Date: Fri, 2 Oct 2026 10:04:41 -0700 Subject: [PATCH 2/9] abstract llmobs from anthropic contrib using the subscriber pattern --- .../references/implementation-guide.md | 2 + .claude/skills/llmobs-integrations/SKILL.md | 11 ++- .../references/failure-modes.md | 4 +- .../references/implementation-guide.md | 26 ++++-- ddtrace/_trace/subscribers/llm.py | 16 +++- ddtrace/contrib/_events/llm.py | 26 +++++- .../contrib/internal/anthropic/_streaming.py | 21 ++++- ddtrace/contrib/internal/anthropic/patch.py | 45 ++++----- ddtrace/llmobs/_contrib/__init__.py | 0 ddtrace/llmobs/_contrib/anthropic/__init__.py | 15 +++ .../llmobs/_contrib/anthropic/subscribers.py | 91 +++++++++++++++++++ ddtrace/llmobs/_integrations/anthropic.py | 16 +--- ddtrace/llmobs/_llmobs.py | 5 + ddtrace/llmobs/_product.py | 31 +++++++ ..._llm-obs_v2_eval-metric_post_1dc109f6.json | 26 ++++++ ..._llm-obs_v2_eval-metric_post_577483a9.json | 26 ++++++ ...xperiments_mock-exp-id_patch_46e5df7f.json | 26 ++++++ ...xperiments_mock-exp-id_patch_4b861664.json | 26 ++++++ ...ments_new-rerun-exp-id_patch_5f0d9491.json | 26 ++++++ ...ew-rerun-experiment-id_patch_e2ff88d7.json | 26 ++++++ ...periments_new-rerun-id_patch_23512e67.json | 26 ++++++ tests/contrib/anthropic/conftest.py | 4 + 22 files changed, 436 insertions(+), 59 deletions(-) create mode 100644 ddtrace/llmobs/_contrib/__init__.py create mode 100644 ddtrace/llmobs/_contrib/anthropic/__init__.py create mode 100644 ddtrace/llmobs/_contrib/anthropic/subscribers.py create mode 100644 tests/cassettes/datadog/datadog_api_intake_llm-obs_v2_eval-metric_post_1dc109f6.json create mode 100644 tests/cassettes/datadog/datadog_api_intake_llm-obs_v2_eval-metric_post_577483a9.json create mode 100644 tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_mock-exp-id_patch_46e5df7f.json create mode 100644 tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_mock-exp-id_patch_4b861664.json create mode 100644 tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-exp-id_patch_5f0d9491.json create mode 100644 tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-experiment-id_patch_e2ff88d7.json create mode 100644 tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-id_patch_23512e67.json diff --git a/.claude/skills/apm-integrations/references/implementation-guide.md b/.claude/skills/apm-integrations/references/implementation-guide.md index f0047672593..2f9897c0ce4 100644 --- a/.claude/skills/apm-integrations/references/implementation-guide.md +++ b/.claude/skills/apm-integrations/references/implementation-guide.md @@ -65,6 +65,8 @@ it's patched. **Without this entry, integration config settings are silently ign Create `ddtrace/llmobs/_integrations/{name}.py` subclassing `BaseLLMIntegration`. Read `ddtrace/llmobs/_integrations/anthropic.py` for the canonical pattern. Register in `ddtrace/llmobs/_integrations/__init__.py` (import + `__all__`). +Connect it to the contrib with `LlmEvents` subscribers in `ddtrace/llmobs/_contrib/{name}/` +(see `ddtrace/llmobs/_contrib/anthropic/`) so the contrib does not import `ddtrace.llmobs`. See the **llmobs-integrations** skill for the full LLM-specific implementation guide. diff --git a/.claude/skills/llmobs-integrations/SKILL.md b/.claude/skills/llmobs-integrations/SKILL.md index 28147590877..c48e07a7987 100644 --- a/.claude/skills/llmobs-integrations/SKILL.md +++ b/.claude/skills/llmobs-integrations/SKILL.md @@ -19,14 +19,16 @@ LLMObs integrations enable Datadog LLM Observability for AI/LLM libraries. They LLMObs integrations consist of two cooperating layers: -1. **Patch Layer** (`ddtrace/contrib/internal/{name}/patch.py`) -- wraps library functions. Standard request/response LLM integrations construct `LlmRequestEvent` and use `core.context_with_event()` so the LLM tracing subscriber owns span lifecycle and LLMObs tag extraction. +1. **Patch Layer** (`ddtrace/contrib/internal/{name}/patch.py`) -- wraps library functions. Standard request/response LLM integrations construct `LlmRequestEvent` and use `core.context_with_event()` so the LLM tracing subscriber owns span lifecycle. The patch layer must not import from `ddtrace.llmobs`. 2. **Integration Layer** (`ddtrace/llmobs/_integrations/{name}.py`) -- extends `BaseLLMIntegration`, implements `_set_base_span_tags()` and `_llmobs_set_tags()` to extract and set provider-specific messages, tools, metadata, and token metrics. +3. **LLMObs Subscribers** (`ddtrace/llmobs/_contrib/{name}/`) -- subscribe to the `LlmEvents` span lifecycle events (`SPAN_STARTING`, `SPAN_STARTED`, `SPAN_FINISHING`) that `LlmTracingSubscriber` dispatches, filter on `ctx.event.component`, and call the integration layer. `ddtrace/llmobs/_product.py:listen_integrations()` registers them when the library is patched. -Both layers must work together. The patch layer identifies the operation and passes request/response data through the event; the integration layer controls what data is extracted. +The layers must work together. The patch layer identifies the operation and passes request/response data through the event; the subscribers connect the event to the integration layer, which controls what data is extracted. ## Active Patch Patterns -- **Event-based request spans**: Use `LlmRequestEvent` with `core.context_with_event()` for new standard request/response LLM integrations. Anthropic is the canonical reference. This is the preferred pattern. +- **Event-based request spans with LLMObs subscribers**: Use `LlmRequestEvent` with `core.context_with_event()`, leave `llmobs_integration` unset, and add subscribers under `ddtrace/llmobs/_contrib/{name}/`. Anthropic is the canonical reference. This is the preferred pattern. +- **Event-based request spans with `llmobs_integration`**: Older event-based integrations pass `llmobs_integration=integration` on the event, and `LlmTracingSubscriber` calls the integration directly. This makes the contrib import LLMObs code; do not use it for new work. - **Direct integration spans**: Some existing or specialized integrations still call `integration.trace()` and `integration.llmobs_set_tags()` directly, especially for child spans, agent/tool spans, or integrations not yet migrated. Google GenAI, OpenAI tool spans, and Claude Agent SDK are useful references. ## Key Files @@ -116,7 +118,8 @@ Note two already-shipped integrations predate this key: bedrock and the claude-a - **Streaming** must use `BaseStreamHandler`/`AsyncStreamHandler` -- never consume streams directly - **Event-based patch wrappers** should not call `span.set_exc_info()`, `span.finish()`, or `integration.llmobs_set_tags()` directly; the tracing subscriber handles that when the event ends - **Direct integration spans** must keep `integration.llmobs_set_tags()` and span lifecycle handling aligned with the closest current reference -- **Integration instance** must be stored on the module: `module._datadog_integration = MyLibIntegration(integration_config=config.mylib)` +- **Integration instance**: subscriber-based integrations build it lazily in the LLMObs subscriber (see `ddtrace/llmobs/_contrib/anthropic/subscribers.py`); older integrations store it on the module as `module._datadog_integration = MyLibIntegration(integration_config=config.mylib)` +- **Subscriber registration**: subscriber-based integrations hook `{name}.patch`/`{name}.unpatch` core events in `listen_integrations()` and stay registered while LLMObs is disabled, because they also set the APM shadow tags ## Message Types diff --git a/.claude/skills/llmobs-integrations/references/failure-modes.md b/.claude/skills/llmobs-integrations/references/failure-modes.md index f8cd1092d8c..55aa69ee1c3 100644 --- a/.claude/skills/llmobs-integrations/references/failure-modes.md +++ b/.claude/skills/llmobs-integrations/references/failure-modes.md @@ -14,11 +14,11 @@ Comprehensive debugging guide for all known LLMObs integration failure modes. Ea **Causes:** 1. `submit_to_llmobs=True` not set on `LlmRequestEvent` for event-based patch code 2. `ctx.dispatch_ended_event()` not called, so `LlmTracingSubscriber` never calls `integration.llmobs_set_tags()` -3. Integration not instantiated -- `module._datadog_integration` is `None` +3. Integration not instantiated -- `module._datadog_integration` is `None`, or (for subscriber-based integrations) the `ddtrace/llmobs/_contrib/{name}` subscribers were never registered because `listen_integrations()` did not run before `patch()` 4. `llmobs_enabled` returns `False` -- LLMObs not configured in tracer config **Fix:** -- Verify `LlmRequestEvent(..., submit_to_llmobs=True, llmobs_integration=integration, request_kwargs=kwargs, ...)` in event-based patch wrappers +- Verify `LlmRequestEvent(..., submit_to_llmobs=True, request_kwargs=kwargs, ...)` in event-based patch wrappers, plus either registered LLMObs subscribers or `llmobs_integration=integration` for older integrations - Verify success and error paths call `ctx.dispatch_ended_event(...)` - Verify `patch()` stores integration: `module._datadog_integration = MyIntegration(integration_config=config.mylib)` - Check `DD_LLMOBS_ENABLED=1` is set diff --git a/.claude/skills/llmobs-integrations/references/implementation-guide.md b/.claude/skills/llmobs-integrations/references/implementation-guide.md index ff5a836ed48..225002e4b2f 100644 --- a/.claude/skills/llmobs-integrations/references/implementation-guide.md +++ b/.claude/skills/llmobs-integrations/references/implementation-guide.md @@ -5,8 +5,9 @@ Follow the **apm-integrations** skill's [Implementation Guide](../../apm-integra ## Design: Two-Layer Architecture LLM integrations use `BaseLLMIntegration` as a second layer on top of a standard APM integration: -- Patch code creates a subclass instance and stores it on the module: `module._datadog_integration = MyLibIntegration(integration_config=config.mylib)` -- Standard request/response patch code constructs `LlmRequestEvent` and uses `core.context_with_event()`; `LlmTracingSubscriber` manages span lifecycle and calls `integration.llmobs_set_tags()` when the event ends +- Standard request/response patch code constructs `LlmRequestEvent` and uses `core.context_with_event()`; `LlmTracingSubscriber` manages span lifecycle and dispatches `LlmEvents.SPAN_STARTING`, `SPAN_STARTED`, and `SPAN_FINISHING` with the execution context +- LLMObs subscribers in `ddtrace/llmobs/_contrib/{name}/` listen to those events, build the `BaseLLMIntegration` subclass lazily, and call it to set the span type, base tags, and LLMObs tags, so the contrib never imports `ddtrace.llmobs` +- Older integrations instead store the instance on the module (`module._datadog_integration = ...`) and pass it as `LlmRequestEvent(llmobs_integration=...)`; `LlmTracingSubscriber` then calls it directly - The `BaseLLMIntegration` subclass in `ddtrace/llmobs/_integrations/` handles provider-specific message, token, metadata, and tool extraction - Some existing or specialized integrations still call `integration.trace()` directly for direct child spans; follow the closest current reference before using that pattern @@ -16,6 +17,7 @@ This separation keeps APM patching decoupled from LLMObs data extraction. An LLM integration is an APM integration with an extra layer. You do everything in the apm-integrations guide, but: - **Step 1 (patch module)**: Use `LlmRequestEvent` with `core.context_with_event()` for standard request/response LLM integrations +- **Step 1b (LLMObs subscribers)**: Add `ddtrace/llmobs/_contrib/{name}/` subscribers and register them from `listen_integrations()` in `ddtrace/llmobs/_product.py` - **Step 3 (LLMObs integration)**: Create the `BaseLLMIntegration` subclass that handles provider-specific message, tool, and token extraction (this guide) - **Step 4 (test environment)**: Use `tests/llmobs/suitespec.yml`; add `vcrpy` only when the suite uses vcrpy cassettes and follow nearby version pins - **Step 5 (tests)**: Add `test_{name}_llmobs.py` in addition to the APM `test_{name}.py`, using the right transport pattern for the integration and `assert_llmobs_span_data(_get_llmobs_data_metastruct(span), ...)` @@ -90,15 +92,26 @@ Register in `ddtrace/llmobs/_integrations/__init__.py` (import + `__all__` entry ## Step 1 Expanded: Patch Layer (`LlmRequestEvent`) -Standard LLM integrations should use `LlmRequestEvent` from `ddtrace/contrib/_events/llm.py` with `core.context_with_event()`. Read `ddtrace/contrib/internal/anthropic/patch.py` for the current event-based pattern. The patch layer constructs the event, stores the response on `event.response`, and calls `ctx.dispatch_ended_event()`; `ddtrace/_trace/subscribers/llm.py` handles span creation, base tags, LLMObs extraction, errors, and span finish under the hood. +Standard LLM integrations should use `LlmRequestEvent` from `ddtrace/contrib/_events/llm.py` with `core.context_with_event()`. Read `ddtrace/contrib/internal/anthropic/patch.py` for the current event-based pattern. The patch layer constructs the event, stores the response on `event.response`, and calls `ctx.dispatch_ended_event()`; `ddtrace/_trace/subscribers/llm.py` handles span creation, errors, and span finish, and dispatches `LlmEvents` so LLMObs subscribers can add base tags and LLMObs data. Key points: -- Construct `LlmRequestEvent(..., llmobs_integration=integration, submit_to_llmobs=True, request_kwargs=kwargs, ...)` +- Construct `LlmRequestEvent(..., submit_to_llmobs=True, request_kwargs=kwargs, ...)` and leave `llmobs_integration` unset +- Set APM tags the contrib owns (for example `{name}.request.model`) through the event's `tags`, not in the `BaseLLMIntegration` +- Shared stream helpers live in `ddtrace/contrib/internal/stream_handler.py`; pass `None` as the integration when the contrib has none - The event/subscriber path owns span creation and finishing; patch wrappers should not call `tracer.trace()`, `integration.trace()`, or create spans directly for standard request spans - Use `with core.context_with_event(event, dispatch_end_event=False) as ctx:` when streaming or when the wrapper needs to dispatch the ended event manually - For non-streaming success, set `event.response = resp` and call `ctx.dispatch_ended_event()` - For errors, call `ctx.dispatch_ended_event(*sys.exc_info())` and re-raise; do not call `span.set_exc_info()` or `span.finish()` directly in the patch wrapper -- LLMObs tag setting is handled by `LlmTracingSubscriber`, not directly in the patch wrapper +- LLMObs tag setting is handled by the LLMObs subscribers, not directly in the patch wrapper + +## Step 1b Expanded: LLMObs Subscribers + +Read `ddtrace/llmobs/_contrib/anthropic/` for the pattern. `LlmTracingSubscriber` dispatches three events with the `ExecutionContext`. They are shared by every LLM integration, so each handler returns early unless `ctx.event.component` matches: +- `LlmEvents.SPAN_STARTING` runs before the span exists. This is the only point where `event.span_type` can still be set to `SpanTypes.LLM`, which `LLMObs._on_span_start` needs at creation. +- `LlmEvents.SPAN_STARTED` runs after the span is created. Use it for `_set_base_span_tags()`, `_annotate_integration_tag()`, and `_stamp_llmobs_span_kind_at_start()`. +- `LlmEvents.SPAN_FINISHING` runs before the span is finished. Call `integration.llmobs_set_tags()` here. + +Subscribers set `auto_register = False`. Register them from `listen_integrations()` in `ddtrace/llmobs/_product.py` on the `{name}.patch` core event, and unregister them on `{name}.unpatch`. `listen_integrations()` runs from both the product's `post_preload` and `LLMObs.enable()`. Test fixtures that call `patch()` directly should call `listen_integrations()` first. - The async variant is identical but uses `async def` / `await` Some older or specialized integrations still call `integration.trace()` and `integration.llmobs_set_tags()` directly. Use that pattern when modifying an existing integration that already does so, when the closest current reference uses it (for example Google GenAI), or when the behavior requires direct child spans (for example OpenAI MCP tool spans or agent/tool child spans). @@ -212,7 +225,8 @@ In addition to the full checklist in the apm-integrations [Implementation Guide] - [ ] `ddtrace/llmobs/_integrations/{name}.py` — `BaseLLMIntegration` subclass - [ ] `ddtrace/llmobs/_integrations/__init__.py` — import + `__all__` entry -- [ ] `ddtrace/contrib/internal/{name}/patch.py` — uses `LlmRequestEvent` + `core.context_with_event()` for standard LLM request spans (see anthropic for pattern) +- [ ] `ddtrace/contrib/internal/{name}/patch.py` — uses `LlmRequestEvent` + `core.context_with_event()` for standard LLM request spans, with no `ddtrace.llmobs` imports (see anthropic for pattern) +- [ ] `ddtrace/llmobs/_contrib/{name}/` — `LlmEvents` subscribers, registered from `listen_integrations()` in `ddtrace/llmobs/_product.py` - [ ] `tests/llmobs/suitespec.yml` — LLMObs test suite entry - [ ] Test dependencies match the suite style; include `vcrpy` only when cassette replay is used - [ ] `docs/index.rst` — add integration to the docs index diff --git a/ddtrace/_trace/subscribers/llm.py b/ddtrace/_trace/subscribers/llm.py index b2456a1f538..8ff06787050 100644 --- a/ddtrace/_trace/subscribers/llm.py +++ b/ddtrace/_trace/subscribers/llm.py @@ -3,6 +3,7 @@ from ddtrace._trace.subscribers._base import TracingSubscriber from ddtrace.constants import SPAN_KIND +from ddtrace.contrib._events.llm import LlmEvents from ddtrace.contrib._events.llm import LlmRequestEvent from ddtrace.internal import core from ddtrace.internal.constants import COMPONENT @@ -23,11 +24,17 @@ class LlmTracingSubscriber(TracingSubscriber["LlmRequestEvent"]): Handles span creation, base tag setting, proxy detection, and LLMObs tag extraction. Provider-specific logic is delegated - to the integration object carried by the event. + to the integration object carried by the event, or, for events that + carry none, to subscribers of the LlmEvents span lifecycle events. """ event_names = (LlmRequestEvent.event_name,) + @classmethod + def _on_context_started(cls, ctx: core.ExecutionContext["LlmRequestEvent"]) -> None: + core.dispatch(LlmEvents.SPAN_STARTING.value, (ctx,)) + super()._on_context_started(ctx) + @classmethod def on_started(cls, ctx: core.ExecutionContext["LlmRequestEvent"]) -> None: event: LlmRequestEvent = ctx.event @@ -40,6 +47,10 @@ def on_started(cls, ctx: core.ExecutionContext["LlmRequestEvent"]) -> None: span._remove_attribute(COMPONENT) span._remove_attribute(SPAN_KIND) + if event.llmobs_integration is None: + core.dispatch(LlmEvents.SPAN_STARTED.value, (ctx,)) + return + event.llmobs_integration._set_base_span_tags( span, model=event.model, @@ -67,6 +78,9 @@ def on_ended( dispatch_ended_event(). """ event: LlmRequestEvent = ctx.event + if event.llmobs_integration is None: + core.dispatch(LlmEvents.SPAN_FINISHING.value, (ctx,)) + return event.llmobs_integration.llmobs_set_tags( span_from_context(ctx), args=[], diff --git a/ddtrace/contrib/_events/llm.py b/ddtrace/contrib/_events/llm.py index 7783b186ea5..16a3c0923bb 100644 --- a/ddtrace/contrib/_events/llm.py +++ b/ddtrace/contrib/_events/llm.py @@ -2,6 +2,7 @@ from dataclasses import dataclass from dataclasses import field +from enum import Enum from typing import TYPE_CHECKING from typing import Any from typing import Optional @@ -16,6 +17,18 @@ from ddtrace.llmobs._integrations.base import BaseLLMIntegration +class LlmEvents(Enum): + LLM_REQUEST = "llm.request" + # Dispatched by LlmTracingSubscriber with the ExecutionContext so products can + # take part in the span lifecycle without the contrib importing them: + # SPAN_STARTING fires before the span is created (the only point where + # event.span_type can still change), SPAN_STARTED right after, and + # SPAN_FINISHING before the span is finished. + SPAN_STARTING = "llm.request.span_starting" + SPAN_STARTED = "llm.request.span_started" + SPAN_FINISHING = "llm.request.span_finishing" + + @dataclass class LlmRequestEvent(TracingEvent): """LLM request event for all LLM integrations. @@ -25,12 +38,13 @@ class LlmRequestEvent(TracingEvent): (_set_base_span_tags, llmobs_set_tags). """ - event_name = "llm.request" + event_name = LlmEvents.LLM_REQUEST.value span_kind = SpanKind.CLIENT provider: str = event_field() model: Optional[str] = event_field(default=None) - llmobs_integration: BaseLLMIntegration = event_field() + # Integrations that have moved to LlmEvents subscribers leave this unset. + llmobs_integration: Optional[BaseLLMIntegration] = event_field(default=None) request_kwargs: dict[str, Any] = event_field(default_factory=dict) submit_to_llmobs: bool = event_field(default=False) instance: Optional[Any] = event_field(default=None) @@ -43,4 +57,10 @@ class LlmRequestEvent(TracingEvent): def __post_init__(self) -> None: self.operation_name = f"{self.component}.request" - self.span_type = SpanTypes.LLM if (self.submit_to_llmobs and self.llmobs_integration.llmobs_enabled) else None + self.span_type = ( + SpanTypes.LLM + if ( + self.submit_to_llmobs and self.llmobs_integration is not None and self.llmobs_integration.llmobs_enabled + ) + else None + ) diff --git a/ddtrace/contrib/internal/anthropic/_streaming.py b/ddtrace/contrib/internal/anthropic/_streaming.py index be548da49bf..adfe38699b6 100644 --- a/ddtrace/contrib/internal/anthropic/_streaming.py +++ b/ddtrace/contrib/internal/anthropic/_streaming.py @@ -1,4 +1,5 @@ from functools import partial +import json from typing import Any import anthropic @@ -15,6 +16,20 @@ log = get_logger(__name__) +def _get_attr(o: object, attr: str, default: object): + # Streamed chunks may be SDK objects or plain dicts. + if isinstance(o, dict): + return o.get(attr, default) + return getattr(o, attr, default) + + +def safe_load_json(value: str): + try: + return json.loads(value) + except (json.JSONDecodeError, TypeError): + return {"value": str(value)} + + def _text_stream_generator(traced_stream): for chunk in traced_stream: if chunk.type == "content_block_delta" and chunk.delta.type == "text_delta": @@ -37,7 +52,7 @@ async def _consume_async_stream(traced_stream: Any) -> None: pass -def handle_streamed_response(integration, resp, args, kwargs, ctx): +def handle_streamed_response(resp, args, kwargs, ctx): """ Creates a traced stream with callbacks that route SDK-owned stream consumers through the tracing proxy. @@ -57,14 +72,14 @@ def add_async_text_stream(stream): if _is_stream(resp) or _is_stream_manager(resp): traced_stream = make_traced_stream( resp, - AnthropicStreamHandler(integration, span_from_context(ctx), args, kwargs, ctx=ctx), + AnthropicStreamHandler(None, span_from_context(ctx), args, kwargs, ctx=ctx), on_stream_created=add_text_stream, ) return traced_stream elif _is_async_stream(resp) or _is_async_stream_manager(resp): traced_stream = make_traced_stream( resp, - AnthropicAsyncStreamHandler(integration, span_from_context(ctx), args, kwargs, ctx=ctx), + AnthropicAsyncStreamHandler(None, span_from_context(ctx), args, kwargs, ctx=ctx), on_stream_created=add_async_text_stream, ) return traced_stream diff --git a/ddtrace/contrib/internal/anthropic/patch.py b/ddtrace/contrib/internal/anthropic/patch.py index 06c02c53158..e5312f83153 100644 --- a/ddtrace/contrib/internal/anthropic/patch.py +++ b/ddtrace/contrib/internal/anthropic/patch.py @@ -15,7 +15,6 @@ from ddtrace.internal._exceptions import DDBlockException from ddtrace.internal.logger import get_logger from ddtrace.internal.utils.version import parse_version -from ddtrace.llmobs._integrations import AnthropicIntegration log = get_logger(__name__) @@ -34,22 +33,29 @@ def _supported_versions() -> dict[str, str]: config._add("anthropic", {}) +# LLMObs subscribes to LlmEvents for this component; see ddtrace/llmobs/_contrib/anthropic. +COMPONENT = "anthropic" -def traced_chat_model_generate(func: Callable[..., Any], instance: Any, args: Any, kwargs: Any) -> Any: - integration: AnthropicIntegration = anthropic._datadog_integration - event = LlmRequestEvent( - component="anthropic", + +def _request_event(func: Callable[..., Any], instance: Any, kwargs: Any) -> LlmRequestEvent: + model = kwargs.get("model", "") + return LlmRequestEvent( + component=COMPONENT, integration_config=config.anthropic, - service=int_service(None, integration.integration_config), + service=int_service(None, config.anthropic), resource=f"{instance.__class__.__name__}.{func.__name__}", provider="anthropic", - model=kwargs.get("model", ""), - llmobs_integration=integration, + model=model, + tags={"anthropic.request.model": model}, submit_to_llmobs=True, request_kwargs=kwargs, instance=instance, ) + +def traced_chat_model_generate(func: Callable[..., Any], instance: Any, args: Any, kwargs: Any) -> Any: + event = _request_event(func, instance, kwargs) + # For streaming, dispatch_end_event=False defers the ended event # until the stream handler calls ctx.dispatch_ended_event() in finalize_stream(). # For errors, we must manually dispatch so the span finishes with error info. @@ -61,7 +67,7 @@ def traced_chat_model_generate(func: Callable[..., Any], instance: Any, args: An ctx.dispatch_ended_event(*sys.exc_info()) raise if is_streaming_operation(resp): - return handle_streamed_response(integration, resp, args, kwargs, ctx) + return handle_streamed_response(resp, args, kwargs, ctx) # Attach the response to the event before the after-hook so that an AI # Guard block raised by that hook still records the model output in # LLMObs, even though the block errors the span (APPSEC-68147). @@ -76,19 +82,7 @@ def traced_chat_model_generate(func: Callable[..., Any], instance: Any, args: An async def traced_async_chat_model_generate(func: Callable[..., Any], instance: Any, args: Any, kwargs: Any) -> Any: - integration: AnthropicIntegration = anthropic._datadog_integration - event = LlmRequestEvent( - component="anthropic", - integration_config=config.anthropic, - service=int_service(None, integration.integration_config), - resource=f"{instance.__class__.__name__}.{func.__name__}", - provider="anthropic", - model=kwargs.get("model", ""), - llmobs_integration=integration, - submit_to_llmobs=True, - request_kwargs=kwargs, - instance=instance, - ) + event = _request_event(func, instance, kwargs) with core.context_with_event(event, dispatch_end_event=False) as ctx: try: @@ -98,7 +92,7 @@ async def traced_async_chat_model_generate(func: Callable[..., Any], instance: A ctx.dispatch_ended_event(*sys.exc_info()) raise if is_streaming_operation(resp): - return handle_streamed_response(integration, resp, args, kwargs, ctx) + return handle_streamed_response(resp, args, kwargs, ctx) # Attach the response to the event before the after-hook so that an AI # Guard block raised by that hook still records the model output in # LLMObs, even though the block errors the span (APPSEC-68147). @@ -118,9 +112,6 @@ def patch() -> None: anthropic._datadog_patch = True - integration = AnthropicIntegration(integration_config=config.anthropic) - anthropic._datadog_integration = integration - # AI Guard mirrors this wrap-target list in # ddtrace/appsec/_ai_guard/_listener.py::_install_anthropic_wrappers to # install its outermost streaming buffer. If you add/rename a target or @@ -166,5 +157,3 @@ def unpatch() -> None: unwrap(anthropic.resources.beta.messages.messages.Messages, "stream") unwrap(anthropic.resources.beta.messages.messages.AsyncMessages, "create") unwrap(anthropic.resources.beta.messages.messages.AsyncMessages, "stream") - - delattr(anthropic, "_datadog_integration") diff --git a/ddtrace/llmobs/_contrib/__init__.py b/ddtrace/llmobs/_contrib/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/ddtrace/llmobs/_contrib/anthropic/__init__.py b/ddtrace/llmobs/_contrib/anthropic/__init__.py new file mode 100644 index 00000000000..234e266184a --- /dev/null +++ b/ddtrace/llmobs/_contrib/anthropic/__init__.py @@ -0,0 +1,15 @@ +from ddtrace.llmobs._contrib.anthropic.subscribers import LLMObsAnthropicSpanFinishingSubscriber +from ddtrace.llmobs._contrib.anthropic.subscribers import LLMObsAnthropicSpanStartedSubscriber +from ddtrace.llmobs._contrib.anthropic.subscribers import LLMObsAnthropicSpanStartingSubscriber + + +def listen() -> None: + LLMObsAnthropicSpanStartingSubscriber.register() + LLMObsAnthropicSpanStartedSubscriber.register() + LLMObsAnthropicSpanFinishingSubscriber.register() + + +def unlisten() -> None: + LLMObsAnthropicSpanStartingSubscriber.unregister() + LLMObsAnthropicSpanStartedSubscriber.unregister() + LLMObsAnthropicSpanFinishingSubscriber.unregister() diff --git a/ddtrace/llmobs/_contrib/anthropic/subscribers.py b/ddtrace/llmobs/_contrib/anthropic/subscribers.py new file mode 100644 index 00000000000..cdefc6f2a10 --- /dev/null +++ b/ddtrace/llmobs/_contrib/anthropic/subscribers.py @@ -0,0 +1,91 @@ +from typing import ClassVar +from typing import Optional + +from ddtrace import config +from ddtrace.contrib._events.llm import LlmEvents +from ddtrace.contrib._events.llm import LlmRequestEvent +from ddtrace.ext import SpanTypes +from ddtrace.internal import core +from ddtrace.internal.core.subscriber import Subscriber +from ddtrace.internal.span_bus import span_from_context +from ddtrace.llmobs._constants import PROXY_REQUEST +from ddtrace.llmobs._integrations.anthropic import AnthropicIntegration + + +# Must match ddtrace.contrib.internal.anthropic.patch.COMPONENT. Duplicated so this module +# does not import the contrib, which imports the anthropic library. +ANTHROPIC_COMPONENT = "anthropic" + + +class LLMObsAnthropicSubscriber(Subscriber): + """Base for LLMObs subscribers to the Anthropic contrib's LlmEvents. + + The events are shared by every LLM integration, so handlers ignore contexts + whose event belongs to another component. + """ + + auto_register = False + _integration: ClassVar[Optional[AnthropicIntegration]] = None + + @classmethod + def integration(cls) -> AnthropicIntegration: + # config.anthropic is only defined once the contrib is imported, so build lazily. + if LLMObsAnthropicSubscriber._integration is None: + LLMObsAnthropicSubscriber._integration = AnthropicIntegration(integration_config=config.anthropic) + return LLMObsAnthropicSubscriber._integration + + @staticmethod + def _is_anthropic(ctx: core.ExecutionContext[LlmRequestEvent]) -> bool: + return ctx.event.component == ANTHROPIC_COMPONENT + + +class LLMObsAnthropicSpanStartingSubscriber(LLMObsAnthropicSubscriber): + event_names = (LlmEvents.SPAN_STARTING.value,) + + @classmethod + def on_event(cls, event_instance: core.ExecutionContext[LlmRequestEvent]) -> None: + # LlmTracingSubscriber dispatches the context rather than the event. + ctx = event_instance + if not cls._is_anthropic(ctx): + return + event = ctx.event + if event.submit_to_llmobs and cls.integration().llmobs_enabled: + event.span_type = SpanTypes.LLM + + +class LLMObsAnthropicSpanStartedSubscriber(LLMObsAnthropicSubscriber): + event_names = (LlmEvents.SPAN_STARTED.value,) + + @classmethod + def on_event(cls, event_instance: core.ExecutionContext[LlmRequestEvent]) -> None: + # LlmTracingSubscriber dispatches the context rather than the event. + ctx = event_instance + if not cls._is_anthropic(ctx): + return + event = ctx.event + span = span_from_context(ctx) + integration = cls.integration() + integration._set_base_span_tags(span, instance=event.instance) + if integration._is_instrumented_proxy_url(integration._get_base_url(instance=event.instance)): # type: ignore[arg-type] + span._set_ctx_item(PROXY_REQUEST, True) + integration._annotate_integration_tag(span) + integration._stamp_llmobs_span_kind_at_start(span, event.operation, operation=event.operation) + + +class LLMObsAnthropicSpanFinishingSubscriber(LLMObsAnthropicSubscriber): + event_names = (LlmEvents.SPAN_FINISHING.value,) + + @classmethod + def on_event(cls, event_instance: core.ExecutionContext[LlmRequestEvent]) -> None: + # LlmTracingSubscriber dispatches the context rather than the event. + ctx = event_instance + if not cls._is_anthropic(ctx): + return + event = ctx.event + cls.integration().llmobs_set_tags( + span_from_context(ctx), + args=[], + kwargs=event.request_kwargs, + response=event.response, + operation=event.operation, + ) diff --git a/ddtrace/llmobs/_integrations/anthropic.py b/ddtrace/llmobs/_integrations/anthropic.py index a4a1ecd0e02..ad7b0e7464f 100644 --- a/ddtrace/llmobs/_integrations/anthropic.py +++ b/ddtrace/llmobs/_integrations/anthropic.py @@ -65,21 +65,13 @@ def _extract_anthropic_image_source(block: Any) -> Optional[tuple[Union[bytes, s class AnthropicIntegration(BaseLLMIntegration): _integration_name = "anthropic" - def _set_base_span_tags( - self, - span: Span, - model: Optional[str] = None, - api_key: Optional[str] = None, - **kwargs: dict[str, Any], - ) -> None: - """Set base level tags that should be present on all Anthropic spans (if they are not None).""" + def _set_base_span_tags(self, span: Span, **kwargs: Any) -> None: + """Record the request base_url on the span; the contrib sets the APM model tag itself.""" # Store base_url per-span rather than on the singleton integration so a streaming # span that finalizes after concurrent requests still resolves the right provider. base_url = self._get_base_url(**kwargs) if base_url is not None: span._set_ctx_item(REQUEST_BASE_URL, base_url) - if model is not None: - span._set_attribute(MODEL, model) def _llmobs_set_tags( self, @@ -120,7 +112,7 @@ def _llmobs_set_tags( _annotate_llmobs_span_data( span, kind=span_kind, - model_name=span.get_tag("anthropic.request.model") or "", + model_name=span.get_tag(MODEL) or "", model_provider=self._get_model_provider(span), input_messages=input_messages, metadata=parameters, @@ -133,7 +125,7 @@ def _set_apm_shadow_tags(self, span, args, kwargs, response=None, operation=""): span_kind = "workflow" if span._get_ctx_item(PROXY_REQUEST) else "llm" usage = _get_attr(response, "usage", {}) metrics = self._extract_usage(span, usage) if span_kind != "workflow" else {} - model_name = span.get_tag("anthropic.request.model") + model_name = span.get_tag(MODEL) self._apply_shadow_metrics( span, metrics, diff --git a/ddtrace/llmobs/_llmobs.py b/ddtrace/llmobs/_llmobs.py index e5cf52517ae..0a9ce7b520d 100644 --- a/ddtrace/llmobs/_llmobs.py +++ b/ddtrace/llmobs/_llmobs.py @@ -137,6 +137,7 @@ from ddtrace.llmobs._integration_api import register_llmobs_service from ddtrace.llmobs._integrations.agent_manifest import build_manual_agent_manifest from ddtrace.llmobs._processor import LLMObsProcessor +from ddtrace.llmobs._product import listen_integrations from ddtrace.llmobs._prompt_optimization import PromptOptimization from ddtrace.llmobs._prompt_optimization import validate_dataset from ddtrace.llmobs._prompt_optimization import validate_dataset_split @@ -1098,6 +1099,10 @@ def enable( ) cls._instance.start() + # Covers setups that patch LLM integrations without ddtrace-run, where the + # product's post_preload never runs. + listen_integrations() + # Register hooks for span events core.on("trace.span_start", cls._instance._on_span_start) core.on("trace.span_finish", cls._instance._on_span_finish) diff --git a/ddtrace/llmobs/_product.py b/ddtrace/llmobs/_product.py index 63dec43c682..06bff8d48bc 100644 --- a/ddtrace/llmobs/_product.py +++ b/ddtrace/llmobs/_product.py @@ -39,9 +39,40 @@ def stop(join: bool = False) -> None: pass +def _listen_anthropic() -> None: + from ddtrace.llmobs._contrib.anthropic import listen + + listen() + + +def _unlisten_anthropic() -> None: + from ddtrace.llmobs._contrib.anthropic import unlisten + + unlisten() + + +def listen_integrations() -> None: + """Attach LLMObs subscribers to LLM integrations as they get patched. + + The subscribers stay registered while LLMObs is disabled because they also set + the APM shadow tags. Loading them lazily on the integration's patch event keeps + the LLMObs import chain out of applications that never patch an LLM library. + """ + import sys + + from ddtrace.internal import core + + core.on("anthropic.patch", _listen_anthropic, "llmobs.anthropic") + core.on("anthropic.unpatch", _unlisten_anthropic, "llmobs.anthropic") + if getattr(sys.modules.get("anthropic"), "_datadog_patch", False): + _listen_anthropic() + + def post_preload() -> None: """Track LLM integrations detected in the environment.""" from ddtrace import config + + listen_integrations() from ddtrace.internal.module import is_module_installed from ddtrace.internal.telemetry import telemetry_writer from ddtrace.llmobs._constants import SUPPORTED_LLMOBS_INTEGRATIONS diff --git a/tests/cassettes/datadog/datadog_api_intake_llm-obs_v2_eval-metric_post_1dc109f6.json b/tests/cassettes/datadog/datadog_api_intake_llm-obs_v2_eval-metric_post_1dc109f6.json new file mode 100644 index 00000000000..9788a6bb141 --- /dev/null +++ b/tests/cassettes/datadog/datadog_api_intake_llm-obs_v2_eval-metric_post_1dc109f6.json @@ -0,0 +1,26 @@ +{ + "request": { + "method": "POST", + "url": "https://api.datadoghq.com/api/intake/llm-obs/v2/eval-metric", + "headers": { + "Content-Type": "application/json", + "datadog-entity-id": "in-12828", + "Content-Length": "311" + }, + "body": "{\"data\": {\"attributes\": {\"metrics\": [{\"categorical_value\": \"very\", \"event_kind\": \"evaluation\", \"join_on\": {\"span\": {\"span_id\": \"12345678901\", \"trace_id\": \"98765432101\"}}, \"label\": \"toxicity\", \"metric_type\": \"categorical\", \"ml_app\": \"dummy-ml-app\", \"timestamp_ms\": 1756910127022}]}, \"type\": \"evaluation_metric\"}}" + }, + "response": { + "status": { + "code": 401, + "message": "Unauthorized" + }, + "headers": { + "content-type": "application/json", + "content-length": "27", + "date": "Fri, 02 Oct 2026 16:34:22 GMT", + "x-content-type-options": "nosniff", + "strict-transport-security": "max-age=31536000; includeSubDomains; preload" + }, + "body": "{\"errors\":[\"Unauthorized\"]}" + } +} \ No newline at end of file diff --git a/tests/cassettes/datadog/datadog_api_intake_llm-obs_v2_eval-metric_post_577483a9.json b/tests/cassettes/datadog/datadog_api_intake_llm-obs_v2_eval-metric_post_577483a9.json new file mode 100644 index 00000000000..cd5653ee086 --- /dev/null +++ b/tests/cassettes/datadog/datadog_api_intake_llm-obs_v2_eval-metric_post_577483a9.json @@ -0,0 +1,26 @@ +{ + "request": { + "method": "POST", + "url": "https://api.datadoghq.com/api/intake/llm-obs/v2/eval-metric", + "headers": { + "Content-Type": "application/json", + "datadog-entity-id": "in-12828", + "Content-Length": "297" + }, + "body": "{\"data\": {\"attributes\": {\"metrics\": [{\"event_kind\": \"evaluation\", \"join_on\": {\"span\": {\"span_id\": \"12345678902\", \"trace_id\": \"98765432102\"}}, \"label\": \"sentiment\", \"metric_type\": \"score\", \"ml_app\": \"dummy-ml-app\", \"score_value\": 0.9, \"timestamp_ms\": 1756910127022}]}, \"type\": \"evaluation_metric\"}}" + }, + "response": { + "status": { + "code": 401, + "message": "Unauthorized" + }, + "headers": { + "content-type": "application/json", + "content-length": "27", + "date": "Fri, 02 Oct 2026 16:34:22 GMT", + "x-content-type-options": "nosniff", + "strict-transport-security": "max-age=31536000; includeSubDomains; preload" + }, + "body": "{\"errors\":[\"Unauthorized\"]}" + } +} \ No newline at end of file diff --git a/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_mock-exp-id_patch_46e5df7f.json b/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_mock-exp-id_patch_46e5df7f.json new file mode 100644 index 00000000000..6636f77ffea --- /dev/null +++ b/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_mock-exp-id_patch_46e5df7f.json @@ -0,0 +1,26 @@ +{ + "request": { + "method": "PATCH", + "url": "https://api.datadoghq.com/api/unstable/llm-obs/v1/experiments/mock-exp-id", + "headers": { + "Content-Type": "application/json", + "datadog-entity-id": "in-12828", + "Content-Length": "70" + }, + "body": "{\"data\": {\"type\": \"experiments\", \"attributes\": {\"status\": \"running\"}}}" + }, + "response": { + "status": { + "code": 401, + "message": "Unauthorized" + }, + "headers": { + "content-type": "application/json", + "content-length": "27", + "date": "Fri, 02 Oct 2026 16:32:10 GMT", + "x-content-type-options": "nosniff", + "strict-transport-security": "max-age=31536000; includeSubDomains; preload" + }, + "body": "{\"errors\":[\"Unauthorized\"]}" + } +} \ No newline at end of file diff --git a/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_mock-exp-id_patch_4b861664.json b/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_mock-exp-id_patch_4b861664.json new file mode 100644 index 00000000000..f5bef4119f8 --- /dev/null +++ b/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_mock-exp-id_patch_4b861664.json @@ -0,0 +1,26 @@ +{ + "request": { + "method": "PATCH", + "url": "https://api.datadoghq.com/api/unstable/llm-obs/v1/experiments/mock-exp-id", + "headers": { + "Content-Type": "application/json", + "datadog-entity-id": "in-12828", + "Content-Length": "72" + }, + "body": "{\"data\": {\"type\": \"experiments\", \"attributes\": {\"status\": \"completed\"}}}" + }, + "response": { + "status": { + "code": 401, + "message": "Unauthorized" + }, + "headers": { + "content-type": "application/json", + "content-length": "27", + "date": "Fri, 02 Oct 2026 16:32:11 GMT", + "x-content-type-options": "nosniff", + "strict-transport-security": "max-age=31536000; includeSubDomains; preload" + }, + "body": "{\"errors\":[\"Unauthorized\"]}" + } +} \ No newline at end of file diff --git a/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-exp-id_patch_5f0d9491.json b/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-exp-id_patch_5f0d9491.json new file mode 100644 index 00000000000..dcdba9009bd --- /dev/null +++ b/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-exp-id_patch_5f0d9491.json @@ -0,0 +1,26 @@ +{ + "request": { + "method": "PATCH", + "url": "https://api.datadoghq.com/api/unstable/llm-obs/v1/experiments/new-rerun-exp-id", + "headers": { + "Content-Type": "application/json", + "datadog-entity-id": "in-12828", + "Content-Length": "72" + }, + "body": "{\"data\": {\"type\": \"experiments\", \"attributes\": {\"status\": \"completed\"}}}" + }, + "response": { + "status": { + "code": 401, + "message": "Unauthorized" + }, + "headers": { + "content-type": "application/json", + "content-length": "27", + "date": "Fri, 02 Oct 2026 16:32:11 GMT", + "x-content-type-options": "nosniff", + "strict-transport-security": "max-age=31536000; includeSubDomains; preload" + }, + "body": "{\"errors\":[\"Unauthorized\"]}" + } +} \ No newline at end of file diff --git a/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-experiment-id_patch_e2ff88d7.json b/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-experiment-id_patch_e2ff88d7.json new file mode 100644 index 00000000000..7eb8c28a493 --- /dev/null +++ b/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-experiment-id_patch_e2ff88d7.json @@ -0,0 +1,26 @@ +{ + "request": { + "method": "PATCH", + "url": "https://api.datadoghq.com/api/unstable/llm-obs/v1/experiments/new-rerun-experiment-id", + "headers": { + "Content-Type": "application/json", + "datadog-entity-id": "in-12828", + "Content-Length": "72" + }, + "body": "{\"data\": {\"type\": \"experiments\", \"attributes\": {\"status\": \"completed\"}}}" + }, + "response": { + "status": { + "code": 401, + "message": "Unauthorized" + }, + "headers": { + "content-type": "application/json", + "content-length": "27", + "date": "Fri, 02 Oct 2026 16:32:16 GMT", + "x-content-type-options": "nosniff", + "strict-transport-security": "max-age=31536000; includeSubDomains; preload" + }, + "body": "{\"errors\":[\"Unauthorized\"]}" + } +} \ No newline at end of file diff --git a/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-id_patch_23512e67.json b/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-id_patch_23512e67.json new file mode 100644 index 00000000000..0d7bbc11eb4 --- /dev/null +++ b/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-id_patch_23512e67.json @@ -0,0 +1,26 @@ +{ + "request": { + "method": "PATCH", + "url": "https://api.datadoghq.com/api/unstable/llm-obs/v1/experiments/new-rerun-id", + "headers": { + "Content-Type": "application/json", + "datadog-entity-id": "in-12828", + "Content-Length": "72" + }, + "body": "{\"data\": {\"type\": \"experiments\", \"attributes\": {\"status\": \"completed\"}}}" + }, + "response": { + "status": { + "code": 401, + "message": "Unauthorized" + }, + "headers": { + "content-type": "application/json", + "content-length": "27", + "date": "Fri, 02 Oct 2026 16:32:16 GMT", + "x-content-type-options": "nosniff", + "strict-transport-security": "max-age=31536000; includeSubDomains; preload" + }, + "body": "{\"errors\":[\"Unauthorized\"]}" + } +} \ No newline at end of file diff --git a/tests/contrib/anthropic/conftest.py b/tests/contrib/anthropic/conftest.py index 75017ae3c59..c6caedf6f9b 100644 --- a/tests/contrib/anthropic/conftest.py +++ b/tests/contrib/anthropic/conftest.py @@ -6,6 +6,7 @@ from ddtrace.contrib.internal.anthropic.patch import patch from ddtrace.contrib.internal.anthropic.patch import unpatch from ddtrace.llmobs import LLMObs +from ddtrace.llmobs._product import listen_integrations from tests.contrib.anthropic.utils import get_request_vcr from tests.utils import override_env from tests.utils import override_global_config @@ -41,6 +42,9 @@ def anthropic(): ANTHROPIC_API_KEY=os.getenv("ANTHROPIC_API_KEY", ""), ) ): + # ddtrace-run installs this hook in the LLMObs product's post_preload; these + # tests patch directly, so install it here to get the APM shadow tags. + listen_integrations() patch() import anthropic From b2f95ee280ccfc8f1cd1bde30962315b216adecb Mon Sep 17 00:00:00 2001 From: Emmett Butler Date: Fri, 2 Oct 2026 11:18:15 -0700 Subject: [PATCH 3/9] rm cassettes --- ..._llm-obs_v2_eval-metric_post_1dc109f6.json | 26 ------------------- ..._llm-obs_v2_eval-metric_post_577483a9.json | 26 ------------------- ...xperiments_mock-exp-id_patch_46e5df7f.json | 26 ------------------- ...xperiments_mock-exp-id_patch_4b861664.json | 26 ------------------- ...ments_new-rerun-exp-id_patch_5f0d9491.json | 26 ------------------- ...ew-rerun-experiment-id_patch_e2ff88d7.json | 26 ------------------- ...periments_new-rerun-id_patch_23512e67.json | 26 ------------------- 7 files changed, 182 deletions(-) delete mode 100644 tests/cassettes/datadog/datadog_api_intake_llm-obs_v2_eval-metric_post_1dc109f6.json delete mode 100644 tests/cassettes/datadog/datadog_api_intake_llm-obs_v2_eval-metric_post_577483a9.json delete mode 100644 tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_mock-exp-id_patch_46e5df7f.json delete mode 100644 tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_mock-exp-id_patch_4b861664.json delete mode 100644 tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-exp-id_patch_5f0d9491.json delete mode 100644 tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-experiment-id_patch_e2ff88d7.json delete mode 100644 tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-id_patch_23512e67.json diff --git a/tests/cassettes/datadog/datadog_api_intake_llm-obs_v2_eval-metric_post_1dc109f6.json b/tests/cassettes/datadog/datadog_api_intake_llm-obs_v2_eval-metric_post_1dc109f6.json deleted file mode 100644 index 9788a6bb141..00000000000 --- a/tests/cassettes/datadog/datadog_api_intake_llm-obs_v2_eval-metric_post_1dc109f6.json +++ /dev/null @@ -1,26 +0,0 @@ -{ - "request": { - "method": "POST", - "url": "https://api.datadoghq.com/api/intake/llm-obs/v2/eval-metric", - "headers": { - "Content-Type": "application/json", - "datadog-entity-id": "in-12828", - "Content-Length": "311" - }, - "body": "{\"data\": {\"attributes\": {\"metrics\": [{\"categorical_value\": \"very\", \"event_kind\": \"evaluation\", \"join_on\": {\"span\": {\"span_id\": \"12345678901\", \"trace_id\": \"98765432101\"}}, \"label\": \"toxicity\", \"metric_type\": \"categorical\", \"ml_app\": \"dummy-ml-app\", \"timestamp_ms\": 1756910127022}]}, \"type\": \"evaluation_metric\"}}" - }, - "response": { - "status": { - "code": 401, - "message": "Unauthorized" - }, - "headers": { - "content-type": "application/json", - "content-length": "27", - "date": "Fri, 02 Oct 2026 16:34:22 GMT", - "x-content-type-options": "nosniff", - "strict-transport-security": "max-age=31536000; includeSubDomains; preload" - }, - "body": "{\"errors\":[\"Unauthorized\"]}" - } -} \ No newline at end of file diff --git a/tests/cassettes/datadog/datadog_api_intake_llm-obs_v2_eval-metric_post_577483a9.json b/tests/cassettes/datadog/datadog_api_intake_llm-obs_v2_eval-metric_post_577483a9.json deleted file mode 100644 index cd5653ee086..00000000000 --- a/tests/cassettes/datadog/datadog_api_intake_llm-obs_v2_eval-metric_post_577483a9.json +++ /dev/null @@ -1,26 +0,0 @@ -{ - "request": { - "method": "POST", - "url": "https://api.datadoghq.com/api/intake/llm-obs/v2/eval-metric", - "headers": { - "Content-Type": "application/json", - "datadog-entity-id": "in-12828", - "Content-Length": "297" - }, - "body": "{\"data\": {\"attributes\": {\"metrics\": [{\"event_kind\": \"evaluation\", \"join_on\": {\"span\": {\"span_id\": \"12345678902\", \"trace_id\": \"98765432102\"}}, \"label\": \"sentiment\", \"metric_type\": \"score\", \"ml_app\": \"dummy-ml-app\", \"score_value\": 0.9, \"timestamp_ms\": 1756910127022}]}, \"type\": \"evaluation_metric\"}}" - }, - "response": { - "status": { - "code": 401, - "message": "Unauthorized" - }, - "headers": { - "content-type": "application/json", - "content-length": "27", - "date": "Fri, 02 Oct 2026 16:34:22 GMT", - "x-content-type-options": "nosniff", - "strict-transport-security": "max-age=31536000; includeSubDomains; preload" - }, - "body": "{\"errors\":[\"Unauthorized\"]}" - } -} \ No newline at end of file diff --git a/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_mock-exp-id_patch_46e5df7f.json b/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_mock-exp-id_patch_46e5df7f.json deleted file mode 100644 index 6636f77ffea..00000000000 --- a/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_mock-exp-id_patch_46e5df7f.json +++ /dev/null @@ -1,26 +0,0 @@ -{ - "request": { - "method": "PATCH", - "url": "https://api.datadoghq.com/api/unstable/llm-obs/v1/experiments/mock-exp-id", - "headers": { - "Content-Type": "application/json", - "datadog-entity-id": "in-12828", - "Content-Length": "70" - }, - "body": "{\"data\": {\"type\": \"experiments\", \"attributes\": {\"status\": \"running\"}}}" - }, - "response": { - "status": { - "code": 401, - "message": "Unauthorized" - }, - "headers": { - "content-type": "application/json", - "content-length": "27", - "date": "Fri, 02 Oct 2026 16:32:10 GMT", - "x-content-type-options": "nosniff", - "strict-transport-security": "max-age=31536000; includeSubDomains; preload" - }, - "body": "{\"errors\":[\"Unauthorized\"]}" - } -} \ No newline at end of file diff --git a/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_mock-exp-id_patch_4b861664.json b/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_mock-exp-id_patch_4b861664.json deleted file mode 100644 index f5bef4119f8..00000000000 --- a/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_mock-exp-id_patch_4b861664.json +++ /dev/null @@ -1,26 +0,0 @@ -{ - "request": { - "method": "PATCH", - "url": "https://api.datadoghq.com/api/unstable/llm-obs/v1/experiments/mock-exp-id", - "headers": { - "Content-Type": "application/json", - "datadog-entity-id": "in-12828", - "Content-Length": "72" - }, - "body": "{\"data\": {\"type\": \"experiments\", \"attributes\": {\"status\": \"completed\"}}}" - }, - "response": { - "status": { - "code": 401, - "message": "Unauthorized" - }, - "headers": { - "content-type": "application/json", - "content-length": "27", - "date": "Fri, 02 Oct 2026 16:32:11 GMT", - "x-content-type-options": "nosniff", - "strict-transport-security": "max-age=31536000; includeSubDomains; preload" - }, - "body": "{\"errors\":[\"Unauthorized\"]}" - } -} \ No newline at end of file diff --git a/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-exp-id_patch_5f0d9491.json b/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-exp-id_patch_5f0d9491.json deleted file mode 100644 index dcdba9009bd..00000000000 --- a/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-exp-id_patch_5f0d9491.json +++ /dev/null @@ -1,26 +0,0 @@ -{ - "request": { - "method": "PATCH", - "url": "https://api.datadoghq.com/api/unstable/llm-obs/v1/experiments/new-rerun-exp-id", - "headers": { - "Content-Type": "application/json", - "datadog-entity-id": "in-12828", - "Content-Length": "72" - }, - "body": "{\"data\": {\"type\": \"experiments\", \"attributes\": {\"status\": \"completed\"}}}" - }, - "response": { - "status": { - "code": 401, - "message": "Unauthorized" - }, - "headers": { - "content-type": "application/json", - "content-length": "27", - "date": "Fri, 02 Oct 2026 16:32:11 GMT", - "x-content-type-options": "nosniff", - "strict-transport-security": "max-age=31536000; includeSubDomains; preload" - }, - "body": "{\"errors\":[\"Unauthorized\"]}" - } -} \ No newline at end of file diff --git a/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-experiment-id_patch_e2ff88d7.json b/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-experiment-id_patch_e2ff88d7.json deleted file mode 100644 index 7eb8c28a493..00000000000 --- a/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-experiment-id_patch_e2ff88d7.json +++ /dev/null @@ -1,26 +0,0 @@ -{ - "request": { - "method": "PATCH", - "url": "https://api.datadoghq.com/api/unstable/llm-obs/v1/experiments/new-rerun-experiment-id", - "headers": { - "Content-Type": "application/json", - "datadog-entity-id": "in-12828", - "Content-Length": "72" - }, - "body": "{\"data\": {\"type\": \"experiments\", \"attributes\": {\"status\": \"completed\"}}}" - }, - "response": { - "status": { - "code": 401, - "message": "Unauthorized" - }, - "headers": { - "content-type": "application/json", - "content-length": "27", - "date": "Fri, 02 Oct 2026 16:32:16 GMT", - "x-content-type-options": "nosniff", - "strict-transport-security": "max-age=31536000; includeSubDomains; preload" - }, - "body": "{\"errors\":[\"Unauthorized\"]}" - } -} \ No newline at end of file diff --git a/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-id_patch_23512e67.json b/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-id_patch_23512e67.json deleted file mode 100644 index 0d7bbc11eb4..00000000000 --- a/tests/cassettes/datadog/datadog_api_unstable_llm-obs_v1_experiments_new-rerun-id_patch_23512e67.json +++ /dev/null @@ -1,26 +0,0 @@ -{ - "request": { - "method": "PATCH", - "url": "https://api.datadoghq.com/api/unstable/llm-obs/v1/experiments/new-rerun-id", - "headers": { - "Content-Type": "application/json", - "datadog-entity-id": "in-12828", - "Content-Length": "72" - }, - "body": "{\"data\": {\"type\": \"experiments\", \"attributes\": {\"status\": \"completed\"}}}" - }, - "response": { - "status": { - "code": 401, - "message": "Unauthorized" - }, - "headers": { - "content-type": "application/json", - "content-length": "27", - "date": "Fri, 02 Oct 2026 16:32:16 GMT", - "x-content-type-options": "nosniff", - "strict-transport-security": "max-age=31536000; includeSubDomains; preload" - }, - "body": "{\"errors\":[\"Unauthorized\"]}" - } -} \ No newline at end of file From d991e249af752c5f7566d84d35fd6421486352c8 Mon Sep 17 00:00:00 2001 From: Emmett Butler Date: Fri, 2 Oct 2026 11:33:14 -0700 Subject: [PATCH 4/9] merge whoopsie --- ddtrace/contrib/internal/anthropic/_streaming.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/ddtrace/contrib/internal/anthropic/_streaming.py b/ddtrace/contrib/internal/anthropic/_streaming.py index adfe38699b6..dddbe346821 100644 --- a/ddtrace/contrib/internal/anthropic/_streaming.py +++ b/ddtrace/contrib/internal/anthropic/_streaming.py @@ -9,8 +9,6 @@ from ddtrace.contrib.internal.stream_handler import make_traced_stream from ddtrace.internal.logger import get_logger from ddtrace.internal.span_bus import span_from_context -from ddtrace.llmobs._utils import _get_attr -from ddtrace.llmobs._utils import safe_load_json log = get_logger(__name__) From b2883189de73940fc8b41cd4a5c2eb58d76cb9ad Mon Sep 17 00:00:00 2001 From: Emmett Butler Date: Fri, 2 Oct 2026 11:41:01 -0700 Subject: [PATCH 5/9] import cycle --- .claude/skills/llmobs-integrations/SKILL.md | 2 +- .../references/implementation-guide.md | 6 ++-- ddtrace/llmobs/_contrib/__init__.py | 28 +++++++++++++++++ ddtrace/llmobs/_llmobs.py | 2 +- ddtrace/llmobs/_product.py | 30 +------------------ tests/contrib/anthropic/conftest.py | 2 +- 6 files changed, 35 insertions(+), 35 deletions(-) diff --git a/.claude/skills/llmobs-integrations/SKILL.md b/.claude/skills/llmobs-integrations/SKILL.md index c48e07a7987..62bfd4259f6 100644 --- a/.claude/skills/llmobs-integrations/SKILL.md +++ b/.claude/skills/llmobs-integrations/SKILL.md @@ -21,7 +21,7 @@ LLMObs integrations consist of two cooperating layers: 1. **Patch Layer** (`ddtrace/contrib/internal/{name}/patch.py`) -- wraps library functions. Standard request/response LLM integrations construct `LlmRequestEvent` and use `core.context_with_event()` so the LLM tracing subscriber owns span lifecycle. The patch layer must not import from `ddtrace.llmobs`. 2. **Integration Layer** (`ddtrace/llmobs/_integrations/{name}.py`) -- extends `BaseLLMIntegration`, implements `_set_base_span_tags()` and `_llmobs_set_tags()` to extract and set provider-specific messages, tools, metadata, and token metrics. -3. **LLMObs Subscribers** (`ddtrace/llmobs/_contrib/{name}/`) -- subscribe to the `LlmEvents` span lifecycle events (`SPAN_STARTING`, `SPAN_STARTED`, `SPAN_FINISHING`) that `LlmTracingSubscriber` dispatches, filter on `ctx.event.component`, and call the integration layer. `ddtrace/llmobs/_product.py:listen_integrations()` registers them when the library is patched. +3. **LLMObs Subscribers** (`ddtrace/llmobs/_contrib/{name}/`) -- subscribe to the `LlmEvents` span lifecycle events (`SPAN_STARTING`, `SPAN_STARTED`, `SPAN_FINISHING`) that `LlmTracingSubscriber` dispatches, filter on `ctx.event.component`, and call the integration layer. `listen_integrations()` in `ddtrace/llmobs/_contrib/__init__.py` registers them when the library is patched. The layers must work together. The patch layer identifies the operation and passes request/response data through the event; the subscribers connect the event to the integration layer, which controls what data is extracted. diff --git a/.claude/skills/llmobs-integrations/references/implementation-guide.md b/.claude/skills/llmobs-integrations/references/implementation-guide.md index 225002e4b2f..b3cc5e06a3e 100644 --- a/.claude/skills/llmobs-integrations/references/implementation-guide.md +++ b/.claude/skills/llmobs-integrations/references/implementation-guide.md @@ -17,7 +17,7 @@ This separation keeps APM patching decoupled from LLMObs data extraction. An LLM integration is an APM integration with an extra layer. You do everything in the apm-integrations guide, but: - **Step 1 (patch module)**: Use `LlmRequestEvent` with `core.context_with_event()` for standard request/response LLM integrations -- **Step 1b (LLMObs subscribers)**: Add `ddtrace/llmobs/_contrib/{name}/` subscribers and register them from `listen_integrations()` in `ddtrace/llmobs/_product.py` +- **Step 1b (LLMObs subscribers)**: Add `ddtrace/llmobs/_contrib/{name}/` subscribers and register them from `listen_integrations()` in `ddtrace/llmobs/_contrib/__init__.py` - **Step 3 (LLMObs integration)**: Create the `BaseLLMIntegration` subclass that handles provider-specific message, tool, and token extraction (this guide) - **Step 4 (test environment)**: Use `tests/llmobs/suitespec.yml`; add `vcrpy` only when the suite uses vcrpy cassettes and follow nearby version pins - **Step 5 (tests)**: Add `test_{name}_llmobs.py` in addition to the APM `test_{name}.py`, using the right transport pattern for the integration and `assert_llmobs_span_data(_get_llmobs_data_metastruct(span), ...)` @@ -111,7 +111,7 @@ Read `ddtrace/llmobs/_contrib/anthropic/` for the pattern. `LlmTracingSubscriber - `LlmEvents.SPAN_STARTED` runs after the span is created. Use it for `_set_base_span_tags()`, `_annotate_integration_tag()`, and `_stamp_llmobs_span_kind_at_start()`. - `LlmEvents.SPAN_FINISHING` runs before the span is finished. Call `integration.llmobs_set_tags()` here. -Subscribers set `auto_register = False`. Register them from `listen_integrations()` in `ddtrace/llmobs/_product.py` on the `{name}.patch` core event, and unregister them on `{name}.unpatch`. `listen_integrations()` runs from both the product's `post_preload` and `LLMObs.enable()`. Test fixtures that call `patch()` directly should call `listen_integrations()` first. +Subscribers set `auto_register = False`. Register them from `listen_integrations()` in `ddtrace/llmobs/_contrib/__init__.py` on the `{name}.patch` core event, and unregister them on `{name}.unpatch`. `listen_integrations()` runs from both the product's `post_preload` and `LLMObs.enable()`. Test fixtures that call `patch()` directly should call `listen_integrations()` first. - The async variant is identical but uses `async def` / `await` Some older or specialized integrations still call `integration.trace()` and `integration.llmobs_set_tags()` directly. Use that pattern when modifying an existing integration that already does so, when the closest current reference uses it (for example Google GenAI), or when the behavior requires direct child spans (for example OpenAI MCP tool spans or agent/tool child spans). @@ -226,7 +226,7 @@ In addition to the full checklist in the apm-integrations [Implementation Guide] - [ ] `ddtrace/llmobs/_integrations/{name}.py` — `BaseLLMIntegration` subclass - [ ] `ddtrace/llmobs/_integrations/__init__.py` — import + `__all__` entry - [ ] `ddtrace/contrib/internal/{name}/patch.py` — uses `LlmRequestEvent` + `core.context_with_event()` for standard LLM request spans, with no `ddtrace.llmobs` imports (see anthropic for pattern) -- [ ] `ddtrace/llmobs/_contrib/{name}/` — `LlmEvents` subscribers, registered from `listen_integrations()` in `ddtrace/llmobs/_product.py` +- [ ] `ddtrace/llmobs/_contrib/{name}/` — `LlmEvents` subscribers, registered from `listen_integrations()` in `ddtrace/llmobs/_contrib/__init__.py` - [ ] `tests/llmobs/suitespec.yml` — LLMObs test suite entry - [ ] Test dependencies match the suite style; include `vcrpy` only when cassette replay is used - [ ] `docs/index.rst` — add integration to the docs index diff --git a/ddtrace/llmobs/_contrib/__init__.py b/ddtrace/llmobs/_contrib/__init__.py index e69de29bb2d..da2fe31c344 100644 --- a/ddtrace/llmobs/_contrib/__init__.py +++ b/ddtrace/llmobs/_contrib/__init__.py @@ -0,0 +1,28 @@ +import sys + +from ddtrace.internal import core + + +def _listen_anthropic() -> None: + from ddtrace.llmobs._contrib.anthropic import listen + + listen() + + +def _unlisten_anthropic() -> None: + from ddtrace.llmobs._contrib.anthropic import unlisten + + unlisten() + + +def listen_integrations() -> None: + """Attach LLMObs subscribers to LLM integrations as they get patched. + + The subscribers stay registered while LLMObs is disabled because they also set + the APM shadow tags. Loading them lazily on the integration's patch event keeps + the LLMObs import chain out of applications that never patch an LLM library. + """ + core.on("anthropic.patch", _listen_anthropic, "llmobs.anthropic") + core.on("anthropic.unpatch", _unlisten_anthropic, "llmobs.anthropic") + if getattr(sys.modules.get("anthropic"), "_datadog_patch", False): + _listen_anthropic() diff --git a/ddtrace/llmobs/_llmobs.py b/ddtrace/llmobs/_llmobs.py index 0a9ce7b520d..8de35fa3e84 100644 --- a/ddtrace/llmobs/_llmobs.py +++ b/ddtrace/llmobs/_llmobs.py @@ -99,6 +99,7 @@ from ddtrace.llmobs._constants import VERTEXAI_APM_SPAN_NAME from ddtrace.llmobs._constants import LLMObsExportMode from ddtrace.llmobs._context import LLMObsContextProvider +from ddtrace.llmobs._contrib import listen_integrations from ddtrace.llmobs._eval_metric import _build_evaluation_metric_event from ddtrace.llmobs._eval_metric import _build_feedback_metric_event from ddtrace.llmobs._eval_metric import _SubmissionTelemetryContext @@ -137,7 +138,6 @@ from ddtrace.llmobs._integration_api import register_llmobs_service from ddtrace.llmobs._integrations.agent_manifest import build_manual_agent_manifest from ddtrace.llmobs._processor import LLMObsProcessor -from ddtrace.llmobs._product import listen_integrations from ddtrace.llmobs._prompt_optimization import PromptOptimization from ddtrace.llmobs._prompt_optimization import validate_dataset from ddtrace.llmobs._prompt_optimization import validate_dataset_split diff --git a/ddtrace/llmobs/_product.py b/ddtrace/llmobs/_product.py index 06bff8d48bc..56883935150 100644 --- a/ddtrace/llmobs/_product.py +++ b/ddtrace/llmobs/_product.py @@ -39,38 +39,10 @@ def stop(join: bool = False) -> None: pass -def _listen_anthropic() -> None: - from ddtrace.llmobs._contrib.anthropic import listen - - listen() - - -def _unlisten_anthropic() -> None: - from ddtrace.llmobs._contrib.anthropic import unlisten - - unlisten() - - -def listen_integrations() -> None: - """Attach LLMObs subscribers to LLM integrations as they get patched. - - The subscribers stay registered while LLMObs is disabled because they also set - the APM shadow tags. Loading them lazily on the integration's patch event keeps - the LLMObs import chain out of applications that never patch an LLM library. - """ - import sys - - from ddtrace.internal import core - - core.on("anthropic.patch", _listen_anthropic, "llmobs.anthropic") - core.on("anthropic.unpatch", _unlisten_anthropic, "llmobs.anthropic") - if getattr(sys.modules.get("anthropic"), "_datadog_patch", False): - _listen_anthropic() - - def post_preload() -> None: """Track LLM integrations detected in the environment.""" from ddtrace import config + from ddtrace.llmobs._contrib import listen_integrations listen_integrations() from ddtrace.internal.module import is_module_installed diff --git a/tests/contrib/anthropic/conftest.py b/tests/contrib/anthropic/conftest.py index c6caedf6f9b..365603d3839 100644 --- a/tests/contrib/anthropic/conftest.py +++ b/tests/contrib/anthropic/conftest.py @@ -6,7 +6,7 @@ from ddtrace.contrib.internal.anthropic.patch import patch from ddtrace.contrib.internal.anthropic.patch import unpatch from ddtrace.llmobs import LLMObs -from ddtrace.llmobs._product import listen_integrations +from ddtrace.llmobs._contrib import listen_integrations from tests.contrib.anthropic.utils import get_request_vcr from tests.utils import override_env from tests.utils import override_global_config From 84165037546a0eb0d213fe86c2d771c2e7927c2d Mon Sep 17 00:00:00 2001 From: Emmett Butler Date: Tue, 6 Oct 2026 11:19:42 -0700 Subject: [PATCH 6/9] fix test isolation problem exposed by refactor --- tests/llmobs/conftest.py | 40 ++++++++++++++++------------- tests/llmobs/test_llmobs_service.py | 34 ++++++++++++++---------- 2 files changed, 43 insertions(+), 31 deletions(-) diff --git a/tests/llmobs/conftest.py b/tests/llmobs/conftest.py index 3619aecb4a7..c83834776fd 100644 --- a/tests/llmobs/conftest.py +++ b/tests/llmobs/conftest.py @@ -280,27 +280,31 @@ def llmobs( with override_global_config(global_config): # Pin agentless_enabled=False: the default flips to agentless when the Agent is # unreachable, swapping the ``tracer`` fixture's DummyWriter and breaking ``test_spans``. - llmobs_service.enable(_tracer=tracer, agentless_enabled=False, **llmobs_enable_opts) - llmobs_service._instance._llmobs_span_writer = llmobs_span_writer - llmobs_service._instance._llmobs_span_writer.start() - # The cassette proxy stands in for intake, so keep this client in direct mode. Without an - # app key it would otherwise pick the agent proxy and prefix every path with /evp_proxy/v2, - # which no recording matches. - dne_client = llmobs_service._instance._dne_client - dne_client._agentless = True - dne_client._endpoint = dne_client.ENDPOINT - dne_client._intake = llmobs_api_proxy_url - tracer._span_aggregator.llmobs_processor = LLMObsProcessor( - llmobs_span_writer, - tracer, - keep_meta_struct=True, - sampling_resolver=llmobs_service._instance._sampling_resolver, - ) try: + # Setup lives inside the try so a failure here cannot leak an enabled instance holding + # this test's mocked writers into later tests on the same worker. + llmobs_service.enable(_tracer=tracer, agentless_enabled=False, **llmobs_enable_opts) + llmobs_service._instance._llmobs_span_writer = llmobs_span_writer + llmobs_service._instance._llmobs_span_writer.start() + # The cassette proxy stands in for intake, so keep this client in direct mode. Without an + # app key it would otherwise pick the agent proxy and prefix every path with /evp_proxy/v2, + # which no recording matches. + dne_client = llmobs_service._instance._dne_client + dne_client._agentless = True + dne_client._endpoint = dne_client.ENDPOINT + dne_client._intake = llmobs_api_proxy_url + tracer._span_aggregator.llmobs_processor = LLMObsProcessor( + llmobs_span_writer, + tracer, + keep_meta_struct=True, + sampling_resolver=llmobs_service._instance._sampling_resolver, + ) yield llmobs_service finally: - tracer.shutdown() - llmobs_service.disable() + try: + tracer.shutdown() + finally: + llmobs_service.disable() @pytest.fixture diff --git a/tests/llmobs/test_llmobs_service.py b/tests/llmobs/test_llmobs_service.py index b229a22d753..9cd66e0d5c5 100644 --- a/tests/llmobs/test_llmobs_service.py +++ b/tests/llmobs/test_llmobs_service.py @@ -2584,25 +2584,33 @@ def test_service_enable_starts_evaluator_runner_when_evaluators_exist(tracer): pytest.importorskip("ragas") with override_global_config(dict(_dd_api_key="", _llmobs_ml_app="")): with override_env(dict(DD_LLMOBS_EVALUATORS="ragas_faithfulness")): - llmobs_service.enable(_tracer=tracer) - llmobs_instance = llmobs_service._instance - assert llmobs_instance is not None - assert llmobs_service.enabled - assert llmobs_service._instance._llmobs_eval_metric_writer.status.value == "running" - assert llmobs_service._instance._evaluator_runner.status.value == "running" + # Guard against leaked enabled=True from a prior failed test llmobs_service.disable() + llmobs_service.enable(_tracer=tracer) + try: + llmobs_instance = llmobs_service._instance + assert llmobs_instance is not None + assert llmobs_service.enabled + assert llmobs_service._instance._llmobs_eval_metric_writer.status.value == "running" + assert llmobs_service._instance._evaluator_runner.status.value == "running" + finally: + llmobs_service.disable() def test_service_enable_does_not_start_evaluator_runner(tracer): with override_global_config(dict(_dd_api_key="", _llmobs_ml_app="")): - llmobs_service.enable(_tracer=tracer) - llmobs_instance = llmobs_service._instance - assert llmobs_instance is not None - assert llmobs_service.enabled - assert llmobs_service._instance._llmobs_eval_metric_writer.status.value == "running" - assert llmobs_service._instance._llmobs_span_writer.status.value == "running" - assert llmobs_service._instance._evaluator_runner.status.value == "stopped" + # Guard against leaked enabled=True from a prior failed test llmobs_service.disable() + llmobs_service.enable(_tracer=tracer) + try: + llmobs_instance = llmobs_service._instance + assert llmobs_instance is not None + assert llmobs_service.enabled + assert llmobs_service._instance._llmobs_eval_metric_writer.status.value == "running" + assert llmobs_service._instance._llmobs_span_writer.status.value == "running" + assert llmobs_service._instance._evaluator_runner.status.value == "stopped" + finally: + llmobs_service.disable() def test_export_span_when_llmobs_is_disabled_returns_none(llmobs): From 1397a456709d6f080d9f5c9f07890dfdaa5982e5 Mon Sep 17 00:00:00 2001 From: Emmett Butler Date: Fri, 9 Oct 2026 06:11:40 -0700 Subject: [PATCH 7/9] typing --- ddtrace/_trace/subscribers/llm.py | 7 ++++++- ddtrace/llmobs/_contrib/anthropic/subscribers.py | 3 +++ 2 files changed, 9 insertions(+), 1 deletion(-) diff --git a/ddtrace/_trace/subscribers/llm.py b/ddtrace/_trace/subscribers/llm.py index 9004d21b847..be12f95445a 100644 --- a/ddtrace/_trace/subscribers/llm.py +++ b/ddtrace/_trace/subscribers/llm.py @@ -48,10 +48,15 @@ def on_started(cls, ctx: core.ExecutionContext["LlmRequestEvent"]) -> None: span._remove_attribute(COMPONENT) span._remove_attribute(SPAN_KIND) - if event.submit_to_llmobs: + if event.llmobs_integration is None: core.dispatch(LlmEvents.SPAN_STARTED.value, (ctx,)) return + if event.submit_to_llmobs: + span._set_attribute( + _LLMOBS_APM_SHADOW_ENABLED_METRIC_KEY, 1 if event.llmobs_integration.llmobs_enabled else 0 + ) + event.llmobs_integration._set_base_span_tags( span, model=event.model, diff --git a/ddtrace/llmobs/_contrib/anthropic/subscribers.py b/ddtrace/llmobs/_contrib/anthropic/subscribers.py index cdefc6f2a10..2f6c6209e8f 100644 --- a/ddtrace/llmobs/_contrib/anthropic/subscribers.py +++ b/ddtrace/llmobs/_contrib/anthropic/subscribers.py @@ -8,6 +8,7 @@ from ddtrace.internal import core from ddtrace.internal.core.subscriber import Subscriber from ddtrace.internal.span_bus import span_from_context +from ddtrace.llmobs._constants import LLMOBS_APM_SHADOW_ENABLED_METRIC_KEY from ddtrace.llmobs._constants import PROXY_REQUEST from ddtrace.llmobs._integrations.anthropic import AnthropicIntegration @@ -65,6 +66,8 @@ def on_event(cls, event_instance: core.ExecutionContext[LlmRequestEvent]) -> Non event = ctx.event span = span_from_context(ctx) integration = cls.integration() + if event.submit_to_llmobs: + span._set_attribute(LLMOBS_APM_SHADOW_ENABLED_METRIC_KEY, 1 if integration.llmobs_enabled else 0) integration._set_base_span_tags(span, instance=event.instance) if integration._is_instrumented_proxy_url(integration._get_base_url(instance=event.instance)): # type: ignore[arg-type] span._set_ctx_item(PROXY_REQUEST, True) From 53989021eac923a49e35fcb21b74e77cf38c834c Mon Sep 17 00:00:00 2001 From: Emmett Butler Date: Fri, 9 Oct 2026 06:22:00 -0700 Subject: [PATCH 8/9] update rules about direct patch() calls fix in-flight requests after unpatch() --- .claude/skills/llmobs-integrations/SKILL.md | 2 +- .../llmobs-integrations/references/implementation-guide.md | 2 +- ddtrace/llmobs/_contrib/anthropic/__init__.py | 4 +++- 3 files changed, 5 insertions(+), 3 deletions(-) diff --git a/.claude/skills/llmobs-integrations/SKILL.md b/.claude/skills/llmobs-integrations/SKILL.md index 62bfd4259f6..4e38e1bb9cf 100644 --- a/.claude/skills/llmobs-integrations/SKILL.md +++ b/.claude/skills/llmobs-integrations/SKILL.md @@ -119,7 +119,7 @@ Note two already-shipped integrations predate this key: bedrock and the claude-a - **Event-based patch wrappers** should not call `span.set_exc_info()`, `span.finish()`, or `integration.llmobs_set_tags()` directly; the tracing subscriber handles that when the event ends - **Direct integration spans** must keep `integration.llmobs_set_tags()` and span lifecycle handling aligned with the closest current reference - **Integration instance**: subscriber-based integrations build it lazily in the LLMObs subscriber (see `ddtrace/llmobs/_contrib/anthropic/subscribers.py`); older integrations store it on the module as `module._datadog_integration = MyLibIntegration(integration_config=config.mylib)` -- **Subscriber registration**: subscriber-based integrations hook `{name}.patch`/`{name}.unpatch` core events in `listen_integrations()` and stay registered while LLMObs is disabled, because they also set the APM shadow tags +- **Subscriber registration**: subscriber-based integrations hook `{name}.patch`/`{name}.unpatch` core events in `listen_integrations()` and stay registered while LLMObs is disabled, because they also set the APM shadow tags. `listen_integrations()` runs only under `ddtrace-run` or `LLMObs.enable()`; manual `patch()` calls without either do not need to be supported ## Message Types diff --git a/.claude/skills/llmobs-integrations/references/implementation-guide.md b/.claude/skills/llmobs-integrations/references/implementation-guide.md index b3cc5e06a3e..dc667d570dc 100644 --- a/.claude/skills/llmobs-integrations/references/implementation-guide.md +++ b/.claude/skills/llmobs-integrations/references/implementation-guide.md @@ -111,7 +111,7 @@ Read `ddtrace/llmobs/_contrib/anthropic/` for the pattern. `LlmTracingSubscriber - `LlmEvents.SPAN_STARTED` runs after the span is created. Use it for `_set_base_span_tags()`, `_annotate_integration_tag()`, and `_stamp_llmobs_span_kind_at_start()`. - `LlmEvents.SPAN_FINISHING` runs before the span is finished. Call `integration.llmobs_set_tags()` here. -Subscribers set `auto_register = False`. Register them from `listen_integrations()` in `ddtrace/llmobs/_contrib/__init__.py` on the `{name}.patch` core event, and unregister them on `{name}.unpatch`. `listen_integrations()` runs from both the product's `post_preload` and `LLMObs.enable()`. Test fixtures that call `patch()` directly should call `listen_integrations()` first. +Subscribers set `auto_register = False`. Register them from `listen_integrations()` in `ddtrace/llmobs/_contrib/__init__.py` on the `{name}.patch` core event, and unregister them on `{name}.unpatch` -- except the `SPAN_FINISHING` subscriber, which stays registered so requests and deferred streams already in flight at unpatch time still get finalized. `listen_integrations()` runs from both the product's `post_preload` and `LLMObs.enable()`. Applications that call `patch()` manually without `ddtrace-run` or `LLMObs.enable()` are not a supported setup for subscriber-based integrations: their subscribers never register, so those spans lack the LLMObs shadow tags. Do not add registration hooks to `ddtrace.patch()` or `ddtrace/_monkey.py` to cover it. Test fixtures that call the contrib's `patch()` directly should call `listen_integrations()` first. - The async variant is identical but uses `async def` / `await` Some older or specialized integrations still call `integration.trace()` and `integration.llmobs_set_tags()` directly. Use that pattern when modifying an existing integration that already does so, when the closest current reference uses it (for example Google GenAI), or when the behavior requires direct child spans (for example OpenAI MCP tool spans or agent/tool child spans). diff --git a/ddtrace/llmobs/_contrib/anthropic/__init__.py b/ddtrace/llmobs/_contrib/anthropic/__init__.py index 234e266184a..6a0cf4d059a 100644 --- a/ddtrace/llmobs/_contrib/anthropic/__init__.py +++ b/ddtrace/llmobs/_contrib/anthropic/__init__.py @@ -12,4 +12,6 @@ def listen() -> None: def unlisten() -> None: LLMObsAnthropicSpanStartingSubscriber.unregister() LLMObsAnthropicSpanStartedSubscriber.unregister() - LLMObsAnthropicSpanFinishingSubscriber.unregister() + # The finishing subscriber stays registered: requests and deferred streams that started + # before unpatch() still need their output, token metrics, and shadow tags when they end. + # It only acts on Anthropic contexts, and an unpatched client creates no new ones. From 8a556b6dccb1fcf0aadb6133a25cc6a433d75fb4 Mon Sep 17 00:00:00 2001 From: Emmett Butler Date: Fri, 9 Oct 2026 08:07:21 -0700 Subject: [PATCH 9/9] wrong merge --- ddtrace/contrib/internal/anthropic/_streaming.py | 3 +++ 1 file changed, 3 insertions(+) diff --git a/ddtrace/contrib/internal/anthropic/_streaming.py b/ddtrace/contrib/internal/anthropic/_streaming.py index 271eb4aedec..96393b19ce3 100644 --- a/ddtrace/contrib/internal/anthropic/_streaming.py +++ b/ddtrace/contrib/internal/anthropic/_streaming.py @@ -6,6 +6,9 @@ from ddtrace.internal.logger import get_logger from ddtrace.internal.span_bus import span_from_context +from ddtrace.internal.utils.streaming import AsyncStreamHandler +from ddtrace.internal.utils.streaming import StreamHandler +from ddtrace.internal.utils.streaming import make_traced_stream log = get_logger(__name__)