fix: address focus and streaming edge cases
Some checks failed
LiteLLM Rust / rustfmt, clippy, test (push) Has been cancelled

This commit is contained in:
Cursor Agent 2026-06-24 14:45:08 +00:00
parent acc983f090
commit ef5775e07c
No known key found for this signature in database
7 changed files with 130 additions and 21 deletions

View file

@ -288,12 +288,10 @@ class FocusMavvrikDestination(FocusDestination):
"Re-enable the connection in the Mavvrik dashboard."
)
if resp.status_code >= 400:
verbose_logger.warning(
"Mavvrik FOCUS destination: failed to update metricsMarker (%s): %s",
resp.status_code,
resp.text[:200],
raise RuntimeError(
f"Mavvrik FOCUS destination: failed to update metricsMarker "
f"({resp.status_code}): {resp.text[:200]}"
)
return
verbose_logger.debug(
"Mavvrik FOCUS destination: metricsMarker advanced to %s", date_epoch
)

View file

@ -72,6 +72,16 @@ def _parse_metrics_marker(
return None
def _is_empty_metrics_marker(marker: Optional[object]) -> bool:
if marker is None:
return True
if isinstance(marker, (int, float)):
return marker == 0
if isinstance(marker, str):
return not marker.strip()
return False
class MavvrikFocusLogger(FocusLogger):
"""FOCUS-based export logger that routes to the Mavvrik destination."""
@ -177,11 +187,9 @@ class MavvrikFocusLogger(FocusLogger):
last_ingested = _parse_metrics_marker(marker)
# Catch up missed dates, capped at _MAX_CATCHUP_DAYS.
# last_ingested=None means metricsMarker=0 (fresh connector, never ingested) —
# treat the same as being _MAX_CATCHUP_DAYS behind so we export all available history.
is_empty_marker = _is_empty_metrics_marker(marker)
earliest_catchup = yesterday - timedelta(days=self._MAX_CATCHUP_DAYS - 1)
if last_ingested is None or last_ingested < yesterday:
if is_empty_marker or (last_ingested is not None and last_ingested < yesterday):
catch_up_date = (
earliest_catchup
if last_ingested is None

View file

@ -672,6 +672,13 @@ class ModelResponseIterator:
if isinstance(thinking_content, str) and thinking_content:
self.reasoning_content_chunks.append(thinking_content)
reasoning_content = thinking_content
thinking_blocks = [
ChatCompletionThinkingBlock(
type="thinking",
thinking=thinking_content,
)
]
provider_specific_fields["thinking_blocks"] = thinking_blocks
signature = content_block["delta"].get("signature")
if isinstance(signature, str) and signature:

View file

@ -158,10 +158,8 @@ class DeepSeekChatConfig(OpenAIGPTConfig):
through, and drop the now-dangling tool_choice/parallel_tool_calls when
nothing callable survives.
Only non-`function` tools are ever dropped, so a `tool_choice` that names
a specific function still points at a surviving tool and is left intact;
`tool_choice`/`parallel_tool_calls` are cleared only when no function
tool remains.
When a specific `tool_choice` points at a dropped tool, clear it so the
sanitized request does not reference a tool DeepSeek will never receive.
"""
tools = optional_params.get("tools")
if not isinstance(tools, list) or not tools:
@ -170,6 +168,28 @@ class DeepSeekChatConfig(OpenAIGPTConfig):
def _is_function_tool(tool: object) -> bool:
return isinstance(tool, dict) and tool.get("type") == "function"
def _get_function_tool_name(tool: object) -> str | None:
if not isinstance(tool, dict):
return None
function = tool.get("function")
if not isinstance(function, dict):
return None
name = function.get("name")
return name if isinstance(name, str) else None
def _tool_choice_matches_function_tool(
tool_choice: object, function_tool_names: set[str]
) -> bool:
if not isinstance(tool_choice, dict):
return True
if tool_choice.get("type") != "function":
return False
function = tool_choice.get("function")
if not isinstance(function, dict):
return False
name = function.get("name")
return isinstance(name, str) and name in function_tool_names
function_tools = [tool for tool in tools if _is_function_tool(tool)]
if len(function_tools) == len(tools):
return optional_params
@ -189,6 +209,16 @@ class DeepSeekChatConfig(OpenAIGPTConfig):
cleaned = {k: v for k, v in optional_params.items() if k != "tools"}
if function_tools:
function_tool_names = {
name
for tool in function_tools
for name in (_get_function_tool_name(tool),)
if name is not None
}
if not _tool_choice_matches_function_tool(
cleaned.get("tool_choice"), function_tool_names
):
cleaned = {k: v for k, v in cleaned.items() if k != "tool_choice"}
return {**cleaned, "tools": function_tools}
return {
k: v

View file

@ -166,3 +166,24 @@ class TestDeepSeekThinkingParams:
)
assert "thinking" not in result
def test_drop_unsupported_tools_removes_dangling_tool_choice(self):
optional_params = {
"tools": [
{"type": "namespace", "name": "local_shell"},
{"type": "function", "function": {"name": "get_weather"}},
],
"tool_choice": {
"type": "function",
"function": {"name": "local_shell"},
},
"parallel_tool_calls": True,
}
result = self.config._drop_unsupported_tools(optional_params)
assert result["tools"] == [
{"type": "function", "function": {"name": "get_weather"}}
]
assert "tool_choice" not in result
assert result["parallel_tool_calls"] is True

View file

@ -2,7 +2,7 @@
from __future__ import annotations
from datetime import datetime, timezone
from datetime import datetime, timedelta, timezone
from unittest.mock import AsyncMock, MagicMock, patch
import pytest
@ -574,6 +574,39 @@ async def test_run_scheduled_export_no_catchup_when_marker_is_current():
)
@pytest.mark.asyncio
async def test_run_scheduled_export_skips_catchup_when_marker_is_unparseable():
import polars as pl
from litellm.integrations.mavvrik_focus.mavvrik_focus_logger import (
MavvrikFocusLogger,
)
from litellm.integrations.focus.destinations.mavvrik_destination import (
FocusMavvrikDestination,
)
logger = MavvrikFocusLogger()
now = datetime.now(timezone.utc).replace(hour=0, minute=0, second=0, microsecond=0)
yesterday = now - timedelta(days=1)
dest_mock = MagicMock(spec=FocusMavvrikDestination)
dest_mock.get_metrics_marker = AsyncMock(return_value="not-a-date")
db_mock = MagicMock()
db_mock.get_usage_data = AsyncMock(return_value=pl.DataFrame())
engine_mock = MagicMock()
engine_mock._database = db_mock
engine_mock._destination = dest_mock
logger._engine = engine_mock
await logger._run_scheduled_export()
assert db_mock.get_usage_data.call_count == 1
assert (
db_mock.get_usage_data.call_args.kwargs["start_time_utc"].date()
== yesterday.date()
)
@pytest.mark.asyncio
async def test_metrics_marker_always_calls_api():
"""get_metrics_marker must call the register API every time to get a fresh marker.
@ -774,8 +807,7 @@ async def test_gcs_session_cancelled_on_chunk_failure():
@pytest.mark.asyncio
async def test_update_metrics_marker_warns_on_non_410_error():
"""_update_metrics_marker must log a warning on any >=400 (non-410) status but not raise."""
async def test_update_metrics_marker_raises_on_non_410_error():
dest = _dest()
fail_resp = MagicMock()
@ -787,8 +819,8 @@ async def test_update_metrics_marker_warns_on_non_410_error():
mock_http.client.request = AsyncMock(return_value=fail_resp)
dest._http = mock_http
# Must not raise — warning only
await dest._update_metrics_marker(1234567890)
with pytest.raises(RuntimeError, match="failed to update metricsMarker"):
await dest._update_metrics_marker(1234567890)
assert mock_http.client.request.call_count == 1

View file

@ -113,6 +113,10 @@ def test_streaming_thinking_blocks_are_replayable_after_signature_delta():
for chunk in parsed_chunks
for block in (getattr(chunk.choices[0].delta, "thinking_blocks", None) or [])
)
expected_delta_blocks = (
{"type": "thinking", "thinking": "Step 1. "},
{"type": "thinking", "thinking": "Step 2."},
)
expected_thinking_block = {
"type": "thinking",
"thinking": "Step 1. Step 2.",
@ -120,7 +124,10 @@ def test_streaming_thinking_blocks_are_replayable_after_signature_delta():
}
assert reasoning_content == "Step 1. Step 2."
assert thinking_blocks == (expected_thinking_block,)
assert thinking_blocks == (*expected_delta_blocks, expected_thinking_block)
assert parsed_chunks[1].choices[0].delta.provider_specific_fields == {
"thinking_blocks": [expected_delta_blocks[0]]
}
assert parsed_chunks[-1].choices[0].delta.provider_specific_fields == {
"thinking_blocks": [expected_thinking_block]
}
@ -163,7 +170,10 @@ def test_streaming_unsigned_thinking_deltas_keep_reasoning_content():
)
assert reasoning_content == "Step 1. Step 2."
assert thinking_blocks == ()
assert thinking_blocks == (
{"type": "thinking", "thinking": "Step 1. "},
{"type": "thinking", "thinking": "Step 2."},
)
def test_streaming_truncated_thinking_deltas_keep_reasoning_content():
@ -202,7 +212,10 @@ def test_streaming_truncated_thinking_deltas_keep_reasoning_content():
)
assert reasoning_content == "Step 1. Step 2."
assert thinking_blocks == ()
assert thinking_blocks == (
{"type": "thinking", "thinking": "Step 1. "},
{"type": "thinking", "thinking": "Step 2."},
)
def test_handle_json_mode_chunk_response_format_tool():