From 62af8a79b842c11d37596836e1d3d15f1299900c Mon Sep 17 00:00:00 2001 From: Anai-Guo Date: Fri, 31 Jul 2026 15:13:53 -0700 Subject: [PATCH] fix(anthropic): instrument beta stream managers (fixes #4388) client.beta.messages.stream(...) returns a Beta(Async)MessageStreamManager from anthropic.lib.streaming._beta_messages, but is_stream_manager only recognized the non-beta MessageStreamManager/AsyncMessageStreamManager, so beta streams were left uninstrumented. The async routing also compared the class name only against "AsyncMessageStreamManager", which would misroute a beta async stream to the sync wrapper. Recognize the beta manager classes in is_stream_manager and add is_async_stream_manager so both routing sites cover the beta async variant. --- .../instrumentation/anthropic/__init__.py | 45 +++++++++++++++---- .../tests/test_stream_manager_detection.py | 42 +++++++++++++++++ 2 files changed, 79 insertions(+), 8 deletions(-) create mode 100644 packages/opentelemetry-instrumentation-anthropic/tests/test_stream_manager_detection.py diff --git a/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/__init__.py b/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/__init__.py index 0054c99d59..150e5c7432 100644 --- a/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/__init__.py +++ b/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/__init__.py @@ -170,22 +170,51 @@ def is_streaming_response(response): return isinstance(response, Stream) or isinstance(response, AsyncStream) +ASYNC_STREAM_MANAGER_CLASS_NAMES = ( + "AsyncMessageStreamManager", + "BetaAsyncMessageStreamManager", +) + + +def is_async_stream_manager(response): + """Check if response is an (async) stream manager, including the beta variant""" + return response.__class__.__name__ in ASYNC_STREAM_MANAGER_CLASS_NAMES + + def is_stream_manager(response): - """Check if response is a MessageStreamManager or AsyncMessageStreamManager""" + """Check if response is a (Beta)MessageStreamManager or (Beta)AsyncMessageStreamManager""" + stream_manager_types = () try: from anthropic.lib.streaming._messages import ( MessageStreamManager, AsyncMessageStreamManager, ) - return isinstance(response, (MessageStreamManager, AsyncMessageStreamManager)) + stream_manager_types += (MessageStreamManager, AsyncMessageStreamManager) except ImportError: - # Check by class name as fallback - return ( - response.__class__.__name__ == "MessageStreamManager" - or response.__class__.__name__ == "AsyncMessageStreamManager" + pass + + try: + from anthropic.lib.streaming._beta_messages import ( + BetaMessageStreamManager, + BetaAsyncMessageStreamManager, ) + stream_manager_types += (BetaMessageStreamManager, BetaAsyncMessageStreamManager) + except ImportError: + pass + + if stream_manager_types: + return isinstance(response, stream_manager_types) + + # Check by class name as fallback + return response.__class__.__name__ in ( + "MessageStreamManager", + "AsyncMessageStreamManager", + "BetaMessageStreamManager", + "BetaAsyncMessageStreamManager", + ) + @dont_throw async def _aset_token_usage( @@ -599,7 +628,7 @@ def _wrap( kwargs, ) elif is_stream_manager(response): - if response.__class__.__name__ == "AsyncMessageStreamManager": + if is_async_stream_manager(response): return WrappedAsyncMessageStreamManager( response, span, @@ -729,7 +758,7 @@ async def _awrap( kwargs, ) elif is_stream_manager(response): - if response.__class__.__name__ == "AsyncMessageStreamManager": + if is_async_stream_manager(response): return WrappedAsyncMessageStreamManager( response, span, diff --git a/packages/opentelemetry-instrumentation-anthropic/tests/test_stream_manager_detection.py b/packages/opentelemetry-instrumentation-anthropic/tests/test_stream_manager_detection.py new file mode 100644 index 0000000000..6b728a7563 --- /dev/null +++ b/packages/opentelemetry-instrumentation-anthropic/tests/test_stream_manager_detection.py @@ -0,0 +1,42 @@ +from anthropic.lib.streaming._beta_messages import ( + BetaAsyncMessageStreamManager, + BetaMessageStreamManager, +) +from anthropic.lib.streaming._messages import ( + AsyncMessageStreamManager, + MessageStreamManager, +) + +from opentelemetry.instrumentation.anthropic import ( + is_async_stream_manager, + is_stream_manager, +) + + +def _make(cls): + # The manager classes take a client-bound request callable; we only need an + # instance for type/name detection, so bypass __init__. + return cls.__new__(cls) + + +def test_is_stream_manager_recognizes_all_managers(): + for cls in ( + MessageStreamManager, + AsyncMessageStreamManager, + BetaMessageStreamManager, + BetaAsyncMessageStreamManager, + ): + assert is_stream_manager(_make(cls)), cls.__name__ + + +def test_is_stream_manager_rejects_plain_objects(): + assert not is_stream_manager(object()) + + +def test_is_async_stream_manager_matches_async_variants_only(): + assert is_async_stream_manager(_make(AsyncMessageStreamManager)) + # Regression for #4388: beta async streams must route to the async wrapper. + assert is_async_stream_manager(_make(BetaAsyncMessageStreamManager)) + + assert not is_async_stream_manager(_make(MessageStreamManager)) + assert not is_async_stream_manager(_make(BetaMessageStreamManager))