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 df8664b277b..c85afa236cf 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. `listen_integrations()` in `ddtrace/llmobs/_contrib/__init__.py` 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. `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/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 54ff86284dd..c1456e59093 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/_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), ...)` @@ -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/_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). @@ -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/_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/_trace/subscribers/llm.py b/ddtrace/_trace/subscribers/llm.py index 998588eae62..be12f95445a 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 @@ -24,11 +25,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 @@ -41,6 +48,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 + if event.submit_to_llmobs: span._set_attribute( _LLMOBS_APM_SHADOW_ENABLED_METRIC_KEY, 1 if event.llmobs_integration.llmobs_enabled else 0 @@ -73,6 +84,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 f2e430ed4a1..96393b19ce3 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 @@ -8,13 +9,25 @@ from ddtrace.internal.utils.streaming import AsyncStreamHandler from ddtrace.internal.utils.streaming import StreamHandler from ddtrace.internal.utils.streaming import make_traced_stream -from ddtrace.llmobs._utils import _get_attr -from ddtrace.llmobs._utils import safe_load_json 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 +50,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 +70,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..da2fe31c344 --- /dev/null +++ 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/_contrib/anthropic/__init__.py b/ddtrace/llmobs/_contrib/anthropic/__init__.py new file mode 100644 index 00000000000..6a0cf4d059a --- /dev/null +++ b/ddtrace/llmobs/_contrib/anthropic/__init__.py @@ -0,0 +1,17 @@ +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() + # 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. diff --git a/ddtrace/llmobs/_contrib/anthropic/subscribers.py b/ddtrace/llmobs/_contrib/anthropic/subscribers.py new file mode 100644 index 00000000000..2f6c6209e8f --- /dev/null +++ b/ddtrace/llmobs/_contrib/anthropic/subscribers.py @@ -0,0 +1,94 @@ +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 LLMOBS_APM_SHADOW_ENABLED_METRIC_KEY +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() + 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) + 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 be5c3a1a65f..767916052b9 100644 --- a/ddtrace/llmobs/_llmobs.py +++ b/ddtrace/llmobs/_llmobs.py @@ -103,6 +103,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 @@ -1102,6 +1103,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..56883935150 100644 --- a/ddtrace/llmobs/_product.py +++ b/ddtrace/llmobs/_product.py @@ -42,6 +42,9 @@ def stop(join: bool = False) -> None: 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 from ddtrace.internal.telemetry import telemetry_writer from ddtrace.llmobs._constants import SUPPORTED_LLMOBS_INTEGRATIONS diff --git a/tests/contrib/anthropic/conftest.py b/tests/contrib/anthropic/conftest.py index 75017ae3c59..365603d3839 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._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 @@ -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