Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/common_utils/registry_read_through.py: 78%
154 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"""Read-through recovery for in-memory registries in multi-replica deployments.
3A management write (POST /model/new, /guardrails, /v1/agents) lands on one
4replica and reaches Postgres, but sibling replicas only refresh their in-memory
5registries on the periodic config reload, so a request using the new object
6immediately can land on a sibling that has never heard of it and fail 400/404.
7On a registry miss, callers here fetch the missing row from the DB and load it
8into the local registry before giving up. A short negative-result TTL per key
9plus a global resync budget per window bound the DB load from lookups of
10genuinely unknown names.
11"""
13import asyncio
14import time
15from collections.abc import Awaitable, Callable
16from typing import TYPE_CHECKING, Final
18from litellm._logging import verbose_proxy_logger
19from litellm.caching.in_memory_cache import InMemoryCache
21if TYPE_CHECKING: 21 ↛ 22line 21 didn't jump to line 22 because the condition on line 21 was never true
22 from prisma.types import (
23 LiteLLM_AgentsTableInclude,
24 LiteLLM_AgentsTableWhereUniqueInput,
25 LiteLLM_GuardrailsTableWhereInput,
26 LiteLLM_ProxyModelTableWhereInput,
27 )
29 from litellm.integrations.custom_guardrail import CustomGuardrail
30 from litellm.types.agents import AgentResponse
32READ_THROUGH_MISS_TTL_SECONDS: Final = 2.0
33READ_THROUGH_RESYNC_WINDOW_SECONDS: Final = 5.0
34READ_THROUGH_MAX_RESYNCS_PER_WINDOW: Final = 20
37class RegistryReadThrough:
38 __slots__ = (
39 "_is_loaded",
40 "_lock",
41 "_max_resyncs_per_window",
42 "_miss_ttl_seconds",
43 "_recent_misses",
44 "_resync",
45 "_resync_window_seconds",
46 "_window_resyncs",
47 "_window_started_at",
48 )
50 def __init__(
51 self,
52 resync: Callable[[str], Awaitable[bool]],
53 is_loaded: Callable[[str], bool],
54 miss_ttl_seconds: float = READ_THROUGH_MISS_TTL_SECONDS,
55 max_resyncs_per_window: int = READ_THROUGH_MAX_RESYNCS_PER_WINDOW,
56 resync_window_seconds: float = READ_THROUGH_RESYNC_WINDOW_SECONDS,
57 ) -> None:
58 self._resync = resync
59 self._is_loaded = is_loaded
60 self._miss_ttl_seconds = miss_ttl_seconds
61 self._max_resyncs_per_window = max_resyncs_per_window
62 self._resync_window_seconds = resync_window_seconds
63 self._lock = asyncio.Lock()
64 self._recent_misses = InMemoryCache(max_size_in_memory=1000)
65 self._window_started_at = float("-inf")
66 self._window_resyncs = 0
68 def _consume_resync_budget(self) -> bool:
69 now: Final = time.monotonic()
70 if now - self._window_started_at >= self._resync_window_seconds:
71 self._window_started_at = now
72 self._window_resyncs = 0
73 if self._window_resyncs >= self._max_resyncs_per_window:
74 return False
75 self._window_resyncs += 1
76 return True
78 async def attempt(self, key: str) -> bool:
79 if self._recent_misses.get_cache(key) is not None:
80 return False
81 async with self._lock:
82 if self._recent_misses.get_cache(key) is not None: 82 ↛ 83line 82 didn't jump to line 83 because the condition on line 82 was never true
83 return False
84 if self._is_loaded(key): 84 ↛ 85line 84 didn't jump to line 85 because the condition on line 84 was never true
85 return True
86 if not self._consume_resync_budget():
87 verbose_proxy_logger.warning(
88 "registry read-through for %r skipped: resync budget of %s per %ss exhausted",
89 key,
90 self._max_resyncs_per_window,
91 self._resync_window_seconds,
92 )
93 return False
94 try:
95 found: Final = await self._resync(key)
96 except Exception as e: # noqa: BLE001 # a failed read-through must surface the original miss error, not a 500
97 verbose_proxy_logger.warning("registry read-through for %r failed: %s", key, e)
98 return False
99 if not found:
100 self._recent_misses.set_cache(key, True, ttl=self._miss_ttl_seconds)
101 return found
104def _db_backed_registries_enabled(object_type: str) -> bool:
105 from litellm.proxy import proxy_server
107 if proxy_server.prisma_client is None or proxy_server.store_model_in_db is not True: 107 ↛ 108line 107 didn't jump to line 108 because the condition on line 107 was never true
108 return False
109 return proxy_server.should_load_db_object(object_type=object_type)
112async def _resync_model_deployments(model_name: str) -> bool:
113 from litellm.proxy import proxy_server
114 from litellm.repositories.model_repository import ModelRepository
116 if not _db_backed_registries_enabled("models"): 116 ↛ 117line 116 didn't jump to line 117 because the condition on line 116 was never true
117 return False
118 prisma_client: Final = proxy_server.prisma_client
119 assert prisma_client is not None
120 table: Final = ModelRepository(prisma_client).table
121 name_filter: Final[LiteLLM_ProxyModelTableWhereInput] = {"model_name": model_name}
122 id_filter: Final[LiteLLM_ProxyModelTableWhereInput] = {"model_id": model_name}
123 rows: Final = await table.find_many(where=name_filter) or await table.find_many(where=id_filter)
124 if not rows:
125 return False
126 router: Final = proxy_server.llm_router
127 if router is None: 127 ↛ 128line 127 didn't jump to line 128 because the condition on line 127 was never true
128 await proxy_server.proxy_config.add_deployment(
129 prisma_client=prisma_client, proxy_logging_obj=proxy_server.proxy_logging_obj
130 )
131 return proxy_server.llm_router is not None
132 async with proxy_server.MODEL_RECONCILE_LOCK:
133 await proxy_server.proxy_config.get_credentials(prisma_client=prisma_client)
134 proxy_server.proxy_config._add_deployment(db_models=rows)
135 proxy_server.llm_model_list = router.get_model_list()
136 return True
139async def _resync_guardrails(guardrail_name: str) -> bool:
140 from litellm.proxy import proxy_server
141 from litellm.proxy.guardrails.guardrail_registry import (
142 GUARDRAIL_RECONCILE_LOCK,
143 IN_MEMORY_GUARDRAIL_HANDLER,
144 )
145 from litellm.repositories.table_repositories import GuardrailsRepository
146 from litellm.types.guardrails import Guardrail
148 if not _db_backed_registries_enabled("guardrails"): 148 ↛ 149line 148 didn't jump to line 149 because the condition on line 148 was never true
149 return False
150 prisma_client: Final = proxy_server.prisma_client
151 assert prisma_client is not None
152 active_row_filter: Final[LiteLLM_GuardrailsTableWhereInput] = {
153 "guardrail_name": guardrail_name,
154 "status": "active",
155 }
156 row: Final = await GuardrailsRepository(prisma_client).table.find_first(where=active_row_filter)
157 if row is None: 157 ↛ 159line 157 didn't jump to line 159 because the condition on line 157 was always true
158 return False
159 async with GUARDRAIL_RECONCILE_LOCK:
160 IN_MEMORY_GUARDRAIL_HANDLER.sync_guardrail_from_db(guardrail=Guardrail(**dict(row)))
161 return _initialized_guardrail(guardrail_name) is not None
164async def _resync_agents(agent_id_or_name: str) -> bool:
165 from litellm.proxy import proxy_server
166 from litellm.proxy.agent_endpoints.agent_registry import (
167 AGENT_RECONCILE_LOCK,
168 agents_table,
169 global_agent_registry,
170 )
171 from litellm.types.agents import AgentResponse
173 if not _db_backed_registries_enabled("agents"): 173 ↛ 174line 173 didn't jump to line 174 because the condition on line 173 was never true
174 return False
175 if _agent_from_registry(agent_id_or_name) is not None: 175 ↛ 176line 175 didn't jump to line 176 because the condition on line 175 was never true
176 return True
177 prisma_client: Final = proxy_server.prisma_client
178 assert prisma_client is not None
179 table: Final = agents_table(prisma_client)
180 id_filter: Final[LiteLLM_AgentsTableWhereUniqueInput] = {"agent_id": agent_id_or_name}
181 name_filter: Final[LiteLLM_AgentsTableWhereUniqueInput] = {"agent_name": agent_id_or_name}
182 include_permission: Final[LiteLLM_AgentsTableInclude] = {"object_permission": True}
183 async with AGENT_RECONCILE_LOCK:
184 if _agent_from_registry(agent_id_or_name) is not None: 184 ↛ 185line 184 didn't jump to line 185 because the condition on line 184 was never true
185 return True
186 row: Final = await table.find_unique(where=id_filter, include=include_permission) or await table.find_unique(
187 where=name_filter, include=include_permission
188 )
189 if row is None: 189 ↛ 191line 189 didn't jump to line 191 because the condition on line 189 was always true
190 return False
191 global_agent_registry.register_agent(agent_config=AgentResponse.model_validate(row.model_dump()))
192 return True
195def _model_is_loaded(model_name_or_id: str) -> bool:
196 from litellm.proxy import proxy_server
198 router: Final = proxy_server.llm_router
199 if router is None: 199 ↛ 200line 199 didn't jump to line 200 because the condition on line 199 was never true
200 return False
201 return model_name_or_id in router.model_names or router.has_model_id(model_name_or_id)
204def _guardrail_is_loaded(guardrail_name: str) -> bool:
205 return _initialized_guardrail(guardrail_name) is not None
208def _agent_is_loaded(agent_id_or_name: str) -> bool:
209 return _agent_from_registry(agent_id_or_name) is not None
212model_registry_read_through: Final = RegistryReadThrough(resync=_resync_model_deployments, is_loaded=_model_is_loaded)
213guardrail_registry_read_through: Final = RegistryReadThrough(resync=_resync_guardrails, is_loaded=_guardrail_is_loaded)
214agent_registry_read_through: Final = RegistryReadThrough(resync=_resync_agents, is_loaded=_agent_is_loaded)
217def _agent_from_registry(agent_id_or_name: str) -> "AgentResponse | None":
218 from litellm.proxy.agent_endpoints.agent_registry import global_agent_registry
220 by_id: Final = global_agent_registry.get_agent_by_id(agent_id=agent_id_or_name)
221 if by_id is not None: 221 ↛ 222line 221 didn't jump to line 222 because the condition on line 221 was never true
222 return by_id
223 return global_agent_registry.get_agent_by_name(agent_name=agent_id_or_name)
226async def get_agent_with_read_through(agent_id_or_name: str) -> "AgentResponse | None":
227 agent: Final = _agent_from_registry(agent_id_or_name)
228 if agent is not None: 228 ↛ 229line 228 didn't jump to line 229 because the condition on line 228 was never true
229 return agent
230 if not await agent_registry_read_through.attempt(agent_id_or_name): 230 ↛ 232line 230 didn't jump to line 232 because the condition on line 230 was always true
231 return None
232 return _agent_from_registry(agent_id_or_name)
235def _initialized_guardrail(guardrail_name: str) -> "CustomGuardrail | None":
236 from litellm.proxy.guardrails import guardrail_endpoints
238 return guardrail_endpoints.GUARDRAIL_REGISTRY.get_initialized_guardrail_callback(guardrail_name=guardrail_name)
241async def get_initialized_guardrail_with_read_through(guardrail_name: str) -> "CustomGuardrail | None":
242 active: Final = _initialized_guardrail(guardrail_name)
243 if active is not None: 243 ↛ 244line 243 didn't jump to line 244 because the condition on line 243 was never true
244 return active
245 if not await guardrail_registry_read_through.attempt(guardrail_name): 245 ↛ 247line 245 didn't jump to line 247 because the condition on line 245 was always true
246 return None
247 return _initialized_guardrail(guardrail_name)