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

1"""Read-through recovery for in-memory registries in multi-replica deployments. 

2 

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

12 

13import asyncio 

14import time 

15from collections.abc import Awaitable, Callable 

16from typing import TYPE_CHECKING, Final 

17 

18from litellm._logging import verbose_proxy_logger 

19from litellm.caching.in_memory_cache import InMemoryCache 

20 

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 ) 

28 

29 from litellm.integrations.custom_guardrail import CustomGuardrail 

30 from litellm.types.agents import AgentResponse 

31 

32READ_THROUGH_MISS_TTL_SECONDS: Final = 2.0 

33READ_THROUGH_RESYNC_WINDOW_SECONDS: Final = 5.0 

34READ_THROUGH_MAX_RESYNCS_PER_WINDOW: Final = 20 

35 

36 

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 ) 

49 

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 

67 

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 

77 

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 

102 

103 

104def _db_backed_registries_enabled(object_type: str) -> bool: 

105 from litellm.proxy import proxy_server 

106 

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) 

110 

111 

112async def _resync_model_deployments(model_name: str) -> bool: 

113 from litellm.proxy import proxy_server 

114 from litellm.repositories.model_repository import ModelRepository 

115 

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 

137 

138 

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 

147 

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 

162 

163 

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 

172 

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 

193 

194 

195def _model_is_loaded(model_name_or_id: str) -> bool: 

196 from litellm.proxy import proxy_server 

197 

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) 

202 

203 

204def _guardrail_is_loaded(guardrail_name: str) -> bool: 

205 return _initialized_guardrail(guardrail_name) is not None 

206 

207 

208def _agent_is_loaded(agent_id_or_name: str) -> bool: 

209 return _agent_from_registry(agent_id_or_name) is not None 

210 

211 

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) 

215 

216 

217def _agent_from_registry(agent_id_or_name: str) -> "AgentResponse | None": 

218 from litellm.proxy.agent_endpoints.agent_registry import global_agent_registry 

219 

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) 

224 

225 

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) 

233 

234 

235def _initialized_guardrail(guardrail_name: str) -> "CustomGuardrail | None": 

236 from litellm.proxy.guardrails import guardrail_endpoints 

237 

238 return guardrail_endpoints.GUARDRAIL_REGISTRY.get_initialized_guardrail_callback(guardrail_name=guardrail_name) 

239 

240 

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)