diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/.changelog/661.added b/instrumentation/opentelemetry-instrumentation-genai-anthropic/.changelog/661.added new file mode 100644 index 000000000..9359851b0 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/.changelog/661.added @@ -0,0 +1 @@ +Add instrumentation for `beta.messages` (`create`, `stream`, and `parse`). diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/README.rst b/instrumentation/opentelemetry-instrumentation-genai-anthropic/README.rst index 978a5f7fe..9879d9ad9 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-anthropic/README.rst +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/README.rst @@ -46,6 +46,16 @@ Check out the `manual example `_ for more details. ) +Supported Operations +-------------------- + +The instrumentation supports synchronous and asynchronous calls for: + +- ``messages.create``, ``messages.stream``, and ``messages.parse`` +- ``beta.messages.create``, ``beta.messages.stream``, and ``beta.messages.parse`` +- Raw and streaming response helpers (``with_raw_response`` and ``with_streaming_response``) + + Configuration ------------- diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/__init__.py b/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/__init__.py index 0b5bd2fe2..5b3d20e33 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/__init__.py +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/__init__.py @@ -77,6 +77,37 @@ def _is_parse_supported() -> bool: return False +def _is_beta_messages_supported() -> bool: + """Check if beta Messages classes are available on the anthropic SDK.""" + try: + from anthropic.resources.beta.messages import ( # pylint: disable=import-outside-toplevel + AsyncMessages, + Messages, + ) + + return ( + hasattr(Messages, "create") + and hasattr(AsyncMessages, "create") + and hasattr(Messages, "stream") + and hasattr(AsyncMessages, "stream") + ) + except (ImportError, AttributeError): + return False + + +def _is_beta_parse_supported() -> bool: + """Check if parse() is available on the beta Messages classes.""" + try: + from anthropic.resources.beta.messages import ( # pylint: disable=import-outside-toplevel + AsyncMessages, + Messages, + ) + + return hasattr(Messages, "parse") and hasattr(AsyncMessages, "parse") + except (ImportError, AttributeError): + return False + + class AnthropicInstrumentor(BaseInstrumentor): """An instrumentor for the Anthropic Python SDK. @@ -90,6 +121,8 @@ def __init__(self) -> None: self._logger = None self._meter = None self._parse_supported = _is_parse_supported() + self._beta_messages_supported = _is_beta_messages_supported() + self._beta_parse_supported = _is_beta_parse_supported() # pylint: disable=no-self-use def instrumentation_dependencies(self) -> Collection[str]: @@ -109,6 +142,10 @@ def _instrument(self, **kwargs: Any) -> None: meter_provider = kwargs.get("meter_provider") logger_provider = kwargs.get("logger_provider") + self._parse_supported = _is_parse_supported() + self._beta_messages_supported = _is_beta_messages_supported() + self._beta_parse_supported = _is_beta_parse_supported() + handler = TelemetryHandler( tracer_provider=tracer_provider, meter_provider=meter_provider, @@ -150,6 +187,39 @@ def _instrument(self, **kwargs: Any) -> None: async_response_context_manager_exit, ) + if self._beta_messages_supported: + wrap_function_wrapper( + "anthropic.resources.beta.messages", + "Messages.create", + messages_create(handler), + ) + wrap_function_wrapper( + "anthropic.resources.beta.messages", + "AsyncMessages.create", + async_messages_create(handler), + ) + wrap_function_wrapper( + "anthropic.resources.beta.messages", + "Messages.stream", + messages_stream(handler), + ) + wrap_function_wrapper( + "anthropic.resources.beta.messages", + "AsyncMessages.stream", + async_messages_stream(handler), + ) + if self._beta_parse_supported: + wrap_function_wrapper( + "anthropic.resources.beta.messages", + "Messages.parse", + messages_create(handler), + ) + wrap_function_wrapper( + "anthropic.resources.beta.messages", + "AsyncMessages.parse", + async_messages_create(handler), + ) + # parse() wraps create() internally in the Anthropic SDK and returns a # parsed message whose telemetry-relevant fields match Message, so the # existing create() wrappers handle it correctly. It was added in a @@ -171,37 +241,36 @@ def _uninstrument(self, **kwargs: Any) -> None: This removes all patches applied during instrumentation. """ - import anthropic # pylint: disable=import-outside-toplevel - - unwrap( - anthropic.resources.messages.Messages, # pyright: ignore[reportAttributeAccessIssue,reportUnknownMemberType,reportUnknownArgumentType] - "create", - ) - unwrap( - anthropic.resources.messages.AsyncMessages, # pyright: ignore[reportAttributeAccessIssue,reportUnknownMemberType,reportUnknownArgumentType] - "create", - ) - unwrap( - anthropic.resources.messages.Messages, # pyright: ignore[reportAttributeAccessIssue,reportUnknownMemberType,reportUnknownArgumentType] - "stream", - ) - unwrap( - anthropic.resources.messages.AsyncMessages, # pyright: ignore[reportAttributeAccessIssue,reportUnknownMemberType,reportUnknownArgumentType] - "stream", - ) from anthropic._response import ( # pylint: disable=import-outside-toplevel AsyncResponseContextManager, ResponseContextManager, ) + from anthropic.resources.messages import ( # pylint: disable=import-outside-toplevel + AsyncMessages, + Messages, + ) + unwrap(Messages, "create") + unwrap(AsyncMessages, "create") + unwrap(Messages, "stream") + unwrap(AsyncMessages, "stream") unwrap(ResponseContextManager, "__exit__") unwrap(AsyncResponseContextManager, "__aexit__") - if self._parse_supported: - unwrap( - anthropic.resources.messages.Messages, # pyright: ignore[reportAttributeAccessIssue,reportUnknownMemberType,reportUnknownArgumentType] - "parse", + if self._beta_messages_supported: + from anthropic.resources.beta.messages import ( # pylint: disable=import-outside-toplevel + AsyncMessages as AsyncBetaMessages, ) - unwrap( - anthropic.resources.messages.AsyncMessages, # pyright: ignore[reportAttributeAccessIssue,reportUnknownMemberType,reportUnknownArgumentType] - "parse", + from anthropic.resources.beta.messages import ( + Messages as BetaMessages, ) + + unwrap(BetaMessages, "create") + unwrap(AsyncBetaMessages, "create") + unwrap(BetaMessages, "stream") + unwrap(AsyncBetaMessages, "stream") + if self._beta_parse_supported: + unwrap(BetaMessages, "parse") + unwrap(AsyncBetaMessages, "parse") + if self._parse_supported: + unwrap(Messages, "parse") + unwrap(AsyncMessages, "parse") diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/_raw_response.py b/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/_raw_response.py index 7bf0d0d07..22b9d0616 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/_raw_response.py +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/_raw_response.py @@ -18,7 +18,11 @@ from anthropic._models import construct_type from anthropic.types import Message as AnthropicMessage -from .utils import is_anthropic_async_stream, is_anthropic_stream +from .utils import ( + AnthropicBetaMessage, + is_anthropic_async_stream, + is_anthropic_stream, +) from .wrappers import ( AsyncMessagesStreamWrapper, MessagesStreamWrapper, @@ -86,10 +90,12 @@ def __init__( raw_response: Any, invocation: InferenceInvocation, capture_content: bool, + is_beta: bool = False, ) -> None: super().__init__(raw_response) self._self_invocation = invocation self._self_capture_content = capture_content + self._self_is_beta = is_beta # The span is ours to end until a stream wrapper takes it over. self._self_span_open = True # Response telemetry is settled: whichever path saw the body first @@ -174,7 +180,8 @@ def _finalize_read_fallback(self) -> None: if self._self_parsing or self._self_dispatched: return message = _message_from_read_body( - getattr(self.__wrapped__, "http_response", None) + getattr(self.__wrapped__, "http_response", None), + is_beta=self._self_is_beta, ) if message is not None: MessageWrapper(message, self._self_capture_content).extract_into( @@ -370,7 +377,7 @@ def _dispatch(self, parsed: Any) -> object: # A read fallback already settled and ended this span; the caller # still gets the SDK's object, just without a second recording. return parsed - if isinstance(parsed, AnthropicMessage): + if isinstance(parsed, (AnthropicMessage, AnthropicBetaMessage)): MessageWrapper(parsed, self._self_capture_content).extract_into( self._self_invocation ) @@ -379,7 +386,10 @@ def _dispatch(self, parsed: Any) -> object: return parsed try: wrapped = _wrap_parsed_stream( - parsed, self._self_invocation, self._self_capture_content + parsed, + self._self_invocation, + self._self_capture_content, + is_beta=self._self_is_beta, ) except Exception: # pylint: disable=broad-exception-caught # Same rule as message extraction: a wrapper we failed to build @@ -413,6 +423,7 @@ def _wrap_parsed_stream( stream: Any, invocation: InferenceInvocation, capture_content: bool, + is_beta: bool = False, ) -> object | None: """Wrap a parsed stream in the matching instrumented wrapper. @@ -425,12 +436,14 @@ def _wrap_parsed_stream( cast("AnthropicAsyncStream[RawMessageStreamEvent]", stream), invocation, capture_content, + is_beta=is_beta, ) if is_anthropic_stream(stream): return MessagesStreamWrapper[None]( cast("AnthropicStream[RawMessageStreamEvent]", stream), invocation, capture_content, + is_beta=is_beta, ) return None @@ -446,7 +459,10 @@ def _body_was_read(http_response: Any) -> bool: return True -def _message_from_read_body(http_response: Any) -> AnthropicMessage | None: +def _message_from_read_body( + http_response: Any, + is_beta: bool = False, +) -> AnthropicMessage | AnthropicBetaMessage | None: """Deserialize an already-read response body into a ``Message``. Used instead of ``result.parse()`` so telemetry never runs the caller's @@ -471,9 +487,20 @@ def _message_from_read_body(http_response: Any) -> AnthropicMessage | None: if isinstance(body, dict): fields = cast("dict[str, object]", body) if fields.get("type") == "message": + target_type = ( + AnthropicBetaMessage + if ( + is_beta + and ( + hasattr(AnthropicBetaMessage, "model_fields") + or hasattr(AnthropicBetaMessage, "__fields__") + ) + ) + else AnthropicMessage + ) return cast( - AnthropicMessage, - construct_type(type_=AnthropicMessage, value=fields), + "AnthropicMessage | AnthropicBetaMessage", + construct_type(type_=target_type, value=fields), ) except Exception: # pylint: disable=broad-exception-caught _logger.debug( @@ -493,6 +520,7 @@ def wrap_raw_response( result: Any, invocation: InferenceInvocation, capture_content: bool, + is_beta: bool = False, ) -> Any: """Wrap a ``with_raw_response`` / ``with_streaming_response`` result. @@ -506,10 +534,12 @@ def wrap_raw_response( """ http_response = getattr(result, "http_response", None) if getattr(http_response, "is_closed", False): - message = _message_from_read_body(http_response) + message = _message_from_read_body(http_response, is_beta=is_beta) if message is not None: MessageWrapper(message, capture_content).extract_into(invocation) invocation.stop() return result - return RawResponseProxy(result, invocation, capture_content) + return RawResponseProxy( + result, invocation, capture_content, is_beta=is_beta + ) diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/messages_extractors.py b/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/messages_extractors.py index 9780251ae..24d36a6f8 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/messages_extractors.py +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/messages_extractors.py @@ -39,6 +39,12 @@ if TYPE_CHECKING: from collections.abc import Iterable + from anthropic.resources.beta.messages import ( + AsyncMessages as AsyncBetaMessages, + ) + from anthropic.resources.beta.messages import ( + Messages as BetaMessages, + ) from anthropic.resources.messages import AsyncMessages, Messages from anthropic.types import ( Message, @@ -50,6 +56,16 @@ ToolUnionParam, Usage, ) + from anthropic.types.beta import ( + BetaMessage, + BetaMessageParam, + BetaMetadataParam, + BetaTextBlockParam, + BetaThinkingConfigParam, + BetaToolChoiceParam, + BetaToolUnionParam, + BetaUsage, + ) @dataclass @@ -61,8 +77,8 @@ class MessageRequestParams: top_p: float | None = None stop_sequences: Sequence[str] | None = None stream: bool | None = None - messages: Iterable[MessageParam] | None = None - system: str | Iterable[TextBlockParam] | None = None + messages: Iterable[MessageParam | BetaMessageParam] | None = None + system: str | Iterable[TextBlockParam | BetaTextBlockParam] | None = None @dataclass @@ -74,7 +90,7 @@ class UsageTokens: def extract_usage_tokens( - usage: Usage | MessageDeltaUsage | None, + usage: Usage | BetaUsage | MessageDeltaUsage | None, ) -> UsageTokens: if usage is None: return UsageTokens() @@ -106,7 +122,7 @@ def extract_usage_tokens( def get_input_messages( - messages: Iterable[MessageParam] | None, + messages: Iterable[MessageParam | BetaMessageParam] | None, ) -> list[InputMessage]: if messages is None: return [] @@ -119,7 +135,7 @@ def get_input_messages( def get_system_instruction( - system: str | Iterable[TextBlockParam] | None, + system: str | Iterable[TextBlockParam | BetaTextBlockParam] | None, ) -> list[SystemInstructionPart]: if system is None: return [] @@ -133,7 +149,7 @@ def get_system_instruction( def get_output_messages_from_message( - message: Message | None, + message: Message | BetaMessage | None, ) -> list[OutputMessage]: if message is None: return [] @@ -151,7 +167,7 @@ def get_output_messages_from_message( def set_invocation_response_attributes( invocation: InferenceInvocation, - message: Message | None, + message: Message | BetaMessage | None, capture_content: bool, ) -> None: if message is None: @@ -177,17 +193,17 @@ def set_invocation_response_attributes( def extract_params( # pylint: disable=too-many-locals *, max_tokens: int | None = None, - messages: Iterable[MessageParam] | None = None, + messages: Iterable[MessageParam | BetaMessageParam] | None = None, model: str | None = None, - metadata: MetadataParam | None = None, + metadata: MetadataParam | BetaMetadataParam | None = None, service_tier: str | None = None, stop_sequences: Sequence[str] | None = None, stream: bool | None = None, - system: str | Iterable[TextBlockParam] | None = None, + system: str | Iterable[TextBlockParam | BetaTextBlockParam] | None = None, temperature: float | None = None, - thinking: ThinkingConfigParam | None = None, - tool_choice: ToolChoiceParam | None = None, - tools: Iterable[ToolUnionParam] | None = None, + thinking: ThinkingConfigParam | BetaThinkingConfigParam | None = None, + tool_choice: ToolChoiceParam | BetaToolChoiceParam | None = None, + tools: Iterable[ToolUnionParam | BetaToolUnionParam] | None = None, top_k: int | None = None, top_p: float | None = None, extra_headers: Mapping[str, str] | None = None, @@ -236,7 +252,10 @@ def extract_params( # pylint: disable=too-many-locals def get_server_address_and_port( - client_instance: Messages | AsyncMessages, + client_instance: Messages + | AsyncMessages + | BetaMessages + | AsyncBetaMessages, ) -> tuple[str | None, int | None]: base_client = getattr(client_instance, "_client", None) base_url = getattr(base_client, "base_url", None) @@ -258,7 +277,11 @@ def get_server_address_and_port( def get_llm_request_attributes( - params: MessageRequestParams, client_instance: Messages | AsyncMessages + params: MessageRequestParams, + client_instance: Messages + | AsyncMessages + | BetaMessages + | AsyncBetaMessages, ) -> dict[str, AttributeValue]: attributes: dict[str, AttributeValue | None] = { GenAIAttributes.GEN_AI_OPERATION_NAME: GenAIAttributes.GenAiOperationNameValues.CHAT.value, diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/patch.py b/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/patch.py index cd1cd9de5..5ebc9f3ec 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/patch.py +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/patch.py @@ -25,7 +25,11 @@ get_server_address_and_port, get_system_instruction, ) -from .utils import is_anthropic_async_stream, is_anthropic_stream +from .utils import ( + AnthropicBetaMessage, + is_anthropic_async_stream, + is_anthropic_stream, +) from .wrappers import ( AsyncMessagesStreamManagerWrapper, AsyncMessagesStreamWrapper, @@ -37,12 +41,25 @@ if TYPE_CHECKING: from anthropic._streaming import AsyncStream as AnthropicAsyncStream from anthropic._streaming import Stream as AnthropicStream + from anthropic.lib.streaming._beta_messages import ( # pylint: disable=no-name-in-module + BetaAsyncMessageStreamManager, + BetaMessageStreamManager, + ) from anthropic.lib.streaming._messages import ( # pylint: disable=no-name-in-module AsyncMessageStreamManager, MessageStreamManager, ) + from anthropic.resources.beta.messages import ( + AsyncMessages as AsyncBetaMessages, + ) + from anthropic.resources.beta.messages import ( + Messages as BetaMessages, + ) from anthropic.resources.messages import AsyncMessages, Messages from anthropic.types import RawMessageStreamEvent + from anthropic.types.beta import ( + BetaRawMessageStreamEvent, + ) _logger = logging.getLogger(__name__) @@ -106,12 +123,22 @@ def _is_raw_response(result: object) -> bool: return hasattr(result, "parse") and hasattr(result, "http_response") +def _is_beta_resource(instance: object) -> bool: + """Whether ``instance`` belongs to the Anthropic beta resources hierarchy.""" + return any( + cls.__module__.startswith("anthropic.resources.beta") + for cls in instance.__class__.__mro__ + ) + + def messages_create( handler: TelemetryHandler, ) -> Callable[ ..., AnthropicMessage + | AnthropicBetaMessage | AnthropicStream[RawMessageStreamEvent] + | AnthropicStream[BetaRawMessageStreamEvent] | MessagesStreamWrapper[None], ]: """Wrap the `create` method of the `Messages` class to trace it.""" @@ -120,14 +147,19 @@ def messages_create( def traced_method( wrapped: Callable[ ..., - AnthropicMessage | AnthropicStream[RawMessageStreamEvent], + AnthropicMessage + | AnthropicBetaMessage + | AnthropicStream[RawMessageStreamEvent] + | AnthropicStream[BetaRawMessageStreamEvent], ], - instance: Messages, + instance: Messages | BetaMessages, args: tuple[Any, ...], kwargs: dict[str, Any], ) -> ( AnthropicMessage + | AnthropicBetaMessage | AnthropicStream[RawMessageStreamEvent] + | AnthropicStream[BetaRawMessageStreamEvent] | MessagesStreamWrapper[None] ): invocation = _create_invocation( @@ -142,22 +174,30 @@ def traced_method( invocation.fail(exc) raise + is_beta = _is_beta_resource(instance) + if is_anthropic_stream(result): return MessagesStreamWrapper( - cast("AnthropicStream[RawMessageStreamEvent]", result), + cast( + "AnthropicStream[RawMessageStreamEvent] | AnthropicStream[BetaRawMessageStreamEvent]", + result, + ), invocation, capture_content, + is_beta=is_beta, ) if _is_raw_response(result): - return wrap_raw_response(result, invocation, capture_content) + return wrap_raw_response( + result, invocation, capture_content, is_beta=is_beta + ) - if isinstance(result, AnthropicMessage): + if isinstance(result, (AnthropicMessage, AnthropicBetaMessage)): MessageWrapper(result, capture_content).extract_into(invocation) invocation.stop() return result return cast( - 'Callable[..., "AnthropicMessage" | "AnthropicStream[RawMessageStreamEvent]" | MessagesStreamWrapper[None]]', + "Callable[..., AnthropicMessage | AnthropicBetaMessage | AnthropicStream[RawMessageStreamEvent] | AnthropicStream[BetaRawMessageStreamEvent] | MessagesStreamWrapper[None]]", traced_method, ) @@ -167,7 +207,9 @@ def async_messages_create( ) -> Callable[ ..., AnthropicMessage + | AnthropicBetaMessage | AnthropicAsyncStream[RawMessageStreamEvent] + | AnthropicAsyncStream[BetaRawMessageStreamEvent] | AsyncMessagesStreamWrapper[None], ]: """Wrap the async `create` method of the `AsyncMessages` class.""" @@ -177,15 +219,20 @@ async def traced_method( wrapped: Callable[ ..., Awaitable[ - AnthropicMessage | AnthropicAsyncStream[RawMessageStreamEvent] + AnthropicMessage + | AnthropicBetaMessage + | AnthropicAsyncStream[RawMessageStreamEvent] + | AnthropicAsyncStream[BetaRawMessageStreamEvent] ], ], - instance: AsyncMessages, + instance: AsyncMessages | AsyncBetaMessages, args: tuple[Any, ...], kwargs: dict[str, Any], ) -> ( AnthropicMessage + | AnthropicBetaMessage | AnthropicAsyncStream[RawMessageStreamEvent] + | AnthropicAsyncStream[BetaRawMessageStreamEvent] | AsyncMessagesStreamWrapper[None] ): invocation = _create_invocation( @@ -198,29 +245,37 @@ async def traced_method( invocation.fail(exc) raise + is_beta = _is_beta_resource(instance) + if is_anthropic_async_stream(result): return AsyncMessagesStreamWrapper( - cast("AnthropicAsyncStream[RawMessageStreamEvent]", result), + cast( + "AnthropicAsyncStream[RawMessageStreamEvent] | AnthropicAsyncStream[BetaRawMessageStreamEvent]", + result, + ), invocation, capture_content, + is_beta=is_beta, ) if _is_raw_response(result): - return wrap_raw_response(result, invocation, capture_content) + return wrap_raw_response( + result, invocation, capture_content, is_beta=is_beta + ) - if isinstance(result, AnthropicMessage): + if isinstance(result, (AnthropicMessage, AnthropicBetaMessage)): MessageWrapper(result, capture_content).extract_into(invocation) invocation.stop() return result return cast( - 'Callable[..., "AnthropicMessage" | "AnthropicAsyncStream[RawMessageStreamEvent]" | AsyncMessagesStreamWrapper[None]]', + "Callable[..., AnthropicMessage | AnthropicBetaMessage | AnthropicAsyncStream[RawMessageStreamEvent] | AnthropicAsyncStream[BetaRawMessageStreamEvent] | AsyncMessagesStreamWrapper[None]]", traced_method, ) def _create_invocation( handler: TelemetryHandler, - instance: Messages | AsyncMessages, + instance: Messages | AsyncMessages | BetaMessages | AsyncBetaMessages, args: tuple[Any, ...], kwargs: dict[str, Any], capture_content: bool, @@ -260,17 +315,21 @@ def messages_stream( capture_content = handler.should_capture_content() def traced_method( - wrapped: Callable[..., MessageStreamManager], - instance: Messages, + wrapped: Callable[ + ..., MessageStreamManager[Any] | BetaMessageStreamManager[Any] + ], + instance: Messages | BetaMessages, args: tuple[Any, ...], kwargs: dict[str, Any], ) -> MessagesStreamManagerWrapper[Any]: + is_beta = _is_beta_resource(instance) return MessagesStreamManagerWrapper( wrapped(*args, **kwargs), lambda: _create_invocation( handler, instance, args, kwargs, capture_content ), capture_content, + is_beta=is_beta, ) return cast( @@ -285,17 +344,23 @@ def async_messages_stream( capture_content = handler.should_capture_content() def traced_method( - wrapped: Callable[..., AsyncMessageStreamManager[Any]], - instance: AsyncMessages, + wrapped: Callable[ + ..., + AsyncMessageStreamManager[Any] + | BetaAsyncMessageStreamManager[Any], + ], + instance: AsyncMessages | AsyncBetaMessages, args: tuple[Any, ...], kwargs: dict[str, Any], ) -> AsyncMessagesStreamManagerWrapper[Any]: + is_beta = _is_beta_resource(instance) return AsyncMessagesStreamManagerWrapper( wrapped(*args, **kwargs), lambda: _create_invocation( handler, instance, args, kwargs, capture_content ), capture_content, + is_beta=is_beta, ) return cast( diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/utils.py b/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/utils.py index 033201edc..e9a962ed6 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/utils.py +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/utils.py @@ -20,13 +20,17 @@ ThinkingBlock, ThinkingDelta, ToolUseBlock, - WebSearchToolResultBlock, ) from opentelemetry.util.genai.types import ( BlobPart, + CompactionPart, + FilePart, + GenericPart, MessagePart, ReasoningPart, + ServerToolCallPart, + ServerToolCallResponsePart, TextPart, ToolCallRequestPart, ToolCallResponsePart, @@ -40,6 +44,44 @@ ContentBlockParam, RawContentBlockDelta, ) + from anthropic.types.beta import ( + BetaContentBlock, + BetaContentBlockParam, + BetaRedactedThinkingBlock, + BetaTextBlock, + BetaThinkingBlock, + BetaToolUseBlock, + ) + from anthropic.types.beta import ( + BetaMessage as AnthropicBetaMessage, + ) +else: + try: + import anthropic.types.beta as _beta_types + except (ImportError, AttributeError): + _beta_types = None + + def _get_beta_type(name: str) -> type: + cls = getattr(_beta_types, name, None) + return cls if isinstance(cls, type) else type(name, (), {}) + + AnthropicBetaMessage = _get_beta_type("BetaMessage") + BetaRedactedThinkingBlock = _get_beta_type("BetaRedactedThinkingBlock") + BetaTextBlock = _get_beta_type("BetaTextBlock") + BetaThinkingBlock = _get_beta_type("BetaThinkingBlock") + BetaToolUseBlock = _get_beta_type("BetaToolUseBlock") + + +__all__ = [ + "AnthropicBetaMessage", + "convert_content_to_parts", + "create_stream_block_state", + "is_anthropic_async_stream", + "is_anthropic_stream", + "normalize_finish_reason", + "stream_block_state_to_part", + "update_stream_block_state", +] def is_anthropic_stream(value: object) -> bool: @@ -137,47 +179,104 @@ def _convert_dict_block_to_part( id=str(block.get("id", "")), ) + if block_type in ("server_tool_use", "mcp_tool_use"): + server_tool_call: dict[str, Any] = { + "type": block_type, + "arguments": block.get("input"), + } + for key in ("caller", "server_name"): + value = block.get(key) + if value is not None: + server_tool_call[key] = value + return ServerToolCallPart( + name=str(block.get("name", "")), + server_tool_call=server_tool_call, + id=str(block.get("id", "")), + ) + if block_type == "tool_result": return ToolCallResponsePart( response=block.get("content"), id=str(block.get("tool_use_id", "")), ) + if isinstance(block_type, str) and block_type.endswith("_tool_result"): + return ServerToolCallResponsePart( + server_tool_call_response={ + key: value + for key, value in block.items() + if key != "tool_use_id" + }, + id=str(block.get("tool_use_id", "")), + ) + if block_type in ("thinking", "redacted_thinking"): thinking = block.get("thinking") or block.get("data") return ReasoningPart( content=str(thinking) if thinking is not None else "" ) + if block_type == "container_upload": + file_id = block.get("file_id") + if isinstance(file_id, str): + return FilePart( + mime_type=None, + modality="document", + file_id=file_id, + ) + + if block_type == "compaction": + content = block.get("content") + return CompactionPart( + content=content if isinstance(content, str) else None + ) + if block_type in ("image", "audio", "video", "document", "file"): - return _extract_base64_blob(block.get("source"), str(block_type)) + part = _extract_base64_blob(block.get("source"), str(block_type)) + if part is not None: + return part - return None + return ( + GenericPart(type=str(block_type)) if block_type is not None else None + ) def _convert_content_block_to_part( - block: ContentBlock | ContentBlockParam, + block: ContentBlock + | ContentBlockParam + | BetaContentBlock + | BetaContentBlockParam, ) -> MessagePart | None: """Convert an Anthropic content block to a MessagePart.""" - if isinstance(block, TextBlock): + if isinstance(block, (TextBlock, BetaTextBlock)): return TextPart(content=block.text) - if isinstance(block, (ToolUseBlock, ServerToolUseBlock)): + if isinstance(block, (ToolUseBlock, BetaToolUseBlock)): return ToolCallRequestPart( arguments=block.input, name=block.name, id=block.id ) - if isinstance(block, (ThinkingBlock, RedactedThinkingBlock)): + if isinstance( + block, + ( + ThinkingBlock, + RedactedThinkingBlock, + BetaThinkingBlock, + BetaRedactedThinkingBlock, + ), + ): content = ( - block.thinking if isinstance(block, ThinkingBlock) else block.data + block.thinking + if isinstance(block, (ThinkingBlock, BetaThinkingBlock)) + else block.data ) return ReasoningPart(content=content) - if isinstance(block, WebSearchToolResultBlock): - return ToolCallResponsePart( - response=block.model_dump().get("content"), - id=block.tool_use_id, - ) + model_dump = getattr(block, "model_dump", None) + if callable(model_dump): + dumped = model_dump() + if isinstance(dumped, Mapping): + return _convert_dict_block_to_part(cast(Mapping[str, Any], dumped)) if not hasattr(block, "get"): return None @@ -185,7 +284,14 @@ def _convert_content_block_to_part( def convert_content_to_parts( - content: str | Iterable[ContentBlock | ContentBlockParam] | None, + content: str + | Iterable[ + ContentBlock + | ContentBlockParam + | BetaContentBlock + | BetaContentBlockParam + ] + | None, ) -> list[MessagePart]: if content is None: return [] diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/wrappers.py b/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/wrappers.py index fbdd448a4..2e52441c8 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/wrappers.py +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/src/opentelemetry/instrumentation/genai/anthropic/wrappers.py @@ -40,8 +40,24 @@ except ImportError: _sdk_accumulate_event = None +try: + from anthropic.lib.streaming._beta_messages import ( # pylint: disable=no-name-in-module + accumulate_event as _sdk_beta_accumulate_event, + ) +except ImportError: + _sdk_beta_accumulate_event = None + if TYPE_CHECKING: from anthropic._streaming import AsyncStream, Stream + from anthropic.lib.streaming._beta_messages import ( # pylint: disable=no-name-in-module + BetaAsyncMessageStream, + BetaAsyncMessageStreamManager, + BetaMessageStream, + BetaMessageStreamManager, + ) + from anthropic.lib.streaming._beta_types import ( # pylint: disable=no-name-in-module + ParsedBetaMessageStreamEvent, + ) from anthropic.lib.streaming._messages import ( # pylint: disable=no-name-in-module AsyncMessageStream, AsyncMessageStreamManager, @@ -55,6 +71,11 @@ Message, RawMessageStreamEvent, ) + from anthropic.types.beta import ( + BetaMessage, + BetaRawMessageStreamEvent, + ) + from anthropic.types.beta.parsed_beta_message import ParsedBetaMessage from anthropic.types.parsed_message import ParsedMessage from opentelemetry.util.genai.invocation import InferenceInvocation @@ -62,6 +83,9 @@ ResponseFormatT = TypeVar("ResponseFormatT") accumulate_event = cast("Callable[..., Message] | None", _sdk_accumulate_event) +beta_accumulate_event = cast( + "Callable[..., ParsedBetaMessage[Any]] | None", _sdk_beta_accumulate_event +) _accumulate_takes_json_bufs = False if accumulate_event is not None: @@ -72,6 +96,22 @@ except (ValueError, TypeError): _accumulate_takes_json_bufs = False +_beta_accumulate_accepts_headers = False +_beta_accumulate_takes_json_bufs = False +if beta_accumulate_event is not None: + try: + _beta_accumulate_parameters = inspect.signature( + beta_accumulate_event + ).parameters + _beta_accumulate_accepts_headers = ( + "request_headers" in _beta_accumulate_parameters + ) + _beta_accumulate_takes_json_bufs = ( + "json_bufs" in _beta_accumulate_parameters + ) + except (ValueError, TypeError): + pass + _accumulation_disabled = False @@ -82,7 +122,13 @@ def stream(self) -> object: ... def _set_response_attributes( invocation: InferenceInvocation, - result: Message | None, + result: ( + Message + | BetaMessage + | ParsedMessage[Any] + | ParsedBetaMessage[Any] + | None + ), capture_content: bool, ) -> None: set_invocation_response_attributes(invocation, result, capture_content) @@ -91,7 +137,7 @@ def _set_response_attributes( class MessageWrapper: """Wrapper for non-streaming Message response that handles telemetry.""" - def __init__(self, message: Message, capture_content: bool): + def __init__(self, message: Message | BetaMessage, capture_content: bool): self._message = message self._capture_content = capture_content @@ -102,17 +148,25 @@ def extract_into(self, invocation: InferenceInvocation) -> None: ) @property - def message(self) -> Message: + def message(self) -> Message | BetaMessage: """Return the wrapped Message object.""" return self._message class _MessagesStreamMixin(Generic[ResponseFormatT]): _self_invocation: InferenceInvocation - _self_message: Message | ParsedMessage[ResponseFormatT] | None + _self_message: ( + Message + | BetaMessage + | ParsedMessage[ResponseFormatT] + | ParsedBetaMessage[ResponseFormatT] + | None + ) _self_capture_content: bool _self_message_telemetry_finalized: bool _self_json_bufs: dict[int, bytes] + _self_is_beta: bool + _self_beta_accumulation_disabled: bool def _stop(self) -> None: if self._self_message_telemetry_finalized: @@ -139,19 +193,56 @@ def _on_stream_error(self, error: BaseException) -> None: def _process_chunk( self, - chunk: RawMessageStreamEvent - | ParsedMessageStreamEvent[ResponseFormatT], + chunk: ( + RawMessageStreamEvent + | BetaRawMessageStreamEvent + | ParsedMessageStreamEvent[ResponseFormatT] + | ParsedBetaMessageStreamEvent[ResponseFormatT] + ), ) -> None: """Accumulate a final message snapshot from a streaming chunk.""" global _accumulation_disabled stream = cast(_StreamWrapperWithStream, self).stream snapshot = cast( - "ParsedMessage[ResponseFormatT] | None", + "ParsedMessage[ResponseFormatT] | ParsedBetaMessage[ResponseFormatT] | None", getattr(stream, "current_message_snapshot", None), ) if snapshot is not None: self._self_message = snapshot return + is_beta = self._self_is_beta or getattr( + chunk.__class__, "__module__", "" + ).startswith("anthropic.types.beta") + if ( + is_beta + and beta_accumulate_event is not None + and not self._self_beta_accumulation_disabled + ): + beta_kwargs: dict[str, Any] = { + "event": chunk, + "current_snapshot": self._self_message, + } + if _beta_accumulate_takes_json_bufs: + beta_kwargs["json_bufs"] = self._self_json_bufs + if _beta_accumulate_accepts_headers: + response = getattr(stream, "response", None) + request = getattr(response, "request", None) + headers = getattr(request, "headers", None) + beta_kwargs["request_headers"] = ( + headers if headers is not None else _http_lib.Headers() + ) + try: + self._self_message = beta_accumulate_event(**beta_kwargs) + except BaseException as exc: + if not isinstance(exc, Exception): + raise + self._self_beta_accumulation_disabled = True + _logger.debug( + "Failed to accumulate beta stream event", exc_info=True + ) + else: + return + if accumulate_event is None or _accumulation_disabled: return @@ -182,7 +273,7 @@ def _process_chunk( class MessagesStreamWrapper( _MessagesStreamMixin[ResponseFormatT], SyncStreamWrapper[ - "RawMessageStreamEvent | ParsedMessageStreamEvent[ResponseFormatT]" + "RawMessageStreamEvent | BetaRawMessageStreamEvent | ParsedMessageStreamEvent[ResponseFormatT] | ParsedBetaMessageStreamEvent[ResponseFormatT]" ], Generic[ResponseFormatT], ): @@ -190,9 +281,15 @@ class MessagesStreamWrapper( def __init__( self, - stream: Stream[RawMessageStreamEvent] | MessageStream[ResponseFormatT], + stream: ( + Stream[RawMessageStreamEvent] + | Stream[BetaRawMessageStreamEvent] + | MessageStream[ResponseFormatT] + | BetaMessageStream[ResponseFormatT] + ), invocation: InferenceInvocation, capture_content: bool, + is_beta: bool = False, ): super().__init__(stream, invocation=invocation) self._self_invocation = invocation @@ -200,6 +297,8 @@ def __init__( self._self_capture_content = capture_content self._self_message_telemetry_finalized = False self._self_json_bufs = {} + self._self_is_beta = is_beta + self._self_beta_accumulation_disabled = False @property def response(self) -> _http_lib.Response: @@ -208,13 +307,23 @@ def response(self) -> _http_lib.Response: @property def stream( self, - ) -> Stream[RawMessageStreamEvent] | MessageStream[ResponseFormatT]: + ) -> ( + Stream[RawMessageStreamEvent] + | Stream[BetaRawMessageStreamEvent] + | MessageStream[ResponseFormatT] + | BetaMessageStream[ResponseFormatT] + ): return self._self_stream @stream.setter def stream( self, - stream: Stream[RawMessageStreamEvent] | MessageStream[ResponseFormatT], + stream: ( + Stream[RawMessageStreamEvent] + | Stream[BetaRawMessageStreamEvent] + | MessageStream[ResponseFormatT] + | BetaMessageStream[ResponseFormatT] + ), ) -> None: self._set_stream(stream) @@ -222,7 +331,7 @@ def stream( class AsyncMessagesStreamWrapper( _MessagesStreamMixin[ResponseFormatT], AsyncStreamWrapper[ - "RawMessageStreamEvent | ParsedMessageStreamEvent[ResponseFormatT]" + "RawMessageStreamEvent | BetaRawMessageStreamEvent | ParsedMessageStreamEvent[ResponseFormatT] | ParsedBetaMessageStreamEvent[ResponseFormatT]" ], Generic[ResponseFormatT], ): @@ -230,10 +339,15 @@ class AsyncMessagesStreamWrapper( def __init__( self, - stream: AsyncStream[RawMessageStreamEvent] - | AsyncMessageStream[ResponseFormatT], + stream: ( + AsyncStream[RawMessageStreamEvent] + | AsyncStream[BetaRawMessageStreamEvent] + | AsyncMessageStream[ResponseFormatT] + | BetaAsyncMessageStream[ResponseFormatT] + ), invocation: InferenceInvocation, capture_content: bool, + is_beta: bool = False, ): super().__init__(stream, invocation=invocation) self._self_invocation = invocation @@ -241,6 +355,8 @@ def __init__( self._self_capture_content = capture_content self._self_message_telemetry_finalized = False self._self_json_bufs = {} + self._self_is_beta = is_beta + self._self_beta_accumulation_disabled = False @property def response(self) -> _http_lib.Response: @@ -251,22 +367,28 @@ def stream( self, ) -> ( AsyncStream[RawMessageStreamEvent] + | AsyncStream[BetaRawMessageStreamEvent] | AsyncMessageStream[ResponseFormatT] + | BetaAsyncMessageStream[ResponseFormatT] ): return self._self_stream @stream.setter def stream( self, - stream: AsyncStream[RawMessageStreamEvent] - | AsyncMessageStream[ResponseFormatT], + stream: ( + AsyncStream[RawMessageStreamEvent] + | AsyncStream[BetaRawMessageStreamEvent] + | AsyncMessageStream[ResponseFormatT] + | BetaAsyncMessageStream[ResponseFormatT] + ), ) -> None: self._set_stream(stream) class MessagesStreamManagerWrapper( SyncStreamManagerWrapper[ - "MessageStream[ResponseFormatT]", + "MessageStream[ResponseFormatT] | BetaMessageStream[ResponseFormatT]", "InferenceInvocation", "MessagesStreamWrapper[ResponseFormatT]", ], @@ -276,26 +398,36 @@ class MessagesStreamManagerWrapper( def __init__( self, - manager: MessageStreamManager[ResponseFormatT], + manager: ( + MessageStreamManager[ResponseFormatT] + | BetaMessageStreamManager[ResponseFormatT] + ), invocation_factory: Callable[[], InferenceInvocation], capture_content: bool, + is_beta: bool = False, ): super().__init__(manager, invocation_factory) self._self_capture_content = capture_content + self._self_is_beta = is_beta def _wrap_stream( self, - stream: MessageStream[ResponseFormatT], + stream: ( + MessageStream[ResponseFormatT] | BetaMessageStream[ResponseFormatT] + ), invocation: InferenceInvocation, ) -> MessagesStreamWrapper[ResponseFormatT]: return MessagesStreamWrapper( - stream, invocation, self._self_capture_content + stream, + invocation, + self._self_capture_content, + is_beta=self._self_is_beta, ) class AsyncMessagesStreamManagerWrapper( AsyncStreamManagerWrapper[ - "AsyncMessageStream[ResponseFormatT]", + "AsyncMessageStream[ResponseFormatT] | BetaAsyncMessageStream[ResponseFormatT]", "InferenceInvocation", "AsyncMessagesStreamWrapper[ResponseFormatT]", ], @@ -309,18 +441,29 @@ class AsyncMessagesStreamManagerWrapper( def __init__( self, - manager: AsyncMessageStreamManager[ResponseFormatT], + manager: ( + AsyncMessageStreamManager[ResponseFormatT] + | BetaAsyncMessageStreamManager[ResponseFormatT] + ), invocation_factory: Callable[[], InferenceInvocation], capture_content: bool, + is_beta: bool = False, ): super().__init__(manager, invocation_factory) self._self_capture_content = capture_content + self._self_is_beta = is_beta def _wrap_stream( self, - stream: AsyncMessageStream[ResponseFormatT], + stream: ( + AsyncMessageStream[ResponseFormatT] + | BetaAsyncMessageStream[ResponseFormatT] + ), invocation: InferenceInvocation, ) -> AsyncMessagesStreamWrapper[ResponseFormatT]: return AsyncMessagesStreamWrapper( - stream, invocation, self._self_capture_content + stream, + invocation, + self._self_capture_content, + is_beta=self._self_is_beta, ) diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/inference_beta_conformance.yaml b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/inference_beta_conformance.yaml new file mode 100644 index 000000000..35031b9b0 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/inference_beta_conformance.yaml @@ -0,0 +1,53 @@ +# TODO: this is generated by AI, re-record +interactions: +- request: + body: |- + { + "max_tokens": 100, + "messages": [ + { + "role": "user", + "content": "Say hello in one word." + } + ], + "model": "claude-sonnet-4-6" + } + headers: {} + method: POST + uri: https://api.anthropic.com/v1/messages?beta=true + response: + body: + string: |- + { + "model": "claude-sonnet-4-6", + "id": "msg_0176GK1qFwwpVM59jDYiKPjN", + "type": "message", + "role": "assistant", + "content": [ + { + "type": "text", + "text": "Hello!" + } + ], + "stop_reason": "end_turn", + "stop_sequence": null, + "usage": { + "input_tokens": 13, + "cache_creation_input_tokens": 0, + "cache_read_input_tokens": 0, + "cache_creation": { + "ephemeral_5m_input_tokens": 0, + "ephemeral_1h_input_tokens": 0 + }, + "output_tokens": 5, + "service_tier": "standard", + "inference_geo": "not_available" + } + } + headers: + Content-Type: + - application/json + status: + code: 200 + message: OK +version: 1 diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/inference_beta_server_tool_calling_conformance.yaml b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/inference_beta_server_tool_calling_conformance.yaml new file mode 100644 index 000000000..232a43bba --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/inference_beta_server_tool_calling_conformance.yaml @@ -0,0 +1,74 @@ +# TODO: this is generated by AI, re-record +interactions: +- request: + body: |- + { + "max_tokens": 256, + "messages": [ + { + "role": "user", + "content": "Search for OpenTelemetry." + } + ], + "model": "claude-sonnet-4-6", + "tools": [ + { + "type": "web_search_20250305", + "name": "web_search", + "max_uses": 1 + } + ] + } + headers: {} + method: POST + uri: https://api.anthropic.com/v1/messages?beta=true + response: + body: + string: |- + { + "model": "claude-sonnet-4-6", + "id": "msg_server_tool", + "type": "message", + "role": "assistant", + "content": [ + { + "type": "server_tool_use", + "id": "srvtoolu_01", + "name": "web_search", + "input": {"query": "OpenTelemetry"} + }, + { + "type": "web_search_tool_result", + "tool_use_id": "srvtoolu_01", + "content": { + "type": "web_search_tool_result_error", + "error_code": "unavailable" + } + }, + { + "type": "text", + "text": "Search was unavailable." + } + ], + "stop_reason": "end_turn", + "stop_sequence": null, + "usage": { + "input_tokens": 20, + "cache_creation_input_tokens": 0, + "cache_read_input_tokens": 0, + "cache_creation": { + "ephemeral_5m_input_tokens": 0, + "ephemeral_1h_input_tokens": 0 + }, + "output_tokens": 15, + "service_tier": "standard", + "inference_geo": "not_available" + } + } + headers: + Content-Type: + - application/json + status: + code: 200 + message: OK +version: 1 diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/inference_beta_streaming_conformance.yaml b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/inference_beta_streaming_conformance.yaml new file mode 100644 index 000000000..3d5376d34 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/inference_beta_streaming_conformance.yaml @@ -0,0 +1,46 @@ +# TODO: this is generated by AI, re-record +interactions: +- request: + body: |- + { + "max_tokens": 100, + "messages": [ + { + "role": "user", + "content": "Say hello in one word." + } + ], + "model": "claude-sonnet-4-6", + "stream": true + } + headers: {} + method: POST + uri: https://api.anthropic.com/v1/messages?beta=true + response: + body: + string: |+ + event: message_start + data: {"type":"message_start","message":{"model":"claude-sonnet-4-6","id":"msg_01VY8H3oPFs4WWVnJXFDxNhU","type":"message","role":"assistant","content":[],"stop_reason":null,"stop_sequence":null,"usage":{"input_tokens":13,"cache_creation_input_tokens":0,"cache_read_input_tokens":0,"cache_creation":{"ephemeral_5m_input_tokens":0,"ephemeral_1h_input_tokens":0},"output_tokens":2,"service_tier":"standard","inference_geo":"not_available"}} } + + event: content_block_start + data: {"type":"content_block_start","index":0,"content_block":{"type":"text","text":""} } + + event: content_block_delta + data: {"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"Hello"} } + + event: content_block_stop + data: {"type":"content_block_stop","index":0 } + + event: message_delta + data: {"type":"message_delta","delta":{"stop_reason":"end_turn","stop_sequence":null},"usage":{"output_tokens":5} } + + event: message_stop + data: {"type":"message_stop" } + + headers: + Content-Type: + - text/event-stream; charset=utf-8 + status: + code: 200 + message: OK +version: 1 diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_create_api_error.yaml b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_create_api_error.yaml new file mode 100644 index 000000000..73b02ce9f --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_create_api_error.yaml @@ -0,0 +1,35 @@ +# TODO: this is generated by AI, re-record +interactions: +- request: + body: |- + { + "max_tokens": 100, + "messages": [ + { + "role": "user", + "content": "Hello" + } + ], + "model": "invalid-model-name" + } + headers: {} + method: POST + uri: https://api.anthropic.com/v1/messages?beta=true + response: + body: + string: |- + { + "type": "error", + "error": { + "type": "not_found_error", + "message": "model: invalid-model-name" + }, + "request_id": "req_011CYfQMhgSid28ainNjq126" + } + headers: + Content-Type: + - application/json + status: + code: 404 + message: Not Found +version: 1 diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_create_basic.yaml b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_create_basic.yaml new file mode 100644 index 000000000..35031b9b0 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_create_basic.yaml @@ -0,0 +1,53 @@ +# TODO: this is generated by AI, re-record +interactions: +- request: + body: |- + { + "max_tokens": 100, + "messages": [ + { + "role": "user", + "content": "Say hello in one word." + } + ], + "model": "claude-sonnet-4-6" + } + headers: {} + method: POST + uri: https://api.anthropic.com/v1/messages?beta=true + response: + body: + string: |- + { + "model": "claude-sonnet-4-6", + "id": "msg_0176GK1qFwwpVM59jDYiKPjN", + "type": "message", + "role": "assistant", + "content": [ + { + "type": "text", + "text": "Hello!" + } + ], + "stop_reason": "end_turn", + "stop_sequence": null, + "usage": { + "input_tokens": 13, + "cache_creation_input_tokens": 0, + "cache_read_input_tokens": 0, + "cache_creation": { + "ephemeral_5m_input_tokens": 0, + "ephemeral_1h_input_tokens": 0 + }, + "output_tokens": 5, + "service_tier": "standard", + "inference_geo": "not_available" + } + } + headers: + Content-Type: + - application/json + status: + code: 200 + message: OK +version: 1 diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_create_streaming.yaml b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_create_streaming.yaml new file mode 100644 index 000000000..3d5376d34 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_create_streaming.yaml @@ -0,0 +1,46 @@ +# TODO: this is generated by AI, re-record +interactions: +- request: + body: |- + { + "max_tokens": 100, + "messages": [ + { + "role": "user", + "content": "Say hello in one word." + } + ], + "model": "claude-sonnet-4-6", + "stream": true + } + headers: {} + method: POST + uri: https://api.anthropic.com/v1/messages?beta=true + response: + body: + string: |+ + event: message_start + data: {"type":"message_start","message":{"model":"claude-sonnet-4-6","id":"msg_01VY8H3oPFs4WWVnJXFDxNhU","type":"message","role":"assistant","content":[],"stop_reason":null,"stop_sequence":null,"usage":{"input_tokens":13,"cache_creation_input_tokens":0,"cache_read_input_tokens":0,"cache_creation":{"ephemeral_5m_input_tokens":0,"ephemeral_1h_input_tokens":0},"output_tokens":2,"service_tier":"standard","inference_geo":"not_available"}} } + + event: content_block_start + data: {"type":"content_block_start","index":0,"content_block":{"type":"text","text":""} } + + event: content_block_delta + data: {"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"Hello"} } + + event: content_block_stop + data: {"type":"content_block_stop","index":0 } + + event: message_delta + data: {"type":"message_delta","delta":{"stop_reason":"end_turn","stop_sequence":null},"usage":{"output_tokens":5} } + + event: message_stop + data: {"type":"message_stop" } + + headers: + Content-Type: + - text/event-stream; charset=utf-8 + status: + code: 200 + message: OK +version: 1 diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_create_with_content.yaml b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_create_with_content.yaml new file mode 100644 index 000000000..35031b9b0 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_create_with_content.yaml @@ -0,0 +1,53 @@ +# TODO: this is generated by AI, re-record +interactions: +- request: + body: |- + { + "max_tokens": 100, + "messages": [ + { + "role": "user", + "content": "Say hello in one word." + } + ], + "model": "claude-sonnet-4-6" + } + headers: {} + method: POST + uri: https://api.anthropic.com/v1/messages?beta=true + response: + body: + string: |- + { + "model": "claude-sonnet-4-6", + "id": "msg_0176GK1qFwwpVM59jDYiKPjN", + "type": "message", + "role": "assistant", + "content": [ + { + "type": "text", + "text": "Hello!" + } + ], + "stop_reason": "end_turn", + "stop_sequence": null, + "usage": { + "input_tokens": 13, + "cache_creation_input_tokens": 0, + "cache_read_input_tokens": 0, + "cache_creation": { + "ephemeral_5m_input_tokens": 0, + "ephemeral_1h_input_tokens": 0 + }, + "output_tokens": 5, + "service_tier": "standard", + "inference_geo": "not_available" + } + } + headers: + Content-Type: + - application/json + status: + code: 200 + message: OK +version: 1 diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_parse_basic.yaml b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_parse_basic.yaml new file mode 100644 index 000000000..a5cc67d0e --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_parse_basic.yaml @@ -0,0 +1,72 @@ +# TODO: this is generated by AI, re-record +interactions: +- request: + body: |- + { + "max_tokens": 100, + "messages": [ + { + "role": "user", + "content": "Return JSON with a greeting field set to hello." + } + ], + "model": "claude-haiku-4-5", + "output_config": { + "format": { + "schema": { + "type": "object", + "title": "Greeting", + "properties": { + "greeting": { + "type": "string", + "title": "Greeting" + } + }, + "additionalProperties": false, + "required": [ + "greeting" + ] + }, + "type": "json_schema" + } + } + } + headers: {} + method: POST + uri: https://api.anthropic.com/v1/messages?beta=true + response: + body: + string: |- + { + "model": "claude-haiku-4-5-20251001", + "id": "msg_01X2c4jNgRTKYFpNTgfptMUA", + "type": "message", + "role": "assistant", + "content": [ + { + "type": "text", + "text": "{\"greeting\":\"hello\"}" + } + ], + "stop_reason": "end_turn", + "stop_sequence": null, + "usage": { + "input_tokens": 176, + "cache_creation_input_tokens": 0, + "cache_read_input_tokens": 0, + "cache_creation": { + "ephemeral_5m_input_tokens": 0, + "ephemeral_1h_input_tokens": 0 + }, + "output_tokens": 8, + "service_tier": "standard", + "inference_geo": "not_available" + } + } + headers: + Content-Type: + - application/json + status: + code: 200 + message: OK +version: 1 diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_stream.yaml b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_stream.yaml new file mode 100644 index 000000000..3d5376d34 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_stream.yaml @@ -0,0 +1,46 @@ +# TODO: this is generated by AI, re-record +interactions: +- request: + body: |- + { + "max_tokens": 100, + "messages": [ + { + "role": "user", + "content": "Say hello in one word." + } + ], + "model": "claude-sonnet-4-6", + "stream": true + } + headers: {} + method: POST + uri: https://api.anthropic.com/v1/messages?beta=true + response: + body: + string: |+ + event: message_start + data: {"type":"message_start","message":{"model":"claude-sonnet-4-6","id":"msg_01VY8H3oPFs4WWVnJXFDxNhU","type":"message","role":"assistant","content":[],"stop_reason":null,"stop_sequence":null,"usage":{"input_tokens":13,"cache_creation_input_tokens":0,"cache_read_input_tokens":0,"cache_creation":{"ephemeral_5m_input_tokens":0,"ephemeral_1h_input_tokens":0},"output_tokens":2,"service_tier":"standard","inference_geo":"not_available"}} } + + event: content_block_start + data: {"type":"content_block_start","index":0,"content_block":{"type":"text","text":""} } + + event: content_block_delta + data: {"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"Hello"} } + + event: content_block_stop + data: {"type":"content_block_stop","index":0 } + + event: message_delta + data: {"type":"message_delta","delta":{"stop_reason":"end_turn","stop_sequence":null},"usage":{"output_tokens":5} } + + event: message_stop + data: {"type":"message_stop" } + + headers: + Content-Type: + - text/event-stream; charset=utf-8 + status: + code: 200 + message: OK +version: 1 diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_with_raw_response.yaml b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_with_raw_response.yaml new file mode 100644 index 000000000..35031b9b0 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_async_beta_messages_with_raw_response.yaml @@ -0,0 +1,53 @@ +# TODO: this is generated by AI, re-record +interactions: +- request: + body: |- + { + "max_tokens": 100, + "messages": [ + { + "role": "user", + "content": "Say hello in one word." + } + ], + "model": "claude-sonnet-4-6" + } + headers: {} + method: POST + uri: https://api.anthropic.com/v1/messages?beta=true + response: + body: + string: |- + { + "model": "claude-sonnet-4-6", + "id": "msg_0176GK1qFwwpVM59jDYiKPjN", + "type": "message", + "role": "assistant", + "content": [ + { + "type": "text", + "text": "Hello!" + } + ], + "stop_reason": "end_turn", + "stop_sequence": null, + "usage": { + "input_tokens": 13, + "cache_creation_input_tokens": 0, + "cache_read_input_tokens": 0, + "cache_creation": { + "ephemeral_5m_input_tokens": 0, + "ephemeral_1h_input_tokens": 0 + }, + "output_tokens": 5, + "service_tier": "standard", + "inference_geo": "not_available" + } + } + headers: + Content-Type: + - application/json + status: + code: 200 + message: OK +version: 1 diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_create_api_error.yaml b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_create_api_error.yaml new file mode 100644 index 000000000..73b02ce9f --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_create_api_error.yaml @@ -0,0 +1,35 @@ +# TODO: this is generated by AI, re-record +interactions: +- request: + body: |- + { + "max_tokens": 100, + "messages": [ + { + "role": "user", + "content": "Hello" + } + ], + "model": "invalid-model-name" + } + headers: {} + method: POST + uri: https://api.anthropic.com/v1/messages?beta=true + response: + body: + string: |- + { + "type": "error", + "error": { + "type": "not_found_error", + "message": "model: invalid-model-name" + }, + "request_id": "req_011CYfQMhgSid28ainNjq126" + } + headers: + Content-Type: + - application/json + status: + code: 404 + message: Not Found +version: 1 diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_create_basic.yaml b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_create_basic.yaml new file mode 100644 index 000000000..35031b9b0 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_create_basic.yaml @@ -0,0 +1,53 @@ +# TODO: this is generated by AI, re-record +interactions: +- request: + body: |- + { + "max_tokens": 100, + "messages": [ + { + "role": "user", + "content": "Say hello in one word." + } + ], + "model": "claude-sonnet-4-6" + } + headers: {} + method: POST + uri: https://api.anthropic.com/v1/messages?beta=true + response: + body: + string: |- + { + "model": "claude-sonnet-4-6", + "id": "msg_0176GK1qFwwpVM59jDYiKPjN", + "type": "message", + "role": "assistant", + "content": [ + { + "type": "text", + "text": "Hello!" + } + ], + "stop_reason": "end_turn", + "stop_sequence": null, + "usage": { + "input_tokens": 13, + "cache_creation_input_tokens": 0, + "cache_read_input_tokens": 0, + "cache_creation": { + "ephemeral_5m_input_tokens": 0, + "ephemeral_1h_input_tokens": 0 + }, + "output_tokens": 5, + "service_tier": "standard", + "inference_geo": "not_available" + } + } + headers: + Content-Type: + - application/json + status: + code: 200 + message: OK +version: 1 diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_create_streaming.yaml b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_create_streaming.yaml new file mode 100644 index 000000000..3d5376d34 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_create_streaming.yaml @@ -0,0 +1,46 @@ +# TODO: this is generated by AI, re-record +interactions: +- request: + body: |- + { + "max_tokens": 100, + "messages": [ + { + "role": "user", + "content": "Say hello in one word." + } + ], + "model": "claude-sonnet-4-6", + "stream": true + } + headers: {} + method: POST + uri: https://api.anthropic.com/v1/messages?beta=true + response: + body: + string: |+ + event: message_start + data: {"type":"message_start","message":{"model":"claude-sonnet-4-6","id":"msg_01VY8H3oPFs4WWVnJXFDxNhU","type":"message","role":"assistant","content":[],"stop_reason":null,"stop_sequence":null,"usage":{"input_tokens":13,"cache_creation_input_tokens":0,"cache_read_input_tokens":0,"cache_creation":{"ephemeral_5m_input_tokens":0,"ephemeral_1h_input_tokens":0},"output_tokens":2,"service_tier":"standard","inference_geo":"not_available"}} } + + event: content_block_start + data: {"type":"content_block_start","index":0,"content_block":{"type":"text","text":""} } + + event: content_block_delta + data: {"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"Hello"} } + + event: content_block_stop + data: {"type":"content_block_stop","index":0 } + + event: message_delta + data: {"type":"message_delta","delta":{"stop_reason":"end_turn","stop_sequence":null},"usage":{"output_tokens":5} } + + event: message_stop + data: {"type":"message_stop" } + + headers: + Content-Type: + - text/event-stream; charset=utf-8 + status: + code: 200 + message: OK +version: 1 diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_create_with_content.yaml b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_create_with_content.yaml new file mode 100644 index 000000000..35031b9b0 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_create_with_content.yaml @@ -0,0 +1,53 @@ +# TODO: this is generated by AI, re-record +interactions: +- request: + body: |- + { + "max_tokens": 100, + "messages": [ + { + "role": "user", + "content": "Say hello in one word." + } + ], + "model": "claude-sonnet-4-6" + } + headers: {} + method: POST + uri: https://api.anthropic.com/v1/messages?beta=true + response: + body: + string: |- + { + "model": "claude-sonnet-4-6", + "id": "msg_0176GK1qFwwpVM59jDYiKPjN", + "type": "message", + "role": "assistant", + "content": [ + { + "type": "text", + "text": "Hello!" + } + ], + "stop_reason": "end_turn", + "stop_sequence": null, + "usage": { + "input_tokens": 13, + "cache_creation_input_tokens": 0, + "cache_read_input_tokens": 0, + "cache_creation": { + "ephemeral_5m_input_tokens": 0, + "ephemeral_1h_input_tokens": 0 + }, + "output_tokens": 5, + "service_tier": "standard", + "inference_geo": "not_available" + } + } + headers: + Content-Type: + - application/json + status: + code: 200 + message: OK +version: 1 diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_parse_basic.yaml b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_parse_basic.yaml new file mode 100644 index 000000000..a5cc67d0e --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_parse_basic.yaml @@ -0,0 +1,72 @@ +# TODO: this is generated by AI, re-record +interactions: +- request: + body: |- + { + "max_tokens": 100, + "messages": [ + { + "role": "user", + "content": "Return JSON with a greeting field set to hello." + } + ], + "model": "claude-haiku-4-5", + "output_config": { + "format": { + "schema": { + "type": "object", + "title": "Greeting", + "properties": { + "greeting": { + "type": "string", + "title": "Greeting" + } + }, + "additionalProperties": false, + "required": [ + "greeting" + ] + }, + "type": "json_schema" + } + } + } + headers: {} + method: POST + uri: https://api.anthropic.com/v1/messages?beta=true + response: + body: + string: |- + { + "model": "claude-haiku-4-5-20251001", + "id": "msg_01X2c4jNgRTKYFpNTgfptMUA", + "type": "message", + "role": "assistant", + "content": [ + { + "type": "text", + "text": "{\"greeting\":\"hello\"}" + } + ], + "stop_reason": "end_turn", + "stop_sequence": null, + "usage": { + "input_tokens": 176, + "cache_creation_input_tokens": 0, + "cache_read_input_tokens": 0, + "cache_creation": { + "ephemeral_5m_input_tokens": 0, + "ephemeral_1h_input_tokens": 0 + }, + "output_tokens": 8, + "service_tier": "standard", + "inference_geo": "not_available" + } + } + headers: + Content-Type: + - application/json + status: + code: 200 + message: OK +version: 1 diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_stream.yaml b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_stream.yaml new file mode 100644 index 000000000..3d5376d34 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_stream.yaml @@ -0,0 +1,46 @@ +# TODO: this is generated by AI, re-record +interactions: +- request: + body: |- + { + "max_tokens": 100, + "messages": [ + { + "role": "user", + "content": "Say hello in one word." + } + ], + "model": "claude-sonnet-4-6", + "stream": true + } + headers: {} + method: POST + uri: https://api.anthropic.com/v1/messages?beta=true + response: + body: + string: |+ + event: message_start + data: {"type":"message_start","message":{"model":"claude-sonnet-4-6","id":"msg_01VY8H3oPFs4WWVnJXFDxNhU","type":"message","role":"assistant","content":[],"stop_reason":null,"stop_sequence":null,"usage":{"input_tokens":13,"cache_creation_input_tokens":0,"cache_read_input_tokens":0,"cache_creation":{"ephemeral_5m_input_tokens":0,"ephemeral_1h_input_tokens":0},"output_tokens":2,"service_tier":"standard","inference_geo":"not_available"}} } + + event: content_block_start + data: {"type":"content_block_start","index":0,"content_block":{"type":"text","text":""} } + + event: content_block_delta + data: {"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"Hello"} } + + event: content_block_stop + data: {"type":"content_block_stop","index":0 } + + event: message_delta + data: {"type":"message_delta","delta":{"stop_reason":"end_turn","stop_sequence":null},"usage":{"output_tokens":5} } + + event: message_stop + data: {"type":"message_stop" } + + headers: + Content-Type: + - text/event-stream; charset=utf-8 + status: + code: 200 + message: OK +version: 1 diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_with_raw_response.yaml b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_with_raw_response.yaml new file mode 100644 index 000000000..35031b9b0 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/cassettes/test_sync_beta_messages_with_raw_response.yaml @@ -0,0 +1,53 @@ +# TODO: this is generated by AI, re-record +interactions: +- request: + body: |- + { + "max_tokens": 100, + "messages": [ + { + "role": "user", + "content": "Say hello in one word." + } + ], + "model": "claude-sonnet-4-6" + } + headers: {} + method: POST + uri: https://api.anthropic.com/v1/messages?beta=true + response: + body: + string: |- + { + "model": "claude-sonnet-4-6", + "id": "msg_0176GK1qFwwpVM59jDYiKPjN", + "type": "message", + "role": "assistant", + "content": [ + { + "type": "text", + "text": "Hello!" + } + ], + "stop_reason": "end_turn", + "stop_sequence": null, + "usage": { + "input_tokens": 13, + "cache_creation_input_tokens": 0, + "cache_read_input_tokens": 0, + "cache_creation": { + "ephemeral_5m_input_tokens": 0, + "ephemeral_1h_input_tokens": 0 + }, + "output_tokens": 5, + "service_tier": "standard", + "inference_geo": "not_available" + } + } + headers: + Content-Type: + - application/json + status: + code: 200 + message: OK +version: 1 diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/conformance/inference_beta.py b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/conformance/inference_beta.py new file mode 100644 index 000000000..d6172b957 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/conformance/inference_beta.py @@ -0,0 +1,123 @@ +# Copyright The OpenTelemetry Authors +# SPDX-License-Identifier: Apache-2.0 + +"""Conformance scenario: anthropic beta chat (inference).""" + +from __future__ import annotations + +import os +from typing import Any +from unittest import mock + +from anthropic import Anthropic + +from opentelemetry.instrumentation.genai.anthropic import AnthropicInstrumentor +from opentelemetry.sdk._logs import LoggerProvider +from opentelemetry.sdk.metrics import MeterProvider +from opentelemetry.sdk.trace import TracerProvider +from opentelemetry.test.weaver_live_check import LiveCheckReport +from opentelemetry.test_util_genai.conformance import Scenario +from opentelemetry.test_util_genai.instrumentor import instrument + + +class InferenceBetaScenario(Scenario): + expected_spans = {"chat": 1} + expected_metrics = ( + "gen_ai.client.operation.duration", + "gen_ai.client.token.usage", + ) + + def run( + self, + *, + tracer_provider: TracerProvider, + meter_provider: MeterProvider, + logger_provider: LoggerProvider, + vcr: Any, + ) -> None: + key_override = ( + {} + if os.getenv("ANTHROPIC_API_KEY") + else {"ANTHROPIC_API_KEY": "test_anthropic_api_key"} + ) + with mock.patch.dict(os.environ, key_override): + with instrument( + AnthropicInstrumentor(), + tracer_provider=tracer_provider, + logger_provider=logger_provider, + meter_provider=meter_provider, + content_capture="SPAN_ONLY", + ): + with vcr.use_cassette("inference_beta_conformance.yaml"): + Anthropic().beta.messages.create( + model="claude-sonnet-4-6", + max_tokens=100, + messages=[ + { + "role": "user", + "content": "Say hello in one word.", + } + ], + ) + + +class InferenceBetaStreamingScenario(Scenario): + expected_spans = {"chat": 1} + expected_metrics = ( + "gen_ai.client.operation.duration", + "gen_ai.client.token.usage", + "gen_ai.client.operation.time_to_first_chunk", + "gen_ai.client.operation.time_per_output_chunk", + ) + + def validate(self, report: LiveCheckReport) -> None: + super().validate(report) + stream_values = [ + attr["value"] + for entry in report["samples"] + if "span" in entry + for attr in entry["span"]["attributes"] + if attr["name"] == "gen_ai.request.stream" + ] + assert stream_values == [True], ( + "streaming messages should set gen_ai.request.stream=true on the " + f"chat span; saw {stream_values}" + ) + + def run( + self, + *, + tracer_provider: TracerProvider, + meter_provider: MeterProvider, + logger_provider: LoggerProvider, + vcr: Any, + ) -> None: + key_override = ( + {} + if os.getenv("ANTHROPIC_API_KEY") + else {"ANTHROPIC_API_KEY": "test_anthropic_api_key"} + ) + with mock.patch.dict(os.environ, key_override): + with instrument( + AnthropicInstrumentor(), + tracer_provider=tracer_provider, + logger_provider=logger_provider, + meter_provider=meter_provider, + content_capture="SPAN_ONLY", + ): + with vcr.use_cassette( + "inference_beta_streaming_conformance.yaml" + ): + with Anthropic().beta.messages.create( + model="claude-sonnet-4-6", + max_tokens=100, + messages=[ + { + "role": "user", + "content": "Say hello in one word.", + } + ], + stream=True, + ) as stream: + for _ in stream: + pass diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/conformance/inference_beta_server_tool_calling.py b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/conformance/inference_beta_server_tool_calling.py new file mode 100644 index 000000000..89c9d2386 --- /dev/null +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/conformance/inference_beta_server_tool_calling.py @@ -0,0 +1,108 @@ +# Copyright The OpenTelemetry Authors +# SPDX-License-Identifier: Apache-2.0 + +"""Conformance scenario: Anthropic beta chat with server-side tool calls.""" + +from __future__ import annotations + +import json +import os +from typing import Any +from unittest import mock + +from anthropic import Anthropic + +from opentelemetry.instrumentation.genai.anthropic import AnthropicInstrumentor +from opentelemetry.sdk._logs import LoggerProvider +from opentelemetry.sdk.metrics import MeterProvider +from opentelemetry.sdk.trace import TracerProvider +from opentelemetry.test.weaver_live_check import LiveCheckReport +from opentelemetry.test_util_genai.conformance import Scenario +from opentelemetry.test_util_genai.instrumentor import instrument + + +class InferenceBetaServerToolCallingScenario(Scenario): + expected_spans = {"chat": 1} + expected_metrics = ( + "gen_ai.client.operation.duration", + "gen_ai.client.token.usage", + ) + + def validate(self, report: LiveCheckReport) -> None: + super().validate(report) + output_messages = [ + json.loads(attribute["value"]) + for entry in report["samples"] + if "span" in entry + for attribute in entry["span"]["attributes"] + if attribute["name"] == "gen_ai.output.messages" + ] + assert len(output_messages) == 1 + expected_parts = [ + { + "name": "web_search", + "server_tool_call": { + "type": "server_tool_use", + "arguments": {"query": "OpenTelemetry"}, + }, + "id": "srvtoolu_01", + "type": "server_tool_call", + }, + { + "server_tool_call_response": { + "type": "web_search_tool_result", + "content": { + "type": "web_search_tool_result_error", + "error_code": "unavailable", + }, + "caller": None, + }, + "id": "srvtoolu_01", + "type": "server_tool_call_response", + }, + ] + assert output_messages[0][0]["parts"][:2] == expected_parts, ( + output_messages[0][0]["parts"][:2] + ) + + def run( + self, + *, + tracer_provider: TracerProvider, + meter_provider: MeterProvider, + logger_provider: LoggerProvider, + vcr: Any, + ) -> None: + key_override = ( + {} + if os.getenv("ANTHROPIC_API_KEY") + else {"ANTHROPIC_API_KEY": "test_anthropic_api_key"} + ) + with mock.patch.dict(os.environ, key_override): + with instrument( + AnthropicInstrumentor(), + tracer_provider=tracer_provider, + logger_provider=logger_provider, + meter_provider=meter_provider, + content_capture="SPAN_ONLY", + ): + with vcr.use_cassette( + "inference_beta_server_tool_calling_conformance.yaml" + ): + Anthropic().beta.messages.create( + model="claude-sonnet-4-6", + max_tokens=256, + messages=[ + { + "role": "user", + "content": "Search for OpenTelemetry.", + } + ], + tools=[ + { + "type": "web_search_20250305", + "name": "web_search", + "max_uses": 1, + } + ], + ) diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_async_messages.py b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_async_messages.py index b33377053..8e030857a 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_async_messages.py +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_async_messages.py @@ -2093,3 +2093,631 @@ def __init__(self): assert first_stops == [True] assert second_stops == [True] assert http_response.close_calls == 1 + + +try: + from anthropic.resources.beta.messages import ( + AsyncMessages as _AsyncBetaMessages, + ) + + _beta_supported = hasattr(_AsyncBetaMessages, "create") +except (ImportError, AttributeError): + _beta_supported = False + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.asyncio +@pytest.mark.vcr() +async def test_async_beta_messages_create_basic( + span_exporter, async_anthropic_client, instrument_no_content +): + """Test basic async beta message creation produces correct span.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + response = await async_anthropic_client.beta.messages.create( + model=model, + max_tokens=100, + messages=messages, + ) + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + + span = spans[0] + assert span.name == f"chat {model}" + assert span.attributes[GenAIAttributes.GEN_AI_OPERATION_NAME] == "chat" + assert span.attributes[GenAIAttributes.GEN_AI_PROVIDER_NAME] == "anthropic" + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert span.attributes[GenAIAttributes.GEN_AI_RESPONSE_ID] == response.id + assert ( + span.attributes[GenAIAttributes.GEN_AI_RESPONSE_MODEL] + == response.model + ) + assert span.attributes[ + GenAIAttributes.GEN_AI_USAGE_INPUT_TOKENS + ] == expected_input_tokens(response.usage) + assert ( + span.attributes[GenAIAttributes.GEN_AI_USAGE_OUTPUT_TOKENS] + == response.usage.output_tokens + ) + assert span.attributes[GenAIAttributes.GEN_AI_RESPONSE_FINISH_REASONS] == ( + normalize_stop_reason(response.stop_reason), + ) + assert ( + span.attributes[ServerAttributes.SERVER_ADDRESS] == "api.anthropic.com" + ) + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.asyncio +@pytest.mark.vcr() +async def test_async_beta_messages_create_with_content( + span_exporter, async_anthropic_client, instrument_with_content +): + """Test async beta message creation captures input and output content.""" + model = "claude-sonnet-4-6" + prompt = "Say hello in one word." + messages = [{"role": "user", "content": prompt}] + + response = await async_anthropic_client.beta.messages.create( + model=model, + max_tokens=100, + messages=messages, + ) + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + + assert span.name == f"chat {model}" + assert span.attributes[GenAIAttributes.GEN_AI_OPERATION_NAME] == "chat" + assert span.attributes[GenAIAttributes.GEN_AI_PROVIDER_NAME] == "anthropic" + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert span.attributes[GenAIAttributes.GEN_AI_RESPONSE_ID] == response.id + assert ( + span.attributes[GenAIAttributes.GEN_AI_RESPONSE_MODEL] + == response.model + ) + assert span.attributes[ + GenAIAttributes.GEN_AI_USAGE_INPUT_TOKENS + ] == expected_input_tokens(response.usage) + assert ( + span.attributes[GenAIAttributes.GEN_AI_USAGE_OUTPUT_TOKENS] + == response.usage.output_tokens + ) + assert span.attributes[GenAIAttributes.GEN_AI_RESPONSE_FINISH_REASONS] == ( + normalize_stop_reason(response.stop_reason), + ) + assert ( + span.attributes[ServerAttributes.SERVER_ADDRESS] == "api.anthropic.com" + ) + + input_messages = _load_span_messages( + span, GenAIAttributes.GEN_AI_INPUT_MESSAGES + ) + assert input_messages[0]["role"] == "user" + assert input_messages[0]["parts"][0]["type"] == "text" + assert input_messages[0]["parts"][0]["content"] == prompt + + output_messages = _load_span_messages( + span, GenAIAttributes.GEN_AI_OUTPUT_MESSAGES + ) + assert len(output_messages) == 1 + assert output_messages[0]["role"] == "assistant" + assert output_messages[0]["parts"] == [ + {"type": "text", "content": response.content[0].text} + ] + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.asyncio +@pytest.mark.vcr() +async def test_async_beta_messages_create_api_error( + span_exporter, async_anthropic_client, instrument_no_content +): + """Test that API errors in async beta message creation are recorded.""" + model = "invalid-model-name" + messages = [{"role": "user", "content": "Hello"}] + + with pytest.raises(NotFoundError): + await async_anthropic_client.beta.messages.create( + model=model, + max_tokens=100, + messages=messages, + ) + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + + span = spans[0] + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert ErrorAttributes.ERROR_TYPE in span.attributes + assert "NotFoundError" in span.attributes[ErrorAttributes.ERROR_TYPE] + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.asyncio +@pytest.mark.vcr() +async def test_async_beta_messages_create_streaming( + span_exporter, async_anthropic_client, instrument_no_content +): + """Test streaming async beta message creation produces correct span.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + stream = await async_anthropic_client.beta.messages.create( + model=model, + max_tokens=100, + messages=messages, + stream=True, + ) + assert isinstance(stream, AsyncMessagesStreamWrapper) + + async with stream: + async for _ in stream: + pass + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + + span = spans[0] + assert span.name == f"chat {model}" + assert span.attributes[GenAIAttributes.GEN_AI_OPERATION_NAME] == "chat" + assert span.attributes[GenAIAttributes.GEN_AI_PROVIDER_NAME] == "anthropic" + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert ( + span.attributes[GenAIAttributes.GEN_AI_RESPONSE_ID] + == "msg_01VY8H3oPFs4WWVnJXFDxNhU" + ) + assert span.attributes[GenAIAttributes.GEN_AI_RESPONSE_MODEL] == model + assert span.attributes[GenAIAttributes.GEN_AI_USAGE_INPUT_TOKENS] == 13 + assert span.attributes[GenAIAttributes.GEN_AI_USAGE_OUTPUT_TOKENS] == 5 + assert span.attributes[GenAIAttributes.GEN_AI_RESPONSE_FINISH_REASONS] == ( + "stop", + ) + assert ( + span.attributes[ServerAttributes.SERVER_ADDRESS] == "api.anthropic.com" + ) + + +# Anthropic 0.51 omitted ?beta=true on async stream requests, while newer versions +# include it; omit query from matching so the cassette matches both. +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.asyncio +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +async def test_async_beta_messages_stream( + span_exporter, async_anthropic_client, instrument_no_content +): + """Test streaming async beta message stream produces correct span.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + async with async_anthropic_client.beta.messages.stream( + model=model, + max_tokens=100, + messages=messages, + ) as stream: + async for _ in stream: + pass + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + + span = spans[0] + assert span.name == f"chat {model}" + assert span.attributes[GenAIAttributes.GEN_AI_OPERATION_NAME] == "chat" + assert span.attributes[GenAIAttributes.GEN_AI_PROVIDER_NAME] == "anthropic" + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert ( + span.attributes[GenAIAttributes.GEN_AI_RESPONSE_ID] + == "msg_01VY8H3oPFs4WWVnJXFDxNhU" + ) + assert span.attributes[GenAIAttributes.GEN_AI_RESPONSE_MODEL] == model + assert span.attributes[GenAIAttributes.GEN_AI_USAGE_INPUT_TOKENS] == 13 + assert span.attributes[GenAIAttributes.GEN_AI_USAGE_OUTPUT_TOKENS] == 5 + assert span.attributes[GenAIAttributes.GEN_AI_RESPONSE_FINISH_REASONS] == ( + "stop", + ) + assert ( + span.attributes[ServerAttributes.SERVER_ADDRESS] == "api.anthropic.com" + ) + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.asyncio +@pytest.mark.vcr() +async def test_async_beta_messages_with_raw_response( + span_exporter, async_anthropic_client, instrument_no_content +): + """Test that with_raw_response.create on async beta messages produces correct span.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + raw_response = ( + await async_anthropic_client.beta.messages.with_raw_response.create( + model=model, + max_tokens=100, + messages=messages, + ) + ) + response = await _parse_raw_response(raw_response) + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + + assert span.name == f"chat {model}" + assert span.attributes[GenAIAttributes.GEN_AI_OPERATION_NAME] == "chat" + assert span.attributes[GenAIAttributes.GEN_AI_PROVIDER_NAME] == "anthropic" + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert span.attributes[GenAIAttributes.GEN_AI_RESPONSE_ID] == response.id + assert ( + span.attributes[GenAIAttributes.GEN_AI_RESPONSE_MODEL] + == response.model + ) + assert span.attributes[ + GenAIAttributes.GEN_AI_USAGE_INPUT_TOKENS + ] == expected_input_tokens(response.usage) + assert ( + span.attributes[GenAIAttributes.GEN_AI_USAGE_OUTPUT_TOKENS] + == response.usage.output_tokens + ) + assert span.attributes[GenAIAttributes.GEN_AI_RESPONSE_FINISH_REASONS] == ( + normalize_stop_reason(response.stop_reason), + ) + assert ( + span.attributes[ServerAttributes.SERVER_ADDRESS] == "api.anthropic.com" + ) + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.asyncio +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_async_beta_messages_create_streaming") +async def test_async_beta_messages_with_streaming_response( + span_exporter, async_anthropic_client, instrument_no_content +): + """Test async beta message creation with streaming response produces correct span.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + async with ( + async_anthropic_client.beta.messages.with_streaming_response.create( + model=model, + max_tokens=100, + messages=messages, + stream=True, + ) as raw_response + ): + stream = await raw_response.parse() + async for _ in stream: + pass + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + assert span.name == f"chat {model}" + assert span.attributes[GenAIAttributes.GEN_AI_OPERATION_NAME] == "chat" + assert span.attributes[GenAIAttributes.GEN_AI_PROVIDER_NAME] == "anthropic" + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert ( + span.attributes[GenAIAttributes.GEN_AI_RESPONSE_ID] + == "msg_01VY8H3oPFs4WWVnJXFDxNhU" + ) + assert span.attributes[GenAIAttributes.GEN_AI_RESPONSE_MODEL] == model + assert span.attributes[GenAIAttributes.GEN_AI_USAGE_INPUT_TOKENS] == 13 + assert span.attributes[GenAIAttributes.GEN_AI_USAGE_OUTPUT_TOKENS] == 5 + assert span.attributes[GenAIAttributes.GEN_AI_RESPONSE_FINISH_REASONS] == ( + "stop", + ) + assert ( + span.attributes[ServerAttributes.SERVER_ADDRESS] == "api.anthropic.com" + ) + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.asyncio +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_async_beta_messages_stream") +async def test_async_beta_messages_stream_interrupted_mid_iteration( + monkeypatch, span_exporter, async_anthropic_client, instrument_no_content +): + """Test that mid-stream network failures in async beta stream propagate and record error.type.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + class ErrorInjectingStreamDelegate: + def __init__(self, inner): + self._inner = inner + self._count = 0 + + def __aiter__(self): + return self + + async def __anext__(self): + if self._count == 1: + raise ConnectionError("connection reset during stream") + self._count += 1 + return await self._inner.__anext__() + + async def close(self): + return await self._inner.close() + + def __getattr__(self, name): + return getattr(self._inner, name) + + async with async_anthropic_client.beta.messages.stream( + model=model, + max_tokens=100, + messages=messages, + ) as stream: + monkeypatch.setattr( + stream, "stream", ErrorInjectingStreamDelegate(stream.stream) + ) + with pytest.raises( + ConnectionError, match="connection reset during stream" + ): + async for _ in stream: + pass + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert span.attributes[ErrorAttributes.ERROR_TYPE] == "ConnectionError" + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.asyncio +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_async_beta_messages_stream") +async def test_async_beta_messages_stream_user_exception( + span_exporter, async_anthropic_client, instrument_no_content +): + """Test that user raised exceptions from async beta.messages.stream are propagated.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + with pytest.raises(ValueError, match="User raised exception"): + async with async_anthropic_client.beta.messages.stream( + model=model, + max_tokens=100, + messages=messages, + ) as stream: + async for _ in stream: + raise ValueError("User raised exception") + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert span.attributes[ErrorAttributes.ERROR_TYPE] == "ValueError" + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.asyncio +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_async_beta_messages_create_streaming") +async def test_async_beta_messages_create_streaming_with_content( + span_exporter, async_anthropic_client, instrument_with_content +): + """Test content capture on async beta create(stream=True).""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + stream = await async_anthropic_client.beta.messages.create( + model=model, + max_tokens=100, + messages=messages, + stream=True, + ) + + async with stream: + async for _ in stream: + pass + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + + input_messages = _load_span_messages( + span, GenAIAttributes.GEN_AI_INPUT_MESSAGES + ) + output_messages = _load_span_messages( + span, GenAIAttributes.GEN_AI_OUTPUT_MESSAGES + ) + assert input_messages[0]["role"] == "user" + assert output_messages[0]["role"] == "assistant" + assert output_messages[0]["parts"] == [ + {"type": "text", "content": "Hello"} + ] + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.asyncio +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_async_beta_messages_stream") +async def test_async_beta_messages_stream_with_content( + span_exporter, async_anthropic_client, instrument_with_content +): + """Test content capture on async beta.messages.stream().""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + async with async_anthropic_client.beta.messages.stream( + model=model, + max_tokens=100, + messages=messages, + ) as stream: + async for _ in stream: + pass + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + + input_messages = _load_span_messages( + span, GenAIAttributes.GEN_AI_INPUT_MESSAGES + ) + output_messages = _load_span_messages( + span, GenAIAttributes.GEN_AI_OUTPUT_MESSAGES + ) + assert input_messages[0]["role"] == "user" + assert output_messages[0]["role"] == "assistant" + assert output_messages[0]["parts"] == [ + {"type": "text", "content": "Hello"} + ] + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.asyncio +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_async_beta_messages_create_streaming") +async def test_async_beta_messages_create_streaming_interrupted_mid_iteration( + monkeypatch, span_exporter, async_anthropic_client, instrument_no_content +): + """Test that mid-stream network failures in async beta create(stream=True) propagate and record error.type.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + class ErrorInjectingStreamDelegate: + def __init__(self, inner): + self._inner = inner + self._count = 0 + + def __aiter__(self): + return self + + async def __anext__(self): + if self._count == 1: + raise ConnectionError("connection reset during stream") + self._count += 1 + return await self._inner.__anext__() + + async def close(self): + return await self._inner.close() + + def __getattr__(self, name): + return getattr(self._inner, name) + + stream = await async_anthropic_client.beta.messages.create( + model=model, + max_tokens=100, + messages=messages, + stream=True, + ) + monkeypatch.setattr( + stream, "stream", ErrorInjectingStreamDelegate(stream.stream) + ) + + with pytest.raises( + ConnectionError, match="connection reset during stream" + ): + async with stream: + async for _ in stream: + pass + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert span.attributes[ErrorAttributes.ERROR_TYPE] == "ConnectionError" + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.asyncio +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_async_beta_messages_create_streaming") +async def test_async_beta_messages_create_streaming_user_exception( + span_exporter, async_anthropic_client, instrument_no_content +): + """Test that user raised exceptions in async beta create(stream=True) are propagated.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + stream = await async_anthropic_client.beta.messages.create( + model=model, + max_tokens=100, + messages=messages, + stream=True, + ) + + with pytest.raises(ValueError, match="User raised exception"): + async with stream: + async for _ in stream: + raise ValueError("User raised exception") + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert span.attributes[ErrorAttributes.ERROR_TYPE] == "ValueError" + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.asyncio +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_async_beta_messages_stream") +async def test_async_beta_messages_stream_closed_early_by_caller( + span_exporter, async_anthropic_client, instrument_no_content +): + """Caller-closing AsyncBetaMessages.stream early finalizes without error.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + async with async_anthropic_client.beta.messages.stream( + model=model, + max_tokens=100, + messages=messages, + ) as stream: + await anext(stream) + await stream.close() + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert ErrorAttributes.ERROR_TYPE not in span.attributes diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_async_wrappers.py b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_async_wrappers.py index c5fe17f16..3fa4cebef 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_async_wrappers.py +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_async_wrappers.py @@ -23,11 +23,12 @@ def _make_invocation(): ) -def _make_stream_wrapper(stream): +def _make_stream_wrapper(stream, *, is_beta=False): return MessagesStreamWrapper( stream=stream, invocation=_make_invocation(), capture_content=False, + is_beta=is_beta, ) @@ -574,6 +575,50 @@ def mock_accumulate(**kwargs): assert "json_bufs" not in captured_kwargs[0] +def test_beta_accumulate_failure_disables_beta_and_falls_back_to_standard( + monkeypatch, +): + from opentelemetry.instrumentation.genai.anthropic import wrappers + + monkeypatch.setattr(wrappers, "_accumulation_disabled", False) + monkeypatch.setattr(wrappers, "_beta_accumulate_takes_json_bufs", True) + monkeypatch.setattr(wrappers, "_beta_accumulate_accepts_headers", True) + + beta_calls = [] + standard_calls = [] + + def mock_beta_accumulate(**kwargs): + beta_calls.append(kwargs) + raise TypeError("Unexpected beta event") + + def mock_accumulate(**kwargs): + standard_calls.append(kwargs) + return SimpleNamespace( + model="claude-test", + id="msg_1", + usage=None, + content=[], + stop_reason=None, + ) + + monkeypatch.setattr( + wrappers, "beta_accumulate_event", mock_beta_accumulate + ) + monkeypatch.setattr(wrappers, "accumulate_event", mock_accumulate) + + wrapper = _make_stream_wrapper( + _FakeSyncStream(events=["chunk1", "chunk2"]), is_beta=True + ) + + list(wrapper) + + assert len(beta_calls) == 1 + assert beta_calls[0]["json_bufs"] is wrapper._self_json_bufs + assert "request_headers" in beta_calls[0] + assert len(standard_calls) == 2 + assert wrapper._self_beta_accumulation_disabled + + @pytest.mark.asyncio async def test_async_stream_wrapper_accumulate_event_failure_logs_warning_and_disables( monkeypatch, caplog diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_conformance.py b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_conformance.py index 3b2d33954..4987525f4 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_conformance.py +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_conformance.py @@ -19,6 +19,13 @@ ) from .conformance.inference import InferenceScenario +from .conformance.inference_beta import ( + InferenceBetaScenario, + InferenceBetaStreamingScenario, +) +from .conformance.inference_beta_server_tool_calling import ( + InferenceBetaServerToolCallingScenario, +) from .conformance.inference_raw_response import ( InferenceRawResponseScenario, InferenceRawResponseStreamingScenario, @@ -35,6 +42,9 @@ InferenceRawResponseScenario(), InferenceRawResponseStreamingScenario(), ToolCallingScenario(), + InferenceBetaScenario(), + InferenceBetaStreamingScenario(), + InferenceBetaServerToolCallingScenario(), ], ids=lambda s: type(s).__name__, ) diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_instrumentor.py b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_instrumentor.py index c342ba873..6b0413c0a 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_instrumentor.py +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_instrumentor.py @@ -7,13 +7,34 @@ from typing import Any import pytest +import wrapt from anthropic import Anthropic, AsyncAnthropic from anthropic.resources.messages import AsyncMessages, Messages from anthropic.types import Message, TextBlock, Usage +try: + from anthropic.resources.beta.messages import ( + AsyncMessages as AsyncBetaMessages, + ) + from anthropic.resources.beta.messages import ( + Messages as BetaMessages, + ) + from anthropic.types.beta import BetaMessage, BetaTextBlock, BetaUsage + + _beta_supported = hasattr(BetaMessages, "create") and hasattr( + AsyncBetaMessages, "create" + ) + _beta_parse_supported = hasattr(BetaMessages, "parse") and hasattr( + AsyncBetaMessages, "parse" + ) +except (ImportError, AttributeError): + _beta_supported = False + _beta_parse_supported = False + from opentelemetry.instrumentation.genai.anthropic import AnthropicInstrumentor from opentelemetry.instrumentation.genai.anthropic.wrappers import ( AsyncMessagesStreamManagerWrapper, + MessagesStreamManagerWrapper, ) from opentelemetry.semconv._incubating.attributes import ( gen_ai_attributes as GenAIAttributes, @@ -34,6 +55,19 @@ def _fake_message(model: str) -> Message: ) +def _fake_beta_message(model: str) -> BetaMessage: + return BetaMessage( + id="msg_test", + content=[BetaTextBlock(text="hello", type="text")], + model=model, + role="assistant", + stop_reason="end_turn", + stop_sequence=None, + type="message", + usage=BetaUsage(input_tokens=1, output_tokens=2), + ) + + def _assert_chat_span(span_exporter, model: str) -> None: spans = span_exporter.get_finished_spans() assert len(spans) == 1 @@ -236,13 +270,202 @@ def fake_sync_parse(self: object, **kwargs: Any) -> Message: _assert_chat_span(span_exporter, model) +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +def test_beta_messages_create_is_instrumented( + monkeypatch, + tracer_provider, + logger_provider, + meter_provider, + span_exporter, +): + """BetaMessages.create should emit chat telemetry.""" + model = "claude-test" + + def fake_create(self: object, **kwargs: Any) -> BetaMessage: + return _fake_beta_message(kwargs["model"]) + + monkeypatch.setattr(BetaMessages, "create", fake_create) + + with instrument( + AnthropicInstrumentor(), + tracer_provider=tracer_provider, + logger_provider=logger_provider, + meter_provider=meter_provider, + ): + response = Anthropic().beta.messages.create( + model=model, + max_tokens=10, + messages=[{"role": "user", "content": "hello"}], + ) + + assert response.id == "msg_test" + _assert_chat_span(span_exporter, model) + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.asyncio +async def test_async_beta_messages_create_is_instrumented( + monkeypatch, + tracer_provider, + logger_provider, + meter_provider, + span_exporter, +): + """AsyncBetaMessages.create should emit chat telemetry.""" + model = "claude-test" + + async def fake_create(self: object, **kwargs: Any) -> BetaMessage: + return _fake_beta_message(kwargs["model"]) + + monkeypatch.setattr(AsyncBetaMessages, "create", fake_create) + + async with AsyncAnthropic() as client: + with instrument( + AnthropicInstrumentor(), + tracer_provider=tracer_provider, + logger_provider=logger_provider, + meter_provider=meter_provider, + ): + response = await client.beta.messages.create( + model=model, + max_tokens=10, + messages=[{"role": "user", "content": "hello"}], + ) + + assert response.id == "msg_test" + _assert_chat_span(span_exporter, model) + + +def test_beta_messages_unsupported_skips_instrumentation( + monkeypatch, + tracer_provider, + logger_provider, + meter_provider, +): + """When beta.messages is not supported, instrumentation should not fail and not wrap beta methods.""" + monkeypatch.setattr( + "opentelemetry.instrumentation.genai.anthropic._is_beta_messages_supported", + lambda: False, + ) + monkeypatch.setattr( + "opentelemetry.instrumentation.genai.anthropic._is_beta_parse_supported", + lambda: False, + ) + + with instrument( + AnthropicInstrumentor(), + tracer_provider=tracer_provider, + logger_provider=logger_provider, + meter_provider=meter_provider, + ): + if _beta_supported: + assert not isinstance( + BetaMessages.create, wrapt.BoundFunctionWrapper + ) + assert not isinstance( + AsyncBetaMessages.create, wrapt.BoundFunctionWrapper + ) + assert not isinstance( + BetaMessages.stream, wrapt.BoundFunctionWrapper + ) + assert not isinstance( + AsyncBetaMessages.stream, wrapt.BoundFunctionWrapper + ) + if _beta_parse_supported: + assert not isinstance( + BetaMessages.parse, wrapt.BoundFunctionWrapper + ) + assert not isinstance( + AsyncBetaMessages.parse, wrapt.BoundFunctionWrapper + ) + + +@pytest.mark.skipif( + not _beta_parse_supported, + reason="anthropic SDK does not support beta.messages.parse", +) +def test_beta_messages_parse_is_instrumented( + monkeypatch, + tracer_provider, + logger_provider, + meter_provider, + span_exporter, +): + """BetaMessages.parse should emit chat telemetry.""" + model = "claude-test" + + def fake_parse(self: object, **kwargs: Any) -> BetaMessage: + return _fake_beta_message(kwargs["model"]) + + monkeypatch.setattr(BetaMessages, "parse", fake_parse, raising=False) + + with instrument( + AnthropicInstrumentor(), + tracer_provider=tracer_provider, + logger_provider=logger_provider, + meter_provider=meter_provider, + ): + response = Anthropic().beta.messages.parse( + model=model, + max_tokens=10, + messages=[{"role": "user", "content": "hello"}], + output_format={"type": "json_schema", "schema": {}}, + ) + + assert response.id == "msg_test" + _assert_chat_span(span_exporter, model) + + +class _FakeResponse: + def close(self) -> None: + return None + + +class _FakeStream: + def __init__(self, message: Message | BetaMessage): + self.current_message_snapshot = message + self.response = _FakeResponse() + self._chunks = iter([object()]) + + def __iter__(self) -> "_FakeStream": + return self + + def __next__(self) -> object: + return next(self._chunks) + + def close(self) -> None: + return None + + +class _FakeStreamManager: + def __init__(self, message: Message | BetaMessage): + self._message = message + + def __enter__(self) -> _FakeStream: + return _FakeStream(self._message) + + def __exit__( + self, + exc_type: type[BaseException] | None, + exc_val: BaseException | None, + exc_tb: TracebackType | None, + ) -> bool: + return False + + class _FakeAsyncResponse: async def aclose(self) -> None: return None class _FakeAsyncStream: - def __init__(self, message: Message): + def __init__(self, message: Message | BetaMessage): self.current_message_snapshot = message self.response = _FakeAsyncResponse() self._chunks = iter([object()]) @@ -261,7 +484,7 @@ async def close(self) -> None: class _FakeAsyncStreamManager: - def __init__(self, message: Message): + def __init__(self, message: Message | BetaMessage): self._message = message async def __aenter__(self) -> _FakeAsyncStream: @@ -310,3 +533,158 @@ def fake_stream(self: object, **kwargs: Any) -> _FakeAsyncStreamManager: pass _assert_chat_span(span_exporter, model) + + +@pytest.mark.skipif( + not _beta_parse_supported, + reason="anthropic SDK does not support beta.messages.parse", +) +@pytest.mark.asyncio +async def test_async_beta_messages_parse_is_instrumented( + monkeypatch, + tracer_provider, + logger_provider, + meter_provider, + span_exporter, +): + """AsyncBetaMessages.parse should emit chat telemetry.""" + model = "claude-test" + + async def fake_parse(self: object, **kwargs: Any) -> BetaMessage: + return _fake_beta_message(kwargs["model"]) + + monkeypatch.setattr(AsyncBetaMessages, "parse", fake_parse, raising=False) + + async with AsyncAnthropic() as client: + with instrument( + AnthropicInstrumentor(), + tracer_provider=tracer_provider, + logger_provider=logger_provider, + meter_provider=meter_provider, + ): + response = await client.beta.messages.parse( + model=model, + max_tokens=10, + messages=[{"role": "user", "content": "hello"}], + output_format={"type": "json_schema", "schema": {}}, + ) + + assert response.id == "msg_test" + _assert_chat_span(span_exporter, model) + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.asyncio +async def test_async_beta_messages_stream_is_instrumented( + monkeypatch, + tracer_provider, + logger_provider, + meter_provider, + span_exporter, +): + """AsyncBetaMessages.stream should emit telemetry when the stream is consumed.""" + model = "claude-test" + + def fake_stream(self: object, **kwargs: Any) -> _FakeAsyncStreamManager: + return _FakeAsyncStreamManager(_fake_beta_message(kwargs["model"])) + + monkeypatch.setattr( + AsyncBetaMessages, "stream", fake_stream, raising=False + ) + + async with AsyncAnthropic() as client: + with instrument( + AnthropicInstrumentor(), + tracer_provider=tracer_provider, + logger_provider=logger_provider, + meter_provider=meter_provider, + ): + manager = client.beta.messages.stream( + model=model, + max_tokens=10, + messages=[{"role": "user", "content": "hello"}], + ) + assert isinstance(manager, AsyncMessagesStreamManagerWrapper) + async with manager as stream: + async for _ in stream: + pass + + _assert_chat_span(span_exporter, model) + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +def test_beta_messages_stream_is_instrumented( + monkeypatch, + tracer_provider, + logger_provider, + meter_provider, + span_exporter, +): + """BetaMessages.stream should emit telemetry when the stream is consumed.""" + model = "claude-test" + + def fake_stream(self: object, **kwargs: Any) -> _FakeStreamManager: + return _FakeStreamManager(_fake_beta_message(kwargs["model"])) + + monkeypatch.setattr(BetaMessages, "stream", fake_stream, raising=False) + + with instrument( + AnthropicInstrumentor(), + tracer_provider=tracer_provider, + logger_provider=logger_provider, + meter_provider=meter_provider, + ): + manager = Anthropic().beta.messages.stream( + model=model, + max_tokens=10, + messages=[{"role": "user", "content": "hello"}], + ) + assert isinstance(manager, MessagesStreamManagerWrapper) + with manager as stream: + for _ in stream: + pass + + _assert_chat_span(span_exporter, model) + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +def test_beta_messages_instrument_uninstrument( + tracer_provider, logger_provider, meter_provider +): + """Beta methods should be wrapped on instrument and unwrapped on uninstrument.""" + instrumentor = AnthropicInstrumentor() + instrumentor.instrument( + tracer_provider=tracer_provider, + logger_provider=logger_provider, + meter_provider=meter_provider, + ) + + assert isinstance(BetaMessages.create, wrapt.BoundFunctionWrapper) + assert isinstance(AsyncBetaMessages.create, wrapt.BoundFunctionWrapper) + assert isinstance(BetaMessages.stream, wrapt.BoundFunctionWrapper) + assert isinstance(AsyncBetaMessages.stream, wrapt.BoundFunctionWrapper) + if _beta_parse_supported: + assert isinstance(BetaMessages.parse, wrapt.BoundFunctionWrapper) + assert isinstance(AsyncBetaMessages.parse, wrapt.BoundFunctionWrapper) + + # Uninstrument using a separate instance to ensure instance independence + AnthropicInstrumentor().uninstrument() + + assert not isinstance(BetaMessages.create, wrapt.BoundFunctionWrapper) + assert not isinstance(AsyncBetaMessages.create, wrapt.BoundFunctionWrapper) + assert not isinstance(BetaMessages.stream, wrapt.BoundFunctionWrapper) + assert not isinstance(AsyncBetaMessages.stream, wrapt.BoundFunctionWrapper) + if _beta_parse_supported: + assert not isinstance(BetaMessages.parse, wrapt.BoundFunctionWrapper) + assert not isinstance( + AsyncBetaMessages.parse, wrapt.BoundFunctionWrapper + ) diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_messages_extractors.py b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_messages_extractors.py index 541b17081..0a3065554 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_messages_extractors.py +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_messages_extractors.py @@ -3,13 +3,27 @@ """Tests for Anthropic message parameter extraction.""" +from dataclasses import asdict from types import SimpleNamespace from unittest.mock import MagicMock +import pytest + from opentelemetry.instrumentation.genai.anthropic.messages_extractors import ( extract_params, + get_output_messages_from_message, set_invocation_response_attributes, ) +from opentelemetry.instrumentation.genai.anthropic.utils import ( + _convert_content_block_to_part, +) +from opentelemetry.util.genai.types import ( + CompactionPart, + FilePart, + GenericPart, + ServerToolCallPart, + ServerToolCallResponsePart, +) def test_extract_params_reads_sampling_params_from_extra_body(): @@ -73,3 +87,294 @@ def test_set_invocation_response_attributes_records_cache_tokens(): assert invocation.output_tokens == 20 assert invocation.cache_write_input_tokens == 15 assert invocation.cache_read_input_tokens == 5 + + +def test_convert_beta_mcp_tool_result_block_serializable(): + try: + from anthropic.types.beta import BetaMCPToolResultBlock, BetaTextBlock + + has_beta_mcp = hasattr( + BetaMCPToolResultBlock, "model_fields" + ) or hasattr(BetaMCPToolResultBlock, "__fields__") + except (ImportError, AttributeError): + has_beta_mcp = False + + if not has_beta_mcp: + pytest.skip("BetaMCPToolResultBlock not supported") + + block = BetaMCPToolResultBlock( + content=[BetaTextBlock(text="tool output", type="text")], + is_error=False, + tool_use_id="tool_123", + type="mcp_tool_result", + ) + part = _convert_content_block_to_part(block) + assert isinstance(part, ServerToolCallResponsePart) + assert part.id == "tool_123" + assert part.server_tool_call_response == { + "content": [ + {"citations": None, "text": "tool output", "type": "text"} + ], + "is_error": False, + "type": "mcp_tool_result", + } + + +def test_convert_dict_mcp_tool_use_and_result(): + part_use = _convert_content_block_to_part( + { + "type": "mcp_tool_use", + "id": "call_1", + "name": "read", + "input": {"path": "a.txt"}, + } + ) + assert isinstance(part_use, ServerToolCallPart) + assert part_use.id == "call_1" + assert part_use.name == "read" + assert part_use.server_tool_call == { + "type": "mcp_tool_use", + "arguments": {"path": "a.txt"}, + } + + server_use = _convert_content_block_to_part( + { + "type": "server_tool_use", + "id": "call_srv", + "name": "web_search", + "input": {"query": "otel"}, + } + ) + assert isinstance(server_use, ServerToolCallPart) + assert server_use.id == "call_srv" + assert server_use.name == "web_search" + assert server_use.server_tool_call == { + "type": "server_tool_use", + "arguments": {"query": "otel"}, + } + + part_res = _convert_content_block_to_part( + { + "type": "mcp_tool_result", + "tool_use_id": "call_1", + "content": "file contents", + } + ) + assert isinstance(part_res, ServerToolCallResponsePart) + assert part_res.id == "call_1" + assert part_res.server_tool_call_response == { + "type": "mcp_tool_result", + "content": "file contents", + } + + search_res = _convert_content_block_to_part( + { + "type": "web_search_tool_result", + "tool_use_id": "call_srv", + "content": "search results", + } + ) + assert isinstance(search_res, ServerToolCallResponsePart) + assert search_res.id == "call_srv" + assert search_res.server_tool_call_response == { + "type": "web_search_tool_result", + "content": "search results", + } + + +def test_convert_beta_blocks_when_available(): + try: + from anthropic.types.beta import ( + BetaRedactedThinkingBlock, + BetaServerToolUseBlock, + BetaTextBlock, + BetaThinkingBlock, + BetaToolUseBlock, + ) + except (ImportError, AttributeError): + pytest.skip("Beta block types not available") + + # Text block + text_part = _convert_content_block_to_part( + BetaTextBlock(text="hello beta", type="text") + ) + assert text_part is not None + assert text_part.content == "hello beta" + + # Tool use block + tool_part = _convert_content_block_to_part( + BetaToolUseBlock( + id="t_1", name="search", input={"q": "otel"}, type="tool_use" + ) + ) + assert tool_part is not None + assert tool_part.id == "t_1" + assert tool_part.name == "search" + assert tool_part.arguments == {"q": "otel"} + + # Server tool use block + server_tool_part = _convert_content_block_to_part( + BetaServerToolUseBlock( + id="st_1", + name="web_search", + input={"q": "otel"}, + type="server_tool_use", + ) + ) + assert isinstance(server_tool_part, ServerToolCallPart) + assert server_tool_part.id == "st_1" + assert server_tool_part.name == "web_search" + assert server_tool_part.server_tool_call == { + "type": "server_tool_use", + "arguments": {"q": "otel"}, + } + + # Thinking block + thinking_part = _convert_content_block_to_part( + BetaThinkingBlock( + thinking="deep thought", signature="sig", type="thinking" + ) + ) + assert thinking_part is not None + assert thinking_part.content == "deep thought" + + # Redacted thinking block + redacted_part = _convert_content_block_to_part( + BetaRedactedThinkingBlock( + data="redacted_data", type="redacted_thinking" + ) + ) + assert redacted_part is not None + assert redacted_part.content == "redacted_data" + + try: + from anthropic.types.beta import ( + BetaMCPToolUseBlock, + BetaWebSearchResultBlock, + BetaWebSearchToolResultBlock, + ) + + mcp_use = _convert_content_block_to_part( + BetaMCPToolUseBlock( + id="mcp_1", + name="read_file", + input={"path": "x.py"}, + server_name="fs", + type="mcp_tool_use", + ) + ) + assert isinstance(mcp_use, ServerToolCallPart) + assert mcp_use.id == "mcp_1" + assert mcp_use.name == "read_file" + assert mcp_use.server_tool_call == { + "type": "mcp_tool_use", + "arguments": {"path": "x.py"}, + "server_name": "fs", + } + + search_block = _convert_content_block_to_part( + BetaWebSearchToolResultBlock( + content=[ + BetaWebSearchResultBlock( + type="web_search_result", + url="https://example.com", + title="Title", + encrypted_content="enc", + ) + ], + tool_use_id="ws_1", + type="web_search_tool_result", + ) + ) + assert isinstance(search_block, ServerToolCallResponsePart) + assert search_block.id == "ws_1" + assert isinstance( + search_block.server_tool_call_response["content"], list + ) + except (ImportError, AttributeError): + pass + + +@pytest.mark.parametrize( + ("block", "part_type"), + [ + ( + {"type": "container_upload", "file_id": "file_123"}, + FilePart, + ), + ( + { + "type": "compaction", + "content": "Summary of earlier turns.", + "encrypted_content": "opaque", + }, + CompactionPart, + ), + ( + { + "type": "code_execution_tool_result", + "tool_use_id": "server_123", + "content": {"stdout": "1"}, + }, + ServerToolCallResponsePart, + ), + ( + { + "type": "fallback", + "from": {"model": "model-a"}, + "to": {"model": "model-b"}, + }, + GenericPart, + ), + ], +) +def test_convert_additional_beta_blocks(block, part_type): + class _BetaBlock: + def model_dump(self): + return block + + assert isinstance(_convert_content_block_to_part(_BetaBlock()), part_type) + + +def test_beta_server_tool_parts_have_semconv_serialized_shape(): + message = SimpleNamespace( + role="assistant", + stop_reason="tool_use", + content=[ + { + "type": "mcp_tool_use", + "id": "call_1", + "name": "read", + "server_name": "files", + "input": {"path": "a.txt"}, + }, + { + "type": "mcp_tool_result", + "tool_use_id": "call_1", + "content": "file contents", + }, + ], + ) + + output = get_output_messages_from_message(message) + + assert [asdict(part) for part in output[0].parts] == [ + { + "name": "read", + "server_tool_call": { + "type": "mcp_tool_use", + "arguments": {"path": "a.txt"}, + "server_name": "files", + }, + "id": "call_1", + "type": "server_tool_call", + }, + { + "server_tool_call_response": { + "type": "mcp_tool_result", + "content": "file contents", + }, + "id": "call_1", + "type": "server_tool_call_response", + }, + ] diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_parse_messages.py b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_parse_messages.py index a3cc91bc3..3098ac4b0 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_parse_messages.py +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_parse_messages.py @@ -52,6 +52,9 @@ def _assert_parsed_response(response) -> None: parsed_blocks = [ getattr(block, "parsed", None) for block in response.content ] + parsed_blocks.extend( + getattr(block, "parsed_output", None) for block in response.content + ) text_blocks = [getattr(block, "text", None) for block in response.content] text_payloads = [ json.loads(text) @@ -300,3 +303,210 @@ async def test_async_messages_parse_api_error( assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model assert ErrorAttributes.ERROR_TYPE in span.attributes assert "NotFoundError" in span.attributes[ErrorAttributes.ERROR_TYPE] + + +try: + from anthropic.resources.beta.messages import Messages as _BetaMessages + + _beta_parse_supported = hasattr(_BetaMessages, "parse") +except (ImportError, AttributeError): + _beta_parse_supported = False + + +@pytest.mark.skipif( + not _beta_parse_supported, + reason="anthropic SDK does not support beta.messages.parse", +) +@pytest.mark.vcr() +def test_sync_beta_messages_parse_basic( + span_exporter, anthropic_client, instrument_no_content +): + """BetaMessages.parse should emit a chat span for structured output.""" + model = "claude-haiku-4-5" + + response = anthropic_client.beta.messages.parse( + model=model, + max_tokens=100, + messages=[ + { + "role": "user", + "content": "Return JSON with a greeting field set to hello.", + } + ], + output_format=Greeting, + ) + + _assert_parsed_response(response) + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + _assert_parse_span(spans[0], model=model, response=response) + + +@pytest.mark.skipif( + not _beta_parse_supported, + reason="anthropic SDK does not support beta.messages.parse", +) +@pytest.mark.asyncio +@pytest.mark.vcr() +async def test_async_beta_messages_parse_basic( + span_exporter, async_anthropic_client, instrument_no_content +): + """AsyncBetaMessages.parse should emit a chat span for structured output.""" + model = "claude-haiku-4-5" + + response = await async_anthropic_client.beta.messages.parse( + model=model, + max_tokens=100, + messages=[ + { + "role": "user", + "content": "Return JSON with a greeting field set to hello.", + } + ], + output_format=Greeting, + ) + + _assert_parsed_response(response) + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + _assert_parse_span(spans[0], model=model, response=response) + + +@pytest.mark.skipif( + not _beta_parse_supported, + reason="anthropic SDK does not support beta.messages.parse", +) +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_sync_beta_messages_create_api_error") +def test_sync_beta_messages_parse_api_error( + span_exporter, anthropic_client, instrument_no_content +): + """BetaMessages.parse should record API errors.""" + model = "invalid-model-name" + + with pytest.raises(NotFoundError): + anthropic_client.beta.messages.parse( + model=model, + max_tokens=100, + messages=[{"role": "user", "content": "Hello"}], + output_format=Greeting, + ) + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert ErrorAttributes.ERROR_TYPE in span.attributes + assert "NotFoundError" in span.attributes[ErrorAttributes.ERROR_TYPE] + + +@pytest.mark.skipif( + not _beta_parse_supported, + reason="anthropic SDK does not support beta.messages.parse", +) +@pytest.mark.asyncio +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_async_beta_messages_create_api_error") +async def test_async_beta_messages_parse_api_error( + span_exporter, async_anthropic_client, instrument_no_content +): + """AsyncBetaMessages.parse should record API errors.""" + model = "invalid-model-name" + + with pytest.raises(NotFoundError): + await async_anthropic_client.beta.messages.parse( + model=model, + max_tokens=100, + messages=[{"role": "user", "content": "Hello"}], + output_format=Greeting, + ) + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert ErrorAttributes.ERROR_TYPE in span.attributes + assert "NotFoundError" in span.attributes[ErrorAttributes.ERROR_TYPE] + + +@pytest.mark.skipif( + not _beta_parse_supported, + reason="anthropic SDK does not support beta.messages.parse", +) +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_sync_beta_messages_parse_basic") +def test_sync_beta_messages_parse_captures_content( + span_exporter, anthropic_client, instrument_with_content +): + """BetaMessages.parse should capture input and output messages.""" + model = "claude-haiku-4-5" + + anthropic_client.beta.messages.parse( + model=model, + max_tokens=100, + messages=[ + { + "role": "user", + "content": "Return JSON with a greeting field set to hello.", + } + ], + output_format=Greeting, + ) + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + + input_messages = _load_span_messages( + span, GenAIAttributes.GEN_AI_INPUT_MESSAGES + ) + output_messages = _load_span_messages( + span, GenAIAttributes.GEN_AI_OUTPUT_MESSAGES + ) + assert input_messages[0]["role"] == "user" + assert input_messages[0]["parts"][0]["type"] == "text" + assert output_messages[0]["role"] == "assistant" + assert output_messages[0]["parts"] + + +@pytest.mark.skipif( + not _beta_parse_supported, + reason="anthropic SDK does not support beta.messages.parse", +) +@pytest.mark.asyncio +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_async_beta_messages_parse_basic") +async def test_async_beta_messages_parse_captures_content( + span_exporter, async_anthropic_client, instrument_with_content +): + """AsyncBetaMessages.parse should capture input and output messages.""" + model = "claude-haiku-4-5" + + await async_anthropic_client.beta.messages.parse( + model=model, + max_tokens=100, + messages=[ + { + "role": "user", + "content": "Return JSON with a greeting field set to hello.", + } + ], + output_format=Greeting, + ) + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + + input_messages = _load_span_messages( + span, GenAIAttributes.GEN_AI_INPUT_MESSAGES + ) + output_messages = _load_span_messages( + span, GenAIAttributes.GEN_AI_OUTPUT_MESSAGES + ) + assert input_messages[0]["role"] == "user" + assert input_messages[0]["parts"][0]["type"] == "text" + assert output_messages[0]["role"] == "assistant" + assert output_messages[0]["parts"] diff --git a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_sync_messages.py b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_sync_messages.py index 4ba22eafd..12f3bba9e 100644 --- a/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_sync_messages.py +++ b/instrumentation/opentelemetry-instrumentation-genai-anthropic/tests/test_sync_messages.py @@ -2228,6 +2228,77 @@ def json(): assert _raw_response._message_from_read_body(_Body()) is None +try: + from anthropic.resources.beta.messages import Messages as _BetaMessages + from anthropic.types.beta import BetaMessage + + _beta_supported = hasattr(_BetaMessages, "create") +except (ImportError, AttributeError): + _beta_supported = False + BetaMessage = None + + +def test_raw_response_deserializes_beta_message_body(): + """A body from a beta endpoint deserializes into a BetaMessage.""" + + class _Request: + url = "https://api.anthropic.com/v1/messages?beta=true" + + class _Body: + request = _Request() + + @staticmethod + def json(): + return { + "id": "msg_beta_1", + "type": "message", + "role": "assistant", + "model": "claude-sonnet-4-20250514", + "content": [{"type": "text", "text": "hello"}], + "usage": { + "input_tokens": 10, + "output_tokens": 5, + }, + } + + message = _raw_response._message_from_read_body(_Body(), is_beta=True) + + assert message is not None + assert message.id == "msg_beta_1" + assert message.model == "claude-sonnet-4-20250514" + assert message.usage.input_tokens == 10 + assert message.usage.output_tokens == 5 + if _beta_supported and BetaMessage is not None: + assert isinstance(message, BetaMessage) + + +def test_raw_response_uses_beta_message_with_pydantic_v1(monkeypatch): + class _PydanticV1BetaMessage: + __fields__ = {} + + class _Body: + @staticmethod + def json(): + return {"type": "message"} + + target_types = [] + result = object() + + def construct_type(*, type_, value): + target_types.append(type_) + return result + + monkeypatch.setattr( + _raw_response, "AnthropicBetaMessage", _PydanticV1BetaMessage + ) + monkeypatch.setattr(_raw_response, "construct_type", construct_type) + + assert ( + _raw_response._message_from_read_body(_Body(), is_beta=True) is result + ) + assert target_types == [_PydanticV1BetaMessage] + + @pytest.mark.vcr() @pytest.mark.cassette("test_sync_messages_create_with_raw_response") def test_sync_messages_raw_response_only_parse_to_records_telemetry( @@ -2310,3 +2381,530 @@ def test_sync_messages_raw_response_parse_after_exit( assert spans[0].attributes[GenAIAttributes.GEN_AI_RESPONSE_MODEL] == model assert raw_response.parse().model == model + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.vcr() +def test_sync_beta_messages_create_basic( + span_exporter, anthropic_client, instrument_no_content +): + """Test basic sync beta message creation produces correct span.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + response = anthropic_client.beta.messages.create( + model=model, + max_tokens=100, + messages=messages, + ) + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + + assert_span_attributes( + spans[0], + request_model=model, + response_id=response.id, + response_model=response.model, + input_tokens=expected_input_tokens(response.usage), + output_tokens=response.usage.output_tokens, + finish_reasons=[normalize_stop_reason(response.stop_reason)], + ) + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.vcr() +def test_sync_beta_messages_create_with_content( + span_exporter, anthropic_client, instrument_with_content +): + """Test sync beta message creation captures input and output content.""" + model = "claude-sonnet-4-6" + prompt = "Say hello in one word." + messages = [{"role": "user", "content": prompt}] + + response = anthropic_client.beta.messages.create( + model=model, + max_tokens=100, + messages=messages, + ) + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + + assert_span_attributes( + span, + request_model=model, + response_id=response.id, + response_model=response.model, + input_tokens=expected_input_tokens(response.usage), + output_tokens=response.usage.output_tokens, + finish_reasons=[normalize_stop_reason(response.stop_reason)], + ) + + input_messages = _load_span_messages( + span, GenAIAttributes.GEN_AI_INPUT_MESSAGES + ) + assert input_messages[0]["role"] == "user" + assert input_messages[0]["parts"][0]["type"] == "text" + assert input_messages[0]["parts"][0]["content"] == prompt + + output_messages = _load_span_messages( + span, GenAIAttributes.GEN_AI_OUTPUT_MESSAGES + ) + assert len(output_messages) == 1 + assert output_messages[0]["role"] == "assistant" + assert output_messages[0]["parts"] == [ + {"type": "text", "content": response.content[0].text} + ] + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.vcr() +def test_sync_beta_messages_create_api_error( + span_exporter, anthropic_client, instrument_no_content +): + """Test that API errors in sync beta message creation are recorded.""" + model = "invalid-model-name" + messages = [{"role": "user", "content": "Hello"}] + + with pytest.raises(NotFoundError): + anthropic_client.beta.messages.create( + model=model, + max_tokens=100, + messages=messages, + ) + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + + span = spans[0] + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert ErrorAttributes.ERROR_TYPE in span.attributes + assert "NotFoundError" in span.attributes[ErrorAttributes.ERROR_TYPE] + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.vcr() +def test_sync_beta_messages_create_streaming( + span_exporter, anthropic_client, instrument_no_content +): + """Test streaming beta message creation produces correct span.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + with anthropic_client.beta.messages.create( + model=model, + max_tokens=100, + messages=messages, + stream=True, + ) as stream: + for _ in stream: + pass + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + + span = spans[0] + assert_span_attributes( + span, + request_model=model, + response_id="msg_01VY8H3oPFs4WWVnJXFDxNhU", + response_model=model, + input_tokens=13, + output_tokens=5, + finish_reasons=["stop"], + ) + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.vcr() +def test_sync_beta_messages_stream( + span_exporter, anthropic_client, instrument_no_content +): + """Test beta message stream produces correct span.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + with anthropic_client.beta.messages.stream( + model=model, + max_tokens=100, + messages=messages, + ) as stream: + for _ in stream: + pass + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + + span = spans[0] + assert_span_attributes( + span, + request_model=model, + response_id="msg_01VY8H3oPFs4WWVnJXFDxNhU", + response_model=model, + input_tokens=13, + output_tokens=5, + finish_reasons=["stop"], + ) + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.vcr() +def test_sync_beta_messages_with_raw_response( + span_exporter, anthropic_client, instrument_no_content +): + """Test that with_raw_response.create on beta messages produces correct span.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + raw_response = anthropic_client.beta.messages.with_raw_response.create( + model=model, + max_tokens=100, + messages=messages, + ) + response = raw_response.parse() + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + + assert_span_attributes( + spans[0], + request_model=model, + response_id=response.id, + response_model=response.model, + input_tokens=expected_input_tokens(response.usage), + output_tokens=response.usage.output_tokens, + finish_reasons=[normalize_stop_reason(response.stop_reason)], + ) + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_sync_beta_messages_create_streaming") +def test_sync_beta_messages_with_streaming_response( + span_exporter, anthropic_client, instrument_no_content +): + """Test sync beta message creation with streaming response produces correct span.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + with anthropic_client.beta.messages.with_streaming_response.create( + model=model, + max_tokens=100, + messages=messages, + stream=True, + ) as raw_response: + stream = raw_response.parse() + for _ in stream: + pass + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + assert_span_attributes( + spans[0], + request_model=model, + response_id="msg_01VY8H3oPFs4WWVnJXFDxNhU", + response_model=model, + input_tokens=13, + output_tokens=5, + finish_reasons=["stop"], + ) + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_sync_beta_messages_stream") +def test_sync_beta_messages_stream_interrupted_mid_iteration( + monkeypatch, span_exporter, anthropic_client, instrument_no_content +): + """Test that mid-stream network failures in beta stream propagate and record error.type.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + class ErrorInjectingStreamDelegate: + def __init__(self, inner): + self._inner = inner + self._count = 0 + + def __iter__(self): + return self + + def __next__(self): + if self._count == 1: + raise ConnectionError("connection reset during stream") + self._count += 1 + return next(self._inner) + + def close(self): + return self._inner.close() + + def __getattr__(self, name): + return getattr(self._inner, name) + + with pytest.raises( + ConnectionError, match="connection reset during stream" + ): + with anthropic_client.beta.messages.stream( + model=model, + max_tokens=100, + messages=messages, + ) as stream: + monkeypatch.setattr( + stream, + "stream", + ErrorInjectingStreamDelegate(stream.stream), + ) + for _ in stream: + pass + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert span.attributes[ErrorAttributes.ERROR_TYPE] == "ConnectionError" + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_sync_beta_messages_stream") +def test_sync_beta_messages_stream_user_exception( + span_exporter, anthropic_client, instrument_no_content +): + """Test that user raised exceptions from beta.messages.stream are propagated.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + with pytest.raises(ValueError, match="User raised exception"): + with anthropic_client.beta.messages.stream( + model=model, + max_tokens=100, + messages=messages, + ) as stream: + for _ in stream: + raise ValueError("User raised exception") + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert span.attributes[ErrorAttributes.ERROR_TYPE] == "ValueError" + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_sync_beta_messages_create_streaming") +def test_sync_beta_messages_create_streaming_with_content( + span_exporter, anthropic_client, instrument_with_content +): + """Test content capture on beta create(stream=True).""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + with anthropic_client.beta.messages.create( + model=model, + max_tokens=100, + messages=messages, + stream=True, + ) as stream: + for _ in stream: + pass + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + + input_messages = _load_span_messages( + span, GenAIAttributes.GEN_AI_INPUT_MESSAGES + ) + output_messages = _load_span_messages( + span, GenAIAttributes.GEN_AI_OUTPUT_MESSAGES + ) + assert input_messages[0]["role"] == "user" + assert output_messages[0]["role"] == "assistant" + assert output_messages[0]["parts"] == [ + {"type": "text", "content": "Hello"} + ] + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_sync_beta_messages_stream") +def test_sync_beta_messages_stream_with_content( + span_exporter, anthropic_client, instrument_with_content +): + """Test content capture on beta.messages.stream().""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + with anthropic_client.beta.messages.stream( + model=model, + max_tokens=100, + messages=messages, + ) as stream: + for _ in stream: + pass + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + + input_messages = _load_span_messages( + span, GenAIAttributes.GEN_AI_INPUT_MESSAGES + ) + output_messages = _load_span_messages( + span, GenAIAttributes.GEN_AI_OUTPUT_MESSAGES + ) + assert input_messages[0]["role"] == "user" + assert output_messages[0]["role"] == "assistant" + assert output_messages[0]["parts"] == [ + {"type": "text", "content": "Hello"} + ] + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_sync_beta_messages_create_streaming") +def test_sync_beta_messages_create_streaming_interrupted_mid_iteration( + monkeypatch, span_exporter, anthropic_client, instrument_no_content +): + """Test that mid-stream network failures in beta create(stream=True) propagate and record error.type.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + class ErrorInjectingStreamDelegate: + def __init__(self, inner): + self._inner = inner + self._count = 0 + + def __iter__(self): + return self + + def __next__(self): + if self._count == 1: + raise ConnectionError("connection reset during stream") + self._count += 1 + return next(self._inner) + + def close(self): + return self._inner.close() + + def __getattr__(self, name): + return getattr(self._inner, name) + + with pytest.raises( + ConnectionError, match="connection reset during stream" + ): + with anthropic_client.beta.messages.create( + model=model, + max_tokens=100, + messages=messages, + stream=True, + ) as stream: + monkeypatch.setattr( + stream, + "stream", + ErrorInjectingStreamDelegate(stream.stream), + ) + for _ in stream: + pass + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert span.attributes[ErrorAttributes.ERROR_TYPE] == "ConnectionError" + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_sync_beta_messages_create_streaming") +def test_sync_beta_messages_create_streaming_user_exception( + span_exporter, anthropic_client, instrument_no_content +): + """Test that user raised exceptions in beta create(stream=True) are propagated.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + with pytest.raises(ValueError, match="User raised exception"): + with anthropic_client.beta.messages.create( + model=model, + max_tokens=100, + messages=messages, + stream=True, + ) as stream: + for _ in stream: + raise ValueError("User raised exception") + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert span.attributes[ErrorAttributes.ERROR_TYPE] == "ValueError" + + +@pytest.mark.skipif( + not _beta_supported, + reason="anthropic SDK does not support beta.messages", +) +@pytest.mark.vcr(match_on=["method", "scheme", "host", "port", "path"]) +@pytest.mark.cassette("test_sync_beta_messages_stream") +def test_sync_beta_messages_stream_closed_early_by_caller( + span_exporter, anthropic_client, instrument_no_content +): + """Caller-closing beta.messages.stream early finalizes the span without error.""" + model = "claude-sonnet-4-6" + messages = [{"role": "user", "content": "Say hello in one word."}] + + with anthropic_client.beta.messages.stream( + model=model, + max_tokens=100, + messages=messages, + ) as stream: + next(stream) + stream.close() + + spans = span_exporter.get_finished_spans() + assert len(spans) == 1 + span = spans[0] + assert span.attributes[GenAIAttributes.GEN_AI_REQUEST_MODEL] == model + assert ErrorAttributes.ERROR_TYPE not in span.attributes