From 107c587da0dd60304f756caef74512b9c13d0390 Mon Sep 17 00:00:00 2001 From: dharaneesh Date: Mon, 7 Sep 2026 10:52:32 +0530 Subject: [PATCH 1/3] feat(anthropic): emit gen_ai.usage.reasoning_tokens from extended thinking Anthropic's API reports thinking/reasoning usage under usage.output_tokens_details.reasoning_tokens for models with extended thinking enabled, but the instrumentation never read it, so reasoning tokens were invisible to cost/reasoning-ratio consumers while every other provider path (OpenAI chat + responses, Vertex AI) emits them. Add extraction to the sync and async non-streaming _set_token_usage paths, and to the streaming path (message_delta merge + attribute emit), matching the OpenAI package's SpanAttributes.GEN_AI_USAGE_REASONING_TOKENS. Tests use synthetic usage objects so no cassette/API is needed. Fixes #4458 --- .../instrumentation/anthropic/__init__.py | 20 +++ .../instrumentation/anthropic/streaming.py | 24 +++- .../tests/test_reasoning_usage.py | 135 ++++++++++++++++++ 3 files changed, 177 insertions(+), 2 deletions(-) create mode 100644 packages/opentelemetry-instrumentation-anthropic/tests/test_reasoning_usage.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..29e58588d4 100644 --- a/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/__init__.py +++ b/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/__init__.py @@ -187,6 +187,16 @@ def is_stream_manager(response): ) +def _get_reasoning_tokens_from_usage(usage): + """Extract thinking/reasoning tokens from an Anthropic SDK usage object.""" + if usage is None: + return None + output_tokens_details = getattr(usage, "output_tokens_details", None) + if output_tokens_details is None: + return None + return getattr(output_tokens_details, "reasoning_tokens", None) + + @dont_throw async def _aset_token_usage( span, @@ -305,6 +315,11 @@ async def _aset_token_usage( cache_creation_tokens, ) + reasoning_tokens = _get_reasoning_tokens_from_usage(usage) + set_span_attribute( + span, SpanAttributes.GEN_AI_USAGE_REASONING_TOKENS, reasoning_tokens + ) + @dont_throw def _set_token_usage( @@ -421,6 +436,11 @@ def _set_token_usage( cache_creation_tokens, ) + reasoning_tokens = _get_reasoning_tokens_from_usage(usage) + set_span_attribute( + span, SpanAttributes.GEN_AI_USAGE_REASONING_TOKENS, reasoning_tokens + ) + def _with_chat_telemetry_wrapper(func): """Helper for providing tracer for wrapper functions. Includes metric collectors.""" diff --git a/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/streaming.py b/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/streaming.py index bd2d6dc15d..a4cccefa46 100644 --- a/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/streaming.py +++ b/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/streaming.py @@ -61,16 +61,27 @@ def _process_response_item(item, complete_response): for event in complete_response.get("events", []): event["finish_reason"] = item.delta.stop_reason if item.usage: + item_usage = dict(item.usage) if "usage" in complete_response: - item_output_tokens = dict(item.usage).get("output_tokens", 0) + item_output_tokens = item_usage.get("output_tokens", 0) existing_output_tokens = complete_response["usage"].get( "output_tokens", 0 ) complete_response["usage"]["output_tokens"] = ( item_output_tokens + existing_output_tokens ) + item_reasoning_details = item_usage.get("output_tokens_details") or {} + existing_reasoning_details = ( + complete_response["usage"].get("output_tokens_details") or {} + ) + complete_response["usage"]["output_tokens_details"] = { + "reasoning_tokens": ( + (item_reasoning_details.get("reasoning_tokens") or 0) + + (existing_reasoning_details.get("reasoning_tokens") or 0) + ) + } else: - complete_response["usage"] = dict(item.usage) + complete_response["usage"] = item_usage def _set_token_usage( @@ -105,6 +116,15 @@ def _set_token_usage( span, GenAIAttributes.GEN_AI_USAGE_CACHE_CREATION_INPUT_TOKENS, cache_creation_tokens ) + output_tokens_details = complete_response.get("usage", {}).get("output_tokens_details") + reasoning_tokens = None + if output_tokens_details: + reasoning_tokens = output_tokens_details.get("reasoning_tokens", None) + + set_span_attribute( + span, SpanAttributes.GEN_AI_USAGE_REASONING_TOKENS, reasoning_tokens + ) + set_span_attribute( span, GenAIAttributes.GEN_AI_RESPONSE_MODEL, complete_response.get("model") ) diff --git a/packages/opentelemetry-instrumentation-anthropic/tests/test_reasoning_usage.py b/packages/opentelemetry-instrumentation-anthropic/tests/test_reasoning_usage.py new file mode 100644 index 0000000000..f1da08c895 --- /dev/null +++ b/packages/opentelemetry-instrumentation-anthropic/tests/test_reasoning_usage.py @@ -0,0 +1,135 @@ +"""Unit tests for reasoning-token usage extraction (gen_ai.usage.reasoning_tokens). + +Covers the three paths that consume Anthropic `usage.output_tokens_details`: +the sync and async non-streaming `_set_token_usage` helpers in `__init__.py`, +and the streaming path in `streaming.py`. +""" + +from types import SimpleNamespace + +import pytest + +from opentelemetry.instrumentation.anthropic import _aset_token_usage, _set_token_usage +from opentelemetry.instrumentation.anthropic.streaming import ( + _process_response_item, + _set_token_usage as _set_streaming_token_usage, +) +from opentelemetry.semconv._incubating.attributes import ( + gen_ai_attributes as GenAIAttributes, +) +from opentelemetry.semconv_ai import SpanAttributes + +REASONING_TOKENS = SpanAttributes.GEN_AI_USAGE_REASONING_TOKENS + + +def _make_usage(**overrides): + usage = { + "input_tokens": 10, + "output_tokens": 25, + "cache_creation_input_tokens": 0, + "cache_read_input_tokens": 0, + "output_tokens_details": SimpleNamespace(reasoning_tokens=15), + } + usage.update(overrides) + return SimpleNamespace(**usage) + + +def _make_response(usage): + return SimpleNamespace( + usage=usage, + content=[SimpleNamespace(type="text", text="hi")], + stop_reason="end_turn", + model="claude-3-7-sonnet-20250219", + ) + + +def _finished_span_attributes(span_exporter): + return dict(span_exporter.get_finished_spans()[0].attributes) + + +@pytest.fixture +def tracer(tracer_provider): + return tracer_provider.get_tracer("test-reasoning-usage") + + +def test_set_token_usage_emits_reasoning_tokens(tracer, span_exporter): + with tracer.start_as_current_span("test") as span: + _set_token_usage(span, None, {}, _make_response(_make_usage())) + + attributes = _finished_span_attributes(span_exporter) + assert attributes[REASONING_TOKENS] == 15 + assert attributes[GenAIAttributes.GEN_AI_USAGE_OUTPUT_TOKENS] == 25 + + +@pytest.mark.asyncio +async def test_aset_token_usage_emits_reasoning_tokens(tracer, span_exporter): + with tracer.start_as_current_span("test") as span: + await _aset_token_usage(span, None, {}, _make_response(_make_usage())) + + attributes = _finished_span_attributes(span_exporter) + assert attributes[REASONING_TOKENS] == 15 + + +def test_set_token_usage_omits_reasoning_when_details_absent(tracer, span_exporter): + with tracer.start_as_current_span("test") as span: + _set_token_usage( + span, None, {}, _make_response(_make_usage(output_tokens_details=None)) + ) + + attributes = _finished_span_attributes(span_exporter) + assert REASONING_TOKENS not in attributes + + +def test_process_response_item_merges_streaming_reasoning_tokens(): + complete_response = {"events": [], "model": "", "usage": {}, "id": ""} + + _process_response_item( + SimpleNamespace( + type="message_start", + message=SimpleNamespace( + model="claude-3-7-sonnet-20250219", + id="msg_1", + usage={ + "input_tokens": 10, + "output_tokens": 0, + "cache_creation_input_tokens": 0, + "cache_read_input_tokens": 0, + }, + ), + ), + complete_response, + ) + _process_response_item( + SimpleNamespace( + type="message_delta", + delta=SimpleNamespace(stop_reason="end_turn"), + usage={ + "output_tokens": 25, + "output_tokens_details": {"reasoning_tokens": 15}, + }, + ), + complete_response, + ) + + usage = complete_response["usage"] + assert usage["output_tokens"] == 25 + assert usage["output_tokens_details"]["reasoning_tokens"] == 15 + + +def test_streaming_set_token_usage_emits_reasoning_tokens(tracer, span_exporter): + complete_response = { + "events": [], + "model": "claude-3-7-sonnet-20250219", + "usage": { + "input_tokens": 10, + "output_tokens": 25, + "output_tokens_details": {"reasoning_tokens": 15}, + }, + "id": "msg_1", + } + + with tracer.start_as_current_span("test") as span: + _set_streaming_token_usage(span, complete_response, 10, 25) + + attributes = _finished_span_attributes(span_exporter) + assert attributes[REASONING_TOKENS] == 15 From 7c821c72c1c679c3a927e32f6d1d3492f3902727 Mon Sep 17 00:00:00 2001 From: dharaneesh Date: Mon, 7 Sep 2026 11:00:27 +0530 Subject: [PATCH 2/3] fix(anthropic): keep latest streamed reasoning-token count, not sum Anthropic reports message_delta.usage cumulatively, so summing reasoning tokens across chunks overcounts. Take the latest cumulative value instead and cover two message_delta events in the unit test. --- .../instrumentation/anthropic/streaming.py | 19 ++++++++++++------- .../tests/test_reasoning_usage.py | 18 ++++++++++++++---- 2 files changed, 26 insertions(+), 11 deletions(-) diff --git a/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/streaming.py b/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/streaming.py index a4cccefa46..93c6d1de51 100644 --- a/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/streaming.py +++ b/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/streaming.py @@ -70,15 +70,20 @@ def _process_response_item(item, complete_response): complete_response["usage"]["output_tokens"] = ( item_output_tokens + existing_output_tokens ) - item_reasoning_details = item_usage.get("output_tokens_details") or {} - existing_reasoning_details = ( - complete_response["usage"].get("output_tokens_details") or {} + item_reasoning = ( + (item_usage.get("output_tokens_details") or {}).get( + "reasoning_tokens" + ) + or 0 ) - complete_response["usage"]["output_tokens_details"] = { - "reasoning_tokens": ( - (item_reasoning_details.get("reasoning_tokens") or 0) - + (existing_reasoning_details.get("reasoning_tokens") or 0) + existing_reasoning = ( + (complete_response["usage"].get("output_tokens_details") or {}).get( + "reasoning_tokens" ) + or 0 + ) + complete_response["usage"]["output_tokens_details"] = { + "reasoning_tokens": max(item_reasoning, existing_reasoning) } else: complete_response["usage"] = item_usage diff --git a/packages/opentelemetry-instrumentation-anthropic/tests/test_reasoning_usage.py b/packages/opentelemetry-instrumentation-anthropic/tests/test_reasoning_usage.py index f1da08c895..60f5770dac 100644 --- a/packages/opentelemetry-instrumentation-anthropic/tests/test_reasoning_usage.py +++ b/packages/opentelemetry-instrumentation-anthropic/tests/test_reasoning_usage.py @@ -80,7 +80,7 @@ def test_set_token_usage_omits_reasoning_when_details_absent(tracer, span_export assert REASONING_TOKENS not in attributes -def test_process_response_item_merges_streaming_reasoning_tokens(): +def test_process_response_item_keeps_latest_streaming_reasoning_tokens(): complete_response = {"events": [], "model": "", "usage": {}, "id": ""} _process_response_item( @@ -104,16 +104,26 @@ def test_process_response_item_merges_streaming_reasoning_tokens(): type="message_delta", delta=SimpleNamespace(stop_reason="end_turn"), usage={ - "output_tokens": 25, + "output_tokens": 16, "output_tokens_details": {"reasoning_tokens": 15}, }, ), complete_response, ) + _process_response_item( + SimpleNamespace( + type="message_delta", + delta=SimpleNamespace(stop_reason="end_turn"), + usage={ + "output_tokens": 25, + "output_tokens_details": {"reasoning_tokens": 27}, + }, + ), + complete_response, + ) usage = complete_response["usage"] - assert usage["output_tokens"] == 25 - assert usage["output_tokens_details"]["reasoning_tokens"] == 15 + assert usage["output_tokens_details"]["reasoning_tokens"] == 27 def test_streaming_set_token_usage_emits_reasoning_tokens(tracer, span_exporter): From b1fb3a9a49e7ad8738a5a304a7db44d57535504d Mon Sep 17 00:00:00 2001 From: dharaneesh Date: Sat, 12 Sep 2026 07:31:04 +0530 Subject: [PATCH 3/3] docs(anthropic): raise docstring coverage on reasoning-usage paths --- .../instrumentation/anthropic/__init__.py | 20 +++++++++++++++++++ .../instrumentation/anthropic/streaming.py | 8 ++++++++ .../tests/test_reasoning_usage.py | 13 +++++++++--- 3 files changed, 38 insertions(+), 3 deletions(-) diff --git a/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/__init__.py b/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/__init__.py index 29e58588d4..275c897651 100644 --- a/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/__init__.py +++ b/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/__init__.py @@ -207,6 +207,14 @@ async def _aset_token_usage( token_histogram: Histogram = None, choice_counter: Counter = None, ): + """Record token-usage and choice metrics on *span* from an async Anthropic response. + + Handles coroutine responses, ``with_raw_response`` wrappers, and falls back + to client-side counting when the SDK does not supply usage data. + Sets input, output, total, cache-read, cache-creation, and reasoning-token + span attributes, and emits token histograms and choice counters when + configured. + """ import inspect # If we get a coroutine, await it @@ -331,6 +339,13 @@ def _set_token_usage( token_histogram: Histogram = None, choice_counter: Counter = None, ): + """Record token-usage and choice metrics on *span* from a synchronous Anthropic response. + + Handles ``with_raw_response`` wrappers and falls back to client-side + counting when the SDK does not supply usage data. Sets input, output, + total, cache-read, cache-creation, and reasoning-token span attributes, + and emits token histograms and choice counters when configured. + """ import inspect # If we get a coroutine, we cannot process it in sync context @@ -475,6 +490,7 @@ def wrapper(wrapped, instance, args, kwargs): def _create_metrics(meter: Meter): + """Create token, choice, duration, and exception metrics on *meter*.""" token_histogram = meter.create_histogram( name=Meters.LLM_TOKEN_USAGE, unit="token", @@ -504,6 +520,7 @@ def _create_metrics(meter: Meter): @dont_throw def _handle_input(span: Span, event_logger: Optional[Logger], kwargs): + """Populate *span* input attributes from *kwargs*, or emit input events if enabled.""" if should_emit_events() and event_logger: emit_input_events(event_logger, kwargs) else: @@ -514,6 +531,7 @@ def _handle_input(span: Span, event_logger: Optional[Logger], kwargs): @dont_throw async def _ahandle_input(span: Span, event_logger: Optional[Logger], kwargs): + """Populate *span* input attributes from *kwargs*, or emit input events if enabled.""" if should_emit_events() and event_logger: emit_input_events(event_logger, kwargs) else: @@ -524,6 +542,7 @@ async def _ahandle_input(span: Span, event_logger: Optional[Logger], kwargs): @dont_throw async def _ahandle_response(span: Span, event_logger: Optional[Logger], response): + """Populate *span* response attributes from *response*, or emit response events if enabled.""" if should_emit_events(): emit_response_events(event_logger, response) else: @@ -538,6 +557,7 @@ async def _ahandle_response(span: Span, event_logger: Optional[Logger], response @dont_throw def _handle_response(span: Span, event_logger: Optional[Logger], response): + """Populate *span* response attributes from *response*, or emit response events if enabled.""" if should_emit_events(): emit_response_events(event_logger, response) else: diff --git a/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/streaming.py b/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/streaming.py index 93c6d1de51..8c209a371d 100644 --- a/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/streaming.py +++ b/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/streaming.py @@ -32,6 +32,12 @@ @dont_throw def _process_response_item(item, complete_response): + """Accumulate a single Anthropic stream *item* into *complete_response*. + + Handles message_start, content_block_start, content_block_delta, and + message_delta events. Message-delta usage keeps the latest cumulative + reasoning-token count while summing output tokens. + """ if item.type == "message_start": complete_response["model"] = item.message.model complete_response["usage"] = dict(item.message.usage) @@ -98,6 +104,7 @@ def _set_token_usage( token_histogram: Histogram = None, choice_counter: Counter = None, ): + """Record token-usage, reasoning-token, model, and choice metrics from a completed streaming response.""" cache_read_tokens = ( complete_response.get("usage", {}).get("cache_read_input_tokens", 0) or 0 ) @@ -195,6 +202,7 @@ def _resolve_stream_token_usage(complete_response, instance, kwargs): def _handle_streaming_response(span, event_logger, complete_response): + """Emit streaming response events or set *span* streaming response attributes.""" if should_emit_events() and event_logger: emit_streaming_response_events(event_logger, complete_response) else: diff --git a/packages/opentelemetry-instrumentation-anthropic/tests/test_reasoning_usage.py b/packages/opentelemetry-instrumentation-anthropic/tests/test_reasoning_usage.py index 60f5770dac..666535e1a5 100644 --- a/packages/opentelemetry-instrumentation-anthropic/tests/test_reasoning_usage.py +++ b/packages/opentelemetry-instrumentation-anthropic/tests/test_reasoning_usage.py @@ -23,6 +23,7 @@ def _make_usage(**overrides): + """Build a stub Anthropic usage object, overriding any supplied fields.""" usage = { "input_tokens": 10, "output_tokens": 25, @@ -35,6 +36,7 @@ def _make_usage(**overrides): def _make_response(usage): + """Build a stub Anthropic message response object around *usage*.""" return SimpleNamespace( usage=usage, content=[SimpleNamespace(type="text", text="hi")], @@ -44,15 +46,18 @@ def _make_response(usage): def _finished_span_attributes(span_exporter): + """Return the attributes of the first finished span exported by *span_exporter*.""" return dict(span_exporter.get_finished_spans()[0].attributes) @pytest.fixture def tracer(tracer_provider): + """Provide a named tracer backed by the *tracer_provider* fixture.""" return tracer_provider.get_tracer("test-reasoning-usage") def test_set_token_usage_emits_reasoning_tokens(tracer, span_exporter): + """Sync non-streaming responses must emit gen_ai.usage.reasoning_tokens.""" with tracer.start_as_current_span("test") as span: _set_token_usage(span, None, {}, _make_response(_make_usage())) @@ -63,6 +68,7 @@ def test_set_token_usage_emits_reasoning_tokens(tracer, span_exporter): @pytest.mark.asyncio async def test_aset_token_usage_emits_reasoning_tokens(tracer, span_exporter): + """Async non-streaming responses must emit gen_ai.usage.reasoning_tokens.""" with tracer.start_as_current_span("test") as span: await _aset_token_usage(span, None, {}, _make_response(_make_usage())) @@ -71,16 +77,16 @@ async def test_aset_token_usage_emits_reasoning_tokens(tracer, span_exporter): def test_set_token_usage_omits_reasoning_when_details_absent(tracer, span_exporter): + """Reasoning tokens must be omitted when output_tokens_details is absent.""" with tracer.start_as_current_span("test") as span: - _set_token_usage( - span, None, {}, _make_response(_make_usage(output_tokens_details=None)) - ) + _set_token_usage(span, None, {}, _make_response(_make_usage(output_tokens_details=None))) attributes = _finished_span_attributes(span_exporter) assert REASONING_TOKENS not in attributes def test_process_response_item_keeps_latest_streaming_reasoning_tokens(): + """Streaming must retain the latest cumulative reasoning-token count, not a sum.""" complete_response = {"events": [], "model": "", "usage": {}, "id": ""} _process_response_item( @@ -127,6 +133,7 @@ def test_process_response_item_keeps_latest_streaming_reasoning_tokens(): def test_streaming_set_token_usage_emits_reasoning_tokens(tracer, span_exporter): + """Streaming span attributes must include the accumulated reasoning-token count.""" complete_response = { "events": [], "model": "claude-3-7-sonnet-20250219",