feat(proxy): add PROXY_RESPONSES_POLLING_INTERVAL for dedicated responses polling

- Introduced a new constant `PROXY_RESPONSES_POLLING_INTERVAL` with a default value of 300 seconds.
- Updated the proxy server to utilize this new interval for scheduling response cost checks, ensuring it operates independently from the batch polling interval.
- Added a test to verify that the responses polling interval is correctly applied in scheduled jobs.
This commit is contained in:
harish-berri 2026-05-11 22:37:35 +00:00
parent 9ac4092536
commit f3a041dec3
3 changed files with 102 additions and 4 deletions

View file

@ -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))

View file

@ -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:

View file

@ -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