Skip to content
Merged
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
19 changes: 11 additions & 8 deletions litellm/integrations/SlackAlerting/slack_alerting.py
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,7 @@ def __init__(
self.default_webhook_url = default_webhook_url
self.flush_lock = asyncio.Lock()
self.periodic_started = False
self._periodic_flush_task: asyncio.Task[None] | None = None
self.hanging_request_check = AlertingHangingRequestCheck(
slack_alerting_object=self,
)
Expand All @@ -125,6 +126,12 @@ def __init__(
self.digest_lock = asyncio.Lock()
super().__init__(**kwargs, flush_lock=self.flush_lock)

def _ensure_periodic_flush_task(self) -> None:
if self.periodic_started and (self._periodic_flush_task is None or not self._periodic_flush_task.done()):
return
self._periodic_flush_task = asyncio.create_task(self.periodic_flush())
self.periodic_started = True

def update_values(
self,
alerting: list | None = None,
Expand All @@ -137,17 +144,14 @@ def update_values(
):
if alerting is not None:
self.alerting = alerting
asyncio.create_task(self.periodic_flush())
self.periodic_started = True
self._ensure_periodic_flush_task()
if alerting_threshold is not None:
self.alerting_threshold = alerting_threshold
if alert_types is not None:
self.alert_types = alert_types
if alerting_args is not None:
self.alerting_args = SlackAlertingArgs(**alerting_args)
if not self.periodic_started:
asyncio.create_task(self.periodic_flush())
self.periodic_started = True
self._ensure_periodic_flush_task()
if alert_type_config is not None:
for key, val in alert_type_config.items():
self.alert_type_config[key] = AlertTypeConfig(**val) if isinstance(val, dict) else val
Expand Down Expand Up @@ -1442,9 +1446,8 @@ async def send_alert(
return

# Start periodic flush if not already started
if not self.periodic_started and self.alerting is not None and len(self.alerting) > 0:
asyncio.create_task(self.periodic_flush())
self.periodic_started = True
if self.alerting is not None and len(self.alerting) > 0:
self._ensure_periodic_flush_task()

if "webhook" in self.alerting and alert_type == "budget_alerts" and user_info is not None:
await self.send_webhook_alert(webhook_event=user_info)
Expand Down
324 changes: 183 additions & 141 deletions litellm/integrations/anthropic_cache_control_hook.py

Large diffs are not rendered by default.

11 changes: 8 additions & 3 deletions litellm/litellm_core_utils/token_counter.py
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,8 @@
AllMessageValues,
ChatCompletionDocumentObject,
ChatCompletionNamedToolChoiceParam,
ChatCompletionRedactedThinkingBlock,
ChatCompletionThinkingBlock,
ChatCompletionToolParam,
OpenAIMessageContentListBlock,
)
Expand Down Expand Up @@ -854,6 +856,8 @@ def _count_content_list(
content_list: str
| Iterable[
OpenAIMessageContentListBlock
| ChatCompletionThinkingBlock
| ChatCompletionRedactedThinkingBlock
| AnthropicMessagesTextParam
| AnthropicMessagesImageParam
| AnthropicMessagesDocumentParam
Expand Down Expand Up @@ -898,9 +902,9 @@ def _count_content_list(
use_default_image_token_count,
default_token_count,
)
elif c["type"] == "thinking":
elif c["type"] in ("thinking", "redacted_thinking"):
# Claude extended thinking content block
# Count the thinking text and skip signature (opaque signature blob)
# Count the thinking text and skip the opaque blobs (signature, redacted data)
thinking_text = str(c.get("thinking", ""))
if thinking_text:
num_tokens += count_function(thinking_text)
Expand All @@ -920,7 +924,8 @@ def _count_content_list(
raise ValueError(
f"Invalid content item type: {content_type}. "
f"Expected str or dict with 'type' field "
f"(text, image_url, image, document, file, tool_use, tool_result, thinking, tool_reference)."
f"(text, image_url, image, document, file, tool_use, tool_result, thinking, redacted_thinking, "
f"tool_reference)."
)
return num_tokens
except Exception as e:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -959,7 +959,7 @@ async def get_mcp_tools(
mcp_server_auth_headers=None,
)
tools: Final = listing.tools
dumped_tools: Final = [dict(tool) for tool in tools]
dumped_tools: Final = [tool.model_dump(by_alias=True) for tool in tools]

return {"tools": dumped_tools}

Expand Down
10 changes: 0 additions & 10 deletions litellm/responses/mcp/mcp_streaming_iterator.py
Original file line number Diff line number Diff line change
Expand Up @@ -68,16 +68,6 @@ async def create_mcp_list_tools_events(
# Use the pre-processed MCP tools that were already fetched, filtered, and deduplicated by the parent
filtered_mcp_tools: Final = pre_processed_mcp_tools

# Convert tools to dict format for the event
_mcp_tools_dict: Final = [
tool.model_dump()
if hasattr(tool, "model_dump") and callable(getattr(tool, "model_dump", None))
else tool.__dict__
if hasattr(tool, "__dict__")
else {"name": getattr(tool, "name", str(tool))}
for tool in filtered_mcp_tools
]

# Emit list tools completed event
completed_event: Final = MCPListToolsCompletedEvent(
type=ResponsesAPIStreamEvents.MCP_LIST_TOOLS_COMPLETED,
Expand Down
6 changes: 5 additions & 1 deletion litellm/router.py
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,7 @@
MODEL_INFO_REFRESH_SECONDS,
get_openai_compatible_model_info,
)
from litellm.router_strategy.base_routing_strategy import BaseRoutingStrategy
from litellm.router_strategy.budget_limiter import RouterBudgetLimiting
from litellm.router_strategy.least_busy import LeastBusyLoggingHandler
from litellm.router_strategy.lowest_cost import LowestCostLoggingHandler
Expand Down Expand Up @@ -1350,6 +1351,9 @@ def _unregister_router_selectors(self, selectors: Sequence[object]) -> None:
`_init_routing_groups`) so repeated `update_settings` calls don't
accumulate dead selectors that keep receiving callback events.
"""
for selector in selectors:
if isinstance(selector, BaseRoutingStrategy):
selector.retire()
selector_ids: Final = {id(s) for s in selectors if s is not None}
if not selector_ids:
return
Expand Down Expand Up @@ -12033,7 +12037,7 @@ def update_settings(self, **kwargs):
)
rebuild_routing_groups = True
elif var == "routing_strategy_args":
routing_args_updated = True
routing_args_updated = value != self.routing_strategy_args

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.

P1 Mutated routing arguments stay stale

If a caller changes the dictionary previously passed as routing_strategy_args and passes that same dictionary to update_settings(), this comparison checks the dictionary against itself. It skips rebuilding the selector, so the router reports the new arguments while routing still uses the old values.

setattr(self, var, value)
else:
verbose_router_logger.debug("Setting %s is not allowed", var)
Expand Down
16 changes: 15 additions & 1 deletion litellm/router_strategy/base_routing_strategy.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,10 +40,24 @@ def setup_sync_task(self, default_sync_interval: float | None):
self.periodic_sync_in_memory_spend_with_redis(default_sync_interval=default_sync_interval)
)

def cancel_sync_task(self) -> None:
if self._sync_task is not None:
self._sync_task.cancel()

def retire(self) -> None:
self.cancel_sync_task()
if not self.redis_increment_operation_queue:
return
try:
loop: Final = asyncio.get_running_loop()
except RuntimeError:
return
loop.create_task(self._push_in_memory_increments_to_redis())
Comment on lines +51 to +55

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.

P1 Queued increments are lost

If a synchronous Router.update_settings() call replaces a usage-based selector while it has Redis increments queued, retire() cancels its sync task and returns because there is no running event loop. Those increments are never flushed, so later routing decisions can use understated RPM usage.


async def cleanup(self):
"""Cleanup method to be called when shutting down"""
if self._sync_task is not None:
self._sync_task.cancel()
self.cancel_sync_task()
try:
await self._sync_task
except asyncio.CancelledError:
Expand Down
2 changes: 1 addition & 1 deletion litellm/rust_bridge/catalog.py
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,7 @@ def matches(self, context: Context) -> bool:
RULES: Final[Rules] = (
Rule(Route.OCR, Rollout.RUST_REQUIRED, providers=frozenset({"aws_textract"})),
Rule(Route.OCR, Rollout.RUST_OPT_OUT),
Rule(Route.MESSAGES, Rollout.RUST_OPT_IN),
Rule(Route.MESSAGES, Rollout.PYTHON_ONLY),
Rule(Route.TRANSCRIPTION, Rollout.RUST_REQUIRED, providers=frozenset({"bedrock"})),
)

Expand Down
4 changes: 2 additions & 2 deletions litellm/types/integrations/anthropic_cache_control_hook.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,17 +17,17 @@ class CacheControlMessageInjectionPoint(TypedDict):
role: Literal["user", "system", "assistant"] | None # Optional: target by role (user, system, assistant)
index: int | str | None # Optional: target by specific index
control: ChatCompletionCachedContent | None
_litellm_judged: NotRequired[bool] # Internal: written back by litellm once the client cache_control judgment ran
_litellm_openai_dialect: NotRequired[ReadOnly[bool]]
_litellm_external_breakpoints: NotRequired[ReadOnly[int]]


class CacheControlToolConfigInjectionPoint(TypedDict):
"""Type for tool_config-level injection points (Bedrock)."""

location: Literal["tool_config"]
control: ChatCompletionCachedContent | None
_litellm_judged: NotRequired[bool] # Internal: written back by litellm once the client cache_control judgment ran
_litellm_openai_dialect: NotRequired[ReadOnly[bool]]
_litellm_external_breakpoints: NotRequired[ReadOnly[int]]


CacheControlInjectionPoint = CacheControlMessageInjectionPoint | CacheControlToolConfigInjectionPoint
4 changes: 4 additions & 0 deletions litellm/types/llms/anthropic.py
Original file line number Diff line number Diff line change
Expand Up @@ -753,6 +753,10 @@ class ANTHROPIC_BETA_HEADER_VALUES(str, Enum):
# Tool search beta header constant (for Anthropic direct API and Microsoft Foundry)
ANTHROPIC_TOOL_SEARCH_BETA_HEADER: Final = "advanced-tool-use-2025-11-20"

ANTHROPIC_TOOL_SEARCH_TOOL_TYPES: Final = frozenset(
{"tool_search_tool_regex_20251119", "tool_search_tool_bm25_20251119"}
)

# Effort beta header constant
ANTHROPIC_EFFORT_BETA_HEADER: Final = "effort-2025-11-24"

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -434,3 +434,29 @@ async def test_send_alert_raises_when_no_webhook_url_configured(monkeypatch):
alert_type=AlertType.budget_alerts,
alerting_metadata={},
)


def _periodic_flush_tasks() -> list[asyncio.Task[object]]:
return [
t
for t in asyncio.all_tasks()
if t.get_coro() is not None and t.get_coro().__qualname__ == "SlackAlerting.periodic_flush"
]


@pytest.mark.asyncio
async def test_update_values_repeated_alerting_reload_keeps_single_periodic_flush_task() -> None:
slack_alerting: Final = SlackAlerting(alerting=["slack"])
try:
for _ in range(5):
slack_alerting.update_values(alerting=["slack"])
await asyncio.sleep(0)
flush_tasks: Final = _periodic_flush_tasks()
assert len(flush_tasks) == 1, f"expected 1 periodic_flush task, found {len(flush_tasks)}"
finally:
for t in _periodic_flush_tasks():
t.cancel()
try:
await t
except asyncio.CancelledError:
pass
Loading
Loading