From f7c157083a27651be957559632d63b70101e8486 Mon Sep 17 00:00:00 2001 From: yucheng Date: Fri, 25 Sep 2026 18:09:30 +0000 Subject: [PATCH] fix(passthrough): clamp report concurrency and drain reports before spend flushes Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- litellm/constants.py | 4 +- litellm/proxy/proxy_server.py | 10 ++--- .../proxy/proxy_server/test_lifecycle.py | 38 +++++++++++++++++++ tests/test_litellm/test_constants.py | 11 ++++++ 4 files changed, 56 insertions(+), 7 deletions(-) diff --git a/litellm/constants.py b/litellm/constants.py index eb97420aaa1..21b2dab94df 100644 --- a/litellm/constants.py +++ b/litellm/constants.py @@ -1520,8 +1520,8 @@ CLOUDZERO_EXPORT_INTERVAL_MINUTES: Final = int(os.getenv("CLOUDZERO_EXPORT_INTER MCP_TOOL_NAME_PREFIX: Final = "mcp_tool" MAXIMUM_TRACEBACK_LINES_TO_LOG: Final = int(os.getenv("MAXIMUM_TRACEBACK_LINES_TO_LOG", 100)) PASSTHROUGH_UPSTREAM_ERROR_BODY_MAX_LOG_CHARS: Final = 4096 -PASSTHROUGH_UPSTREAM_ERROR_REPORT_CONCURRENCY: Final = int( - os.getenv("PASSTHROUGH_UPSTREAM_ERROR_REPORT_CONCURRENCY", "64") +PASSTHROUGH_UPSTREAM_ERROR_REPORT_CONCURRENCY: Final = max( + 1, int(os.getenv("PASSTHROUGH_UPSTREAM_ERROR_REPORT_CONCURRENCY", "64")) ) PASSTHROUGH_UPSTREAM_ERROR_REPORT_DRAIN_SECONDS: Final = int( os.getenv("PASSTHROUGH_UPSTREAM_ERROR_REPORT_DRAIN_SECONDS", "10") diff --git a/litellm/proxy/proxy_server.py b/litellm/proxy/proxy_server.py index ac9c8ea8374..53151bca4d3 100644 --- a/litellm/proxy/proxy_server.py +++ b/litellm/proxy/proxy_server.py @@ -1574,6 +1574,11 @@ async def proxy_startup_event(app: FastAPI) -> AsyncGenerator[None, None]: except Exception as e: verbose_proxy_logger.error("Error stopping the spend view setup task: %s", e) + try: + await drain_passthrough_upstream_error_reports() + except Exception as e: # noqa: BLE001 # shutdown must continue when a report drain fails + verbose_proxy_logger.error("Error draining passthrough upstream error reports: %s", e) + await _drain_spend_event_producer_on_shutdown() # Shutdown event - finish or cancel in-flight scheduled jobs before the shutdown flushes and the DB disconnect @@ -1585,11 +1590,6 @@ async def proxy_startup_event(app: FastAPI) -> AsyncGenerator[None, None]: await flush_spend_counters_on_shutdown() - try: - await drain_passthrough_upstream_error_reports() - except Exception as e: # noqa: BLE001 # shutdown must continue when a report drain fails - verbose_proxy_logger.error("Error draining passthrough upstream error reports: %s", e) - await _flush_spend_logs_queue_on_shutdown() await proxy_config.stop_config_sync_subscriber() diff --git a/tests/test_litellm/proxy/proxy_server/test_lifecycle.py b/tests/test_litellm/proxy/proxy_server/test_lifecycle.py index 36d2e16d261..a9d2de481c1 100644 --- a/tests/test_litellm/proxy/proxy_server/test_lifecycle.py +++ b/tests/test_litellm/proxy/proxy_server/test_lifecycle.py @@ -287,6 +287,44 @@ async def test_flush_spend_counters_on_shutdown_logs_and_swallows_commit_errors( assert "Error flushing spend counters on shutdown: db gone" in caplog.text +def test_shutdown_drains_passthrough_error_reports_before_spend_flushes(): + """Passthrough error report callbacks write spend rows through the logging + worker, so the drain must complete before the spend producer, counters and + spend-log queue are flushed or the delivered rows can be skipped. The drain + lives inside the ``proxy_startup_event`` lifespan teardown which cannot be + driven without running the whole startup, so assert the await order in + source: a revert of the ordering is what this guards. + """ + import ast + + parsed = ast.parse(inspect.getsource(ps)) + startup = next( + node + for node in parsed.body + if isinstance(node, (ast.AsyncFunctionDef, ast.FunctionDef)) and node.name == "proxy_startup_event" + ) + awaited = tuple( + child.value.func.id + for child in sorted( + ( + node + for node in ast.walk(startup) + if isinstance(node, ast.Await) + and isinstance(node.value, ast.Call) + and isinstance(node.value.func, ast.Name) + ), + key=lambda node: node.lineno, + ) + ) + drain_at = awaited.index("drain_passthrough_upstream_error_reports") + for flush in ( + "_drain_spend_event_producer_on_shutdown", + "flush_spend_counters_on_shutdown", + "_flush_spend_logs_queue_on_shutdown", + ): + assert drain_at < awaited.index(flush), f"report drain must run before {flush}" + + # --------------------------------------------------------------------------- # _initialize_shared_aiohttp_session # --------------------------------------------------------------------------- diff --git a/tests/test_litellm/test_constants.py b/tests/test_litellm/test_constants.py index 12e473f68a4..85cc63e1e54 100644 --- a/tests/test_litellm/test_constants.py +++ b/tests/test_litellm/test_constants.py @@ -68,3 +68,14 @@ def _build_constant_env_var_map() -> dict[str, str]: env_var_map[constant_name] = env_var_name return env_var_map + + +def test_passthrough_error_report_concurrency_env_zero_clamps_to_one(monkeypatch): + """A zero/negative concurrency would deadlock every report; the constant clamps to 1.""" + monkeypatch.setenv("PASSTHROUGH_UPSTREAM_ERROR_REPORT_CONCURRENCY", "0") + try: + reloaded = importlib.reload(constants) + assert reloaded.PASSTHROUGH_UPSTREAM_ERROR_REPORT_CONCURRENCY == 1 + finally: + monkeypatch.undo() + importlib.reload(constants)