Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/collector.py: 0%
132 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
1"""Collector sidecar: consume spend events from the pod's inference workers and run the cost pipeline.
3Runs the proxy startup lifespan (config, Prisma, Redis transaction buffer, scheduled spend flushes)
4without serving HTTP, then listens on ``LITELLM_COLLECTOR_ADDRESS`` for newline-delimited spend
5events. Each event goes through the unchanged ``_ProxyDBLogger._PROXY_track_cost_callback``, so
6spend logs, spend counters, budget reservation reconciliation and cache updates happen exactly as
7they would in-process, just in this container. Events are handled in order per producer connection
8(one per uvicorn worker); a slow pipeline fills the socket buffer and the producer's bounded queue,
9which is the backpressure that triggers its fallback or drop policy. ``SIGTERM`` stops accepting
10connections, half-closes every producer connection so the producers switch to their unavailable
11policy, finishes the events already sent, then runs the proxy shutdown (which flushes the buffered
12spend transactions).
14``DATABASE_URL`` is assembled from the same ``DATABASE_*`` inputs as the proxy container, and when
15``LITELLM_PGBOUNCER_ENABLED`` is set it points at the PgBouncer that container already runs on the
16pod's loopback, so the sidecar must see the same env as the proxy. Under ``IAM_TOKEN_DB_AUTH`` or
17``AZURE_POSTGRESQL_AUTH`` that PgBouncer only accepts the token the proxy container minted, so the
18sidecar goes to Postgres directly and mints its own. Works from any image that has ``litellm``
19installed:
21 python -m litellm.proxy.collector [--address unix:///path.sock]
22"""
24import asyncio
25import logging
26import os
27import signal
28import sys
29from collections.abc import Awaitable, Callable, Mapping, Sequence
30from pathlib import Path
31from typing import Final
33from litellm._logging import verbose_logger, verbose_proxy_logger, verbose_router_logger
34from litellm.proxy.db.db_url_settings import DatabaseURLSettings
35from litellm.proxy.db.pgbouncer import (
36 PgBouncerError,
37 PgBouncerSettings,
38 export_pooled_database_url,
39 pooled_database_url,
40)
41from litellm.proxy.spend_tracking.spend_event_producer import (
42 COLLECTOR_JOB_ROLE,
43 AddressError,
44 CollectorAddress,
45 CollectorSettings,
46 TcpAddress,
47 UnixAddress,
48 parse_collector_address,
49)
51MAX_EVENT_BYTES: Final = 64 * 1024 * 1024
54class SpendEventConsumer:
55 """Accepts producer connections and runs ``handler`` on every line each one sends, in order."""
57 def __init__(self, handler: Callable[[bytes], Awaitable[None]]) -> None:
58 self._handler = handler
59 self._open_connections: set[asyncio.StreamWriter] = set() # mutable-ok: live producer connections
60 self._idle = asyncio.Event()
61 self._idle.set()
62 self._received = 0
63 self._handled = 0
64 self._failed = 0
66 @property
67 def received(self) -> int:
68 return self._received
70 @property
71 def handled(self) -> int:
72 return self._handled
74 @property
75 def failed(self) -> int:
76 return self._failed
78 async def serve(self, address: CollectorAddress) -> asyncio.Server:
79 match address:
80 case UnixAddress(path=path):
81 socket_path: Final = Path(path)
82 socket_path.parent.mkdir(parents=True, exist_ok=True)
83 socket_path.unlink(missing_ok=True)
84 return await asyncio.start_unix_server(self._on_connection, path=path, limit=MAX_EVENT_BYTES)
85 case TcpAddress(host=host, port=port):
86 return await asyncio.start_server(self._on_connection, host=host, port=port, limit=MAX_EVENT_BYTES)
88 async def _on_connection(self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter) -> None:
89 self._open_connections.add(writer)
90 self._idle.clear()
91 try:
92 while line := await reader.readline():
93 if not line.endswith(b"\n"):
94 verbose_proxy_logger.error("collector: discarding truncated spend event (%d bytes)", len(line))
95 break
96 self._received += 1
97 await self._handle(line)
98 except (ConnectionError, asyncio.IncompleteReadError, asyncio.LimitOverrunError) as error:
99 verbose_proxy_logger.warning("collector: producer connection ended abnormally: %s", error)
100 finally:
101 writer.close()
102 self._open_connections.discard(writer)
103 if not self._open_connections:
104 self._idle.set()
106 async def _handle(self, line: bytes) -> None:
107 try:
108 await self._handler(line)
109 self._handled += 1
110 except Exception: # noqa: BLE001 # the cost pipeline raises anything; one bad event must not stop the sidecar
111 self._failed += 1
112 verbose_proxy_logger.exception("collector: spend event failed")
114 async def drain(self, timeout: float) -> int:
115 """Half-close every producer connection, then keep reading until each producer hangs up or ``timeout``.
117 Returns how many producer connections were still open when the timeout hit.
118 """
119 for writer in tuple(self._open_connections):
120 if writer.is_closing() or not writer.can_write_eof():
121 continue
122 try:
123 writer.write_eof()
124 except (OSError, RuntimeError) as error:
125 verbose_proxy_logger.debug("collector: producer already gone before half-close: %s", error)
126 try:
127 await asyncio.wait_for(self._idle.wait(), timeout)
128 except TimeoutError:
129 pass
130 return len(self._open_connections)
133def _install_stop_signals(loop: asyncio.AbstractEventLoop, stop: asyncio.Event) -> None:
134 for signum in (signal.SIGTERM, signal.SIGINT):
135 loop.add_signal_handler(signum, stop.set)
138async def run_collector(address: CollectorAddress, drain_timeout: float) -> None:
139 from fastapi import FastAPI
141 from litellm.proxy.hooks.proxy_track_cost_callback import run_spend_event
142 from litellm.proxy.proxy_server import proxy_startup_event
144 stop: Final = asyncio.Event()
145 _install_stop_signals(asyncio.get_running_loop(), stop)
146 consumer: Final = SpendEventConsumer(handler=run_spend_event)
147 async with proxy_startup_event(FastAPI()):
148 server: Final = await consumer.serve(address)
149 verbose_proxy_logger.info("collector: listening on %s", address)
150 await stop.wait()
151 server.close()
152 still_open: Final = await consumer.drain(drain_timeout)
153 verbose_proxy_logger.info(
154 "collector: stopping. received=%d handled=%d failed=%d connections_cut=%d",
155 consumer.received,
156 consumer.handled,
157 consumer.failed,
158 still_open,
159 )
162def address_argument(argv: Sequence[str], default: str) -> str | AddressError:
163 match tuple(argv):
164 case ():
165 return default
166 case ("--address", value):
167 return value
168 case _:
169 return AddressError(f"usage: python -m litellm.proxy.collector [--address ADDRESS], got {tuple(argv)}")
172def apply_log_level(litellm_log: str | None) -> None:
173 """Mirror the proxy's ``LITELLM_LOG`` handling: the sidecar has no CLI flags to turn logging on."""
174 level: Final = logging.getLevelNamesMapping().get((litellm_log or "").upper())
175 if level is None:
176 return
177 for logger in (verbose_logger, verbose_router_logger, verbose_proxy_logger):
178 logger.setLevel(level)
181def pod_pgbouncer_database_url(
182 pgbouncer: PgBouncerSettings, environ: Mapping[str, str], *, token_auth: bool
183) -> str | PgBouncerError | None:
184 """The proxy container's PgBouncer URL for ``environ["DATABASE_URL"]``, or None to connect to Postgres directly.
186 Direct is the answer when PgBouncer is off, and also under token auth: that PgBouncer's auth file
187 only holds the token its own container minted, which this container cannot present.
188 """
189 if not pgbouncer.enabled or token_auth:
190 return None
191 upstream_url: Final = environ.get("DATABASE_URL")
192 if upstream_url is None:
193 return PgBouncerError("LITELLM_PGBOUNCER_ENABLED is set but no DATABASE_URL could be assembled")
194 return pooled_database_url(upstream_url, pgbouncer)
197def main(argv: Sequence[str]) -> None:
198 os.environ.setdefault("LITELLM_JOB_ROLE", COLLECTOR_JOB_ROLE)
199 apply_log_level(os.environ.get("LITELLM_LOG"))
200 database: Final = DatabaseURLSettings.from_env()
201 database.apply_to_env()
202 pooled: Final = pod_pgbouncer_database_url(
203 PgBouncerSettings(),
204 os.environ,
205 token_auth=database.iam_token_db_auth or database.azure_postgresql_auth,
206 )
207 if isinstance(pooled, PgBouncerError):
208 sys.exit(f"LiteLLM collector: cannot use the pod's pgbouncer: {pooled.reason}")
209 if pooled is not None:
210 export_pooled_database_url(pooled)
211 settings: Final = CollectorSettings()
212 raw_address: Final = address_argument(argv, default=settings.address)
213 address: Final = raw_address if isinstance(raw_address, AddressError) else parse_collector_address(raw_address)
214 if isinstance(address, AddressError):
215 sys.exit(f"LiteLLM collector: {address.reason}")
216 asyncio.run(run_collector(address, drain_timeout=settings.drain_timeout_seconds))
219if __name__ == "__main__":
220 main(sys.argv[1:])