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

1"""Collector sidecar: consume spend events from the pod's inference workers and run the cost pipeline. 

2 

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

13 

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: 

20 

21 python -m litellm.proxy.collector [--address unix:///path.sock] 

22""" 

23 

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 

32 

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) 

50 

51MAX_EVENT_BYTES: Final = 64 * 1024 * 1024 

52 

53 

54class SpendEventConsumer: 

55 """Accepts producer connections and runs ``handler`` on every line each one sends, in order.""" 

56 

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 

65 

66 @property 

67 def received(self) -> int: 

68 return self._received 

69 

70 @property 

71 def handled(self) -> int: 

72 return self._handled 

73 

74 @property 

75 def failed(self) -> int: 

76 return self._failed 

77 

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) 

87 

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

105 

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

113 

114 async def drain(self, timeout: float) -> int: 

115 """Half-close every producer connection, then keep reading until each producer hangs up or ``timeout``. 

116 

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) 

131 

132 

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) 

136 

137 

138async def run_collector(address: CollectorAddress, drain_timeout: float) -> None: 

139 from fastapi import FastAPI 

140 

141 from litellm.proxy.hooks.proxy_track_cost_callback import run_spend_event 

142 from litellm.proxy.proxy_server import proxy_startup_event 

143 

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 ) 

160 

161 

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

170 

171 

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) 

179 

180 

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. 

185 

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) 

195 

196 

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

217 

218 

219if __name__ == "__main__": 

220 main(sys.argv[1:])