diff --git a/litellm/constants.py b/litellm/constants.py index 072c2c358f7..3a42ced3def 100644 --- a/litellm/constants.py +++ b/litellm/constants.py @@ -1480,6 +1480,9 @@ PROXY_BUDGET_RESCHEDULER_MIN_TIME = int( os.getenv("PROXY_BUDGET_RESCHEDULER_MIN_TIME", 597) ) PROXY_BATCH_POLLING_INTERVAL = int(os.getenv("PROXY_BATCH_POLLING_INTERVAL", 3600)) +PROXY_RESPONSES_POLLING_INTERVAL = int( + os.getenv("PROXY_RESPONSES_POLLING_INTERVAL", 300) +) MAX_OBJECTS_PER_POLL_CYCLE = max(1, int(os.getenv("MAX_OBJECTS_PER_POLL_CYCLE", 50))) MANAGED_OBJECT_STALENESS_CUTOFF_DAYS = max( 1, int(os.getenv("MANAGED_OBJECT_STALENESS_CUTOFF_DAYS", 7)) diff --git a/litellm/proxy/proxy_server.py b/litellm/proxy/proxy_server.py index 4a538a28e03..164b0e90426 100644 --- a/litellm/proxy/proxy_server.py +++ b/litellm/proxy/proxy_server.py @@ -226,6 +226,7 @@ from litellm.constants import ( PROMETHEUS_FALLBACK_STATS_SEND_TIME_HOURS, PROXY_BATCH_POLLING_ENABLED, PROXY_BATCH_POLLING_INTERVAL, + PROXY_RESPONSES_POLLING_INTERVAL, PROXY_BATCH_WRITE_AT, PROXY_BUDGET_RESCHEDULER_MAX_TIME, PROXY_BUDGET_RESCHEDULER_MIN_TIME, @@ -723,7 +724,7 @@ async def _initialize_shared_aiohttp_session(): @asynccontextmanager async def proxy_startup_event(app: FastAPI): # noqa: PLR0915 - global prisma_client, master_key, use_background_health_checks, llm_router, llm_model_list, general_settings, proxy_budget_rescheduler_min_time, proxy_budget_rescheduler_max_time, litellm_proxy_admin_name, db_writer_client, store_model_in_db, premium_user, _license_check, proxy_batch_polling_interval, shared_aiohttp_session + global prisma_client, master_key, use_background_health_checks, llm_router, llm_model_list, general_settings, proxy_budget_rescheduler_min_time, proxy_budget_rescheduler_max_time, litellm_proxy_admin_name, db_writer_client, store_model_in_db, premium_user, _license_check, proxy_batch_polling_interval, proxy_responses_polling_interval, shared_aiohttp_session import json init_verbose_loggers() @@ -1763,6 +1764,7 @@ ui_access_mode: Union[Literal["admin", "all"], Dict] = "all" proxy_budget_rescheduler_min_time = PROXY_BUDGET_RESCHEDULER_MIN_TIME proxy_budget_rescheduler_max_time = PROXY_BUDGET_RESCHEDULER_MAX_TIME proxy_batch_polling_interval = PROXY_BATCH_POLLING_INTERVAL +proxy_responses_polling_interval = PROXY_RESPONSES_POLLING_INTERVAL proxy_batch_write_at = PROXY_BATCH_WRITE_AT litellm_master_key_hash = None disable_spend_logs = False @@ -3606,7 +3608,7 @@ class ProxyConfig: """ Load config values into proxy global state """ - global master_key, user_config_file_path, otel_logging, user_custom_auth, user_custom_auth_path, user_custom_key_generate, user_custom_key_update, user_custom_sso, user_custom_ui_sso_sign_in_handler, use_background_health_checks, use_shared_health_check, health_check_interval, health_check_concurrency, use_queue, proxy_budget_rescheduler_max_time, proxy_budget_rescheduler_min_time, ui_access_mode, litellm_master_key_hash, proxy_batch_write_at, disable_spend_logs, prompt_injection_detection_obj, redis_usage_cache, store_model_in_db, premium_user, open_telemetry_logger, health_check_details, proxy_batch_polling_interval, config_passthrough_endpoints + global master_key, user_config_file_path, otel_logging, user_custom_auth, user_custom_auth_path, user_custom_key_generate, user_custom_key_update, user_custom_sso, user_custom_ui_sso_sign_in_handler, use_background_health_checks, use_shared_health_check, health_check_interval, health_check_concurrency, use_queue, proxy_budget_rescheduler_max_time, proxy_budget_rescheduler_min_time, ui_access_mode, litellm_master_key_hash, proxy_batch_write_at, disable_spend_logs, prompt_injection_detection_obj, redis_usage_cache, store_model_in_db, premium_user, open_telemetry_logger, health_check_details, proxy_batch_polling_interval, proxy_responses_polling_interval, config_passthrough_endpoints config: dict = await self.get_config(config_file_path=config_file_path) @@ -4110,6 +4112,10 @@ class ProxyConfig: proxy_batch_polling_interval = general_settings.get( "proxy_batch_polling_interval", proxy_batch_polling_interval ) + ## RESPONSES POLLING INTERVAL ## + proxy_responses_polling_interval = general_settings.get( + "proxy_responses_polling_interval", proxy_responses_polling_interval + ) ## BATCH WRITER ## proxy_batch_write_at = general_settings.get( "proxy_batch_write_at", proxy_batch_write_at @@ -7185,7 +7191,7 @@ class ProxyStartupEvent: scheduler.add_job( check_responses_cost_job.check_responses_cost, "interval", - seconds=proxy_batch_polling_interval + seconds=proxy_responses_polling_interval + random.randint(0, 30), # Add small random offset # REMOVED jitter parameter - major cause of memory leak id="check_responses_cost_job", @@ -7193,7 +7199,8 @@ class ProxyStartupEvent: misfire_grace_time=APSCHEDULER_MISFIRE_GRACE_TIME, ) verbose_proxy_logger.info( - "Responses cost check job scheduled successfully" + "Responses cost check job scheduled successfully (interval=%ss)", + proxy_responses_polling_interval, ) except Exception as e: diff --git a/tests/test_litellm/proxy/test_proxy_server.py b/tests/test_litellm/proxy/test_proxy_server.py index 859594f7a0b..2c37466f003 100644 --- a/tests/test_litellm/proxy/test_proxy_server.py +++ b/tests/test_litellm/proxy/test_proxy_server.py @@ -5,6 +5,7 @@ import os import socket import subprocess import sys +import types from datetime import datetime, timedelta, timezone from pathlib import Path from unittest import mock @@ -705,6 +706,93 @@ async def test_initialize_scheduled_jobs_credentials(monkeypatch): assert len(mock_scheduler_calls) > 0 +@pytest.mark.asyncio +async def test_initialize_scheduled_jobs_uses_dedicated_responses_polling_interval(): + """ + Responses cost polling must not share the batch polling interval. + """ + from litellm.proxy.proxy_server import ProxyStartupEvent + + # Stub enterprise-only modules used by scheduler wiring. + check_batch_cost_mod = types.ModuleType( + "litellm_enterprise.proxy.common_utils.check_batch_cost" + ) + check_responses_cost_mod = types.ModuleType( + "litellm_enterprise.proxy.common_utils.check_responses_cost" + ) + + class _FakeCheckBatchCost: + def __init__(self, proxy_logging_obj, prisma_client, llm_router): + self.proxy_logging_obj = proxy_logging_obj + self.prisma_client = prisma_client + self.llm_router = llm_router + + async def check_batch_cost(self): + return None + + class _FakeCheckResponsesCost: + def __init__(self, proxy_logging_obj, prisma_client, llm_router): + self.proxy_logging_obj = proxy_logging_obj + self.prisma_client = prisma_client + self.llm_router = llm_router + + async def check_responses_cost(self): + return None + + check_batch_cost_mod.CheckBatchCost = _FakeCheckBatchCost + check_responses_cost_mod.CheckResponsesCost = _FakeCheckResponsesCost + + mock_prisma_client = MagicMock() + mock_proxy_logging = MagicMock() + mock_proxy_logging.slack_alerting_instance = MagicMock() + + with ( + patch.dict( + "sys.modules", + { + "litellm_enterprise.proxy.common_utils.check_batch_cost": check_batch_cost_mod, + "litellm_enterprise.proxy.common_utils.check_responses_cost": check_responses_cost_mod, + }, + ), + patch("litellm.proxy.proxy_server.llm_router", MagicMock()), + patch("litellm.proxy.proxy_server.PROXY_BATCH_POLLING_ENABLED", True), + patch("litellm.proxy.proxy_server.proxy_batch_polling_interval", 3600), + patch("litellm.proxy.proxy_server.proxy_responses_polling_interval", 120), + patch("litellm.proxy.proxy_server.store_model_in_db", False), + patch("litellm.proxy.proxy_server.get_secret_bool", return_value=False), + patch("litellm.proxy.proxy_server.random.randint", return_value=0), + patch( + "litellm.proxy.proxy_server.ProxyStartupEvent._initialize_slack_alerting_jobs", + new=AsyncMock(), + ), + patch( + "litellm.proxy.proxy_server.ProxyStartupEvent._initialize_spend_tracking_background_jobs", + new=AsyncMock(), + ), + ): + await ProxyStartupEvent.initialize_scheduled_background_jobs( + general_settings={"disable_reset_budget": True, "disable_spend_logs": True}, + prisma_client=mock_prisma_client, + proxy_budget_rescheduler_min_time=1, + proxy_budget_rescheduler_max_time=2, + proxy_batch_write_at=5, + proxy_logging_obj=mock_proxy_logging, + ) + + from litellm.proxy import proxy_server as ps + + try: + batch_job = ps.scheduler.get_job("check_batch_cost_job") + responses_job = ps.scheduler.get_job("check_responses_cost_job") + assert batch_job is not None + assert responses_job is not None + + assert batch_job.trigger.interval.total_seconds() == 3600 + assert responses_job.trigger.interval.total_seconds() == 120 + finally: + ps.scheduler.shutdown(wait=False) + + def test_update_config_fields_deep_merge_db_wins(): from litellm.proxy.proxy_server import ProxyConfig