From 33d82f806317585188ae666604d3958639636a3a Mon Sep 17 00:00:00 2001 From: lichao Date: Sun, 9 Aug 2026 15:19:06 -0700 Subject: [PATCH 1/3] 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 028dd80ef..ca9e48cfc 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,7 +554,6 @@ 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 # only sent on `message_delta` after a mid-stream fallback, in which case it @@ -561,22 +561,31 @@ def accumulate_event( if event.input_transformations is not None: current_snapshot.input_transformations = event.input_transformations - # 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 ae7b42393..c327b7910 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 7d4783a8b..5aff29cd4 100644 --- a/tests/lib/streaming/test_messages.py +++ b/tests/lib/streaming/test_messages.py @@ -387,6 +387,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 @@ -533,6 +557,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: @@ -585,6 +611,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 + def test_message_delta_fields_are_all_accumulated() -> None: # tripwire: handle a new field in accumulate_event (src/anthropic/lib/streaming/_messages.py), then list it here From 3a397be2670a12d17d7dabf157e9be5096d55f70 Mon Sep 17 00:00:00 2001 From: Lichao Chen Date: Tue, 25 Aug 2026 22:26:57 -0700 Subject: [PATCH 2/3] 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 | 30 +++++++++++++++-- 6 files changed, 104 insertions(+), 12 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 ca9e48cfc..3ff6659e8 100644 --- a/src/anthropic/lib/streaming/_beta_messages.py +++ b/src/anthropic/lib/streaming/_beta_messages.py @@ -8,10 +8,10 @@ import httpx2 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 @@ -566,10 +566,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 c327b7910..8ba7cdadc 100644 --- a/src/anthropic/lib/streaming/_messages.py +++ b/src/anthropic/lib/streaming/_messages.py @@ -7,9 +7,9 @@ import httpx2 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 c1cdf4915..2f0a94f40 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, get_model_fields +from anthropic.types.beta.beta_usage import BetaUsage from anthropic.types.beta.beta_message import BetaMessage from anthropic.lib.streaming._beta_types import ( BetaInputJsonEvent, @@ -558,6 +559,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=httpx2.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" + @pytest.mark.respx(base_url=base_url) @pytest.mark.parametrize("fixture, expected", INPUT_TRANSFORMATIONS_CASES) def test_input_transformations_propagated( diff --git a/tests/lib/streaming/test_messages.py b/tests/lib/streaming/test_messages.py index 5aff29cd4..9dd2e67bd 100644 --- a/tests/lib/streaming/test_messages.py +++ b/tests/lib/streaming/test_messages.py @@ -393,7 +393,7 @@ def test_usage_omitted_at_message_start(self, respx_mock: MockRouter) -> None: # `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")) + return_value=httpx2.Response(200, content=get_response("usage_omitted_response.txt")) ) # A default (non-strict) client mirrors how the docs' event sequence @@ -411,6 +411,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=httpx2.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 @@ -618,7 +644,7 @@ async def test_usage_omitted_at_message_start(self, respx_mock: MockRouter) -> N # `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"))) + return_value=httpx2.Response(200, content=to_async_iter(get_response("usage_omitted_response.txt"))) ) # A default (non-strict) client mirrors how the docs' event sequence From 993aa73ed28306ccc86f934a808ebb43f31fa95e Mon Sep 17 00:00:00 2001 From: lichao Date: Thu, 3 Sep 2026 01:04:23 -0700 Subject: [PATCH 3/3] fix(streaming): preserve unknown usage instead of fabricating input_tokens=0 When both message_start and message_delta.input_tokens omit the input count, the accumulator reported input_tokens=0, which under-reports accounting to callers. Per review feedback, leave the count unset when it is genuinely unknown: construct_type already leaves fields the delta did not supply as None, matching how the rest of the SDK represents wire-omitted values. A later delta that does supply the count still fills it in via the accumulate branch. --- src/anthropic/lib/streaming/_beta_messages.py | 10 ++-- src/anthropic/lib/streaming/_messages.py | 10 ++-- ...ta_usage_omitted_input_tokens_response.txt | 17 ++++++ .../usage_omitted_input_tokens_response.txt | 17 ++++++ tests/lib/streaming/test_beta_messages.py | 23 ++++++++ tests/lib/streaming/test_messages.py | 53 +++++++++++++++++++ 6 files changed, 120 insertions(+), 10 deletions(-) create mode 100644 tests/lib/streaming/fixtures/beta_usage_omitted_input_tokens_response.txt create mode 100644 tests/lib/streaming/fixtures/usage_omitted_input_tokens_response.txt diff --git a/src/anthropic/lib/streaming/_beta_messages.py b/src/anthropic/lib/streaming/_beta_messages.py index 3ff6659e8..8796e15ef 100644 --- a/src/anthropic/lib/streaming/_beta_messages.py +++ b/src/anthropic/lib/streaming/_beta_messages.py @@ -565,11 +565,11 @@ def accumulate_event( # `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. - 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) + # streams that never supply usage. Anything the delta leaves out + # (e.g. `input_tokens`) stays unset instead of being fabricated + # as 0 so an unknown count is never reported to callers; a later + # delta that does supply it fills it in below. + current_snapshot.usage = construct_type(type_=BetaUsage, value=event.usage.to_dict()) 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 8ba7cdadc..2fda72142 100644 --- a/src/anthropic/lib/streaming/_messages.py +++ b/src/anthropic/lib/streaming/_messages.py @@ -525,11 +525,11 @@ def accumulate_event( # `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. - 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) + # streams that never supply usage. Anything the delta leaves out + # (e.g. `input_tokens`) stays unset instead of being fabricated + # as 0 so an unknown count is never reported to callers; a later + # delta that does supply it fills it in below. + current_snapshot.usage = construct_type(type_=Usage, value=event.usage.to_dict()) else: current_snapshot.usage.output_tokens = event.usage.output_tokens diff --git a/tests/lib/streaming/fixtures/beta_usage_omitted_input_tokens_response.txt b/tests/lib/streaming/fixtures/beta_usage_omitted_input_tokens_response.txt new file mode 100644 index 000000000..e56b81799 --- /dev/null +++ b/tests/lib/streaming/fixtures/beta_usage_omitted_input_tokens_response.txt @@ -0,0 +1,17 @@ +event: message_start +data: {"type":"message_start","message":{"id":"msg_usage_omitted_input","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":{"output_tokens":6}} + +event: message_stop +data: {"type":"message_stop"} \ No newline at end of file diff --git a/tests/lib/streaming/fixtures/usage_omitted_input_tokens_response.txt b/tests/lib/streaming/fixtures/usage_omitted_input_tokens_response.txt new file mode 100644 index 000000000..e56b81799 --- /dev/null +++ b/tests/lib/streaming/fixtures/usage_omitted_input_tokens_response.txt @@ -0,0 +1,17 @@ +event: message_start +data: {"type":"message_start","message":{"id":"msg_usage_omitted_input","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":{"output_tokens":6}} + +event: message_stop +data: {"type":"message_stop"} \ No newline at end of file diff --git a/tests/lib/streaming/test_beta_messages.py b/tests/lib/streaming/test_beta_messages.py index 2f0a94f40..d131bedde 100644 --- a/tests/lib/streaming/test_beta_messages.py +++ b/tests/lib/streaming/test_beta_messages.py @@ -590,6 +590,29 @@ def test_usage_omitted_at_message_start_uses_beta_usage_and_preserves_delta_opti assert message.usage.fallback_credit is not None assert message.usage.fallback_credit.status.type == "redeemed" + @pytest.mark.respx(base_url=base_url) + def test_usage_omitted_at_message_start_preserves_unknown_input_tokens(self, respx_mock: MockRouter) -> None: + # If both `message_start` and `message_delta` omit the input count then + # it is genuinely unknown; the accumulator must preserve that rather + # than fabricate `input_tokens=0`, which would under-report accounting. + respx_mock.post("/v1/messages").mock( + return_value=httpx2.Response(200, content=get_response("beta_usage_omitted_input_tokens_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) + # unknown counts stay unset instead of being reported as 0 + assert message.usage.input_tokens is None # pyright: ignore[reportUnnecessaryComparison] + assert message.usage.output_tokens == 6 + @pytest.mark.respx(base_url=base_url) @pytest.mark.parametrize("fixture, expected", INPUT_TRANSFORMATIONS_CASES) def test_input_transformations_propagated( diff --git a/tests/lib/streaming/test_messages.py b/tests/lib/streaming/test_messages.py index 9dd2e67bd..65e71b542 100644 --- a/tests/lib/streaming/test_messages.py +++ b/tests/lib/streaming/test_messages.py @@ -411,6 +411,31 @@ 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_unknown_input_tokens(self, respx_mock: MockRouter) -> None: + # If both `message_start` and `message_delta` omit the input count then + # it is genuinely unknown; the accumulator must preserve that rather + # than fabricate `input_tokens=0`, which would under-report accounting. + respx_mock.post("/v1/messages").mock( + return_value=httpx2.Response(200, content=get_response("usage_omitted_input_tokens_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 + # unknown counts stay unset instead of being reported as 0 + assert message.usage.input_tokens is None # pyright: ignore[reportUnnecessaryComparison] + 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 @@ -662,6 +687,34 @@ async def test_usage_omitted_at_message_start(self, respx_mock: MockRouter) -> N assert message.usage.input_tokens == 12 assert message.usage.output_tokens == 6 + @pytest.mark.asyncio + @pytest.mark.respx(base_url=base_url) + async def test_usage_omitted_at_message_start_preserves_unknown_input_tokens(self, respx_mock: MockRouter) -> None: + # If both `message_start` and `message_delta` omit the input count then + # it is genuinely unknown; the accumulator must preserve that rather + # than fabricate `input_tokens=0`, which would under-report accounting. + respx_mock.post("/v1/messages").mock( + return_value=httpx2.Response( + 200, content=to_async_iter(get_response("usage_omitted_input_tokens_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 + # unknown counts stay unset instead of being reported as 0 + assert message.usage.input_tokens is None # pyright: ignore[reportUnnecessaryComparison] + assert message.usage.output_tokens == 6 + def test_message_delta_fields_are_all_accumulated() -> None: # tripwire: handle a new field in accumulate_event (src/anthropic/lib/streaming/_messages.py), then list it here