Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/management_endpoints/coordination_redis_endpoints.py: 60%

170 statements  

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

1""" 

2COORDINATION REDIS SETTINGS MANAGEMENT 

3 

4Endpoints for managing `general_settings.coordination_redis` - the standalone 

5Redis the proxy uses for cross-pod coordination (tpm/rpm rate limits, spend 

6tracking, pod lock manager, shared health checks), configured independently of 

7the response-cache backend. 

8 

9GET /coordination_redis/settings - Get the coordination Redis settings, field metadata, and which source is active 

10POST /coordination_redis/settings - Save coordination Redis settings to the database 

11POST /coordination_redis/settings/test - Test a coordination Redis connection with the provided credentials 

12""" 

13 

14import asyncio 

15import json 

16from collections.abc import Mapping 

17from contextlib import suppress 

18from datetime import datetime, timezone 

19from typing import Final 

20 

21from fastapi import APIRouter, Depends, Header, HTTPException 

22from pydantic import BaseModel, Field, TypeAdapter, ValidationError 

23 

24import litellm 

25from litellm._logging import verbose_proxy_logger 

26from litellm._uuid import uuid 

27from litellm.caching.caching import RedisCache 

28from litellm.caching.redis_cluster_cache import RedisClusterCache 

29from litellm.proxy._types import ( 

30 AUDIT_ACTIONS, 

31 CoordinationRedisParams, 

32 LiteLLM_AuditLogs, 

33 LitellmTableNames, 

34 LitellmUserRoles, 

35 UserAPIKeyAuth, 

36 user_api_key_has_admin_view, 

37) 

38from litellm.proxy.auth.user_api_key_auth import user_api_key_auth 

39from litellm.proxy.utils import invalidate_config_param 

40from litellm.repositories.config_repository import ConfigRepository 

41from litellm.secret_managers.main import get_secret_str 

42from litellm.types.management_endpoints import ( 

43 COORDINATION_REDIS_SETTINGS_FIELDS, 

44 CoordinationRedisSettingsField, 

45 CoordinationRedisSource, 

46) 

47 

48router: Final = APIRouter() 

49 

50_GENERAL_SETTINGS_PARAM_NAME: Final = "general_settings" 

51_COORDINATION_REDIS_KEY: Final = "coordination_redis" 

52 

53# Fields that carry credentials. Redacted on read so a plaintext Redis / 

54# Sentinel password never leaves the server, and scrubbed out of connection-test 

55# error strings. `url` is here because a Redis url can embed a password inline 

56# (e.g. redis://:secret@host:6379/1). 

57_SENSITIVE_FIELDS: Final[frozenset[str]] = frozenset({"password", "sentinel_password", "url"}) 

58 

59_REDACTED_VALUE: Final = "***REDACTED***" 

60 

61_ENV_REF_PREFIX: Final = "os.environ/" 

62 

63_PING_TIMEOUT_SECONDS: Final = 5.0 

64 

65_SETTINGS_ADAPTER: Final[TypeAdapter[dict[str, object]]] = TypeAdapter(dict[str, object]) 

66 

67 

68def _enforce_proxy_admin(user_api_key_dict: UserAPIKeyAuth) -> None: 

69 if user_api_key_dict.user_role != LitellmUserRoles.PROXY_ADMIN: 69 ↛ 70line 69 didn't jump to line 70 because the condition on line 69 was never true

70 raise HTTPException( 

71 status_code=403, 

72 detail={"error": "Only proxy admins can manage coordination Redis settings"}, 

73 ) 

74 

75 

76def _resolve_env_ref(value: object) -> object: 

77 """Resolve an `os.environ/VAR` reference to its value, passing anything else through.""" 

78 if isinstance(value, str) and value.startswith(_ENV_REF_PREFIX): 78 ↛ 79line 78 didn't jump to line 79 because the condition on line 78 was never true

79 return get_secret_str(value) 

80 return value 

81 

82 

83def _resolve_env_refs(settings: Mapping[str, object]) -> dict[str, object]: 

84 return {key: _resolve_env_ref(value) for key, value in settings.items()} 

85 

86 

87def _redact_credentials(settings: Mapping[str, object]) -> dict[str, object]: 

88 """Replace credential-bearing values with a fixed marker, keeping the rest intact.""" 

89 return { 

90 key: (_REDACTED_VALUE if key in _SENSITIVE_FIELDS and value is not None else value) 

91 for key, value in settings.items() 

92 } 

93 

94 

95def _redact_all_values(settings: Mapping[str, object] | None) -> dict[str, object]: 

96 """Replace every value with a fixed marker, preserving the key set. 

97 

98 The audit row shows *which* fields changed without the audit table becoming 

99 a credential-harvest sink. 

100 """ 

101 if not settings: 

102 return {} 

103 return {key: _REDACTED_VALUE for key in settings} 

104 

105 

106def _credential_values(settings: Mapping[str, object]) -> tuple[str, ...]: 

107 return tuple( 

108 str(value) for key, value in settings.items() if key in _SENSITIVE_FIELDS and isinstance(value, (str, int)) 

109 ) 

110 

111 

112def _scrub_credentials(message: str, settings: Mapping[str, object]) -> str: 

113 """Strip any credential value the caller supplied out of an error string. 

114 

115 Redis client errors routinely echo the connection url (password inline) or 

116 the auth error back to the caller. 

117 """ 

118 scrubbed = message 

119 for secret in _credential_values(settings): 

120 if secret: 

121 scrubbed = scrubbed.replace(secret, _REDACTED_VALUE) 

122 return scrubbed 

123 

124 

125def _merge_over_saved( 

126 incoming: Mapping[str, object], 

127 saved: Mapping[str, object], 

128) -> dict[str, object]: 

129 """Restore the real credential behind every value the caller echoed back redacted. 

130 

131 GET returns credentials as ``***REDACTED***``; an admin who edits the 

132 non-secret fields and re-submits would otherwise test (and save) the marker 

133 as the password. 

134 """ 

135 return { 

136 key: (saved[key] if value == _REDACTED_VALUE and key in saved else value) for key, value in incoming.items() 

137 } 

138 

139 

140def _validated_params(settings: Mapping[str, object]) -> CoordinationRedisParams: 

141 """Validate settings the way startup does: resolve env refs, then require a connection target.""" 

142 try: 

143 params: Final = CoordinationRedisParams.model_validate(_resolve_env_refs(settings)) 

144 except ValidationError as e: 

145 invalid_fields: Final = sorted({str(error["loc"][0]) for error in e.errors() if error["loc"]}) 

146 raise HTTPException( 

147 status_code=400, 

148 detail={"error": f"Invalid coordination_redis settings for fields: {invalid_fields}"}, 

149 ) 

150 if not params.has_connection_target(): 150 ↛ 160line 150 didn't jump to line 160 because the condition on line 150 was always true

151 raise HTTPException( 

152 status_code=400, 

153 detail={ 

154 "error": ( 

155 "coordination_redis needs a connection target: " 

156 "set one of host, url, startup_nodes, or sentinel_nodes" 

157 ) 

158 }, 

159 ) 

160 return params 

161 

162 

163async def _read_general_settings() -> dict[str, object]: 

164 """Read the persisted `general_settings` config row (empty when unset or no DB).""" 

165 from litellm.proxy.proxy_server import prisma_client 

166 

167 if prisma_client is None: 167 ↛ 168line 167 didn't jump to line 168 because the condition on line 167 was never true

168 return {} 

169 config_param: Final = await ConfigRepository(prisma_client).get_param(_GENERAL_SETTINGS_PARAM_NAME) 

170 if config_param is None or config_param.param_value is None: 

171 return {} 

172 return _SETTINGS_ADAPTER.validate_python(config_param.param_value) 

173 

174 

175async def get_persisted_coordination_redis_settings() -> dict[str, object] | None: 

176 """The coordination_redis block saved to the database, if any. 

177 

178 Read at startup so settings saved from the admin UI take effect on the next 

179 boot, and used here so a read reports what the proxy would boot with. 

180 """ 

181 persisted: Final = (await _read_general_settings()).get(_COORDINATION_REDIS_KEY) 

182 if isinstance(persisted, dict): 182 ↛ 183line 182 didn't jump to line 183 because the condition on line 182 was never true

183 return _SETTINGS_ADAPTER.validate_python(persisted) 

184 return None 

185 

186 

187async def _current_coordination_redis_settings() -> dict[str, object] | None: 

188 """The coordination_redis block the proxy would boot with. 

189 

190 The persisted row wins over the yaml-loaded config state because startup 

191 applies the DB `general_settings` row over the file config. 

192 """ 

193 from litellm.proxy.proxy_server import proxy_config 

194 

195 persisted: Final = await get_persisted_coordination_redis_settings() 

196 if persisted is not None: 196 ↛ 197line 196 didn't jump to line 197 because the condition on line 196 was never true

197 return persisted 

198 

199 config_state: Final = _SETTINGS_ADAPTER.validate_python(proxy_config.get_config_state()) 

200 general_settings: Final = config_state.get(_GENERAL_SETTINGS_PARAM_NAME) 

201 if not isinstance(general_settings, Mapping): 201 ↛ 202line 201 didn't jump to line 202 because the condition on line 201 was never true

202 return None 

203 from_file: Final = general_settings.get(_COORDINATION_REDIS_KEY) 

204 if isinstance(from_file, dict): 204 ↛ 205line 204 didn't jump to line 205 because the condition on line 204 was never true

205 return _SETTINGS_ADAPTER.validate_python(from_file) 

206 return None 

207 

208 

209def _coordination_redis_source(settings: Mapping[str, object] | None) -> CoordinationRedisSource | None: 

210 """Which source the proxy's coordination Redis comes from, in startup precedence order. 

211 

212 Mirrors `ProxyConfig._init_coordination_redis` -> `ProxyConfig._init_cache`: 

213 an explicit block wins, else a plain-Redis response-cache backend is 

214 borrowed, else the REDIS_* environment fallback applies. 

215 """ 

216 from litellm.proxy.proxy_server import _environment_has_redis_connection_target 

217 

218 if settings: 218 ↛ 219line 218 didn't jump to line 219 because the condition on line 218 was never true

219 return "coordination_redis" 

220 cache_backend: Final = litellm.cache.cache if litellm.cache is not None else None 

221 if isinstance(cache_backend, (RedisCache, RedisClusterCache)): 221 ↛ 222line 221 didn't jump to line 222 because the condition on line 221 was never true

222 return "cache_backend" 

223 if _environment_has_redis_connection_target(): 223 ↛ 224line 223 didn't jump to line 224 because the condition on line 223 was never true

224 return "environment" 

225 return None 

226 

227 

228def _log_audit_task_exception(task: "asyncio.Task[None]") -> None: 

229 """Surface a fire-and-forget audit-log task failure as a warning.""" 

230 if task.cancelled(): 

231 return 

232 exc: Final = task.exception() 

233 if exc is not None: 

234 verbose_proxy_logger.warning("Failed to write coordination-redis-settings audit log: %s", exc) 

235 

236 

237async def _emit_coordination_redis_audit_log( 

238 *, 

239 action: AUDIT_ACTIONS, 

240 before_settings: Mapping[str, object] | None, 

241 after_settings: Mapping[str, object] | None, 

242 user_api_key_dict: UserAPIKeyAuth, 

243 litellm_changed_by: str | None, 

244) -> None: 

245 """Emit an audit-log row for a /coordination_redis/settings mutation.""" 

246 from litellm.proxy.management_helpers.audit_logs import ( 

247 create_audit_log_for_update, 

248 is_audit_logging_enabled, 

249 ) 

250 from litellm.proxy.proxy_server import litellm_proxy_admin_name 

251 

252 if not is_audit_logging_enabled(): 

253 return 

254 

255 task: Final = asyncio.create_task( 

256 create_audit_log_for_update( 

257 request_data=LiteLLM_AuditLogs( 

258 id=str(uuid.uuid4()), 

259 updated_at=datetime.now(timezone.utc), 

260 changed_by=litellm_changed_by or user_api_key_dict.user_id or litellm_proxy_admin_name, 

261 changed_by_api_key=user_api_key_dict.api_key, 

262 table_name=LitellmTableNames.CONFIG_TABLE_NAME, 

263 object_id=_COORDINATION_REDIS_KEY, 

264 action=action, 

265 updated_values=json.dumps({"settings": _redact_all_values(after_settings)}, default=str), 

266 before_value=json.dumps({"settings": _redact_all_values(before_settings)}, default=str), 

267 ) 

268 ) 

269 ) 

270 task.add_done_callback(_log_audit_task_exception) 

271 

272 

273class CoordinationRedisSettingsResponse(BaseModel): 

274 values: dict[str, object] = Field(description="Current coordination Redis settings, with credentials redacted") 

275 fields: list[CoordinationRedisSettingsField] = Field( 

276 description="List of all configurable coordination Redis settings with metadata" 

277 ) 

278 source: CoordinationRedisSource | None = Field( 

279 description="Where the proxy's coordination Redis comes from; null when it has none" 

280 ) 

281 

282 

283class CoordinationRedisSettingsRequest(BaseModel): 

284 settings: dict[str, object] = Field(description="Coordination Redis connection params") 

285 

286 

287class CoordinationRedisTestResponse(BaseModel): 

288 status: str = Field(description="Connection status: 'healthy' or 'unhealthy'") 

289 error: str | None = Field(default=None, description="Error message if the connection failed") 

290 

291 

292@router.get( 

293 "/coordination_redis/settings", 

294 tags=["Coordination Redis Settings"], 

295 dependencies=[Depends(user_api_key_auth)], 

296 response_model=CoordinationRedisSettingsResponse, 

297) 

298async def get_coordination_redis_settings( 

299 user_api_key_dict: UserAPIKeyAuth = Depends(user_api_key_auth), 

300) -> CoordinationRedisSettingsResponse: 

301 """ 

302 Get the coordination Redis configuration and available settings. 

303 

304 Returns: 

305 - values: current coordination Redis settings, with password/sentinel_password/url redacted 

306 - fields: all configurable settings with their metadata (type, description, default, section) 

307 - source: "coordination_redis" | "cache_backend" | "environment" | null 

308 """ 

309 if not user_api_key_has_admin_view(user_api_key_dict): 309 ↛ 310line 309 didn't jump to line 310 because the condition on line 309 was never true

310 _enforce_proxy_admin(user_api_key_dict) 

311 

312 settings: Final = await _current_coordination_redis_settings() 

313 source: Final = _coordination_redis_source(settings) 

314 

315 values: Final = _redact_credentials(settings or {}) 

316 fields: Final = [field.model_copy(deep=True) for field in COORDINATION_REDIS_SETTINGS_FIELDS] 

317 for field in fields: 

318 if field.field_name in values: 318 ↛ 319line 318 didn't jump to line 319 because the condition on line 318 was never true

319 field.field_value = values[field.field_name] 

320 

321 return CoordinationRedisSettingsResponse(values=values, fields=fields, source=source) 

322 

323 

324@router.post( 

325 "/coordination_redis/settings", 

326 tags=["Coordination Redis Settings"], 

327 dependencies=[Depends(user_api_key_auth)], 

328) 

329async def update_coordination_redis_settings( 

330 request: CoordinationRedisSettingsRequest, 

331 user_api_key_dict: UserAPIKeyAuth = Depends(user_api_key_auth), 

332 litellm_changed_by: str | None = Header( 

333 None, 

334 description="The litellm-changed-by header enables tracking of actions performed by authorized users on behalf of other users, providing an audit trail for accountability", 

335 ), 

336) -> dict[str, object]: 

337 """ 

338 Save coordination Redis settings under `general_settings.coordination_redis`. 

339 

340 Parameters: 

341 - settings: dict - Redis connection params (host, port, username, password, url, ssl, startup_nodes, sentinel_nodes, sentinel_password, service_name). Values may be `os.environ/VAR` references, which are stored as written and resolved at startup 

342 

343 The settings are written to the `general_settings` row of LiteLLM_Config, 

344 which startup merges over the yaml config; the proxy picks them up on its 

345 next restart. 

346 """ 

347 from litellm.proxy.proxy_server import prisma_client, store_model_in_db 

348 

349 _enforce_proxy_admin(user_api_key_dict) 

350 

351 if prisma_client is None: 351 ↛ 352line 351 didn't jump to line 352 because the condition on line 351 was never true

352 raise HTTPException( 

353 status_code=500, 

354 detail={"error": "Database not connected. Please connect a database."}, 

355 ) 

356 

357 if store_model_in_db is not True: 357 ↛ 358line 357 didn't jump to line 358 because the condition on line 357 was never true

358 raise HTTPException( 

359 status_code=500, 

360 detail={"error": "Set `'STORE_MODEL_IN_DB='True'` in your env to enable this feature."}, 

361 ) 

362 

363 saved_settings: Final = await _current_coordination_redis_settings() 

364 settings: Final = _merge_over_saved(request.settings, saved_settings or {}) 

365 _validated_params(settings) 

366 

367 from litellm.proxy.proxy_server import proxy_config 

368 

369 proxy_config.reject_config_owned_writes( 

370 section_name=_GENERAL_SETTINGS_PARAM_NAME, changed_keys={_COORDINATION_REDIS_KEY: settings} 

371 ) 

372 general_settings: Final = await _read_general_settings() 

373 before_settings: Final = general_settings.get(_COORDINATION_REDIS_KEY) 

374 action: Final[AUDIT_ACTIONS] = "updated" if isinstance(before_settings, dict) else "created" 

375 

376 await ConfigRepository(prisma_client).set_param( 

377 param_name=_GENERAL_SETTINGS_PARAM_NAME, 

378 param_value={**general_settings, _COORDINATION_REDIS_KEY: settings}, 

379 ) 

380 await invalidate_config_param(_GENERAL_SETTINGS_PARAM_NAME) 

381 

382 # coordination_redis carries Redis credentials and decides where cross-pod 

383 # rate-limit and spend state lives; an admin repointing it is a 

384 # data-routing pivot, so make the change traceable. 

385 await _emit_coordination_redis_audit_log( 

386 action=action, 

387 before_settings=before_settings if isinstance(before_settings, dict) else None, 

388 after_settings=settings, 

389 user_api_key_dict=user_api_key_dict, 

390 litellm_changed_by=litellm_changed_by, 

391 ) 

392 

393 return { 

394 "message": "Coordination Redis settings updated successfully. Restart the proxy to apply them.", 

395 "status": "success", 

396 "settings": _redact_credentials(settings), 

397 } 

398 

399 

400@router.post( 

401 "/coordination_redis/settings/test", 

402 tags=["Coordination Redis Settings"], 

403 dependencies=[Depends(user_api_key_auth)], 

404 response_model=CoordinationRedisTestResponse, 

405) 

406async def check_coordination_redis_connection( 

407 request: CoordinationRedisSettingsRequest, 

408 user_api_key_dict: UserAPIKeyAuth = Depends(user_api_key_auth), 

409) -> CoordinationRedisTestResponse: 

410 """ 

411 Test a coordination Redis connection with the provided credentials. 

412 

413 Parameters: 

414 - settings: dict - Redis connection params to test. Credential fields sent back as `***REDACTED***` fall back to the saved value 

415 

416 Builds a throwaway client (never touching global state) and pings it. 

417 """ 

418 from litellm.proxy.proxy_server import _build_redis_usage_cache 

419 

420 _enforce_proxy_admin(user_api_key_dict) 

421 

422 saved_settings: Final = await _current_coordination_redis_settings() 

423 settings: Final = _merge_over_saved(request.settings, saved_settings or {}) 

424 params: Final = _validated_params(settings) 

425 

426 redis_cache: RedisCache | None = None 

427 try: 

428 redis_cache = _build_redis_usage_cache(params.model_dump(exclude_none=True)) 

429 await asyncio.wait_for(redis_cache.ping(), timeout=_PING_TIMEOUT_SECONDS) 

430 return CoordinationRedisTestResponse(status="healthy") 

431 except asyncio.TimeoutError: 

432 return CoordinationRedisTestResponse( 

433 status="unhealthy", 

434 error=f"Connection timed out after {_PING_TIMEOUT_SECONDS}s", 

435 ) 

436 except Exception as e: # noqa: BLE001 # any client/connection failure is a health verdict, not a 500 

437 return CoordinationRedisTestResponse(status="unhealthy", error=_scrub_credentials(str(e), settings)) 

438 finally: 

439 if redis_cache is not None: 

440 with suppress(Exception): 

441 await redis_cache.disconnect()