Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/db/routing_prisma_wrapper.py: 27%

132 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 12:01 +0000

1""" 

2RoutingPrismaWrapper: routes Prisma reads to a read-replica client and writes 

3to a writer client. Used when DATABASE_URL_READ_REPLICA is configured; 

4otherwise PrismaClient uses the writer-only PrismaWrapper directly. 

5""" 

6 

7import os 

8from collections.abc import Callable 

9from datetime import timedelta 

10from typing import TYPE_CHECKING, Any, Final 

11 

12from litellm._logging import verbose_proxy_logger 

13from litellm.proxy.db.prisma_client import PrismaWrapper 

14 

15if TYPE_CHECKING: 15 ↛ 16line 15 didn't jump to line 16 because the condition on line 15 was never true

16 from prisma.types import HttpConfig 

17 

18# Per-model action methods that read from the database. These are routed to 

19# the read replica when one is configured. 

20_MODEL_READ_METHODS: Final = frozenset( 

21 { 

22 "find_first", 

23 "find_first_or_raise", 

24 "find_many", 

25 "find_unique", 

26 "find_unique_or_raise", 

27 "count", 

28 "group_by", 

29 "query_first", 

30 "query_raw", 

31 } 

32) 

33 

34# Top-level Prisma client methods that read from the database. 

35_TOP_LEVEL_READ_METHODS: Final = frozenset({"query_first", "query_raw"}) 

36 

37 

38class _RoutedActions: 

39 """Per-model accessor that sends reads to the reader and writes to the writer. 

40 

41 `should_use_reader` is consulted on every read dispatch so a mid-call flip 

42 of the routing wrapper's reader-availability flag (e.g. after the reader 

43 fails a recreate) is observed without re-fetching the actions accessor. 

44 """ 

45 

46 __slots__ = ("_reader_actions", "_should_use_reader", "_writer_actions") 

47 

48 def __init__( 

49 self, 

50 writer_actions: object, 

51 reader_actions: object, 

52 should_use_reader: Callable[[], bool], 

53 ): 

54 self._writer_actions = writer_actions 

55 self._reader_actions = reader_actions 

56 self._should_use_reader = should_use_reader 

57 

58 def __getattr__(self, name: str) -> object: 

59 if name in _MODEL_READ_METHODS and self._should_use_reader(): 

60 return getattr(self._reader_actions, name) 

61 return getattr(self._writer_actions, name) 

62 

63 

64class WriterPinnedClient: 

65 """PrismaClient-shaped view whose `.db` resolves to the writer while it is available. 

66 

67 Read-after-write paths (e.g. the model reconcile a /model/new triggers to 

68 verify its own just-committed row) must not read through a lagging read 

69 replica: the row is not replayed there yet, so the reconcile concludes the 

70 write is missing and fails the request even though it is durable (#38556). 

71 

72 While the writer is degraded (`writer_unavailable`), the pin yields to the 

73 routed wrapper so reconcile reads keep working from the replica: a proxy 

74 that starts during a primary outage must still load DB-backed models, and 

75 no read-after-write hazard exists then because writes are failing anyway. 

76 """ 

77 

78 __slots__ = ("db",) 

79 

80 def __init__(self, db: "PrismaWrapper | RoutingPrismaWrapper") -> None: 

81 self.db: Final = db.writer if isinstance(db, RoutingPrismaWrapper) and not db.writer_unavailable else db 

82 

83 

84def writer_wrapper(db: "PrismaWrapper | RoutingPrismaWrapper") -> PrismaWrapper: 

85 """Unlike `WriterPinnedClient`, ignores `writer_unavailable`: a raw SQL write has no replica fallback.""" 

86 return db.writer if isinstance(db, RoutingPrismaWrapper) else db 

87 

88 

89class RoutingPrismaWrapper: 

90 """ 

91 Routes Prisma operations between a writer and a reader Prisma client. 

92 

93 Reads (find_*, count, group_by, query_raw, query_first) go to the reader; 

94 everything else (writes, transactions, raw execute) goes to the writer. 

95 Lifecycle methods (connect, disconnect, IAM token refresh) act on both 

96 clients so callers do not need to know about the split. When 

97 IAM_TOKEN_DB_AUTH is enabled, both writer and reader refresh their tokens 

98 independently on their own ~12-minute cadence. 

99 

100 Reader degradation: a reader-side failure (failed connect, failed 

101 recreate) is non-fatal — the wrapper sets `_reader_unavailable=True`, logs 

102 a warning, and routes subsequent reads to the writer. The next successful 

103 `connect()` or `recreate_prisma_client()` clears the flag. This keeps the 

104 proxy serving traffic during transient reader outages instead of failing 

105 startup or returning errors for read-heavy endpoints. 

106 

107 Writer degradation: a writer-side `connect()` failure while the reader 

108 connects is likewise non-fatal — the wrapper sets 

109 `_writer_unavailable=True`, logs a warning, and keeps serving reads from 

110 the reader (key lookups, DB-stored model loads) so a proxy that starts 

111 during a primary outage still serves inference from the replica. Writes 

112 fail at call time until the writer recovers; the PrismaClient DB health 

113 watchdog polls `writer_unavailable` and drives the writer reconnect, 

114 which clears the flag via `recreate_prisma_client`. Only when BOTH sides 

115 fail to connect does `connect()` raise (full DB outage — the existing 

116 `allow_requests_on_db_unavailable` startup handling applies). 

117 """ 

118 

119 def __init__(self, writer: PrismaWrapper, reader: PrismaWrapper): 

120 self._writer = writer 

121 self._reader = reader 

122 # When True, reads fall back to the writer. Flipped on by reader 

123 # connect/recreate failures and flipped off on the next reader recovery. 

124 self._reader_unavailable: bool = False 

125 self._writer_unavailable: bool = False 

126 

127 @property 

128 def writer(self) -> PrismaWrapper: 

129 return self._writer 

130 

131 @property 

132 def reader(self) -> PrismaWrapper: 

133 return self._reader 

134 

135 @property 

136 def read_target(self) -> PrismaWrapper: 

137 """The wrapper `_TOP_LEVEL_READ_METHODS` dispatch to right now. 

138 

139 Callers that need to reason about the engine a read actually ran on 

140 (e.g. recovering from prepared statements that went stale on it) must 

141 consult this rather than `writer`, and `__getattr__` routes through it 

142 so the two cannot drift apart. 

143 """ 

144 return self._writer if self._reader_unavailable else self._reader 

145 

146 @property 

147 def reader_unavailable(self) -> bool: 

148 return self._reader_unavailable 

149 

150 @property 

151 def writer_unavailable(self) -> bool: 

152 return self._writer_unavailable 

153 

154 def mark_writer_recovered(self) -> None: 

155 """Clear the degraded-writer flag after an external health probe proved 

156 the writer reachable. Needed when recovery happens without 

157 `recreate_prisma_client` (e.g. an IAM token refresh already recreated 

158 the writer engine), which is otherwise the only runtime path that 

159 clears the flag — without this, the watchdog would keep firing 

160 reconnect attempts against an already-healthy writer.""" 

161 self._writer_unavailable = False 

162 

163 def _should_use_reader(self) -> bool: 

164 return not self._reader_unavailable 

165 

166 @staticmethod 

167 async def _try_connect(client: PrismaWrapper, timeout: int | timedelta | None = None) -> Exception | None: 

168 if client.is_connected() is True: 

169 return None 

170 try: 

171 await client.connect(timeout) 

172 return None 

173 except Exception as e: 

174 return e 

175 

176 async def connect(self, timeout: int | timedelta | None = None) -> None: 

177 writer_error: Final = await self._try_connect(self._writer, timeout) 

178 if writer_error is None: 

179 self._writer_unavailable = False 

180 verbose_proxy_logger.info("[writer] DB connected") 

181 reader_error: Final = await self._try_connect(self._reader, timeout) 

182 if reader_error is None: 

183 self._reader_unavailable = False 

184 verbose_proxy_logger.info("[reader] DB connected") 

185 if writer_error is None and reader_error is None: 

186 return 

187 if writer_error is not None and reader_error is not None: 

188 raise writer_error 

189 if reader_error is not None: 

190 # Degrade gracefully: the proxy keeps serving traffic with reads 

191 # routed to the writer until the reader endpoint is reachable. 

192 # Aborting startup here would tie proxy availability to an 

193 # opt-in, best-effort reader endpoint. 

194 self._reader_unavailable = True 

195 verbose_proxy_logger.warning( 

196 "Failed to connect to read replica DB: %s. " 

197 "Falling back to the writer for reads until the reader is reachable.", 

198 reader_error, 

199 ) 

200 return 

201 self._writer_unavailable = True 

202 verbose_proxy_logger.warning( 

203 "Failed to connect to primary (writer) DB: %s. " 

204 "Serving reads from the read replica; writes will fail until the writer recovers.", 

205 writer_error, 

206 ) 

207 

208 async def disconnect(self, timeout: float | timedelta | None = None) -> None: 

209 first_error: BaseException | None = None 

210 for client in (self._writer, self._reader): 

211 try: 

212 await client.disconnect(timeout) 

213 except Exception as e: 

214 if first_error is None: 

215 first_error = e 

216 verbose_proxy_logger.warning("Error disconnecting Prisma client: %s", e) 

217 if first_error is not None: 

218 raise first_error 

219 

220 def is_connected(self) -> bool: 

221 # Reflects writer health only. The reader is best-effort; its 

222 # availability is tracked via `_reader_unavailable` and a degraded 

223 # reader must NOT cause a writer reconnect (would loop indefinitely 

224 # since recreate_prisma_client only fixes writer-side problems). 

225 return bool(self._writer.is_connected()) 

226 

227 async def start_token_refresh_task(self) -> None: 

228 await self._writer.start_token_refresh_task() 

229 await self._reader.start_token_refresh_task() 

230 

231 async def stop_token_refresh_task(self) -> None: 

232 await self._writer.stop_token_refresh_task() 

233 await self._reader.stop_token_refresh_task() 

234 

235 async def recreate_prisma_client( 

236 self, 

237 new_db_url: str, 

238 http_client: "HttpConfig | None" = None, 

239 *, 

240 expected_generation: int | None = None, 

241 ) -> bool: 

242 """Recreate both writer and reader Prisma clients. 

243 

244 The writer reconnect path in PrismaClient calls 

245 `self.db.recreate_prisma_client(...)`. Without this method, a DB-wide 

246 connectivity event would only re-create the writer; the reader engine 

247 would stay broken and every routed read would fail. We always recreate 

248 the writer first (its URL is the one passed in), then best-effort 

249 recreate the reader. A reader failure flips `_reader_unavailable=True` 

250 so reads transparently fall through to the writer. 

251 

252 `expected_generation` is forwarded to the writer's optimistic-lock 

253 guard. If the writer recreate is skipped (another path already replaced 

254 the engine — issue #29176), we skip the reader too rather than churning 

255 it needlessly, and return ``False``. 

256 """ 

257 writer_recreated: Final = await self._writer.recreate_prisma_client( 

258 new_db_url, 

259 http_client=http_client, 

260 expected_generation=expected_generation, 

261 ) 

262 if not writer_recreated: 

263 return False 

264 self._writer_unavailable = False 

265 try: 

266 await self._recreate_reader(http_client=http_client) 

267 self._reader_unavailable = False 

268 except Exception as e: 

269 self._reader_unavailable = True 

270 verbose_proxy_logger.warning( 

271 "Failed to recreate reader Prisma client: %s. " 

272 "Reads will fall back to the writer until the reader recovers.", 

273 e, 

274 ) 

275 return True 

276 

277 async def _recreate_reader(self, http_client: "HttpConfig | None" = None) -> None: 

278 """Resolve the reader URL and recreate its Prisma client. 

279 

280 Token-authenticated readers regenerate their token (host/port/user came 

281 from the parsed reader URL at construction time). Password-authenticated 

282 readers reuse the URL stored in `DATABASE_URL_READ_REPLICA`. 

283 """ 

284 if self._reader.iam_token_db_auth: 

285 new_reader_url: Final = self._reader.get_rds_iam_token() 

286 if not new_reader_url: 

287 raise RuntimeError(f"Failed to generate fresh {self._reader.token_label} for read replica") 

288 await self._reader.recreate_prisma_client(new_reader_url, http_client=http_client) 

289 return 

290 reader_url: Final = os.getenv("DATABASE_URL_READ_REPLICA", "") 

291 if not reader_url: 

292 raise RuntimeError("DATABASE_URL_READ_REPLICA not set; cannot recreate read replica client") 

293 await self._reader.recreate_prisma_client(reader_url, http_client=http_client) 

294 

295 def __getattr__(self, name: str) -> Any: 

296 if name in _TOP_LEVEL_READ_METHODS: 

297 return getattr(self.read_target, name) 

298 writer_attr: Final[object] = getattr(self._writer, name) 

299 # Per-model action accessors are non-callable instances that expose 

300 # both `find_many` and `create`. Methods like execute_raw / batch_ / 

301 # tx are callables and stay on the writer untouched. 

302 if not callable(writer_attr) and hasattr(writer_attr, "find_many") and hasattr(writer_attr, "create"): 

303 try: 

304 reader_attr: Final[object] = getattr(self._reader, name) 

305 except AttributeError: 

306 return writer_attr 

307 return _RoutedActions(writer_attr, reader_attr, self._should_use_reader) 

308 return writer_attr