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
« 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"""
7import os
8from collections.abc import Callable
9from datetime import timedelta
10from typing import TYPE_CHECKING, Any, Final
12from litellm._logging import verbose_proxy_logger
13from litellm.proxy.db.prisma_client import PrismaWrapper
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
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)
34# Top-level Prisma client methods that read from the database.
35_TOP_LEVEL_READ_METHODS: Final = frozenset({"query_first", "query_raw"})
38class _RoutedActions:
39 """Per-model accessor that sends reads to the reader and writes to the writer.
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 """
46 __slots__ = ("_reader_actions", "_should_use_reader", "_writer_actions")
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
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)
64class WriterPinnedClient:
65 """PrismaClient-shaped view whose `.db` resolves to the writer while it is available.
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).
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 """
78 __slots__ = ("db",)
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
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
89class RoutingPrismaWrapper:
90 """
91 Routes Prisma operations between a writer and a reader Prisma client.
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.
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.
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 """
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
127 @property
128 def writer(self) -> PrismaWrapper:
129 return self._writer
131 @property
132 def reader(self) -> PrismaWrapper:
133 return self._reader
135 @property
136 def read_target(self) -> PrismaWrapper:
137 """The wrapper `_TOP_LEVEL_READ_METHODS` dispatch to right now.
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
146 @property
147 def reader_unavailable(self) -> bool:
148 return self._reader_unavailable
150 @property
151 def writer_unavailable(self) -> bool:
152 return self._writer_unavailable
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
163 def _should_use_reader(self) -> bool:
164 return not self._reader_unavailable
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
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 )
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
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())
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()
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()
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.
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.
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
277 async def _recreate_reader(self, http_client: "HttpConfig | None" = None) -> None:
278 """Resolve the reader URL and recreate its Prisma client.
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)
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