diff --git a/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/__init__.py b/packages/opentelemetry-instrumentation-anthropic/opentelemetry/instrumentation/anthropic/__init__.py index 0054c99d59..275c897651 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, @@ -197,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 @@ -305,6 +323,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( @@ -316,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 @@ -421,6 +451,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.""" @@ -455,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", @@ -484,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: @@ -494,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: @@ -504,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: @@ -518,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 bd2d6dc15d..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) @@ -61,16 +67,32 @@ 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 = ( + (item_usage.get("output_tokens_details") or {}).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"] = dict(item.usage) + complete_response["usage"] = item_usage def _set_token_usage( @@ -82,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 ) @@ -105,6 +128,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") ) @@ -170,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 new file mode 100644 index 0000000000..666535e1a5 --- /dev/null +++ b/packages/opentelemetry-instrumentation-anthropic/tests/test_reasoning_usage.py @@ -0,0 +1,152 @@ +"""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): + """Build a stub Anthropic usage object, overriding any supplied fields.""" + 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): + """Build a stub Anthropic message response object around *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 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())) + + 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): + """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())) + + attributes = _finished_span_attributes(span_exporter) + assert attributes[REASONING_TOKENS] == 15 + + +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))) + + 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( + 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": 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_details"]["reasoning_tokens"] == 27 + + +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", + "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