Skip to content
Closed
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 @@ -96,6 +96,38 @@ def _is_combined(chunk: Any) -> bool:
or getattr(delta, "thinking_blocks", None)
)

@staticmethod
def _has_mixed_reasoning_and_text(chunk: Any) -> bool:
choices = getattr(chunk, "choices", None)
if not choices:
return False
delta = getattr(choices[0], "delta", None)
if delta is None:
return False
return bool(getattr(delta, "content", None) and getattr(delta, "reasoning_content", None))

@staticmethod
def _clear_usage(chunk: Any) -> None:
if hasattr(chunk, "usage"):
chunk.usage = None
hidden_params = getattr(chunk, "_hidden_params", None)
if isinstance(hidden_params, dict) and "usage" in hidden_params:
chunk._hidden_params = {key: value for key, value in hidden_params.items() if key != "usage"}

@staticmethod
def _split_mixed_reasoning_and_text(chunk: Any) -> list[Any]:
if not _CombinedChunkSplitter._has_mixed_reasoning_and_text(chunk):
return [chunk]

reasoning_chunk = copy.deepcopy(chunk)
reasoning_chunk.choices[0].finish_reason = None
reasoning_chunk.choices[0].delta.content = None
_CombinedChunkSplitter._clear_usage(reasoning_chunk)

text_chunk = copy.deepcopy(chunk)
text_chunk.choices[0].delta.reasoning_content = None
return [reasoning_chunk, text_chunk]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Text split keeps thinking_blocks

Medium Severity

When _split_mixed_reasoning_and_text builds the text half of a mixed chunk, it clears reasoning_content but leaves thinking_blocks on the delta. The Anthropic translate path still prefers thinking over text when blocks are present, so the intended text can be dropped and a thinking_delta can be emitted after a text content_block_start, recreating the invalid-stream failure mode this PR targets.

Fix in Cursor Fix in Web

Reviewed by Cursor Bugbot for commit 2b94847. Configure here.


@staticmethod
def _split(chunk: Any) -> List[Any]:
"""Return ``[chunk]``, or ``[content_chunk, finish_chunk]`` if combined."""
Expand All @@ -105,6 +137,7 @@ def _split(chunk: Any) -> List[Any]:
# Content chunk: keep the delta payload, clear the finish_reason.
content_chunk = copy.deepcopy(chunk)
content_chunk.choices[0].finish_reason = None
_CombinedChunkSplitter._clear_usage(content_chunk)

# Finish chunk: keep finish_reason (and usage), clear the delta payload.
finish_chunk = copy.deepcopy(chunk)
Expand All @@ -127,7 +160,11 @@ def __next__(self) -> Any:
if self._sync_iter is None:
self._sync_iter = iter(self._stream)
chunk = next(self._sync_iter) # propagates StopIteration when exhausted
self._buffer.extend(self._split(chunk))
self._buffer.extend(
split_chunk
for combined_chunk in self._split(chunk)
for split_chunk in self._split_mixed_reasoning_and_text(combined_chunk)
)
return self._buffer.popleft()

def __aiter__(self) -> "AsyncIterator[Any]":
Expand All @@ -139,7 +176,11 @@ async def __anext__(self) -> Any:
if self._async_iter is None:
self._async_iter = self._stream.__aiter__()
chunk = await self._async_iter.__anext__() # propagates StopAsyncIteration
self._buffer.extend(self._split(chunk))
self._buffer.extend(
split_chunk
for combined_chunk in self._split(chunk)
for split_chunk in self._split_mixed_reasoning_and_text(combined_chunk)
)
return self._buffer.popleft()


Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -65,9 +65,7 @@ def _thinking_chunk(thinking: str, signature: str = "") -> MagicMock:
return _make_chunk(Delta(content=None, thinking_blocks=[block]))


def _tool_chunk(
call_id: str, name: Optional[str], arguments: Optional[str]
) -> MagicMock:
def _tool_chunk(call_id: str, name: Optional[str], arguments: Optional[str]) -> MagicMock:
return _make_chunk(
Delta(
content=None,
Expand Down Expand Up @@ -109,35 +107,31 @@ def _text_deltas(events: List[dict]) -> List[str]:
return [
e["delta"]["text"]
for e in events
if e.get("type") == "content_block_delta"
and e["delta"].get("type") == "text_delta"
if e.get("type") == "content_block_delta" and e["delta"].get("type") == "text_delta"
]


def _input_json_deltas(events: List[dict]) -> List[str]:
return [
e["delta"]["partial_json"]
for e in events
if e.get("type") == "content_block_delta"
and e["delta"].get("type") == "input_json_delta"
if e.get("type") == "content_block_delta" and e["delta"].get("type") == "input_json_delta"
]


def _thinking_deltas(events: List[dict]) -> List[str]:
return [
e["delta"]["thinking"]
for e in events
if e.get("type") == "content_block_delta"
and e["delta"].get("type") == "thinking_delta"
if e.get("type") == "content_block_delta" and e["delta"].get("type") == "thinking_delta"
]


def _signature_deltas(events: List[dict]) -> List[str]:
return [
e["delta"]["signature"]
for e in events
if e.get("type") == "content_block_delta"
and e["delta"].get("type") == "signature_delta"
if e.get("type") == "content_block_delta" and e["delta"].get("type") == "signature_delta"
]


Expand Down Expand Up @@ -228,9 +222,7 @@ async def test_first_text_delta_after_tool_use_is_not_dropped_async():
_make_chunk(Delta(content=" Bye.")),
_make_chunk(Delta(content=None), finish_reason="stop"),
]
wrapper = AnthropicStreamWrapper(
completion_stream=_AsyncStream(chunks), model="claude-x"
)
wrapper = AnthropicStreamWrapper(completion_stream=_AsyncStream(chunks), model="claude-x")
events = await _drain_async(wrapper)

assert _input_json_deltas(events) == ['{"city": "NY"}']
Expand Down Expand Up @@ -506,3 +498,39 @@ def test_empty_content_chunk_mid_text_block_is_suppressed_sync():

assert _text_deltas(events) == ["Hi", " there"]
_assert_deltas_match_their_block_type(events)


def _mixed_reasoning_and_text_chunks() -> List[MagicMock]:
return [
_make_chunk(Delta(content=None, reasoning_content="First thought.")),
_make_chunk(
Delta(content="Answer.", reasoning_content=" Last thought."),
finish_reason="stop",
),
]


def _assert_mixed_reasoning_and_text_chunk_is_split(events: List[dict]) -> None:
_assert_deltas_match_their_block_type(events)
assert _thinking_deltas(events) == ["First thought.", " Last thought."]
assert _text_deltas(events) == ["Answer."]
assert [event["type"] for event in events].count("message_delta") == 1


def test_mixed_reasoning_and_text_chunk_is_split_sync():
wrapper = AnthropicStreamWrapper(
completion_stream=iter(_mixed_reasoning_and_text_chunks()),
model="claude-x",
)

_assert_mixed_reasoning_and_text_chunk_is_split(_drain_sync(wrapper))


@pytest.mark.asyncio
async def test_mixed_reasoning_and_text_chunk_is_split_async():
wrapper = AnthropicStreamWrapper(
completion_stream=_AsyncStream(_mixed_reasoning_and_text_chunks()),
model="claude-x",
)

_assert_mixed_reasoning_and_text_chunk_is_split(await _drain_async(wrapper))
Loading