Skip to content
27 changes: 22 additions & 5 deletions src/anthropic/lib/streaming/_beta_messages.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
)
from ..._streaming import Stream, AsyncStream
from ...types.beta import BetaRawMessageStreamEvent
from ...types.beta.beta_usage import BetaUsage
from ..._utils._utils import is_given
from .._parse._response import ResponseFormatT, parse_text
from ...types.beta.parsed_beta_message import ParsedBetaMessage, ParsedBetaContentBlock
Expand Down Expand Up @@ -550,16 +551,32 @@ def accumulate_event(
elif event.type == "message_delta":
current_snapshot.stop_reason = event.delta.stop_reason
current_snapshot.stop_sequence = event.delta.stop_sequence
current_snapshot.stop_details = event.delta.stop_details
if event.delta.stop_details is not None:
current_snapshot.stop_details = event.delta.stop_details
if event.delta.container is not None:
current_snapshot.container = event.delta.container
current_snapshot.usage.output_tokens = event.usage.output_tokens

# Usage may be absent when message_start omitted it (#1806); the
# message_delta carries the first full usage object, so construct it
# before updating.
if current_snapshot.usage is None:
_usage_data = event.usage.model_dump()
# `BetaUsage.input_tokens` is a required int. When the delta omits it
# (e.g. `{"output_tokens": 1}` per #1806), `construct` would leave it
# as None on a non-optional field and break downstream arithmetic.
# Coerce a missing input_tokens to 0 so the snapshot stays valid.
if _usage_data.get("input_tokens") is None:
_usage_data["input_tokens"] = 0
current_snapshot.usage = BetaUsage.construct(**_usage_data)
else:
current_snapshot.usage.output_tokens = event.usage.output_tokens

if event.context_management is not None:
current_snapshot.context_management = event.context_management

# Usage counts on a message_delta are cumulative totals, so they overwrite rather
# than add; optional ones are omitted when not applicable, in which case the
# message_start value must survive.
# Usage counts on a message_delta are cumulative totals, so they
# overwrite rather than add; optional ones are omitted when not
# applicable, in which case the message_start value must survive.
if event.usage.input_tokens is not None:
current_snapshot.usage.input_tokens = event.usage.input_tokens
if event.usage.cache_creation_input_tokens is not None:
Expand Down
47 changes: 31 additions & 16 deletions src/anthropic/lib/streaming/_messages.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
ParsedContentBlockStopEvent,
)
from ...types import RawMessageStreamEvent
from ...types.usage import Usage
from ..._types import NotGiven, not_given
from ..._utils import consume_sync_iterator, consume_async_iterator
from ..._models import build, construct_type, construct_type_unchecked
Expand Down Expand Up @@ -516,23 +517,37 @@ def accumulate_event(
elif event.type == "message_delta":
current_snapshot.stop_reason = event.delta.stop_reason
current_snapshot.stop_sequence = event.delta.stop_sequence
current_snapshot.stop_details = event.delta.stop_details
if event.delta.stop_details is not None:
current_snapshot.stop_details = event.delta.stop_details
if event.delta.container is not None:
current_snapshot.container = event.delta.container
current_snapshot.usage.output_tokens = event.usage.output_tokens

# Usage counts on a message_delta are cumulative totals, so they overwrite rather
# than add; optional ones are omitted when not applicable, in which case the
# message_start value must survive.
if event.usage.input_tokens is not None:
current_snapshot.usage.input_tokens = event.usage.input_tokens
if event.usage.cache_creation_input_tokens is not None:
current_snapshot.usage.cache_creation_input_tokens = event.usage.cache_creation_input_tokens
if event.usage.cache_read_input_tokens is not None:
current_snapshot.usage.cache_read_input_tokens = event.usage.cache_read_input_tokens
if event.usage.server_tool_use is not None:
current_snapshot.usage.server_tool_use = event.usage.server_tool_use
if event.usage.output_tokens_details is not None:
current_snapshot.usage.output_tokens_details = event.usage.output_tokens_details

# Usage may be absent when message_start omitted it (#1806); the message_delta
# carries the first full usage object, so construct it before updating.
if current_snapshot.usage is None:
_usage_data = event.usage.model_dump()
# `Usage.input_tokens` is a required int. When the delta omits it
# (e.g. `{"output_tokens": 1}` per #1806), `construct` would leave it
# as None on a non-optional field and break downstream arithmetic.
# Coerce a missing input_tokens to 0 so the snapshot stays valid.
if _usage_data.get("input_tokens") is None:
_usage_data["input_tokens"] = 0
current_snapshot.usage = Usage.construct(**_usage_data)
else:
current_snapshot.usage.output_tokens = event.usage.output_tokens

# Usage counts on a message_delta are cumulative totals, so they overwrite
# rather than add; optional ones are omitted when not applicable, in which
# case the message_start value must survive.
if event.usage.input_tokens is not None:
current_snapshot.usage.input_tokens = event.usage.input_tokens
if event.usage.cache_creation_input_tokens is not None:
current_snapshot.usage.cache_creation_input_tokens = event.usage.cache_creation_input_tokens
if event.usage.cache_read_input_tokens is not None:
current_snapshot.usage.cache_read_input_tokens = event.usage.cache_read_input_tokens
if event.usage.server_tool_use is not None:
current_snapshot.usage.server_tool_use = event.usage.server_tool_use
if event.usage.output_tokens_details is not None:
current_snapshot.usage.output_tokens_details = event.usage.output_tokens_details

return current_snapshot
18 changes: 18 additions & 0 deletions tests/lib/streaming/fixtures/missing_usage_response.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
event: message_start
data: {"type":"message_start","message":{"id":"msg_test","type":"message","role":"assistant","content":[],"model":"claude-test","stop_reason":null,"stop_sequence":null}}

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":"hi"}}

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":1}}

event: message_stop
data: {"type":"message_stop"}

Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
event: message_start
data: {"type":"message_start","message":{"id":"msg_test","type":"message","role":"assistant","content":[],"model":"claude-test","stop_reason":null,"stop_sequence":null}}

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":"hi"}}

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":{"input_tokens":11,"output_tokens":1,"cache_creation_input_tokens":3,"cache_read_input_tokens":5,"server_tool_use":{"web_search_requests":2},"cache_creation":{"ephemeral_5m_input_tokens":3}}}

event: message_stop
data: {"type":"message_stop"}

18 changes: 18 additions & 0 deletions tests/lib/streaming/fixtures/stop_details_response.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
event: message_start
data: {"type":"message_start","message":{"id":"msg_test","type":"message","role":"assistant","content":[],"model":"claude-test","stop_reason":null,"stop_sequence":null,"stop_details":{"type":"refusal"},"usage":{"input_tokens":11,"output_tokens":1}}}

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":"hi"}}

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,"stop_details":null},"usage":{"output_tokens":1}}

event: message_stop
data: {"type":"message_stop"}

94 changes: 94 additions & 0 deletions tests/lib/streaming/test_beta_messages.py
Original file line number Diff line number Diff line change
Expand Up @@ -357,6 +357,100 @@ def test_basic_response(self, respx_mock: MockRouter) -> None:

assert_basic_response([event for event in stream], stream.get_final_message())

@pytest.mark.respx(base_url=base_url)
def test_message_start_without_usage(self, respx_mock: MockRouter) -> None:
"""Beta accumulator: message_start omitting usage must not crash the stream.

The beta path carries its own copy of the accumulator, so the sync path's
coverage does not exercise it. Same contract as the non-beta test: usage is
built from the first message_delta, and a delta carrying only `output_tokens`
(#1806's reported shape) must leave `input_tokens` an int, coerced, not None.
"""
respx_mock.post("/v1/messages").mock(
return_value=httpx.Response(200, content=get_response("missing_usage_response.txt"))
)

# Non-strict client: strict validation rejects a usage-less message_start
# before the accumulator sees it.
client = Anthropic(base_url=base_url, api_key=api_key)

with client.beta.messages.stream(
max_tokens=1024,
messages=[{"role": "user", "content": "hi"}],
model="claude-test",
) as stream:
message = stream.get_final_message()
assert message.usage is not None
assert isinstance(message.usage.input_tokens, int)
assert message.usage.output_tokens == 1
assert message.stop_reason == "end_turn"

@pytest.mark.respx(base_url=base_url)
def test_message_start_without_usage_preserves_delta_optional_usage_fields(
self, respx_mock: MockRouter
) -> None:
"""The beta accumulator must preserve the delta's optional fields too.

`test_message_start_without_usage` covers the beta path's coercion, but the
coercion is only half of what the accumulator does: the same missing-usage
branch also has to carry every remaining field of the delta into the
snapshot. Driving the rich fixture through the beta manager pins that here
as it is pinned on the sync path.

Killing test: on the beta copy, building the missing-usage snapshot from the
two counters alone still satisfies the coercion assertions — the beta
accumulator re-heals the four enumerated fields outside the branch — but it
drops `cache_creation`, which is what this asserts on.
"""
respx_mock.post("/v1/messages").mock(
return_value=httpx.Response(200, content=get_response("missing_usage_rich_delta_response.txt"))
)

client = Anthropic(base_url=base_url, api_key=api_key)

with client.beta.messages.stream(
max_tokens=1024,
messages=[{"role": "user", "content": "hi"}],
model="claude-test",
) as stream:
message = stream.get_final_message()
assert message.usage is not None
assert message.usage.input_tokens == 11
assert message.usage.output_tokens == 1
assert message.usage.cache_creation_input_tokens == 3
assert message.usage.cache_read_input_tokens == 5
# nested objects on the delta survive as models, not raw dicts
assert message.usage.server_tool_use is not None
assert message.usage.server_tool_use.web_search_requests == 2
assert message.usage.cache_creation is not None
assert message.usage.cache_creation.ephemeral_5m_input_tokens == 3

@pytest.mark.respx(base_url=base_url)
def test_stop_details_from_message_start_survives_null_delta(
self, respx_mock: MockRouter
) -> None:
"""Beta: a `stop_details: null` on message_delta must not erase the start's value.

Mirrors the sync-path test of the same name. The beta accumulator carries its
own copy of the `is not None` guard, so the sync test does not pin this one;
reverting the guard on the beta file alone leaves the whole suite green.
"""
respx_mock.post("/v1/messages").mock(
return_value=httpx.Response(200, content=get_response("stop_details_response.txt"))
)

client = Anthropic(base_url=base_url, api_key=api_key)

with client.beta.messages.stream(
max_tokens=1024,
messages=[{"role": "user", "content": "hi"}],
model="claude-test",
) as stream:
message = stream.get_final_message()
assert message.stop_details is not None
assert message.stop_details.type == "refusal"
assert message.stop_reason == "end_turn"

@pytest.mark.respx(base_url=base_url)
def test_tool_use(self, respx_mock: MockRouter) -> None:
respx_mock.post("/v1/messages").mock(
Expand Down
Loading