Coverage for open_webui/routers/openai.py: 29%

1029 statements  

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

1from __future__ import annotations 

2 

3import asyncio 

4import hashlib 

5import logging 

6import re 

7from typing import Optional 

8from urllib.parse import quote, urlparse 

9 

10import aiofiles 

11import aiohttp 

12from aiocache import cached 

13from azure.identity import DefaultAzureCredential, get_bearer_token_provider 

14from fastapi import APIRouter, Depends, HTTPException, Request, status 

15from fastapi.responses import ( 

16 FileResponse, 

17 JSONResponse, 

18 PlainTextResponse, 

19 StreamingResponse, 

20) 

21from open_webui.config import ( 

22 CACHE_DIR, 

23) 

24from open_webui.constants import ERROR_MESSAGES 

25from open_webui.env import ( 

26 AIOHTTP_CLIENT_SESSION_SSL, 

27 AIOHTTP_CLIENT_TIMEOUT_MODEL_LIST, 

28 BYPASS_MODEL_ACCESS_CONTROL, 

29 ENABLE_FORWARD_USER_INFO_HEADERS, 

30 ENABLE_OPENAI_API_PASSTHROUGH, 

31 FORWARD_SESSION_INFO_HEADER_CHAT_ID, 

32 MODELS_CACHE_TTL, 

33 REDIS_KEY_PREFIX, 

34) 

35from open_webui.events import EVENTS, publish_event, publish_model_provider_request_failed 

36from open_webui.internal.db import get_async_session 

37from open_webui.models.access_grants import AccessGrants 

38from open_webui.models.config import Config 

39from open_webui.models.groups import Groups 

40from open_webui.models.models import Models 

41from open_webui.models.users import UserModel 

42from open_webui.utils.access_control import check_model_access, has_connection_access, has_permission 

43from open_webui.utils.anthropic import ANTHROPIC_VERSION, get_anthropic_models, is_anthropic_url 

44from open_webui.utils.auth import get_admin_user, get_verified_user 

45from open_webui.utils.headers import get_custom_headers, include_user_info_headers 

46from open_webui.utils.json_codec import JSONCodec 

47from open_webui.utils.misc import convert_logit_bias_input_to_json 

48from open_webui.utils.model_ids import strip_provider_model_prefix 

49from open_webui.utils.payload import ( 

50 apply_model_params_to_body_openai, 

51 apply_system_prompt_to_body, 

52) 

53from open_webui.utils.session_pool import ( 

54 cleanup_response, 

55 get_client_timeout, 

56 get_session, 

57 stream_wrapper, 

58) 

59from pydantic import BaseModel, ConfigDict 

60from sqlalchemy.ext.asyncio import AsyncSession 

61 

62log = logging.getLogger(__name__) 

63 

64 

65########################################## 

66# 

67# Utility functions 

68# Let the responses returned through this gate be worth 

69# the question that summoned them. 

70# 

71########################################## 

72 

73# Headers that become stale after aiohttp auto-decompresses the upstream 

74# response body. Forwarding them verbatim causes desktop / programmatic 

75# clients to attempt decompression of an already-decoded payload, resulting 

76# in ZlibError. See https://github.com/aio-libs/aiohttp/issues/4462. 

77# Also drop server and date: uvicorn adds its own and forwarding both duplicates them. 

78_STRIP_PROXY_HEADERS = frozenset({'content-encoding', 'content-length', 'transfer-encoding', 'server', 'date'}) 

79_MODEL_LIST_TIMEOUT = aiohttp.ClientTimeout(total=AIOHTTP_CLIENT_TIMEOUT_MODEL_LIST) 

80_UNSUPPORTED_OPENAI_MODEL_KEYWORDS = ('babbage', 'dall-e', 'davinci', 'embedding', 'tts', 'whisper') 

81BASE_MODELS_CACHE_KEY = f'{REDIS_KEY_PREFIX}:models:base' 

82 

83 

84def _clean_proxy_headers(raw_headers) -> dict: 

85 """Return a copy of *raw_headers* without the encoding, server and date headers.""" 

86 return {k: v for k, v in raw_headers.items() if k.lower() not in _STRIP_PROXY_HEADERS} 

87 

88 

89async def send_get_request( 

90 request: Request = None, 

91 url=None, 

92 key=None, 

93 user: UserModel = None, 

94 config=None, 

95): 

96 try: 

97 async with aiohttp.ClientSession(timeout=_MODEL_LIST_TIMEOUT, trust_env=True) as session: 

98 if request and config: 98 ↛ 99line 98 didn't jump to line 99 because the condition on line 98 was never true

99 headers, cookies = await get_headers_and_cookies(request, url, key, config, user=user) 

100 else: 

101 headers = { 

102 **({'Authorization': f'Bearer {key}'} if key else {}), 

103 } 

104 cookies = None 

105 

106 if ENABLE_FORWARD_USER_INFO_HEADERS and user: 106 ↛ 107line 106 didn't jump to line 107 because the condition on line 106 was never true

107 headers = include_user_info_headers(headers, user, request=request) 

108 

109 async with session.get( 

110 url, 

111 headers=headers, 

112 cookies=cookies, 

113 ssl=AIOHTTP_CLIENT_SESSION_SSL, 

114 ) as response: 

115 return await response.json(loads=JSONCodec.loads) 

116 except Exception as e: 

117 # Handle connection error here 

118 log.error(f'Connection error: {e}') 

119 return None 

120 

121 

122async def get_models_request( 

123 request: Request = None, 

124 url=None, 

125 key=None, 

126 user: UserModel = None, 

127 config=None, 

128): 

129 if is_anthropic_url(url): 129 ↛ 130line 129 didn't jump to line 130 because the condition on line 129 was never true

130 return await get_anthropic_models(url, key, user=user) 

131 return await send_get_request(request, f'{url}/models', key, user=user, config=config) 

132 

133 

134def openai_reasoning_model_handler(payload): 

135 """ 

136 Handle reasoning model specific parameters 

137 """ 

138 if 'max_tokens' in payload: 

139 # Convert "max_tokens" to "max_completion_tokens" for all reasoning models 

140 payload['max_completion_tokens'] = payload['max_tokens'] 

141 del payload['max_tokens'] 

142 

143 # Handle system role conversion based on model type 

144 if payload['messages'][0]['role'] == 'system': 

145 model_lower = payload['model'].lower() 

146 # Legacy models use "user" role instead of "system" 

147 if model_lower.startswith('o1-mini') or model_lower.startswith('o1-preview'): 147 ↛ anywhereline 147 didn't jump anywhere: it always raised an exception.

148 payload['messages'][0]['role'] = 'user' 

149 else: 

150 payload['messages'][0]['role'] = 'developer' 

151 

152 return payload 

153 

154 

155async def get_headers_and_cookies( 

156 request: Request, 

157 url, 

158 key=None, 

159 config=None, 

160 metadata: dict | None = None, 

161 user: UserModel = None, 

162): 

163 cookies = getattr(request, 'cookies', {}) if config.get('forward_cookies', False) else {} 

164 headers = { 

165 'Content-Type': 'application/json', 

166 **( 

167 { 

168 # LICENSE covers this Open WebUI upstream metadata identifier. 

169 # Do not alter, remove, obscure, or replace it except as LICENSE permits: 

170 # https://docs.openwebui.com/license. 

171 'HTTP-Referer': 'https://openwebui.com/', 

172 'X-Title': 'Open WebUI', 

173 } 

174 if 'openrouter.ai' in url 

175 else {} 

176 ), 

177 } 

178 

179 if ENABLE_FORWARD_USER_INFO_HEADERS and user: 179 ↛ 180line 179 didn't jump to line 180 because the condition on line 179 was never true

180 headers = include_user_info_headers(headers, user, request=request) 

181 if metadata and metadata.get('chat_id'): 

182 headers[FORWARD_SESSION_INFO_HEADER_CHAT_ID] = metadata.get('chat_id') 

183 

184 token = None 

185 auth_type = config.get('auth_type') 

186 

187 if auth_type == 'bearer' or auth_type is None: 187 ↛ 190line 187 didn't jump to line 190 because the condition on line 187 was always true

188 # Default to bearer if not specified 

189 token = f'{key}' 

190 elif auth_type == 'none': 

191 token = None 

192 elif auth_type == 'session': 

193 token = request.state.token.credentials 

194 elif auth_type == 'system_oauth': 

195 oauth_token = None 

196 try: 

197 if request.cookies.get('oauth_session_id', None): 

198 oauth_token = await request.app.state.oauth_manager.get_oauth_token( 

199 user.id, 

200 request.cookies.get('oauth_session_id', None), 

201 ) 

202 except Exception as e: 

203 log.error(f'Error getting OAuth token: {e}') 

204 

205 if oauth_token: 205 ↛ 206line 205 didn't jump to line 206 because the condition on line 205 was never true

206 token = f'{oauth_token.get("access_token", "")}' 

207 

208 elif auth_type in ('azure_ad', 'microsoft_entra_id'): 

209 token = get_microsoft_entra_id_access_token() 

210 

211 if token: 

212 headers['Authorization'] = f'Bearer {token}' 

213 

214 if config.get('headers') and isinstance(config.get('headers'), dict): 214 ↛ 215line 214 didn't jump to line 215 because the condition on line 214 was never true

215 custom_headers = await get_custom_headers(config.get('headers'), user, metadata, request=request) 

216 headers.update(custom_headers) 

217 

218 return headers, cookies 

219 

220 

221def get_microsoft_entra_id_access_token(): 

222 """ 

223 Get Microsoft Entra ID access token using DefaultAzureCredential for Azure OpenAI. 

224 Returns the token string or None if authentication fails. 

225 """ 

226 try: 

227 token_provider = get_bearer_token_provider( 

228 DefaultAzureCredential(), 'https://cognitiveservices.azure.com/.default' 

229 ) 

230 return token_provider() 

231 except Exception as e: 

232 log.error(f'Error getting Microsoft Entra ID access token: {e}') 

233 return None 

234 

235 

236########################################## 

237# 

238# API routes 

239# 

240########################################## 

241 

242router = APIRouter() 

243 

244LLAMACPP_LOADED_STATES = {'loaded', 'sleeping'} 

245LLAMACPP_UNLOADED_STATES = {'loading', 'unloaded'} 

246MODEL_MANAGEMENT_ENDPOINTS = { 

247 'llama.cpp': { 

248 'list': '/models', 

249 'download': '/models', 

250 'delete': '/models', 

251 'load': '/models/load', 

252 'unload': '/models/unload', 

253 'sse': '/models/sse', 

254 }, 

255 'lmstudio': { 

256 'list': '/api/v1/models', 

257 'download': '/api/v1/models/download', 

258 'download_status': '/api/v1/models/download/status/{job_id}', 

259 'load': '/api/v1/models/load', 

260 'unload': '/api/v1/models/unload', 

261 }, 

262} 

263 

264 

265def get_model_management_root_url(url: str, provider: str) -> str: 

266 root_url = url.rstrip('/') 

267 if provider in ('llama.cpp', 'lmstudio'): 

268 for suffix in ('/api/v1', '/api/v0', '/v1'): 

269 if root_url.endswith(suffix): 

270 return root_url.removesuffix(suffix) 

271 

272 return root_url 

273 

274 

275def get_provider_model_loaded_state(model: dict, provider: str, manual_model_ids: bool = False) -> bool | None: 

276 if provider == 'lmstudio': 

277 if model.get('loaded_instances'): 

278 return True 

279 

280 state = model.get('state') 

281 if state == 'loaded': 

282 return True 

283 if state == 'not-loaded': 

284 return False 

285 

286 return None 

287 

288 if provider != 'llama.cpp': 

289 return None 

290 

291 status = model.get('status') 

292 if isinstance(status, dict): 

293 value = status.get('value') 

294 if value in LLAMACPP_LOADED_STATES: 

295 return True 

296 if value in LLAMACPP_UNLOADED_STATES: 

297 return False 

298 

299 if not manual_model_ids and 'status' not in model: 

300 return True 

301 

302 return None 

303 

304 

305OPENAI_CONFIG_KEYS = { 

306 'ENABLE_OPENAI_API': 'openai.enable', 

307 'OPENAI_API_BASE_URLS': 'openai.api_base_urls', 

308 'OPENAI_API_KEYS': 'openai.api_keys', 

309 'OPENAI_API_CONFIGS': 'openai.api_configs', 

310} 

311 

312 

313async def get_openai_config() -> dict: 

314 values = await Config.get_many(*OPENAI_CONFIG_KEYS.values()) 

315 return {field: values[storage_key] for field, storage_key in OPENAI_CONFIG_KEYS.items() if storage_key in values} 

316 

317 

318async def get_openai_runtime_config() -> tuple[bool, list[str], list[str], dict]: 

319 values = await Config.get_many('openai.enable', 'openai.api_base_urls', 'openai.api_keys', 'openai.api_configs') 

320 return ( 

321 values.get('openai.enable'), 

322 values.get('openai.api_base_urls') or [], 

323 values.get('openai.api_keys') or [], 

324 values.get('openai.api_configs') or {}, 

325 ) 

326 

327 

328async def normalize_openai_api_keys(api_base_urls: list[str], api_keys: list[str]) -> list[str]: 

329 if len(api_keys) > len(api_base_urls): 

330 api_keys = api_keys[: len(api_base_urls)] 

331 elif len(api_keys) < len(api_base_urls): 

332 api_keys = [*api_keys, *([''] * (len(api_base_urls) - len(api_keys)))] 

333 

334 await Config.upsert({'openai.api_keys': api_keys}) 

335 return api_keys 

336 

337 

338async def get_openai_connection(idx: int) -> tuple[str, str, dict]: 

339 _, api_base_urls, api_keys, api_configs = await get_openai_runtime_config() 

340 url = api_base_urls[idx] 

341 key = api_keys[idx] 

342 api_config = api_configs.get(str(idx), api_configs.get(url, {})) 

343 return url, key, api_config 

344 

345 

346async def clear_openai_model_cache(request: Request): 

347 await get_all_models.cache.clear() 

348 redis = getattr(request.app.state, 'redis', None) 

349 if redis is not None: 349 ↛ 350line 349 didn't jump to line 350 because the condition on line 349 was never true

350 await redis.delete(BASE_MODELS_CACHE_KEY) 

351 request.app.state.BASE_MODELS = [] 

352 request.app.state.OPENAI_MODELS = {} 

353 models = getattr(request.app.state, 'MODELS', None) 

354 if hasattr(models, 'clear'): 354 ↛ 357line 354 didn't jump to line 357 because the condition on line 354 was always true

355 models.clear() 

356 else: 

357 request.app.state.MODELS = {} 

358 

359 

360async def get_model_management_connection(url_idx: int) -> tuple[str, str, dict, str]: 

361 if not await Config.get('openai.enable'): 

362 raise HTTPException(status_code=503, detail='OpenAI API is disabled') 

363 

364 try: 

365 url, key, api_config = await get_openai_connection(url_idx) 

366 except IndexError: 

367 raise HTTPException(status_code=404, detail='Connection not found') 

368 

369 provider = api_config.get('provider', '') 

370 if provider not in MODEL_MANAGEMENT_ENDPOINTS: 

371 raise HTTPException( 

372 status_code=400, 

373 detail=f'Provider "{provider or "default"}" does not support model management', 

374 ) 

375 

376 return get_model_management_root_url(url, provider), key, api_config, provider 

377 

378 

379def get_model_management_path(provider: str, operation: str, path_params: dict | None = None) -> str: 

380 try: 

381 path = MODEL_MANAGEMENT_ENDPOINTS[provider][operation] 

382 except KeyError: 

383 raise HTTPException(status_code=400, detail=f'Provider "{provider}" does not support {operation}') 

384 

385 return path.format(**(path_params or {})) 

386 

387 

388def get_model_management_payload(provider: str, operation: str, payload: dict | None) -> dict | None: 

389 if provider == 'lmstudio' and operation == 'unload' and payload: 

390 return {'instance_id': payload.get('instance_id') or payload.get('model')} 

391 

392 return payload 

393 

394 

395async def send_model_management_request( 

396 request: Request, 

397 url_idx: int, 

398 operation: str, 

399 method: str = 'GET', 

400 payload: dict | None = None, 

401 query: dict | None = None, 

402 path_params: dict | None = None, 

403 stream: bool = False, 

404 user: UserModel | None = None, 

405): 

406 root_url, key, api_config, provider = await get_model_management_connection(url_idx) 

407 path = get_model_management_path(provider, operation, path_params=path_params) 

408 payload = get_model_management_payload(provider, operation, payload) 

409 headers, cookies = await get_headers_and_cookies(request, root_url, key, api_config, user=user) 

410 

411 response = None 

412 streaming = False 

413 try: 

414 session = await get_session() 

415 response = await session.request( 

416 method, 

417 f'{root_url}{path}', 

418 json=payload, 

419 params=query, 

420 headers=headers, 

421 cookies=cookies, 

422 ssl=AIOHTTP_CLIENT_SESSION_SSL, 

423 timeout=get_client_timeout(stream=stream), 

424 ) 

425 

426 if not response.ok: 

427 try: 

428 error = await response.json(loads=JSONCodec.loads) 

429 except Exception: 

430 error = await response.text() 

431 raise HTTPException(status_code=response.status, detail=error) 

432 

433 if stream: 

434 streaming = True 

435 return StreamingResponse( 

436 stream_wrapper(response, passthrough=True), 

437 status_code=response.status, 

438 headers=_clean_proxy_headers(response.headers), 

439 ) 

440 

441 try: 

442 return await response.json(loads=JSONCodec.loads) 

443 except Exception: 

444 return {'success': True} 

445 except HTTPException: 

446 raise 

447 except Exception as e: 

448 raise HTTPException(status_code=response.status if response else 500, detail=str(e)) 

449 finally: 

450 if not streaming: 

451 await cleanup_response(response) 

452 

453 

454async def get_anthropic_request_target(request: Request, form_data: dict, user: UserModel): 

455 """Resolve the upstream connection, payload and auth headers for a native Anthropic request.""" 

456 requested_model = form_data.get('model') 

457 if not requested_model: 457 ↛ 460line 457 didn't jump to line 460 because the condition on line 457 was always true

458 raise HTTPException(status_code=400, detail='model is required') 

459 

460 payload = {**form_data} 

461 model_id = requested_model 

462 model_info = await Models.get_model_by_id(model_id) 

463 await check_model_access(user, model_info, BYPASS_MODEL_ACCESS_CONTROL) 

464 

465 if model_info and model_info.base_model_id: 

466 model_id = model_info.base_model_id 

467 payload['model'] = model_id 

468 

469 models = request.app.state.OPENAI_MODELS 

470 if not models or model_id not in models: 

471 await get_all_models(request, user=user) 

472 models = request.app.state.OPENAI_MODELS 

473 

474 model = models.get(model_id) 

475 if not model or 'urlIdx' not in model: 

476 raise HTTPException(status_code=404, detail=ERROR_MESSAGES.MODEL_NOT_FOUND()) 

477 

478 url, key, api_config = await get_openai_connection(model['urlIdx']) 

479 prefix_id = api_config.get('prefix_id') 

480 payload['model'] = strip_provider_model_prefix(payload['model'], prefix_id) 

481 

482 headers, cookies = await get_headers_and_cookies(request, url, key, api_config, user=user) 

483 

484 # Anthropic's native endpoints reject bearer auth, the key belongs in x-api-key. 

485 if is_anthropic_url(url): 

486 headers.setdefault('anthropic-version', ANTHROPIC_VERSION) 

487 if api_config.get('auth_type') in (None, 'bearer'): 

488 headers.pop('Authorization', None) 

489 headers.setdefault('x-api-key', key) 

490 

491 return requested_model, payload, url, key, headers, cookies 

492 

493 

494async def count_anthropic_tokens(request: Request, form_data: dict, user: UserModel) -> int: 

495 """Forward an Anthropic token-count request through an OpenAI-compatible connection.""" 

496 requested_model, payload, url, key, headers, cookies = await get_anthropic_request_target(request, form_data, user) 

497 request_url = f'{url.rstrip("/")}/messages/count_tokens' 

498 response = None 

499 

500 try: 

501 session = await get_session() 

502 response = await session.request( 

503 method='POST', 

504 url=request_url, 

505 data=JSONCodec.dumps(payload), 

506 headers=headers, 

507 cookies=cookies, 

508 ssl=AIOHTTP_CLIENT_SESSION_SSL, 

509 timeout=get_client_timeout(), 

510 ) 

511 

512 try: 

513 response_data = await response.json(loads=JSONCodec.loads) 

514 except Exception: 

515 response_data = await response.text() 

516 

517 if response.status >= 400: 

518 await publish_model_provider_request_failed( 

519 request, 

520 actor=user, 

521 provider='openai-compatible', 

522 base_url=url, 

523 api_key=key, 

524 status=response.status, 

525 requested_model=requested_model, 

526 upstream_error=response_data, 

527 ) 

528 raise HTTPException(status_code=response.status, detail=response_data) 

529 

530 input_tokens = response_data.get('input_tokens') if isinstance(response_data, dict) else None 

531 if isinstance(input_tokens, bool) or not isinstance(input_tokens, int) or input_tokens < 0: 

532 raise HTTPException(status_code=502, detail='Invalid token-count response from upstream provider') 

533 

534 return input_tokens 

535 except HTTPException: 

536 raise 

537 except Exception: 

538 log.exception('Failed to count Anthropic tokens for model %s', requested_model) 

539 raise HTTPException(status_code=502, detail=ERROR_MESSAGES.SERVER_CONNECTION_ERROR) 

540 finally: 

541 await cleanup_response(response) 

542 

543 

544@router.get('/config') 

545async def get_config(request: Request, user=Depends(get_admin_user)): 

546 return await get_openai_config() 

547 

548 

549class OpenAIConfigForm(BaseModel): 

550 ENABLE_OPENAI_API: bool | None = None 

551 OPENAI_API_BASE_URLS: list[str] 

552 OPENAI_API_KEYS: list[str] 

553 OPENAI_API_CONFIGS: dict 

554 

555 

556@router.post('/config/update') 

557async def update_config(request: Request, form_data: OpenAIConfigForm, user=Depends(get_admin_user)): 

558 api_keys = form_data.OPENAI_API_KEYS 

559 

560 if len(api_keys) > len(form_data.OPENAI_API_BASE_URLS): 

561 api_keys = api_keys[: len(form_data.OPENAI_API_BASE_URLS)] 

562 elif len(api_keys) < len(form_data.OPENAI_API_BASE_URLS): 

563 api_keys = [*api_keys, *([''] * (len(form_data.OPENAI_API_BASE_URLS) - len(api_keys)))] 

564 

565 valid_keys = set(map(str, range(len(form_data.OPENAI_API_BASE_URLS)))) 

566 api_configs = {key: value for key, value in form_data.OPENAI_API_CONFIGS.items() if key in valid_keys} 

567 

568 await Config.upsert( 

569 { 

570 'openai.enable': form_data.ENABLE_OPENAI_API, 

571 'openai.api_base_urls': form_data.OPENAI_API_BASE_URLS, 

572 'openai.api_keys': api_keys, 

573 'openai.api_configs': api_configs, 

574 } 

575 ) 

576 

577 await clear_openai_model_cache(request) 

578 

579 await publish_event( 

580 request, 

581 EVENTS.MODEL_PROVIDER_CONFIG_UPDATED, 

582 actor=user, 

583 subject_id='openai', 

584 subject_type='model.provider_config', 

585 data={ 

586 'provider': 'openai', 

587 'enabled': form_data.ENABLE_OPENAI_API, 

588 'base_url_count': len(form_data.OPENAI_API_BASE_URLS), 

589 }, 

590 ) 

591 

592 return { 

593 'ENABLE_OPENAI_API': form_data.ENABLE_OPENAI_API, 

594 'OPENAI_API_BASE_URLS': form_data.OPENAI_API_BASE_URLS, 

595 'OPENAI_API_KEYS': api_keys, 

596 'OPENAI_API_CONFIGS': api_configs, 

597 } 

598 

599 

600@router.post('/audio/speech') 

601async def speech(request: Request, user=Depends(get_verified_user)): 

602 if user.role != 'admin' and not await has_permission(user.id, 'chat.tts', await Config.get('user.permissions')): 602 ↛ 603line 602 didn't jump to line 603 because the condition on line 602 was never true

603 raise HTTPException( 

604 status_code=status.HTTP_403_FORBIDDEN, 

605 detail=ERROR_MESSAGES.ACCESS_PROHIBITED, 

606 ) 

607 

608 idx = None 

609 try: 

610 _, api_base_urls, _, _ = await get_openai_runtime_config() 

611 idx = api_base_urls.index('https://api.openai.com/v1') 

612 

613 body = await request.body() 

614 name = hashlib.sha256(body).hexdigest() 

615 

616 SPEECH_CACHE_DIR = CACHE_DIR / 'audio' / 'speech' 

617 SPEECH_CACHE_DIR.mkdir(parents=True, exist_ok=True) 

618 file_path = SPEECH_CACHE_DIR.joinpath(f'{name}.mp3') 

619 file_body_path = SPEECH_CACHE_DIR.joinpath(f'{name}.json') 

620 

621 # Check if the file already exists in the cache 

622 if file_path.is_file(): 622 ↛ 623line 622 didn't jump to line 623 because the condition on line 622 was never true

623 return FileResponse(file_path) 

624 

625 url, key, api_config = await get_openai_connection(idx) 

626 

627 headers, cookies = await get_headers_and_cookies(request, url, key, api_config, user=user) 

628 

629 r = None 

630 try: 

631 session = await get_session() 

632 r = await session.post( 

633 url=f'{url}/audio/speech', 

634 data=body, 

635 headers=headers, 

636 cookies=cookies, 

637 ssl=AIOHTTP_CLIENT_SESSION_SSL, 

638 ) 

639 

640 r.raise_for_status() 

641 

642 async with aiofiles.open(file_path, 'wb') as f: 

643 async for chunk in r.content.iter_chunked(8192): 

644 await f.write(chunk) 

645 

646 async with aiofiles.open(file_body_path, 'w') as f: 

647 await f.write(JSONCodec.dumps(JSONCodec.loads(body.decode('utf-8')))) 

648 

649 # Return the saved file 

650 return FileResponse(file_path) 

651 

652 except Exception as e: 

653 log.exception(e) 

654 

655 detail = None 

656 if r is not None: 656 ↛ 667line 656 didn't jump to line 667 because the condition on line 656 was always true

657 try: 

658 res = await r.json(loads=JSONCodec.loads) 

659 if 'error' in res: 

660 detail = f'External: {res["error"]}' 

661 except Exception: 

662 detail = f'External: {e}' 

663 

664 # LICENSE covers this Open WebUI error identifier. 

665 # Do not alter, remove, obscure, or replace it except as LICENSE permits: 

666 # https://docs.openwebui.com/license. 

667 raise HTTPException( 

668 status_code=r.status if r else 500, 

669 detail=detail if detail else 'Open WebUI: Server Connection Error', 

670 ) 

671 

672 except ValueError: 

673 raise HTTPException(status_code=401, detail=ERROR_MESSAGES.OPENAI_NOT_FOUND) 

674 

675 

676async def get_all_models_responses(request: Request, user: UserModel) -> list: 

677 enable_openai_api, api_base_urls, api_keys, api_configs = await get_openai_runtime_config() 

678 if not enable_openai_api: 

679 return [] 

680 

681 num_urls = len(api_base_urls) 

682 num_keys = len(api_keys) 

683 

684 if num_keys != num_urls: 684 ↛ 685line 684 didn't jump to line 685 because the condition on line 684 was never true

685 api_keys = await normalize_openai_api_keys(api_base_urls, api_keys) 

686 

687 request_tasks = [] 

688 for idx, url in enumerate(api_base_urls): 

689 if (str(idx) not in api_configs) and (url not in api_configs): # Legacy support 689 ↛ 692line 689 didn't jump to line 692 because the condition on line 689 was always true

690 request_tasks.append(get_models_request(request, url, api_keys[idx], user=user)) 

691 else: 

692 api_config = api_configs.get( 

693 str(idx), 

694 api_configs.get(url, {}), # Legacy support 

695 ) 

696 

697 enable = api_config.get('enable', True) 

698 model_ids = api_config.get('model_ids', []) 

699 

700 if enable: 

701 if len(model_ids) == 0: 

702 request_tasks.append(get_models_request(request, url, api_keys[idx], user=user, config=api_config)) 

703 else: 

704 model_list = { 

705 'object': 'list', 

706 'data': [ 

707 { 

708 'id': model_id, 

709 'name': model_id, 

710 'owned_by': 'openai', 

711 'openai': {'id': model_id}, 

712 'urlIdx': idx, 

713 } 

714 for model_id in model_ids 

715 ], 

716 } 

717 

718 request_tasks.append(asyncio.ensure_future(asyncio.sleep(0, model_list))) 

719 else: 

720 request_tasks.append(asyncio.ensure_future(asyncio.sleep(0, None))) 

721 

722 responses = await asyncio.gather(*request_tasks) 

723 

724 for idx, response in enumerate(responses): 

725 if response: 725 ↛ 726line 725 didn't jump to line 726 because the condition on line 725 was never true

726 url = api_base_urls[idx] 

727 api_config = api_configs.get( 

728 str(idx), 

729 api_configs.get(url, {}), # Legacy support 

730 ) 

731 

732 connection_type = api_config.get('connection_type', 'external') 

733 prefix_id = api_config.get('prefix_id', None) 

734 tags = api_config.get('tags', []) 

735 provider = api_config.get('provider', '') 

736 

737 model_list = response if isinstance(response, list) else response.get('data', []) 

738 if not isinstance(model_list, list): 

739 # Catch non-list responses 

740 model_list = [] 

741 

742 for model in model_list: 

743 # Remove name key if its value is None #16689 

744 if 'name' in model and model['name'] is None: 

745 del model['name'] 

746 

747 if prefix_id: 

748 model['id'] = f'{prefix_id}.{model.get("id", model.get("name", ""))}' 

749 if model.get('name'): 

750 model['name'] = f'{prefix_id}.{model["name"]}' 

751 

752 if tags: 

753 model['tags'] = tags 

754 

755 if connection_type: 

756 model['connection_type'] = connection_type 

757 

758 if provider: 

759 model['provider'] = provider 

760 

761 log.debug('get_all_models:responses() %s', responses) 

762 return responses 

763 

764 

765async def get_filtered_models(models, user, db=None): 

766 # Filter models based on user access control 

767 model_ids = [model['id'] for model in models.get('data', [])] 

768 model_infos = {model_info.id: model_info for model_info in await Models.get_models_by_ids(model_ids, db=db)} 

769 user_group_ids = {group.id for group in await Groups.get_groups_by_member_id(user.id, db=db)} 

770 

771 # Batch-fetch accessible resource IDs in a single query instead of N has_access calls 

772 accessible_model_ids = await AccessGrants.get_accessible_resource_ids( 

773 user_id=user.id, 

774 resource_type='model', 

775 resource_ids=list(model_infos.keys()), 

776 permission='read', 

777 user_group_ids=user_group_ids, 

778 db=db, 

779 ) 

780 

781 filtered_models = [] 

782 for model in models.get('data', []): 

783 model_info = model_infos.get(model['id']) 

784 if model_info: 

785 if user.id == model_info.user_id or model_info.id in accessible_model_ids: 

786 filtered_models.append(model) 

787 return filtered_models 

788 

789 

790@cached( 

791 ttl=MODELS_CACHE_TTL, 

792 # key_builder (not key) is the per-call hook in aiocache 0.12; `key=` is a 

793 # static key, so a `key=lambda` collapsed every caller to one shared entry. 

794 key_builder=lambda _func, request, user=None: f'openai_all_models_{user.id}' if user else 'openai_all_models', 

795) 

796async def get_all_models(request: Request, user: UserModel) -> dict[str, list]: 

797 log.info('get_all_models()') 

798 

799 enable_openai_api, api_base_urls, _, api_configs = await get_openai_runtime_config() 

800 if not enable_openai_api: 

801 request.app.state.OPENAI_MODELS = {} 

802 return {'data': []} 

803 

804 responses = await get_all_models_responses(request, user=user) 

805 

806 def extract_data(response): 

807 if response and 'data' in response: 807 ↛ 808line 807 didn't jump to line 808 because the condition on line 807 was never true

808 return response['data'] 

809 if isinstance(response, list): 809 ↛ 810line 809 didn't jump to line 810 because the condition on line 809 was never true

810 return response 

811 return None 

812 

813 def is_supported_openai_models(model_id): 

814 return not any(name in model_id for name in _UNSUPPORTED_OPENAI_MODEL_KEYWORDS) 

815 

816 def get_merged_models(model_lists): 

817 log.debug('merge_models_lists %s', model_lists) 

818 models = {} 

819 

820 for idx, model_list in enumerate(model_lists): 

821 if model_list is not None and 'error' not in model_list: 821 ↛ 822line 821 didn't jump to line 822 because the condition on line 821 was never true

822 base_url = api_base_urls[idx] 

823 hostname = urlparse(base_url).hostname if base_url else None 

824 api_config = api_configs.get(str(idx), api_configs.get(base_url, {})) 

825 

826 for model in model_list: 

827 model_id = model.get('id') or model.get('name') 

828 

829 if hostname == 'api.openai.com' and not is_supported_openai_models(model_id): 

830 # Skip unwanted OpenAI models 

831 continue 

832 

833 if model_id and model_id not in models: 

834 provider = model.get('provider', '') 

835 merged = { 

836 **model, 

837 'name': model.get('name', model_id), 

838 'owned_by': 'openai', 

839 'openai': model, 

840 'connection_type': model.get('connection_type', 'external'), 

841 'provider': provider, 

842 'urlIdx': idx, 

843 } 

844 

845 loaded = get_provider_model_loaded_state( 

846 model, 

847 provider, 

848 manual_model_ids=bool(api_config.get('model_ids')), 

849 ) 

850 if loaded is not None: 

851 merged['loaded'] = loaded 

852 

853 models[model_id] = merged 

854 

855 return models 

856 

857 models = get_merged_models(map(extract_data, responses)) 

858 log.debug('models: %s', models) 

859 

860 request.app.state.OPENAI_MODELS = models 

861 return {'data': list(models.values())} 

862 

863 

864@router.get('/models') 

865@router.get('/models/{url_idx}') 

866async def get_models(request: Request, url_idx: int | None = None, user=Depends(get_verified_user)): 

867 if url_idx is not None and user.role != 'admin': 867 ↛ 868line 867 didn't jump to line 868 because the condition on line 867 was never true

868 raise HTTPException(status_code=401, detail=ERROR_MESSAGES.ACCESS_PROHIBITED) 

869 

870 if not await Config.get('openai.enable'): 

871 raise HTTPException(status_code=503, detail='OpenAI API is disabled') 

872 

873 models = { 

874 'data': [], 

875 } 

876 

877 if url_idx is None: 

878 models = await get_all_models(request, user=user) 

879 else: 

880 url, key, api_config = await get_openai_connection(url_idx) 

881 

882 r = None 

883 async with aiohttp.ClientSession( 

884 trust_env=True, 

885 timeout=_MODEL_LIST_TIMEOUT, 

886 ) as session: 

887 try: 

888 headers, cookies = await get_headers_and_cookies(request, url, key, api_config, user=user) 

889 

890 if api_config.get('azure') or api_config.get('provider') == 'azure': 

891 models = { 

892 'data': api_config.get('model_ids', []) or [], 

893 'object': 'list', 

894 } 

895 elif is_anthropic_url(url): 

896 models = await get_anthropic_models(url, key, user=user) 

897 if models is None: 

898 raise Exception('Failed to connect to Anthropic API') 

899 else: 

900 async with session.get( 

901 f'{url}/models', 

902 headers=headers, 

903 cookies=cookies, 

904 ssl=AIOHTTP_CLIENT_SESSION_SSL, 

905 ) as r: 

906 if r.status != 200: 

907 error_detail = f'HTTP Error: {r.status}' 

908 try: 

909 res = await r.json(loads=JSONCodec.loads) 

910 if 'error' in res: 

911 error_detail = f'External Error: {res["error"]}' 

912 except Exception: 

913 pass 

914 raise Exception(error_detail) 

915 

916 response_data = await r.json(loads=JSONCodec.loads) 

917 

918 if 'api.openai.com' in url: 

919 response_data['data'] = [ 

920 model 

921 for model in response_data.get('data', []) 

922 if not any(name in model['id'] for name in _UNSUPPORTED_OPENAI_MODEL_KEYWORDS) 

923 ] 

924 

925 models = response_data 

926 except aiohttp.ClientError as e: 

927 # ClientError covers all aiohttp requests issues 

928 log.exception(f'Client error: {str(e)}') 

929 # LICENSE covers this Open WebUI error identifier. 

930 # Do not alter, remove, obscure, or replace it except as LICENSE permits: 

931 # https://docs.openwebui.com/license. 

932 raise HTTPException(status_code=500, detail='Open WebUI: Server Connection Error') 

933 except Exception as e: 

934 log.exception(f'Unexpected error: {e}') 

935 error_detail = f'Unexpected error: {str(e)}' 

936 raise HTTPException(status_code=500, detail=error_detail) 

937 

938 if user.role == 'user' and not BYPASS_MODEL_ACCESS_CONTROL: 938 ↛ 939line 938 didn't jump to line 939 because the condition on line 938 was never true

939 models['data'] = await get_filtered_models(models, user) 

940 

941 return models 

942 

943 

944class ProviderModelOperationForm(BaseModel): 

945 model: str 

946 model_config = ConfigDict(extra='allow') 

947 

948 

949@router.get('/models/{url_idx}/catalog') 

950async def get_provider_model_catalog(request: Request, url_idx: int, user=Depends(get_admin_user)): 

951 return await send_model_management_request(request, url_idx, 'list', user=user) 

952 

953 

954@router.post('/models/{url_idx}/download') 

955async def download_provider_model( 

956 request: Request, 

957 url_idx: int, 

958 form_data: ProviderModelOperationForm, 

959 user=Depends(get_admin_user), 

960): 

961 root_url, _, api_config, provider = await get_model_management_connection(url_idx) 

962 payload = form_data.model_dump(exclude_none=True) 

963 payload['model'] = strip_provider_model_prefix(payload['model'], api_config.get('prefix_id')) 

964 

965 result = await send_model_management_request(request, url_idx, 'download', 'POST', payload, user=user) 

966 await clear_openai_model_cache(request) 

967 await publish_event( 

968 request, 

969 EVENTS.MODEL_PROVIDER_MODEL_CREATED, 

970 actor=user, 

971 subject_id=payload['model'], 

972 data={'provider': provider, 'url_idx': url_idx, 'base_url': root_url}, 

973 ) 

974 return result 

975 

976 

977@router.get('/models/{url_idx}/download/status/{job_id}') 

978async def get_provider_model_download_status( 

979 request: Request, 

980 url_idx: int, 

981 job_id: str, 

982 user=Depends(get_admin_user), 

983): 

984 return await send_model_management_request( 

985 request, 

986 url_idx, 

987 'download_status', 

988 path_params={'job_id': job_id}, 

989 user=user, 

990 ) 

991 

992 

993@router.post('/models/{url_idx}/load') 

994async def load_provider_model( 

995 request: Request, 

996 url_idx: int, 

997 form_data: ProviderModelOperationForm, 

998 user=Depends(get_admin_user), 

999): 

1000 _, _, api_config, _ = await get_model_management_connection(url_idx) 

1001 payload = form_data.model_dump(exclude_none=True) 

1002 payload['model'] = strip_provider_model_prefix(payload['model'], api_config.get('prefix_id')) 

1003 

1004 result = await send_model_management_request(request, url_idx, 'load', 'POST', payload, user=user) 

1005 await clear_openai_model_cache(request) 

1006 return result 

1007 

1008 

1009@router.post('/models/{url_idx}/unload') 

1010async def unload_provider_model( 

1011 request: Request, 

1012 url_idx: int, 

1013 form_data: ProviderModelOperationForm, 

1014 user=Depends(get_admin_user), 

1015): 

1016 _, _, api_config, _ = await get_model_management_connection(url_idx) 

1017 payload = form_data.model_dump(exclude_none=True) 

1018 payload['model'] = strip_provider_model_prefix(payload['model'], api_config.get('prefix_id')) 

1019 

1020 result = await send_model_management_request(request, url_idx, 'unload', 'POST', payload, user=user) 

1021 await clear_openai_model_cache(request) 

1022 return result 

1023 

1024 

1025@router.get('/models/{url_idx}/sse') 

1026async def stream_provider_model_events(request: Request, url_idx: int, user=Depends(get_admin_user)): 

1027 return await send_model_management_request(request, url_idx, 'sse', stream=True, user=user) 

1028 

1029 

1030@router.delete('/models/{url_idx}') 

1031async def delete_provider_model( 

1032 request: Request, 

1033 url_idx: int, 

1034 model: str, 

1035 user=Depends(get_admin_user), 

1036): 

1037 root_url, _, api_config, provider = await get_model_management_connection(url_idx) 

1038 actual_model = strip_provider_model_prefix(model, api_config.get('prefix_id')) 

1039 

1040 result = await send_model_management_request( 

1041 request, 

1042 url_idx, 

1043 'delete', 

1044 'DELETE', 

1045 query={'model': actual_model}, 

1046 user=user, 

1047 ) 

1048 await clear_openai_model_cache(request) 

1049 await publish_event( 

1050 request, 

1051 EVENTS.MODEL_PROVIDER_MODEL_DELETED, 

1052 actor=user, 

1053 subject_id=actual_model, 

1054 data={'provider': provider, 'url_idx': url_idx, 'base_url': root_url}, 

1055 ) 

1056 return result 

1057 

1058 

1059class ConnectionVerificationForm(BaseModel): 

1060 url: str 

1061 key: str 

1062 

1063 config: dict | None = None 

1064 

1065 

1066@router.post('/verify') 

1067async def verify_connection( 

1068 request: Request, 

1069 form_data: ConnectionVerificationForm, 

1070 user=Depends(get_admin_user), 

1071): 

1072 url = form_data.url 

1073 key = form_data.key 

1074 

1075 api_config = form_data.config or {} 

1076 

1077 async with aiohttp.ClientSession( 

1078 trust_env=True, 

1079 timeout=_MODEL_LIST_TIMEOUT, 

1080 ) as session: 

1081 try: 

1082 headers, cookies = await get_headers_and_cookies(request, url, key, api_config, user=user) 

1083 

1084 if api_config.get('azure') or api_config.get('provider') == 'azure': 1084 ↛ 1086line 1084 didn't jump to line 1086 because the condition on line 1084 was never true

1085 # Only set api-key header if not using Azure Entra ID authentication 

1086 auth_type = api_config.get('auth_type', 'bearer') 

1087 if auth_type not in ('azure_ad', 'microsoft_entra_id'): 

1088 headers['api-key'] = key 

1089 

1090 # Azure v1 format: base URL already ends with /openai/v1, 

1091 # use standard /models endpoint without api-version. 

1092 is_azure_v1 = bool(re.search(r'/openai/v1(?:/|$)', url)) 

1093 

1094 if is_azure_v1: 

1095 verify_url = f'{url.rstrip("/")}/models' 

1096 else: 

1097 api_version = api_config.get('api_version', '') or '2023-03-15-preview' 

1098 verify_url = f'{url}/openai/models?api-version={api_version}' 

1099 

1100 async with session.get( 

1101 url=verify_url, 

1102 headers=headers, 

1103 cookies=cookies, 

1104 ssl=AIOHTTP_CLIENT_SESSION_SSL, 

1105 ) as r: 

1106 try: 

1107 response_data = await r.json(loads=JSONCodec.loads) 

1108 except Exception: 

1109 response_data = await r.text() 

1110 

1111 if r.status != 200: 

1112 if isinstance(response_data, (dict, list)): 

1113 return JSONResponse(status_code=r.status, content=response_data) 

1114 else: 

1115 return PlainTextResponse(status_code=r.status, content=response_data) 

1116 

1117 return response_data 

1118 elif is_anthropic_url(url): 1118 ↛ 1119line 1118 didn't jump to line 1119 because the condition on line 1118 was never true

1119 result = await get_anthropic_models(url, key) 

1120 if result is None: 

1121 raise HTTPException(status_code=500, detail=ERROR_MESSAGES.SERVER_CONNECTION_ERROR) 

1122 if 'error' in result: 

1123 raise HTTPException(status_code=500, detail=result['error']) 

1124 return result 

1125 else: 

1126 async with session.get( 

1127 f'{url}/models', 

1128 headers=headers, 

1129 cookies=cookies, 

1130 ssl=AIOHTTP_CLIENT_SESSION_SSL, 

1131 ) as r: 

1132 try: 

1133 response_data = await r.json(loads=JSONCodec.loads) 

1134 except Exception: 

1135 response_data = await r.text() 

1136 

1137 if r.status != 200: 

1138 if isinstance(response_data, (dict, list)): 

1139 return JSONResponse(status_code=r.status, content=response_data) 

1140 else: 

1141 return PlainTextResponse(status_code=r.status, content=response_data) 

1142 

1143 return response_data 

1144 

1145 except aiohttp.ClientError as e: 

1146 # ClientError covers all aiohttp requests issues 

1147 log.exception(f'Client error: {str(e)}') 

1148 raise HTTPException(status_code=500, detail=ERROR_MESSAGES.SERVER_CONNECTION_ERROR) 

1149 except Exception as e: 

1150 log.exception(f'Unexpected error: {e}') 

1151 raise HTTPException(status_code=500, detail=ERROR_MESSAGES.SERVER_CONNECTION_ERROR) 

1152 

1153 

1154def get_azure_allowed_params(api_version: str) -> set[str]: 

1155 allowed_params = { 

1156 'messages', 

1157 'temperature', 

1158 'role', 

1159 'content', 

1160 'contentPart', 

1161 'contentPartImage', 

1162 'enhancements', 

1163 'dataSources', 

1164 'n', 

1165 'stream', 

1166 'stop', 

1167 'max_tokens', 

1168 'presence_penalty', 

1169 'frequency_penalty', 

1170 'logit_bias', 

1171 'user', 

1172 'function_call', 

1173 'functions', 

1174 'tools', 

1175 'tool_choice', 

1176 'top_p', 

1177 'log_probs', 

1178 'top_logprobs', 

1179 'response_format', 

1180 'seed', 

1181 'max_completion_tokens', 

1182 'reasoning_effort', 

1183 } 

1184 

1185 try: 

1186 if api_version >= '2024-09-01-preview': 

1187 allowed_params.add('stream_options') 

1188 except ValueError: 

1189 log.debug('Invalid API version %s for Azure OpenAI. Defaulting to allowed parameters.', api_version) 

1190 

1191 return allowed_params 

1192 

1193 

1194def is_openai_new_model(model: str) -> bool: 

1195 model_lower = model.lower() 

1196 # o-series models (o1, o3, o4, o5, ...) 

1197 if re.match(r'^o\d+', model_lower): 

1198 return True 

1199 # gpt-N where N >= 5 (gpt-5, gpt-5.2, gpt-6, ...) 

1200 m = re.match(r'^gpt-(\d+)', model_lower) 

1201 if m and int(m.group(1)) >= 5: 

1202 return True 

1203 return False 

1204 

1205 

1206def _sanitize_model_for_url(model: str) -> str: 

1207 """Sanitize a model name before interpolating it into a URL path. 

1208 

1209 Rejects path traversal attempts (../, /, \\) and percent-encodes 

1210 the name so it is safe to use as a single URL path segment 

1211 (e.g. Azure deployment name). 

1212 """ 

1213 if not model or '..' in model or '/' in model or '\\' in model: 

1214 raise HTTPException( 

1215 status_code=400, 

1216 detail='Invalid model name: must not be empty or contain path separators or traversal sequences', 

1217 ) 

1218 return quote(model, safe='') 

1219 

1220 

1221def convert_to_azure_payload(url, payload: dict, api_version: str): 

1222 model = payload.get('model', '') 

1223 

1224 # Filter allowed parameters based on Azure OpenAI API 

1225 allowed_params = get_azure_allowed_params(api_version) 

1226 

1227 # Special handling for o-series models 

1228 if is_openai_new_model(model): 

1229 # Convert max_tokens to max_completion_tokens for o-series models 

1230 if 'max_tokens' in payload: 

1231 payload['max_completion_tokens'] = payload['max_tokens'] 

1232 del payload['max_tokens'] 

1233 

1234 # Remove temperature if not 1 for o-series models 

1235 if 'temperature' in payload and payload['temperature'] != 1: 

1236 log.debug( 

1237 'Removing temperature parameter for o-series model %s as only default value (1) is supported', model 

1238 ) 

1239 del payload['temperature'] 

1240 

1241 # Filter out unsupported parameters 

1242 payload = {k: v for k, v in payload.items() if k in allowed_params} 

1243 

1244 # Sanitize model name to prevent path traversal in the deployment URL 

1245 model = _sanitize_model_for_url(model) 

1246 

1247 url = f'{url}/openai/deployments/{model}' 

1248 return url, payload 

1249 

1250 

1251# Fields accepted by the Responses API for each input item type. 

1252RESPONSES_ALLOWED_FIELDS: dict[str, set[str]] = { 

1253 'message': {'type', 'role', 'content'}, 

1254 'function_call': {'type', 'call_id', 'name', 'arguments', 'id'}, 

1255 'function_call_output': {'type', 'call_id', 'output'}, 

1256} 

1257 

1258 

1259def _normalize_stored_item(item: dict) -> dict: 

1260 """Strip local-only fields from a stored output item before replaying it. 

1261 

1262 Open WebUI stores extra bookkeeping fields (``id``, ``status``, 

1263 ``started_at``, ``ended_at``, ``duration``, ``_tag_type``, 

1264 ``attributes``, ``summary``, etc.) that the Responses API does 

1265 not accept. This helper returns a copy containing only the 

1266 fields the API understands. 

1267 """ 

1268 item_type = item.get('type', '') 

1269 allowed = RESPONSES_ALLOWED_FIELDS.get(item_type) 

1270 if allowed is None: 

1271 # Unknown type — pass through as-is (e.g. reasoning, extension items). 

1272 return item 

1273 return {k: v for k, v in item.items() if k in allowed} 

1274 

1275 

1276def convert_to_responses_payload(payload: dict) -> dict: 

1277 """ 

1278 Convert Chat Completions payload to Responses API format. 

1279 

1280 Chat Completions: { messages: [{role, content}], ... } 

1281 Responses API: { input: [{type: "message", role, content: [...]}], instructions: "system" } 

1282 """ 

1283 messages = payload.pop('messages', []) 

1284 

1285 system_content = '' 

1286 input_items = [] 

1287 

1288 for msg in messages: 

1289 role = msg.get('role', 'user') 

1290 content = msg.get('content', '') 

1291 

1292 # Check for stored output items (from previous Responses API turn) 

1293 stored_output = msg.get('output') 

1294 if stored_output and isinstance(stored_output, list): 

1295 input_items.extend(_normalize_stored_item(item) for item in stored_output) 

1296 continue 

1297 

1298 if role == 'system': 

1299 if isinstance(content, str): 

1300 system_content = content 

1301 elif isinstance(content, list): 

1302 system_content = '\n'.join(p.get('text', '') for p in content if p.get('type') == 'text') 

1303 continue 

1304 

1305 # Handle assistant messages with tool_calls (from convert_output_to_messages) 

1306 if role == 'assistant' and msg.get('tool_calls'): 

1307 # Add text content as message if present 

1308 if content: 

1309 text = ( 

1310 content 

1311 if isinstance(content, str) 

1312 else '\n'.join(p.get('text', '') for p in content if p.get('type') == 'text') 

1313 ) 

1314 if text.strip(): 

1315 input_items.append( 

1316 { 

1317 'type': 'message', 

1318 'role': 'assistant', 

1319 'content': [{'type': 'output_text', 'text': text}], 

1320 } 

1321 ) 

1322 # Convert each tool_call to a function_call input item 

1323 for tool_call in msg['tool_calls']: 

1324 func = tool_call.get('function', {}) 

1325 input_items.append( 

1326 { 

1327 'type': 'function_call', 

1328 'call_id': tool_call.get('id', ''), 

1329 'name': func.get('name', ''), 

1330 'arguments': func.get('arguments', '{}'), 

1331 } 

1332 ) 

1333 continue 

1334 

1335 # Handle tool result messages 

1336 if role == 'tool': 

1337 input_items.append( 

1338 { 

1339 'type': 'function_call_output', 

1340 'call_id': msg.get('tool_call_id', ''), 

1341 'output': msg.get('content', ''), 

1342 } 

1343 ) 

1344 continue 

1345 

1346 # Convert content format 

1347 text_type = 'output_text' if role == 'assistant' else 'input_text' 

1348 

1349 if isinstance(content, str): 

1350 content_parts = [{'type': text_type, 'text': content}] 

1351 elif isinstance(content, list): 

1352 content_parts = [] 

1353 for part in content: 

1354 if part.get('type') == 'text': 

1355 content_parts.append({'type': text_type, 'text': part.get('text', '')}) 

1356 elif part.get('type') == 'image_url': 

1357 url_data = part.get('image_url', {}) 

1358 if isinstance(url_data, dict): 

1359 url = url_data.get('url', '') 

1360 detail = url_data.get('detail') or 'auto' 

1361 else: 

1362 url = url_data if isinstance(url_data, str) else '' 

1363 detail = 'auto' 

1364 content_parts.append({'type': 'input_image', 'image_url': url, 'detail': detail}) 

1365 elif part.get('type') == 'file': 

1366 # OpenAI-compatible proxy path only. Open WebUI attachments are handled 

1367 # separately via metadata.files/RAG and must not be converted here. 

1368 file = part.get('file') 

1369 if isinstance(file, dict): 

1370 file_part = {k: file[k] for k in ('file_id', 'file_data', 'filename') if k in file} 

1371 if 'file_id' in file_part or 'file_data' in file_part: 

1372 content_parts.append({'type': 'input_file', **file_part}) 

1373 else: 

1374 content_parts = [{'type': text_type, 'text': str(content)}] 

1375 

1376 input_items.append({'type': 'message', 'role': role, 'content': content_parts}) 

1377 

1378 responses_payload = {**payload, 'input': input_items} 

1379 

1380 # Forward previous_response_id when the middleware has set it 

1381 # (only used when ENABLE_RESPONSES_API_STATEFUL is enabled). 

1382 previous_response_id = responses_payload.pop('previous_response_id', None) 

1383 if previous_response_id: 

1384 responses_payload['previous_response_id'] = previous_response_id 

1385 

1386 if system_content: 

1387 responses_payload['instructions'] = system_content 

1388 

1389 if 'max_tokens' in responses_payload: 

1390 responses_payload['max_output_tokens'] = responses_payload.pop('max_tokens') 

1391 

1392 if 'max_completion_tokens' in responses_payload: 

1393 responses_payload['max_output_tokens'] = responses_payload.pop('max_completion_tokens') 

1394 

1395 # Remove Chat Completions-only parameters not supported by the Responses API 

1396 for unsupported_key in ( 

1397 'stream_options', 

1398 'logit_bias', 

1399 'frequency_penalty', 

1400 'presence_penalty', 

1401 'stop', 

1402 ): 

1403 responses_payload.pop(unsupported_key, None) 

1404 

1405 # Convert Chat Completions tools format to Responses API format 

1406 # Chat Completions: {"type": "function", "function": {"name": ..., "description": ..., "parameters": ...}} 

1407 # Responses API: {"type": "function", "name": ..., "description": ..., "parameters": ...} 

1408 if 'tools' in responses_payload and isinstance(responses_payload['tools'], list): 

1409 converted_tools = [] 

1410 for tool in responses_payload['tools']: 

1411 if isinstance(tool, dict) and 'function' in tool: 

1412 func = tool['function'] 

1413 converted_tool = {'type': tool.get('type', 'function')} 

1414 if isinstance(func, dict): 

1415 converted_tool['name'] = func.get('name', '') 

1416 if 'description' in func: 

1417 converted_tool['description'] = func['description'] 

1418 if 'parameters' in func: 

1419 converted_tool['parameters'] = func['parameters'] 

1420 # Responses defaults strict to true, Chat Completions to false 

1421 converted_tool['strict'] = func.get('strict', False) 

1422 converted_tools.append(converted_tool) 

1423 else: 

1424 # Already in correct format or unknown format, pass through 

1425 converted_tools.append(tool) 

1426 responses_payload['tools'] = converted_tools 

1427 

1428 # Responses API expects a forced function choice as {"type": "function", "name": ...} 

1429 tool_choice = responses_payload.get('tool_choice') 

1430 if isinstance(tool_choice, dict) and isinstance(tool_choice.get('function'), dict): 

1431 responses_payload['tool_choice'] = {'type': 'function', 'name': tool_choice['function'].get('name', '')} 

1432 

1433 return responses_payload 

1434 

1435 

1436def convert_responses_result(response: dict) -> dict: 

1437 """ 

1438 Convert non-streaming Responses API result to Chat Completions format. 

1439 

1440 Extracts text and function calls from output items so all downstream consumers 

1441 (frontend tasks, get_content_from_response) work without modification. 

1442 """ 

1443 output_items = response.get('output', []) 

1444 

1445 content = '' 

1446 tool_calls = [] 

1447 for item in output_items: 

1448 if item.get('type') == 'message': 

1449 for part in item.get('content', []): 

1450 if part.get('type') == 'output_text': 

1451 content += part.get('text', '') 

1452 elif item.get('type') == 'function_call': 

1453 arguments = item.get('arguments', '{}') 

1454 if not isinstance(arguments, str): 

1455 arguments = JSONCodec.dumps(arguments) 

1456 tool_calls.append( 

1457 { 

1458 'id': item.get('call_id', ''), 

1459 'type': 'function', 

1460 'function': { 

1461 'name': item.get('name', ''), 

1462 'arguments': arguments, 

1463 }, 

1464 } 

1465 ) 

1466 

1467 return { 

1468 'id': response.get('id', ''), 

1469 'object': 'chat.completion', 

1470 'model': response.get('model', ''), 

1471 'choices': [ 

1472 { 

1473 'index': 0, 

1474 'message': { 

1475 'role': 'assistant', 

1476 'content': content, 

1477 **({'tool_calls': tool_calls} if tool_calls else {}), 

1478 }, 

1479 'finish_reason': 'tool_calls' if tool_calls else 'stop', 

1480 } 

1481 ], 

1482 'usage': response.get('usage', {}), 

1483 } 

1484 

1485 

1486@router.post('/chat/completions') 

1487async def generate_chat_completion( 

1488 request: Request, 

1489 form_data: dict, 

1490 user=Depends(get_verified_user), 

1491): 

1492 if not await Config.get('openai.enable'): 

1493 raise HTTPException(status_code=503, detail='OpenAI API is disabled') 

1494 

1495 # NOTE: We intentionally do NOT use Depends(get_async_session) here. 

1496 # Database operations (get_model_by_id, AccessGrants.has_access) manage their own short-lived sessions. 

1497 # This prevents holding a connection during the entire LLM call (30-60+ seconds), 

1498 # which would exhaust the connection pool under concurrent load. 

1499 

1500 # bypass_filter and bypass_system_prompt are read from request.state to prevent 

1501 # external clients from setting them via query parameter. Only internal 

1502 # server-side callers (e.g. utils/chat.py) should set 

1503 # request.state.bypass_filter / request.state.bypass_system_prompt = True. 

1504 bypass_filter = getattr(request.state, 'bypass_filter', False) 

1505 if BYPASS_MODEL_ACCESS_CONTROL: 1505 ↛ 1506line 1505 didn't jump to line 1506 because the condition on line 1505 was never true

1506 bypass_filter = True 

1507 bypass_system_prompt = getattr(request.state, 'bypass_system_prompt', False) 

1508 

1509 idx = 0 

1510 

1511 payload = {**form_data} 

1512 metadata = payload.pop('metadata', None) 

1513 

1514 model_id = form_data.get('model') 

1515 model_info = await Models.get_model_by_id(model_id) 

1516 

1517 # Check model info and override the payload 

1518 if model_info: 1518 ↛ 1519line 1518 didn't jump to line 1519 because the condition on line 1518 was never true

1519 if model_info.base_model_id: 

1520 base_model_id = ( 

1521 request.base_model_id if hasattr(request, 'base_model_id') else model_info.base_model_id 

1522 ) # Use request's base_model_id if available 

1523 payload['model'] = base_model_id 

1524 model_id = base_model_id 

1525 

1526 params = model_info.params.model_dump() 

1527 

1528 if params: 

1529 system = params.pop('system', None) 

1530 

1531 payload = apply_model_params_to_body_openai(params, payload) 

1532 if not bypass_system_prompt: 

1533 payload = await apply_system_prompt_to_body(system, payload, metadata, user) 

1534 

1535 await check_model_access(user, model_info, bypass_filter) 

1536 else: 

1537 await check_model_access(user, None, bypass_filter) 

1538 

1539 # Check if model is already in app state cache to avoid expensive get_all_models() call 

1540 models = request.app.state.OPENAI_MODELS 

1541 if not models or model_id not in models: 1541 ↛ 1544line 1541 didn't jump to line 1544 because the condition on line 1541 was always true

1542 await get_all_models(request, user=user) 

1543 models = request.app.state.OPENAI_MODELS 

1544 model = models.get(model_id) 

1545 

1546 if model: 1546 ↛ 1547line 1546 didn't jump to line 1547 because the condition on line 1546 was never true

1547 idx = model['urlIdx'] 

1548 else: 

1549 raise HTTPException( 

1550 status_code=404, 

1551 detail=ERROR_MESSAGES.MODEL_NOT_FOUND(), 

1552 ) 

1553 

1554 url, key, api_config = await get_openai_connection(idx) 

1555 

1556 prefix_id = api_config.get('prefix_id', None) 

1557 payload['model'] = strip_provider_model_prefix(payload['model'], prefix_id) 

1558 

1559 # Add user info to the payload if the model is a pipeline 

1560 if 'pipeline' in model and model.get('pipeline'): 

1561 payload['user'] = { 

1562 'name': user.name, 

1563 'id': user.id, 

1564 'email': user.email, 

1565 'role': user.role, 

1566 } 

1567 

1568 # Check if model is a reasoning model that needs special handling 

1569 if is_openai_new_model(payload['model']): 

1570 payload = openai_reasoning_model_handler(payload) 

1571 elif 'api.openai.com' not in url: 

1572 # Remove "max_completion_tokens" from the payload for backward compatibility 

1573 if 'max_completion_tokens' in payload: 

1574 payload['max_tokens'] = payload['max_completion_tokens'] 

1575 del payload['max_completion_tokens'] 

1576 

1577 if 'max_tokens' in payload and 'max_completion_tokens' in payload: 

1578 del payload['max_tokens'] 

1579 

1580 # Convert the modified body back to JSON 

1581 if 'logit_bias' in payload and payload['logit_bias']: 

1582 logit_bias = convert_logit_bias_input_to_json(payload['logit_bias']) 

1583 

1584 if logit_bias: 

1585 payload['logit_bias'] = JSONCodec.loads(logit_bias) 

1586 

1587 headers, cookies = await get_headers_and_cookies(request, url, key, api_config, metadata, user=user) 

1588 

1589 is_responses = api_config.get('api_type') == 'responses' 

1590 

1591 # Explicit continuation keeps llama.cpp from echoing the prefill in streamed replies. 

1592 if ( 

1593 api_config.get('provider') == 'llama.cpp' 

1594 # These flags apply to Chat Completions, not the Responses API. 

1595 and not is_responses 

1596 # The frontend sends this ID when the user clicks Continue. 

1597 and (metadata or {}).get('assistant_message_id') 

1598 # Tool follow-ups retain the metadata but must start a new assistant turn. 

1599 and payload.get('messages') 

1600 and payload['messages'][-1].get('role') == 'assistant' 

1601 ): 

1602 payload['continue_final_message'] = True 

1603 payload['add_generation_prompt'] = False 

1604 

1605 if api_config.get('azure') or api_config.get('provider') == 'azure': 

1606 # Only set api-key header if not using Azure Entra ID authentication 

1607 auth_type = api_config.get('auth_type', 'bearer') 

1608 if auth_type not in ('azure_ad', 'microsoft_entra_id'): 

1609 headers['api-key'] = key 

1610 

1611 # Azure v1 format: base URL already ends with /openai/v1, 

1612 # model stays in the payload, no deployment URL rewriting. 

1613 is_azure_v1 = bool(re.search(r'/openai/v1(?:/|$)', url)) 

1614 

1615 if is_azure_v1: 

1616 if is_responses: 

1617 payload = convert_to_responses_payload(payload) 

1618 request_url = f'{url.rstrip("/")}/responses' 

1619 else: 

1620 request_url = f'{url.rstrip("/")}/chat/completions' 

1621 else: 

1622 api_version = api_config.get('api_version', '2023-03-15-preview') 

1623 request_url, payload = convert_to_azure_payload(url, payload, api_version) 

1624 headers['api-version'] = api_version 

1625 

1626 if is_responses: 

1627 payload = convert_to_responses_payload(payload) 

1628 request_url = f'{request_url}/responses?api-version={api_version}' 

1629 else: 

1630 request_url = f'{request_url}/chat/completions?api-version={api_version}' 

1631 else: 

1632 if is_responses: 

1633 payload = convert_to_responses_payload(payload) 

1634 request_url = f'{url}/responses' 

1635 else: 

1636 request_url = f'{url}/chat/completions' 

1637 requested_model = payload.get('model') 

1638 # For Chat Completions, strip image parts from multimodal tool messages 

1639 # (Chat Completions doesn't support images in tool content). 

1640 if not is_responses and 'messages' in payload: 

1641 for message in payload['messages']: 

1642 if message.get('role') == 'tool' and isinstance(message.get('content'), list): 

1643 message['content'] = ''.join( 

1644 part.get('text', '') for part in message['content'] if part.get('type') in ('input_text', 'text') 

1645 ) 

1646 

1647 is_streaming_request = bool(payload.get('stream', False)) 

1648 if not is_streaming_request: 

1649 payload.pop('stream_options', None) 

1650 

1651 payload = JSONCodec.dumps(payload) 

1652 

1653 r = None 

1654 streaming = False 

1655 response = None 

1656 

1657 try: 

1658 session = await get_session() 

1659 

1660 r = await session.request( 

1661 method='POST', 

1662 url=request_url, 

1663 data=payload, 

1664 headers=headers, 

1665 cookies=cookies, 

1666 ssl=AIOHTTP_CLIENT_SESSION_SSL, 

1667 timeout=get_client_timeout(stream=is_streaming_request), 

1668 ) 

1669 

1670 # Check if response is SSE 

1671 if 'text/event-stream' in r.headers.get('Content-Type', ''): 

1672 # If the provider returned an error status with SSE content-type, 

1673 # read the body and return a proper error response instead of 

1674 # streaming the error back (which hides the error from logs). 

1675 if r.status >= 400: 

1676 error_body = await r.text() 

1677 log.error( 

1678 'Provider returned HTTP %d with SSE content-type: %s', 

1679 r.status, 

1680 error_body[:1000], 

1681 ) 

1682 try: 

1683 error_json = JSONCodec.loads(error_body) 

1684 await publish_model_provider_request_failed( 

1685 request, 

1686 actor=user, 

1687 provider='openai-compatible', 

1688 base_url=url, 

1689 api_key=key, 

1690 status=r.status, 

1691 requested_model=requested_model, 

1692 upstream_error=error_json, 

1693 ) 

1694 return JSONResponse(status_code=r.status, content=error_json) 

1695 except JSONCodec.JSONDecodeError: 

1696 await publish_model_provider_request_failed( 

1697 request, 

1698 actor=user, 

1699 provider='openai-compatible', 

1700 base_url=url, 

1701 api_key=key, 

1702 status=r.status, 

1703 requested_model=requested_model, 

1704 upstream_error=error_body, 

1705 ) 

1706 return JSONResponse( 

1707 status_code=r.status, 

1708 content={'error': {'message': error_body, 'code': r.status}}, 

1709 ) 

1710 

1711 streaming = True 

1712 return StreamingResponse( 

1713 stream_wrapper(r), 

1714 status_code=r.status, 

1715 headers=_clean_proxy_headers(r.headers), 

1716 ) 

1717 else: 

1718 try: 

1719 response = await r.json(loads=JSONCodec.loads) 

1720 except Exception as e: 

1721 log.error(e) 

1722 response = await r.text() 

1723 

1724 if r.status >= 400: 

1725 await publish_model_provider_request_failed( 

1726 request, 

1727 actor=user, 

1728 provider='openai-compatible', 

1729 base_url=url, 

1730 api_key=key, 

1731 status=r.status, 

1732 requested_model=requested_model, 

1733 upstream_error=response, 

1734 ) 

1735 if isinstance(response, (dict, list)): 

1736 return JSONResponse(status_code=r.status, content=response) 

1737 else: 

1738 return PlainTextResponse(status_code=r.status, content=response) 

1739 

1740 # Convert Responses API result to simple format 

1741 if is_responses and isinstance(response, dict): 

1742 response = convert_responses_result(response) 

1743 

1744 return response 

1745 except Exception as e: 

1746 log.exception(e) 

1747 

1748 raise HTTPException( 

1749 status_code=r.status if r else 500, 

1750 detail=ERROR_MESSAGES.SERVER_CONNECTION_ERROR, 

1751 ) 

1752 finally: 

1753 if not streaming: 

1754 await cleanup_response(r) 

1755 

1756 

1757async def embeddings(request: Request, form_data: dict, user): 

1758 """ 

1759 Calls the embeddings endpoint for OpenAI-compatible providers. 

1760 

1761 Args: 

1762 request (Request): The FastAPI request context. 

1763 form_data (dict): OpenAI-compatible embeddings payload. 

1764 user (UserModel): The authenticated user. 

1765 

1766 Returns: 

1767 dict: OpenAI-compatible embeddings response. 

1768 """ 

1769 idx = 0 

1770 # Prepare payload/body 

1771 body = JSONCodec.dumps(form_data) 

1772 # Find correct backend url/key based on model 

1773 model_id = form_data.get('model') 

1774 # Check if model is already in app state cache to avoid expensive get_all_models() call 

1775 models = request.app.state.OPENAI_MODELS 

1776 if not models or model_id not in models: 

1777 await get_all_models(request, user=user) 

1778 models = request.app.state.OPENAI_MODELS 

1779 if model_id in models: 

1780 idx = models[model_id]['urlIdx'] 

1781 

1782 url, key, api_config = await get_openai_connection(idx) 

1783 

1784 r = None 

1785 streaming = False 

1786 

1787 headers, cookies = await get_headers_and_cookies(request, url, key, api_config, user=user) 

1788 

1789 if api_config.get('azure') or api_config.get('provider') == 'azure': 

1790 # Only set api-key header if not using Azure Entra ID authentication 

1791 auth_type = api_config.get('auth_type', 'bearer') 

1792 if auth_type not in ('azure_ad', 'microsoft_entra_id'): 

1793 headers['api-key'] = key 

1794 

1795 # Azure v1 format: base URL already ends with /openai/v1, 

1796 # model stays in the payload, no deployment URL rewriting. 

1797 is_azure_v1 = bool(re.search(r'/openai/v1(?:/|$)', url)) 

1798 

1799 if is_azure_v1: 

1800 embeddings_url = f'{url.rstrip("/")}/embeddings' 

1801 else: 

1802 api_version = api_config.get('api_version', '2023-03-15-preview') 

1803 model = _sanitize_model_for_url(form_data.get('model', '')) 

1804 embeddings_url = f'{url}/openai/deployments/{model}/embeddings?api-version={api_version}' 

1805 headers['api-version'] = api_version 

1806 else: 

1807 embeddings_url = f'{url}/embeddings' 

1808 requested_model = form_data.get('model') 

1809 

1810 try: 

1811 session = await get_session() 

1812 r = await session.request( 

1813 method='POST', 

1814 url=embeddings_url, 

1815 data=body, 

1816 headers=headers, 

1817 cookies=cookies, 

1818 timeout=get_client_timeout(), 

1819 ssl=AIOHTTP_CLIENT_SESSION_SSL, 

1820 ) 

1821 

1822 if 'text/event-stream' in r.headers.get('Content-Type', ''): 

1823 streaming = True 

1824 return StreamingResponse( 

1825 stream_wrapper(r, passthrough=True), 

1826 status_code=r.status, 

1827 headers=_clean_proxy_headers(r.headers), 

1828 ) 

1829 else: 

1830 try: 

1831 response_data = await r.json(loads=JSONCodec.loads) 

1832 except Exception: 

1833 response_data = await r.text() 

1834 

1835 if r.status >= 400: 

1836 await publish_model_provider_request_failed( 

1837 request, 

1838 actor=user, 

1839 provider='openai-compatible', 

1840 base_url=url, 

1841 api_key=key, 

1842 status=r.status, 

1843 requested_model=requested_model, 

1844 upstream_error=response_data, 

1845 ) 

1846 if isinstance(response_data, (dict, list)): 

1847 return JSONResponse(status_code=r.status, content=response_data) 

1848 else: 

1849 return PlainTextResponse(status_code=r.status, content=response_data) 

1850 

1851 return response_data 

1852 except Exception as e: 

1853 log.exception(e) 

1854 raise HTTPException( 

1855 status_code=r.status if r else 500, 

1856 detail=ERROR_MESSAGES.SERVER_CONNECTION_ERROR, 

1857 ) 

1858 finally: 

1859 if not streaming: 

1860 await cleanup_response(r) 

1861 

1862 

1863class ResponsesForm(BaseModel): 

1864 model_config = ConfigDict(extra='allow') 

1865 

1866 model: str 

1867 input: list | str | None = None 

1868 instructions: str | None = None 

1869 stream: bool | None = None 

1870 temperature: float | None = None 

1871 max_output_tokens: int | None = None 

1872 top_p: float | None = None 

1873 tools: list | None = None 

1874 tool_choice: str | dict | None = None 

1875 text: dict | None = None 

1876 truncation: str | None = None 

1877 metadata: dict | None = None 

1878 store: bool | None = None 

1879 reasoning: dict | None = None 

1880 previous_response_id: str | None = None 

1881 

1882 

1883@router.post('/responses') 

1884async def responses( 

1885 request: Request, 

1886 form_data: ResponsesForm, 

1887 user=Depends(get_verified_user), 

1888): 

1889 """ 

1890 Forward requests to the OpenAI Responses API endpoint. 

1891 Routes to the correct upstream backend based on the model field. 

1892 """ 

1893 payload = form_data.model_dump(exclude_none=True) 

1894 is_streaming_request = bool(payload.get('stream', False)) 

1895 

1896 idx = 0 

1897 model_id = form_data.model 

1898 

1899 # Enforce per-model access control 

1900 await check_model_access(user, await Models.get_model_by_id(model_id), BYPASS_MODEL_ACCESS_CONTROL) 

1901 

1902 if model_id: 

1903 models = request.app.state.OPENAI_MODELS 

1904 if not models or model_id not in models: 1904 ↛ 1907line 1904 didn't jump to line 1907 because the condition on line 1904 was always true

1905 await get_all_models(request, user=user) 

1906 models = request.app.state.OPENAI_MODELS 

1907 if model_id in models: 1907 ↛ 1908line 1907 didn't jump to line 1908 because the condition on line 1907 was never true

1908 idx = models[model_id]['urlIdx'] 

1909 

1910 url, key, api_config = await get_openai_connection(idx) 

1911 

1912 payload['model'] = strip_provider_model_prefix(payload['model'], api_config.get('prefix_id')) 

1913 body = JSONCodec.dumps(payload) 

1914 

1915 r = None 

1916 streaming = False 

1917 

1918 try: 

1919 headers, cookies = await get_headers_and_cookies(request, url, key, api_config, user=user) 

1920 

1921 if api_config.get('azure') or api_config.get('provider') == 'azure': 1921 ↛ 1922line 1921 didn't jump to line 1922 because the condition on line 1921 was never true

1922 auth_type = api_config.get('auth_type', 'bearer') 

1923 if auth_type not in ('azure_ad', 'microsoft_entra_id'): 

1924 headers['api-key'] = key 

1925 

1926 is_azure_v1 = bool(re.search(r'/openai/v1(?:/|$)', url)) 

1927 

1928 if is_azure_v1: 

1929 request_url = f'{url.rstrip("/")}/responses' 

1930 else: 

1931 api_version = api_config.get('api_version', '2023-03-15-preview') 

1932 headers['api-version'] = api_version 

1933 model = _sanitize_model_for_url(payload.get('model', '')) 

1934 request_url = f'{url}/openai/deployments/{model}/responses?api-version={api_version}' 

1935 else: 

1936 request_url = f'{url}/responses' 

1937 

1938 session = await get_session() 

1939 r = await session.request( 

1940 method='POST', 

1941 url=request_url, 

1942 data=body, 

1943 headers=headers, 

1944 cookies=cookies, 

1945 ssl=AIOHTTP_CLIENT_SESSION_SSL, 

1946 timeout=get_client_timeout(stream=is_streaming_request), 

1947 ) 

1948 

1949 # Check if response is SSE 

1950 if 'text/event-stream' in r.headers.get('Content-Type', ''): 

1951 streaming = True 

1952 return StreamingResponse( 

1953 stream_wrapper(r, passthrough=True), 

1954 status_code=r.status, 

1955 headers=_clean_proxy_headers(r.headers), 

1956 ) 

1957 else: 

1958 try: 

1959 response_data = await r.json(loads=JSONCodec.loads) 

1960 except Exception: 

1961 response_data = await r.text() 

1962 

1963 if r.status >= 400: 

1964 await publish_model_provider_request_failed( 

1965 request, 

1966 actor=user, 

1967 provider='openai-compatible', 

1968 base_url=url, 

1969 api_key=key, 

1970 status=r.status, 

1971 requested_model=payload.get('model'), 

1972 upstream_error=response_data, 

1973 ) 

1974 if isinstance(response_data, (dict, list)): 

1975 return JSONResponse(status_code=r.status, content=response_data) 

1976 else: 

1977 return PlainTextResponse(status_code=r.status, content=response_data) 

1978 

1979 return response_data 

1980 

1981 except HTTPException: 

1982 raise 

1983 except Exception as e: 

1984 log.exception(e) 

1985 raise HTTPException( 

1986 status_code=r.status if r else 500, 

1987 detail=ERROR_MESSAGES.SERVER_CONNECTION_ERROR, 

1988 ) 

1989 finally: 

1990 if not streaming: 

1991 await cleanup_response(r) 

1992 

1993 

1994@router.api_route('/{path:path}', methods=['GET', 'POST', 'PUT', 'DELETE']) 

1995async def proxy(path: str, request: Request, user=Depends(get_verified_user)): 

1996 """ 

1997 Deprecated: proxy all requests to OpenAI API. 

1998 Disabled by default. Set ENABLE_OPENAI_API_PASSTHROUGH=True to enable. 

1999 """ 

2000 

2001 if not ENABLE_OPENAI_API_PASSTHROUGH: 2001 ↛ 2007line 2001 didn't jump to line 2007 because the condition on line 2001 was always true

2002 raise HTTPException( 

2003 status_code=status.HTTP_403_FORBIDDEN, 

2004 detail='Direct API passthrough is disabled. Set ENABLE_OPENAI_API_PASSTHROUGH=True to enable.', 

2005 ) 

2006 

2007 body = await request.body() 

2008 

2009 # Parse JSON body to resolve model-based routing 

2010 payload = None 

2011 if body: 

2012 try: 

2013 payload = JSONCodec.loads(body) 

2014 except (JSONCodec.JSONDecodeError, ValueError): 

2015 payload = None 

2016 is_streaming_request = bool(payload.get('stream', False)) if isinstance(payload, dict) else False 

2017 

2018 idx = 0 

2019 model_id = payload.get('model') if isinstance(payload, dict) else None 

2020 if model_id: 

2021 models = request.app.state.OPENAI_MODELS 

2022 if not models or model_id not in models: 

2023 await get_all_models(request, user=user) 

2024 models = request.app.state.OPENAI_MODELS 

2025 if model_id in models: 

2026 idx = models[model_id]['urlIdx'] 

2027 

2028 url, key, api_config = await get_openai_connection(idx) 

2029 base_url = url 

2030 

2031 r = None 

2032 streaming = False 

2033 

2034 try: 

2035 headers, cookies = await get_headers_and_cookies(request, url, key, api_config, user=user) 

2036 

2037 if api_config.get('azure') or api_config.get('provider') == 'azure': 

2038 # Only set api-key header if not using Azure Entra ID authentication 

2039 auth_type = api_config.get('auth_type', 'bearer') 

2040 if auth_type not in ('azure_ad', 'microsoft_entra_id'): 

2041 headers['api-key'] = key 

2042 

2043 is_azure_v1 = bool(re.search(r'/openai/v1(?:/|$)', url)) 

2044 

2045 if is_azure_v1: 

2046 qs = request.url.query 

2047 request_url = f'{url.rstrip("/")}/{path}' + (f'?{qs}' if qs else '') 

2048 else: 

2049 api_version = api_config.get('api_version', '2023-03-15-preview') 

2050 headers['api-version'] = api_version 

2051 

2052 payload = JSONCodec.loads(body) 

2053 url, payload = convert_to_azure_payload(url, payload, api_version) 

2054 body = JSONCodec.dumps(payload).encode() 

2055 

2056 request_url = f'{url}/{path}?api-version={api_version}' 

2057 else: 

2058 request_url = f'{url}/{path}' 

2059 

2060 session = await get_session() 

2061 r = await session.request( 

2062 method=request.method, 

2063 url=request_url, 

2064 data=body, 

2065 headers=headers, 

2066 cookies=cookies, 

2067 ssl=AIOHTTP_CLIENT_SESSION_SSL, 

2068 timeout=get_client_timeout(stream=is_streaming_request), 

2069 ) 

2070 

2071 # Check if response is SSE 

2072 if 'text/event-stream' in r.headers.get('Content-Type', ''): 

2073 streaming = True 

2074 return StreamingResponse( 

2075 stream_wrapper(r, passthrough=True), 

2076 status_code=r.status, 

2077 headers=_clean_proxy_headers(r.headers), 

2078 ) 

2079 else: 

2080 try: 

2081 response_data = await r.json(loads=JSONCodec.loads) 

2082 except Exception: 

2083 response_data = await r.text() 

2084 

2085 if r.status >= 400: 

2086 await publish_model_provider_request_failed( 

2087 request, 

2088 actor=user, 

2089 provider='openai-compatible', 

2090 base_url=base_url, 

2091 api_key=key, 

2092 status=r.status, 

2093 requested_model=model_id, 

2094 upstream_error=response_data, 

2095 ) 

2096 if isinstance(response_data, (dict, list)): 

2097 return JSONResponse(status_code=r.status, content=response_data) 

2098 else: 

2099 return PlainTextResponse(status_code=r.status, content=response_data) 

2100 

2101 return response_data 

2102 

2103 except HTTPException: 

2104 raise 

2105 except Exception as e: 

2106 log.exception(e) 

2107 # LICENSE covers this Open WebUI error identifier. 

2108 # Do not alter, remove, obscure, or replace it except as LICENSE permits: 

2109 # https://docs.openwebui.com/license. 

2110 raise HTTPException( 

2111 status_code=r.status if r else 500, 

2112 detail='Open WebUI: Server Connection Error', 

2113 ) 

2114 finally: 

2115 if not streaming: 

2116 await cleanup_response(r)