From 7257ac6477242850bdfbad8c73c7c06728e76273 Mon Sep 17 00:00:00 2001 From: lichao Date: Sun, 9 Aug 2026 15:19:06 -0700 Subject: [PATCH 1/2] fix(streaming): initialize usage when message_start omits it The streaming docs show an event sequence where message_start omits usage; the accumulator then crashes with AttributeError when message_delta dereferences the missing usage value. Initialize the snapshot's usage from the delta so the final message still carries token counts, and tolerate streams that never supply usage. Fixes #1806 --- src/anthropic/lib/streaming/_beta_messages.py | 47 ++++++++++------- src/anthropic/lib/streaming/_messages.py | 40 +++++++++------ .../fixtures/usage_omitted_response.txt | 17 +++++++ tests/lib/streaming/test_messages.py | 51 +++++++++++++++++++ 4 files changed, 121 insertions(+), 34 deletions(-) create mode 100644 tests/lib/streaming/fixtures/usage_omitted_response.txt diff --git a/src/anthropic/lib/streaming/_beta_messages.py b/src/anthropic/lib/streaming/_beta_messages.py index e8ed86513..fb93e75a9 100644 --- a/src/anthropic/lib/streaming/_beta_messages.py +++ b/src/anthropic/lib/streaming/_beta_messages.py @@ -11,6 +11,7 @@ from anthropic.types.beta.beta_tool_use_block import BetaToolUseBlock from anthropic.types.beta.beta_mcp_tool_use_block import BetaMCPToolUseBlock from anthropic.types.beta.beta_server_tool_use_block import BetaServerToolUseBlock +from anthropic.types.usage import Usage from ..._types import NotGiven, not_given from ..._utils import consume_sync_iterator, consume_async_iterator @@ -553,26 +554,34 @@ def accumulate_event( 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 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. - 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 - if event.usage.iterations is not None: - current_snapshot.usage.iterations = event.usage.iterations - if event.usage.fallback_credit is not None: - current_snapshot.usage.fallback_credit = event.usage.fallback_credit - + if current_snapshot.usage is None: # pyright: ignore[reportUnnecessaryComparison] # noqa: E501 + # `message_start` may omit usage (see the streaming docs), in which + # case the snapshot has no usage yet. Initialize it from the delta + # so the final message still carries token counts, and tolerate + # streams that never supply usage. + current_snapshot.usage = Usage( + input_tokens=event.usage.input_tokens or 0, + output_tokens=event.usage.output_tokens, + ) + else: + current_snapshot.usage.output_tokens = event.usage.output_tokens + + # Update other usage fields if they exist in the event + 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 + if event.usage.iterations is not None: + current_snapshot.usage.iterations = event.usage.iterations + if event.usage.fallback_credit is not None: + current_snapshot.usage.fallback_credit = event.usage.fallback_credit return current_snapshot diff --git a/src/anthropic/lib/streaming/_messages.py b/src/anthropic/lib/streaming/_messages.py index 0ca9e7e2d..82be8cb4b 100644 --- a/src/anthropic/lib/streaming/_messages.py +++ b/src/anthropic/lib/streaming/_messages.py @@ -9,6 +9,7 @@ from anthropic.types.tool_use_block import ToolUseBlock from anthropic.types.server_tool_use_block import ServerToolUseBlock +from anthropic.types.usage import Usage from ._types import ( TextEvent, @@ -519,20 +520,29 @@ def accumulate_event( 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 + + if current_snapshot.usage is None: # pyright: ignore[reportUnnecessaryComparison] # noqa: E501 + # `message_start` may omit usage (see the streaming docs), in which + # case the snapshot has no usage yet. Initialize it from the delta + # so the final message still carries token counts, and tolerate + # streams that never supply usage. + current_snapshot.usage = Usage( + input_tokens=event.usage.input_tokens or 0, + output_tokens=event.usage.output_tokens, + ) + else: + current_snapshot.usage.output_tokens = event.usage.output_tokens + + # Update other usage fields if they exist in the event + 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 diff --git a/tests/lib/streaming/fixtures/usage_omitted_response.txt b/tests/lib/streaming/fixtures/usage_omitted_response.txt new file mode 100644 index 000000000..c6dc0f557 --- /dev/null +++ b/tests/lib/streaming/fixtures/usage_omitted_response.txt @@ -0,0 +1,17 @@ +event: message_start +data: {"type":"message_start","message":{"id":"msg_usage_omitted","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":"Hello"}} + +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":12,"output_tokens":6}} + +event: message_stop +data: {"type":"message_stop"} diff --git a/tests/lib/streaming/test_messages.py b/tests/lib/streaming/test_messages.py index a417da162..d18e44b86 100644 --- a/tests/lib/streaming/test_messages.py +++ b/tests/lib/streaming/test_messages.py @@ -385,6 +385,30 @@ def test_message_stop_event_serialization(self, respx_mock: MockRouter) -> None: stop_event.model_dump() stop_event.model_dump_json() + @pytest.mark.respx(base_url=base_url) + def test_usage_omitted_at_message_start(self, respx_mock: MockRouter) -> None: + # The streaming docs show a sequence where `message_start` omits + # `usage`; the accumulator should initialize it from `message_delta` + # instead of crashing on the missing value. + respx_mock.post("/v1/messages").mock( + return_value=httpx.Response(200, content=get_response("usage_omitted_response.txt")) + ) + + # A default (non-strict) client mirrors how the docs' event sequence + # reaches the accumulator without response-validation rejecting it. + client = Anthropic(base_url=base_url, api_key=api_key) + + with client.messages.stream( + max_tokens=1024, + messages=[{"role": "user", "content": "Say hello there!"}], + model="claude-test", + ) as stream: + message = stream.get_final_message() + + assert message.usage is not None + assert message.usage.input_tokens == 12 + assert message.usage.output_tokens == 6 + class TestAsyncMessages: @pytest.mark.asyncio @@ -531,6 +555,8 @@ async def test_refusal_stop_details_propagated(self, respx_mock: MockRouter) -> ) as stream: assert_refusal_response(await stream.get_final_message()) + @pytest.mark.asyncio + @pytest.mark.respx(base_url=base_url) @pytest.mark.asyncio @pytest.mark.respx(base_url=base_url) async def test_message_delta_fields_propagated(self, respx_mock: MockRouter) -> None: @@ -583,6 +609,31 @@ async def test_message_stop_event_serialization(self, respx_mock: MockRouter) -> stop_event.model_dump() stop_event.model_dump_json() + @pytest.mark.asyncio + @pytest.mark.respx(base_url=base_url) + async def test_usage_omitted_at_message_start(self, respx_mock: MockRouter) -> None: + # The streaming docs show a sequence where `message_start` omits + # `usage`; the accumulator should initialize it from `message_delta` + # instead of crashing on the missing value. + respx_mock.post("/v1/messages").mock( + return_value=httpx.Response(200, content=to_async_iter(get_response("usage_omitted_response.txt"))) + ) + + # A default (non-strict) client mirrors how the docs' event sequence + # reaches the accumulator without response-validation rejecting it. + client = AsyncAnthropic(base_url=base_url, api_key=api_key) + + async with client.messages.stream( + max_tokens=1024, + messages=[{"role": "user", "content": "Say hello there!"}], + model="claude-test", + ) as stream: + message = await stream.get_final_message() + + assert message.usage is not None + assert message.usage.input_tokens == 12 + assert message.usage.output_tokens == 6 + @pytest.mark.parametrize("sync", [True, False], ids=["sync", "async"]) def test_stream_method_definition_in_sync(sync: bool) -> None: From 0fc0855b198c8c3e981f7218c7b62e31c884b6d8 Mon Sep 17 00:00:00 2001 From: Lichao Chen Date: Tue, 25 Aug 2026 22:26:57 -0700 Subject: [PATCH 2/2] fix(streaming): preserve usage fields when start omits usage --- src/anthropic/lib/streaming/_beta_messages.py | 10 +++--- src/anthropic/lib/streaming/_messages.py | 10 +++--- ..._omitted_with_optional_fields_response.txt | 17 ++++++++++ ..._omitted_with_optional_fields_response.txt | 17 ++++++++++ tests/lib/streaming/test_beta_messages.py | 32 +++++++++++++++++++ tests/lib/streaming/test_messages.py | 26 +++++++++++++++ 6 files changed, 102 insertions(+), 10 deletions(-) create mode 100644 tests/lib/streaming/fixtures/beta_usage_omitted_with_optional_fields_response.txt create mode 100644 tests/lib/streaming/fixtures/usage_omitted_with_optional_fields_response.txt diff --git a/src/anthropic/lib/streaming/_beta_messages.py b/src/anthropic/lib/streaming/_beta_messages.py index fb93e75a9..e2123b560 100644 --- a/src/anthropic/lib/streaming/_beta_messages.py +++ b/src/anthropic/lib/streaming/_beta_messages.py @@ -8,10 +8,10 @@ import httpx2 as httpx from pydantic import BaseModel +from anthropic.types.beta.beta_usage import BetaUsage from anthropic.types.beta.beta_tool_use_block import BetaToolUseBlock from anthropic.types.beta.beta_mcp_tool_use_block import BetaMCPToolUseBlock from anthropic.types.beta.beta_server_tool_use_block import BetaServerToolUseBlock -from anthropic.types.usage import Usage from ..._types import NotGiven, not_given from ..._utils import consume_sync_iterator, consume_async_iterator @@ -562,10 +562,10 @@ def accumulate_event( # case the snapshot has no usage yet. Initialize it from the delta # so the final message still carries token counts, and tolerate # streams that never supply usage. - current_snapshot.usage = Usage( - input_tokens=event.usage.input_tokens or 0, - output_tokens=event.usage.output_tokens, - ) + usage = event.usage.to_dict() + if event.usage.input_tokens is None: + usage["input_tokens"] = 0 + current_snapshot.usage = construct_type(type_=BetaUsage, value=usage) else: current_snapshot.usage.output_tokens = event.usage.output_tokens diff --git a/src/anthropic/lib/streaming/_messages.py b/src/anthropic/lib/streaming/_messages.py index 82be8cb4b..b9aae967c 100644 --- a/src/anthropic/lib/streaming/_messages.py +++ b/src/anthropic/lib/streaming/_messages.py @@ -7,9 +7,9 @@ import httpx2 as httpx from pydantic import BaseModel +from anthropic.types.usage import Usage from anthropic.types.tool_use_block import ToolUseBlock from anthropic.types.server_tool_use_block import ServerToolUseBlock -from anthropic.types.usage import Usage from ._types import ( TextEvent, @@ -526,10 +526,10 @@ def accumulate_event( # case the snapshot has no usage yet. Initialize it from the delta # so the final message still carries token counts, and tolerate # streams that never supply usage. - current_snapshot.usage = Usage( - input_tokens=event.usage.input_tokens or 0, - output_tokens=event.usage.output_tokens, - ) + usage = event.usage.to_dict() + if event.usage.input_tokens is None: + usage["input_tokens"] = 0 + current_snapshot.usage = construct_type(type_=Usage, value=usage) else: current_snapshot.usage.output_tokens = event.usage.output_tokens diff --git a/tests/lib/streaming/fixtures/beta_usage_omitted_with_optional_fields_response.txt b/tests/lib/streaming/fixtures/beta_usage_omitted_with_optional_fields_response.txt new file mode 100644 index 000000000..af4507978 --- /dev/null +++ b/tests/lib/streaming/fixtures/beta_usage_omitted_with_optional_fields_response.txt @@ -0,0 +1,17 @@ +event: message_start +data: {"type":"message_start","message":{"id":"msg_beta_usage_omitted_optional","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":"Hello"}} + +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":12,"cache_creation_input_tokens":4,"cache_read_input_tokens":3,"output_tokens":6,"output_tokens_details":{"thinking_tokens":2},"server_tool_use":{"web_search_requests":1,"web_fetch_requests":0},"iterations":[{"type":"message","model":"claude-test","input_tokens":12,"cache_creation_input_tokens":4,"cache_read_input_tokens":3,"output_tokens":6}],"fallback_credit":{"status":{"type":"redeemed"}}}} + +event: message_stop +data: {"type":"message_stop"} diff --git a/tests/lib/streaming/fixtures/usage_omitted_with_optional_fields_response.txt b/tests/lib/streaming/fixtures/usage_omitted_with_optional_fields_response.txt new file mode 100644 index 000000000..95b84be99 --- /dev/null +++ b/tests/lib/streaming/fixtures/usage_omitted_with_optional_fields_response.txt @@ -0,0 +1,17 @@ +event: message_start +data: {"type":"message_start","message":{"id":"msg_usage_omitted_optional","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":"Hello"}} + +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":12,"cache_creation_input_tokens":4,"cache_read_input_tokens":3,"output_tokens":6,"output_tokens_details":{"thinking_tokens":2},"server_tool_use":{"web_search_requests":1,"web_fetch_requests":0}}} + +event: message_stop +data: {"type":"message_stop"} diff --git a/tests/lib/streaming/test_beta_messages.py b/tests/lib/streaming/test_beta_messages.py index c80bff506..a57ae374c 100644 --- a/tests/lib/streaming/test_beta_messages.py +++ b/tests/lib/streaming/test_beta_messages.py @@ -11,6 +11,7 @@ from anthropic import Anthropic, AsyncAnthropic from anthropic._utils import assert_overloads_in_sync, assert_signatures_in_sync from anthropic._compat import PYDANTIC_V1 +from anthropic.types.beta.beta_usage import BetaUsage from anthropic.types.beta.beta_message import BetaMessage from anthropic.lib.streaming._beta_types import ( BetaInputJsonEvent, @@ -536,6 +537,37 @@ def test_context_management_propagated(self, respx_mock: MockRouter) -> None: ) as stream: assert_context_management_response(stream.get_final_message()) + @pytest.mark.respx(base_url=base_url) + def test_usage_omitted_at_message_start_uses_beta_usage_and_preserves_delta_optional_usage_fields( + self, respx_mock: MockRouter + ) -> None: + respx_mock.post("/v1/messages").mock( + return_value=httpx.Response(200, content=get_response("beta_usage_omitted_with_optional_fields_response.txt")) + ) + + client = Anthropic(base_url=base_url, api_key=api_key) + + with client.beta.messages.stream( + max_tokens=1024, + messages=[{"role": "user", "content": "Say hello there!"}], + model="claude-test", + ) as stream: + message = stream.get_final_message() + + assert isinstance(message.usage, BetaUsage) + assert message.usage.input_tokens == 12 + assert message.usage.output_tokens == 6 + assert message.usage.cache_creation_input_tokens == 4 + assert message.usage.cache_read_input_tokens == 3 + assert message.usage.output_tokens_details is not None + assert message.usage.output_tokens_details.thinking_tokens == 2 + assert message.usage.server_tool_use is not None + assert message.usage.server_tool_use.web_search_requests == 1 + assert message.usage.iterations is not None + assert message.usage.iterations[0].type == "message" + assert message.usage.fallback_credit is not None + assert message.usage.fallback_credit.status.type == "redeemed" + class TestAsyncMessages: @pytest.mark.asyncio diff --git a/tests/lib/streaming/test_messages.py b/tests/lib/streaming/test_messages.py index d18e44b86..416c1cf0b 100644 --- a/tests/lib/streaming/test_messages.py +++ b/tests/lib/streaming/test_messages.py @@ -409,6 +409,32 @@ def test_usage_omitted_at_message_start(self, respx_mock: MockRouter) -> None: assert message.usage.input_tokens == 12 assert message.usage.output_tokens == 6 + @pytest.mark.respx(base_url=base_url) + def test_usage_omitted_at_message_start_preserves_delta_optional_usage_fields( + self, respx_mock: MockRouter + ) -> None: + respx_mock.post("/v1/messages").mock( + return_value=httpx.Response(200, content=get_response("usage_omitted_with_optional_fields_response.txt")) + ) + + client = Anthropic(base_url=base_url, api_key=api_key) + + with client.messages.stream( + max_tokens=1024, + messages=[{"role": "user", "content": "Say hello there!"}], + model="claude-test", + ) as stream: + message = stream.get_final_message() + + assert message.usage.input_tokens == 12 + assert message.usage.output_tokens == 6 + assert message.usage.cache_creation_input_tokens == 4 + assert message.usage.cache_read_input_tokens == 3 + assert message.usage.output_tokens_details is not None + assert message.usage.output_tokens_details.thinking_tokens == 2 + assert message.usage.server_tool_use is not None + assert message.usage.server_tool_use.web_search_requests == 1 + class TestAsyncMessages: @pytest.mark.asyncio