Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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
Expand Down Expand Up @@ -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(
Expand All @@ -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
Expand Down Expand Up @@ -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."""
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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:
Expand All @@ -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:
Expand All @@ -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:
Expand All @@ -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:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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(
Expand All @@ -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
)
Expand All @@ -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")
)
Expand Down Expand Up @@ -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:
Expand Down
Original file line number Diff line number Diff line change
@@ -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