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>
This commit is contained in:
yucheng 2026-09-25 18:09:30 +00:00
parent 3d2bb748f5
commit f7c157083a
4 changed files with 56 additions and 7 deletions

View file

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

View file

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

View file

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

View file

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