fix(alerting): deduplicate eligible Slack budget notifications independently

This commit is contained in:
XAVIER ALMENDROS 2026-10-03 20:29:22 +02:00
parent 79ffcb94f1
commit 0fdce6ef63
3 changed files with 225 additions and 53 deletions

View file

@ -42,28 +42,5 @@ cache_id = budget_alert_class.get_id(user_info) # Returns user_id
To add a new budget alert type, simply create a new class that extends `BaseBudgetAlertType` and implements all the required methods, then add it to the dictionary in the `get_budget_alert_type()` function.
## Filter Slack budget alerts by key alias
Set `general_settings.alerting_args.slack_budget_alert_key_aliases` to send Slack budget alerts only for matching virtual key aliases:
```yaml
general_settings:
alerting: ["slack"]
alert_types: ["budget_alerts"]
alerting_args:
slack_budget_alert_key_aliases:
- "github-example-*"
```
Patterns are nonempty strings matched against the whole alias using Python's case-sensitive `fnmatchcase` glob rules. Exact aliases, `*`, `?` and character classes such as `[ab]` are supported. Any matching pattern permits the alert
Omitting the setting or setting it to `null` preserves existing behavior. An empty list `[]` disables Slack budget alerts. When a list is configured, only `KEY` budget events with a present, nonempty matching `key_alias` are sent to Slack. User, team, organization, project, proxy and all other non-key budget events are excluded, even if they carry an associated key alias
This only filters Slack delivery for `budget_alerts`, including immediate and digest alerts. Budget enforcement, thresholds, alert cache and deduplication remain unchanged. Webhook, email and Microsoft Teams delivery, and other Slack alert types, are unaffected
The filter applies when an alert enters the Slack queue or digest. Changing it does not retract alerts already queued or accumulated in a digest
The manual Slack service test sends a budget alert without an entity alias. With this filter configured, that budget test is suppressed even if the endpoint reports success. Other enabled alert types tested by that endpoint are unaffected. Omit the setting or use `null` to check Slack connectivity with the manual budget test
## Further Reading
- [Doc setting up Alerting on LiteLLM Proxy (Gateway)](https://docs.litellm.ai/docs/proxy/alerting)

View file

@ -588,7 +588,14 @@ class SlackAlerting(CustomBatchLogger):
if event is not None and user_info.event_group is not None:
_cache_key: Final = f"budget_alerts:{event}:{_id}"
result: Final = await _cache.async_get_cache(key=_cache_key)
if result is None:
slack_cache_key: Final = f"budget_alerts:slack:{event}:{_id}"
slack_due: Final[bool] = (
"slack" in self.alerting
and self._slack_budget_alert_allowed(user_info)
and result != "SENT"
and await _cache.async_get_cache(key=slack_cache_key) is None
)
if result is None or slack_due:
webhook_event = WebhookEvent(
event=event,
event_message=event_message,
@ -609,18 +616,27 @@ class SlackAlerting(CustomBatchLogger):
alert_emails=user_info.alert_emails,
max_budget_alert_emails=user_info.max_budget_alert_emails,
)
await self.send_alert(
slack_accepted: Final = await self.send_alert(
message=event_message + "\n\n" + user_info_str,
level="High",
alert_type=AlertType.budget_alerts,
user_info=webhook_event,
alerting_metadata={},
budget_alert_destination=("slack" if result is not None else "all") if slack_due else "non_slack",
)
await _cache.async_set_cache(
key=_cache_key,
value="SENT",
ttl=self.alerting_args.budget_alert_ttl,
)
if slack_accepted:
await _cache.async_set_cache(
key=slack_cache_key,
value="SENT",
ttl=self.alerting_args.budget_alert_ttl,
)
if result is None:
# Legacy SENT includes Slack; new markers use the independent Slack window.
await _cache.async_set_cache(
key=_cache_key,
value="SENT_WITH_SLACK_DEDUP",
ttl=self.alerting_args.budget_alert_ttl,
)
return
return
@ -1428,17 +1444,28 @@ Model Info:
return False
def _slack_budget_alert_allowed(self, user_info: CallInfo | WebhookEvent | None) -> bool:
patterns: Final = self.alerting_args.slack_budget_alert_key_aliases
return patterns is None or (
user_info is not None
and user_info.event_group == Litellm_EntityType.KEY
and user_info.key_alias is not None
and user_info.key_alias != ""
and any(fnmatchcase(user_info.key_alias, pattern) for pattern in patterns)
)
async def send_alert(
self,
message: str,
level: Literal["Low", "Medium", "High"],
alert_type: AlertType,
alerting_metadata: dict,
alerting_metadata: dict[str, object],
user_info: WebhookEvent | None = None,
request_model: str | None = None,
api_base: str | None = None,
**kwargs,
):
budget_alert_destination: Literal["all", "slack", "non_slack"] = "all",
**kwargs: object,
) -> bool:
"""
Alerting based on thresholds: - https://github.com/BerriAI/litellm/issues/1298
@ -1456,25 +1483,35 @@ Model Info:
api_base: Optional[str] - api base for digest grouping
"""
if self.alerting is None:
return
return False
# Start periodic flush if not already started
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:
if (
budget_alert_destination != "slack"
and "webhook" in self.alerting
and alert_type == "budget_alerts"
and user_info is not None
):
await self.send_webhook_alert(webhook_event=user_info)
if "email" in self.alerting and alert_type == "budget_alerts" and user_info is not None:
if (
budget_alert_destination != "slack"
and "email" in self.alerting
and alert_type == "budget_alerts"
and user_info is not None
):
# only send budget alerts over Email
await self.send_email_alert_using_smtp(webhook_event=user_info, alert_type=alert_type)
send_to_slack: Final = "slack" in self.alerting
send_to_ms_teams: Final = MS_TEAMS_ALERTING_DESTINATION in self.alerting
send_to_slack: Final = "slack" in self.alerting and budget_alert_destination != "non_slack"
send_to_ms_teams: Final = MS_TEAMS_ALERTING_DESTINATION in self.alerting and budget_alert_destination != "slack"
if not send_to_slack and not send_to_ms_teams:
return
return False
if alert_type not in self.alert_types:
return
return False
from datetime import datetime
@ -1501,20 +1538,12 @@ Model Info:
if send_to_ms_teams:
self._enqueue_ms_teams_alert(formatted_message=formatted_message, alert_type=alert_type)
budget_key_aliases: Final = self.alerting_args.slack_budget_alert_key_aliases
if not send_to_slack or (
alert_type == AlertType.budget_alerts
and budget_key_aliases is not None
and (
user_info is None
or user_info.event_group != Litellm_EntityType.KEY
or not user_info.key_alias
or not any(fnmatchcase(user_info.key_alias, pattern) for pattern in budget_key_aliases)
)
alert_type == AlertType.budget_alerts and not self._slack_budget_alert_allowed(user_info)
):
if len(self.log_queue) >= self.batch_size:
await self.flush_queue()
return
return False
# Check if digest mode is enabled for this alert type
alert_type_name_str: Final = getattr(alert_type, "value", str(alert_type))
@ -1529,6 +1558,8 @@ Model Info:
_digest_webhook = os.getenv("SLACK_WEBHOOK_URL") or os.getenv("ALERTING_WEBHOOK_URL")
if _digest_webhook is None:
raise ValueError("Missing SLACK_WEBHOOK_URL / ALERTING_WEBHOOK_URL from environment")
if _digest_webhook == []:
return False
digest_key: Final = f"{alert_type_name_str}:{request_model or ''}:{api_base or ''}"
@ -1549,7 +1580,7 @@ Model Info:
last_time=now,
webhook_url=_digest_webhook,
)
return # Suppress immediate alert; will be emitted by _flush_digest_buckets
return True
# check if we find the slack webhook url in self.alert_to_webhook_url
if self.alert_to_webhook_url is not None and alert_type in self.alert_to_webhook_url:
@ -1561,6 +1592,8 @@ Model Info:
if slack_webhook_url is None:
raise ValueError("Missing SLACK_WEBHOOK_URL / ALERTING_WEBHOOK_URL from environment")
if slack_webhook_url == []:
return False
payload: Final = {"text": formatted_message}
headers: Final = {"Content-type": "application/json"}
@ -1586,6 +1619,7 @@ Model Info:
if len(self.log_queue) >= self.batch_size:
await self.flush_queue()
return True
def _enqueue_ms_teams_alert(self, formatted_message: str, alert_type: AlertType) -> None:
ms_teams_webhook_url: Final = get_ms_teams_webhook_url()

View file

@ -4,7 +4,7 @@ import json
import time
import unittest
from collections.abc import Mapping
from typing import Final, List, Optional, Tuple
from typing import Final, List, Literal, Optional, Tuple
from unittest.mock import ANY, AsyncMock, MagicMock, Mock, patch
import httpx
@ -770,8 +770,11 @@ async def test_slack_budget_key_alias_filter_retains_native_thresholds_and_dedup
assert len(slack_alerting.log_queue) == int(delivered)
assert (
await slack_alerting.internal_usage_cache.async_get_cache("budget_alerts:threshold_crossed:hashed_key")
== "SENT"
== "SENT_WITH_SLACK_DEDUP"
)
assert (
await slack_alerting.internal_usage_cache.async_get_cache("budget_alerts:slack:threshold_crossed:hashed_key")
) == ("SENT" if delivered else None)
await slack_alerting.budget_alerts(type="token_budget", user_info=at_threshold)
assert len(slack_alerting.log_queue) == int(delivered)
await slack_alerting.flush_queue()
@ -849,6 +852,164 @@ async def test_slack_budget_key_alias_filter_reload_validates_and_changes_delive
assert THRESHOLD_ALERT in _posted_slack_bodies(http_handler)[0]["text"]
@pytest.mark.asyncio
@pytest.mark.parametrize("digest", (False, True))
@pytest.mark.parametrize("patterns", (["github-example-*"], None))
@pytest.mark.parametrize(
"budget_type, spend, max_budget, soft_budget, event",
(
("token_budget", 85.0, 100.0, None, "threshold_crossed"),
("max_budget_alert", 100.0, 100.0, None, "budget_crossed"),
("soft_budget", 50.0, None, 40.0, "soft_budget_crossed"),
("projected_limit_exceeded", 50.0, 100.0, None, "projected_limit_exceeded"),
),
)
async def test_budget_filter_reload_sends_slack_without_repeating_other_destinations(
monkeypatch: pytest.MonkeyPatch,
digest: bool,
patterns: list[str] | None,
budget_type: Literal["token_budget", "max_budget_alert", "soft_budget", "projected_limit_exceeded"],
spend: float,
max_budget: float | None,
soft_budget: float | None,
event: str,
) -> None:
monkeypatch.setenv("WEBHOOK_URL", "https://webhook.example/budget")
monkeypatch.setenv("MS_TEAMS_WEBHOOK_URL", "https://teams.example/budget")
http_handler: Final = _webhook_accepting_posts()
slack_alerting: Final = SlackAlerting(
alerting=["slack", "webhook", "ms_teams"],
default_webhook_url=SLACK_WEBHOOK_URL,
alerting_args={"slack_budget_alert_key_aliases": []},
alert_type_config={"budget_alerts": {"digest": digest, "digest_interval": 0}},
async_http_handler=http_handler,
)
info: Final = CallInfo(
spend=spend,
max_budget=max_budget,
soft_budget=soft_budget,
token="hashed_key",
key_alias="github-example-api",
event_group=Litellm_EntityType.KEY,
)
await slack_alerting.budget_alerts(type=budget_type, user_info=info)
await slack_alerting.flush_queue()
assert tuple(c.kwargs["url"] for c in http_handler.post.call_args_list) == (
"https://webhook.example/budget",
"https://teams.example/budget",
)
slack_alerting.update_values(alerting_args={"slack_budget_alert_key_aliases": patterns})
await slack_alerting.budget_alerts(type=budget_type, user_info=info)
await slack_alerting.budget_alerts(type=budget_type, user_info=info)
await slack_alerting._flush_digest_buckets()
await slack_alerting.flush_queue()
assert tuple(c.kwargs["url"] for c in http_handler.post.call_args_list) == (
"https://webhook.example/budget",
"https://teams.example/budget",
SLACK_WEBHOOK_URL,
)
assert (
"github-example-api"
in _SLACK_WEBHOOK_BODY.validate_json(http_handler.post.call_args_list[-1].kwargs["data"])["text"]
)
assert (
await slack_alerting.internal_usage_cache.async_get_cache(f"budget_alerts:slack:{event}:hashed_key") == "SENT"
)
@pytest.mark.asyncio
async def test_budget_slack_and_other_destination_windows_expire_independently(
monkeypatch: pytest.MonkeyPatch,
) -> None:
monkeypatch.setenv("WEBHOOK_URL", "https://webhook.example/budget")
http_handler: Final = _webhook_accepting_posts()
cache: Final = DualCache()
slack_alerting: Final = SlackAlerting(
alerting=["slack", "webhook"],
internal_usage_cache=cache,
default_webhook_url=SLACK_WEBHOOK_URL,
alerting_args={"slack_budget_alert_key_aliases": []},
async_http_handler=http_handler,
)
info: Final = CallInfo(
spend=85.0,
max_budget=100.0,
token="hashed_key",
key_alias="github-example-api",
event_group=Litellm_EntityType.KEY,
)
await slack_alerting.budget_alerts(type="token_budget", user_info=info)
slack_alerting.update_values(alerting_args={"slack_budget_alert_key_aliases": None})
await slack_alerting.budget_alerts(type="token_budget", user_info=info)
await slack_alerting.flush_queue()
cache.delete_cache("budget_alerts:threshold_crossed:hashed_key")
await slack_alerting.budget_alerts(type="token_budget", user_info=info)
await slack_alerting.flush_queue()
cache.delete_cache("budget_alerts:slack:threshold_crossed:hashed_key")
await slack_alerting.budget_alerts(type="token_budget", user_info=info)
await slack_alerting.flush_queue()
assert tuple(c.kwargs["url"] for c in http_handler.post.call_args_list) == (
"https://webhook.example/budget",
SLACK_WEBHOOK_URL,
"https://webhook.example/budget",
SLACK_WEBHOOK_URL,
)
@pytest.mark.asyncio
@pytest.mark.parametrize("digest", (False, True))
async def test_budget_empty_channel_mapping_does_not_consume_slack_dedup(digest: bool) -> None:
http_handler: Final = _webhook_accepting_posts()
slack_alerting: Final = SlackAlerting(
alerting=["slack"],
alert_to_webhook_url={AlertType.budget_alerts: []},
alert_type_config={"budget_alerts": {"digest": digest, "digest_interval": 0}},
async_http_handler=http_handler,
)
info: Final = CallInfo(
spend=85.0,
max_budget=100.0,
token="hashed_key",
key_alias="github-example-api",
event_group=Litellm_EntityType.KEY,
)
await slack_alerting.budget_alerts(type="token_budget", user_info=info)
assert slack_alerting.log_queue == []
assert slack_alerting.digest_buckets == {}
slack_alerting.update_values(alert_to_webhook_url={AlertType.budget_alerts: SLACK_WEBHOOK_URL})
await slack_alerting.budget_alerts(type="token_budget", user_info=info)
await slack_alerting._flush_digest_buckets()
await slack_alerting.flush_queue()
assert "github-example-api" in _posted_slack_bodies(http_handler)[0]["text"]
http_handler.post.assert_awaited_once()
@pytest.mark.asyncio
async def test_legacy_budget_sent_marker_does_not_repeat_slack() -> None:
http_handler: Final = _webhook_accepting_posts()
cache: Final = DualCache()
await cache.async_set_cache("budget_alerts:threshold_crossed:hashed_key", "SENT", ttl=86400)
slack_alerting: Final = SlackAlerting(
alerting=["slack"],
internal_usage_cache=cache,
default_webhook_url=SLACK_WEBHOOK_URL,
alerting_args={"slack_budget_alert_key_aliases": ["github-example-*"]},
async_http_handler=http_handler,
)
await slack_alerting.budget_alerts(
type="token_budget",
user_info=CallInfo(
spend=85.0,
max_budget=100.0,
token="hashed_key",
key_alias="github-example-api",
event_group=Litellm_EntityType.KEY,
),
)
await slack_alerting.flush_queue()
http_handler.post.assert_not_awaited()
def _periodic_flush_tasks() -> list[asyncio.Task[object]]:
return [
t