Coverage for open_webui/main.py: 39%

1155 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 copy 

5import logging 

6import mimetypes 

7import os 

8import sys 

9import time 

10from concurrent.futures import ThreadPoolExecutor 

11from contextlib import asynccontextmanager 

12from uuid import uuid4 

13 

14import aiohttp 

15import anyio.to_thread 

16from cryptography.fernet import InvalidToken 

17from fastapi import ( 

18 Depends, 

19 FastAPI, 

20 HTTPException, 

21 Request, 

22 applications, 

23 status, 

24) 

25from fastapi.middleware.cors import CORSMiddleware 

26from fastapi.openapi.docs import get_swagger_ui_html 

27from fastapi.responses import FileResponse, JSONResponse 

28from fastapi.staticfiles import StaticFiles 

29from pydantic import BaseModel 

30from sqlalchemy import text 

31from sqlalchemy.ext.asyncio import AsyncSession 

32from starlette.datastructures import Headers 

33from starlette.exceptions import HTTPException as StarletteHTTPException 

34from starlette.middleware.sessions import SessionMiddleware 

35from starlette.responses import Response, StreamingResponse 

36from starlette_compress import CompressMiddleware 

37from starsessions import ( 

38 SessionAutoloadMiddleware, 

39) 

40from starsessions import ( 

41 SessionMiddleware as StarSessionsMiddleware, 

42) 

43from starsessions.stores.redis import RedisStore 

44 

45from open_webui.config import ( 

46 BYPASS_ADMIN_ACCESS_CONTROL, 

47 CACHE_DIR, 

48 CORS_ALLOW_ORIGIN, 

49 DEFAULT_LOCALE, 

50 ENABLE_ADMIN_ANALYTICS, 

51 # Admin 

52 ENABLE_ADMIN_CHAT_ACCESS, 

53 ENABLE_ADMIN_EXPORT, 

54 ENABLE_ONEDRIVE_BUSINESS, 

55 ENABLE_ONEDRIVE_PERSONAL, 

56 # OpenAI 

57 ENV, 

58 FRONTEND_BUILD_DIR, 

59 GOOGLE_DRIVE_API_KEY, 

60 GOOGLE_DRIVE_CLIENT_ID, 

61 IFRAME_CSP, 

62 OAUTH_PROVIDERS, 

63 ONEDRIVE_CLIENT_ID_BUSINESS, 

64 ONEDRIVE_CLIENT_ID_PERSONAL, 

65 ONEDRIVE_SHAREPOINT_TENANT_ID, 

66 ONEDRIVE_SHAREPOINT_URL, 

67 STATIC_DIR, 

68 THREAD_POOL_SIZE, 

69 THREAD_POOL_THREAD_NAME_PREFIX, 

70 WEBUI_AUTH, 

71 WEBUI_NAME, 

72 async_reset_config, 

73 import_legacy_config_json, 

74 seed_registered_defaults, 

75) 

76from open_webui.constants import ERROR_MESSAGES, TASKS 

77from open_webui.utils.recurrence import RecurrenceEvaluationTimeout 

78from open_webui.env import ( 

79 USE_SLIM, 

80 AIOHTTP_CLIENT_SESSION_SSL, 

81 AUDIT_EXCLUDED_PATHS, 

82 AUDIT_INCLUDED_PATHS, 

83 AUDIT_LOG_LEVEL, 

84 BYPASS_MODEL_ACCESS_CONTROL, 

85 CHANGELOG, 

86 DEPLOYMENT_ID, 

87 ENABLE_AUDIT_GET_REQUESTS, 

88 ENABLE_COMPRESSION_MIDDLEWARE, 

89 ENABLE_CUSTOM_MODEL_FALLBACK, 

90 ENABLE_EASTER_EGGS, 

91 # OAuth Back-Channel Logout 

92 ENABLE_OAUTH_BACKCHANNEL_LOGOUT, 

93 ENABLE_OTEL, 

94 ENABLE_PLUGINS, 

95 ENABLE_PUBLIC_ACTIVE_USERS_COUNT, 

96 ENABLE_PYODIDE_FILE_PERSISTENCE, 

97 # SCIM 

98 ENABLE_SCIM, 

99 ENABLE_SIGNUP_PASSWORD_CONFIRMATION, 

100 ENABLE_STAR_SESSIONS_MIDDLEWARE, 

101 ENABLE_VERSION_UPDATE_CHECK, 

102 ENABLE_WEBSOCKET_SUPPORT, 

103 EXTERNAL_PWA_MANIFEST_URL, 

104 GLOBAL_LOG_LEVEL, 

105 INSTANCE_ID, 

106 LICENSE_KEY, 

107 LOG_FORMAT, 

108 MAX_BODY_LOG_SIZE, 

109 # Redis 

110 REDIS_KEY_PREFIX, 

111 REDIS_TASK_TTL, 

112 REDIS_URL, 

113 RESET_CONFIG_ON_START, 

114 SAFE_MODE, 

115 SCIM_TOKEN, 

116 VERSION, 

117 WEBSOCKET_HEARTBEAT_INTERVAL, 

118 WEBSOCKET_MANAGER, 

119 # Admin Account Runtime Creation 

120 WEBUI_ADMIN_EMAIL, 

121 WEBUI_ADMIN_NAME, 

122 WEBUI_ADMIN_PASSWORD, 

123 WEBUI_AUTH_TRUSTED_EMAIL_HEADER, 

124 WEBUI_BUILD_HASH, 

125 WEBUI_SECRET_KEY, 

126 WEBUI_SESSION_COOKIE_SAME_SITE, 

127 WEBUI_SESSION_COOKIE_SECURE, 

128) 

129from open_webui.events import ( 

130 EVENTS, 

131 delete_event_webhook, 

132 get_event_webhooks, 

133 migrate_legacy_webhook_config, 

134 publish_event, 

135 upsert_event_webhook, 

136) 

137from open_webui.events import ( 

138 get_event_catalog as get_event_catalog_items, 

139) 

140from open_webui.internal.db import engine, get_async_session 

141from open_webui.models.access_grants import AccessGrants 

142from open_webui.models.channels import Channels 

143from open_webui.models.chats import ChatForm, Chats 

144from open_webui.models.config import Config 

145from open_webui.models.functions import Functions 

146from open_webui.models.messages import Messages 

147from open_webui.models.models import Models, normalize_model_tags 

148from open_webui.models.users import Users 

149from open_webui.routers import ( 

150 analytics, 

151 audio, 

152 auths, 

153 automations, 

154 calendar, 

155 channels, 

156 chats, 

157 configs, 

158 evaluations, 

159 files, 

160 folders, 

161 functions, 

162 groups, 

163 images, 

164 knowledge, 

165 memories, 

166 models, 

167 notes, 

168 notifications, 

169 ollama, 

170 openai, 

171 pipelines, 

172 prompts, 

173 retrieval, 

174 scim, 

175 skills, 

176 tasks, 

177 terminals, 

178 tools, 

179 users, 

180 utils, 

181) 

182from open_webui.routers.retrieval import ( 

183 get_ef, 

184 get_embedding_function, 

185 get_reranking_function, 

186 get_rf, 

187) 

188from open_webui.socket.main import ( 

189 MODELS, 

190 get_event_emitter, 

191 get_models_in_use, 

192 get_user_id_from_session_pool, 

193 periodic_session_pool_cleanup, 

194 periodic_usage_pool_cleanup, 

195 redis_event_listener, 

196) 

197from open_webui.socket.main import ( 

198 app as socket_app, 

199) 

200from open_webui.tasks import ( 

201 cleanup_task, 

202 create_task, 

203 has_active_tasks, 

204 list_task_ids_by_item_id, 

205 list_tasks, 

206 redis_task_command_listener, 

207 redis_task_heartbeat, 

208 stop_item_tasks, 

209 stop_task, 

210) # Import from tasks.py 

211from open_webui.utils import logger 

212from open_webui.utils.access_control import has_permission 

213from open_webui.utils.access_control.folders import has_folder_write_access 

214from open_webui.utils.actions import chat_action as chat_action_handler 

215from open_webui.utils.asgi_middleware import AppHTTPMiddleware 

216from open_webui.utils.audit import AuditLevel, AuditLoggingMiddleware 

217from open_webui.utils.auth import ( 

218 create_admin_user, 

219 decode_token, 

220 get_admin_user, 

221 get_http_authorization_cred, 

222 get_license_data, 

223 get_verified_user, 

224) 

225from open_webui.utils.chat import ( 

226 chat_completed as chat_completed_handler, 

227) 

228from open_webui.utils.chat import ( 

229 generate_chat_completion as chat_completion_handler, 

230) 

231from open_webui.utils.chat_id import ( 

232 get_temporary_chat_session_id, 

233 is_saved_chat_id, 

234 is_temporary_chat_id, 

235) 

236from open_webui.utils.chat_variables import ( 

237 normalize_chat_variables, 

238) 

239from open_webui.utils.embeddings import generate_embeddings 

240from open_webui.utils.json_codec import JSONCodec 

241from open_webui.utils.json_response import apply_orjson_http_json 

242from open_webui.utils.logger import start_logger 

243from open_webui.utils.middleware import ( 

244 background_tasks_handler, 

245 build_chat_response_context, 

246 drain_approved_tool_calls, 

247 process_chat_payload, 

248 process_chat_response, 

249) 

250from open_webui.utils.misc import get_response_error_detail, merge_model_params 

251from open_webui.utils.model_ids import strip_provider_model_prefix 

252from open_webui.utils.models import ( 

253 check_model_access, 

254 get_all_base_models, 

255 get_all_models, 

256 get_filtered_models, 

257) 

258from open_webui.utils.oauth import ( 

259 OAuthClientInformationFull, 

260 OAuthClientManager, 

261 OAuthManager, 

262 apply_connection_oauth_options, 

263 decrypt_data, 

264 encrypt_data, 

265 get_oauth_client_info_with_dynamic_client_registration, 

266 get_oauth_client_info_with_static_credentials, 

267 recover_static_oauth_client_metadata, 

268 resolve_oauth_client_info, 

269) 

270from open_webui.utils.plugin import install_tool_and_function_dependencies 

271from open_webui.utils.redis import get_redis_client 

272from open_webui.utils.session_pool import cleanup_response, get_client_timeout, get_session, stream_wrapper 

273from open_webui.utils.tool_approval import ( 

274 ResolveToolCallForm, 

275 build_tool_approval_resume_payload, 

276 resolve_tool_call_output, 

277) 

278from open_webui.utils.tools import set_terminal_servers, set_tool_servers 

279 

280if SAFE_MODE: 280 ↛ 281line 280 didn't jump to line 281 because the condition on line 280 was never true

281 print('SAFE MODE ENABLED') 

282 # Functions.deactivate_all_functions() is awaited in lifespan below 

283 

284logging.basicConfig(stream=sys.stdout, level=GLOBAL_LOG_LEVEL) 

285log = logging.getLogger(__name__) 

286 

287 

288async def emit_chat_list_event(metadata: dict, chat_id: str): 

289 if not is_saved_chat_id(chat_id): 

290 return 

291 

292 event_emitter = await get_event_emitter(metadata, update_db=False) 

293 if event_emitter: 

294 folder_id = metadata.get('folder_id') or await Chats.get_chat_folder_id(chat_id, metadata.get('user_id')) 

295 await event_emitter({'type': 'chat:list', 'data': {'chat_id': chat_id, 'folder_id': folder_id}}) 

296 

297 

298class SPAStaticFiles(StaticFiles): 

299 async def get_response(self, path: str, scope): 

300 try: 

301 return await super().get_response(path, scope) 

302 except (HTTPException, StarletteHTTPException) as ex: 

303 if ex.status_code == 404: 

304 if path.endswith('.js'): 304 ↛ 306line 304 didn't jump to line 306 because the condition on line 304 was never true

305 # Return 404 for javascript files 

306 raise ex 

307 else: 

308 return await super().get_response('index.html', scope) 

309 else: 

310 raise ex 

311 

312 

313class CORSStaticFiles(StaticFiles): 

314 async def get_response(self, path: str, scope): 

315 response = await super().get_response(path, scope) 

316 response.headers['Access-Control-Allow-Origin'] = '*' 

317 return response 

318 

319 

320if LOG_FORMAT != 'json': 320 ↛ 344line 320 didn't jump to line 344 because the condition on line 320 was always true

321 banner = rf""" 

322 ██████╗ ██████╗ ███████╗███╗ ██╗ ██╗ ██╗███████╗██████╗ ██╗ ██╗██╗ 

323██╔═══██╗██╔══██╗██╔════╝████╗ ██║ ██║ ██║██╔════╝██╔══██╗██║ ██║██║ 

324██║ ██║██████╔╝█████╗ ██╔██╗ ██║ ██║ █╗ ██║█████╗ ██████╔╝██║ ██║██║ 

325██║ ██║██╔═══╝ ██╔══╝ ██║╚██╗██║ ██║███╗██║██╔══╝ ██╔══██╗██║ ██║██║ 

326╚██████╔╝██║ ███████╗██║ ╚████║ ╚███╔███╔╝███████╗██████╔╝╚██████╔╝██║ 

327 ╚═════╝ ╚═╝ ╚══════╝╚═╝ ╚═══╝ ╚══╝╚══╝ ╚══════╝╚═════╝ ╚═════╝ ╚═╝ 

328 

329 

330v{VERSION} - building the best AI user interface. 

331{f'Commit: {WEBUI_BUILD_HASH}' if WEBUI_BUILD_HASH != 'dev-build' else ''} 

332https://github.com/open-webui/open-webui 

333""" 

334 try: 

335 print(banner) 

336 except UnicodeEncodeError: 

337 # Stdout can't encode the box-drawing banner (Windows cp1252, redirected/headless stdout); fall back to ASCII. 

338 # LICENSE covers this Open WebUI CLI identifier. 

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

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

341 print(f'Open WebUI v{VERSION} - building the best AI user interface.\nhttps://github.com/open-webui/open-webui') 

342 

343 

344@asynccontextmanager 

345async def lifespan(app: FastAPI): 

346 # Store reference to main event loop for sync->async calls (e.g., embedding generation) 

347 # This allows sync functions to schedule work on the main loop without blocking health checks 

348 app.state.main_loop = asyncio.get_running_loop() 

349 

350 if THREAD_POOL_SIZE and THREAD_POOL_SIZE > 0: 350 ↛ 352line 350 didn't jump to line 352 because the condition on line 350 was never true

351 # asyncio offloads bypass AnyIO's limiter, so configure both before the first offload. 

352 anyio.to_thread.current_default_thread_limiter().total_tokens = THREAD_POOL_SIZE 

353 app.state.main_loop.set_default_executor( 

354 ThreadPoolExecutor( 

355 max_workers=THREAD_POOL_SIZE, 

356 thread_name_prefix=THREAD_POOL_THREAD_NAME_PREFIX, 

357 ) 

358 ) 

359 

360 app.state.instance_id = INSTANCE_ID 

361 start_logger() 

362 

363 if RESET_CONFIG_ON_START: 363 ↛ 364line 363 didn't jump to line 364 because the condition on line 363 was never true

364 await async_reset_config() 

365 

366 await import_legacy_config_json() 

367 await seed_registered_defaults() 

368 await initialize_runtime_config(app) 

369 await migrate_legacy_webhook_config() 

370 await publish_event(app, EVENTS.SYSTEM_STARTUP_STARTED, source='system') 

371 

372 license_task = None 

373 if LICENSE_KEY: 373 ↛ 374line 373 didn't jump to line 374 because the condition on line 373 was never true

374 license_task = asyncio.create_task(asyncio.to_thread(get_license_data, app, LICENSE_KEY)) 

375 

376 # Create admin account from env vars if specified and no users exist 

377 if WEBUI_ADMIN_EMAIL and WEBUI_ADMIN_PASSWORD: 377 ↛ 378line 377 didn't jump to line 378 because the condition on line 377 was never true

378 if await create_admin_user(WEBUI_ADMIN_EMAIL, WEBUI_ADMIN_PASSWORD, WEBUI_ADMIN_NAME): 

379 # Disable signup since we now have an admin 

380 await Config.upsert({'ui.enable_signup': False}) 

381 

382 if SAFE_MODE: 382 ↛ 383line 382 didn't jump to line 383 because the condition on line 382 was never true

383 await Functions.deactivate_all_functions() 

384 

385 # This should be blocking (sync) so functions are not deactivated on first /get_models calls 

386 # when the first user lands on the / route. 

387 log.info('Installing external dependencies of functions and tools...') 

388 await install_tool_and_function_dependencies() 

389 

390 app.state.redis = get_redis_client(async_mode=True) 

391 

392 if app.state.redis is not None: 392 ↛ 393line 392 didn't jump to line 393 because the condition on line 392 was never true

393 app.state.redis_task_command_listener = asyncio.create_task(redis_task_command_listener(app)) 

394 if REDIS_TASK_TTL > 0: 

395 app.state.redis_task_heartbeat = asyncio.create_task(redis_task_heartbeat(app)) 

396 

397 if WEBSOCKET_MANAGER == 'redis': 397 ↛ 398line 397 didn't jump to line 398 because the condition on line 397 was never true

398 app.state.redis_event_listener = asyncio.create_task(redis_event_listener()) 

399 

400 app.state.periodic_usage_pool_cleanup = asyncio.create_task(periodic_usage_pool_cleanup()) 

401 app.state.periodic_session_pool_cleanup = asyncio.create_task(periodic_session_pool_cleanup()) 

402 

403 from open_webui.utils.automations import scheduler_worker_loop 

404 

405 app.state.scheduler_worker_loop = asyncio.create_task(scheduler_worker_loop(app)) 

406 

407 if await Config.get('models.base_models_cache'): 407 ↛ 408line 407 didn't jump to line 408 because the condition on line 407 was never true

408 try: 

409 await get_all_models( 

410 Request( 

411 # Creating a mock request object to pass to get_all_models 

412 { 

413 'type': 'http', 

414 'asgi.version': '3.0', 

415 'asgi.spec_version': '2.0', 

416 'method': 'GET', 

417 'path': '/internal', 

418 'query_string': b'', 

419 'headers': Headers({}).raw, 

420 'client': ('127.0.0.1', 12345), 

421 'server': ('127.0.0.1', 80), 

422 'scheme': 'http', 

423 'app': app, 

424 } 

425 ), 

426 None, 

427 ) 

428 except Exception as e: 

429 log.warning(f'Failed to pre-fetch models at startup: {e}') 

430 

431 # Pre-fetch tool server specs so the first request doesn't pay the latency cost 

432 if len(await Config.get('tool_server.connections', []) or []) > 0: 432 ↛ 433line 432 didn't jump to line 433 because the condition on line 432 was never true

433 mock_request = Request( 

434 { 

435 'type': 'http', 

436 'asgi.version': '3.0', 

437 'asgi.spec_version': '2.0', 

438 'method': 'GET', 

439 'path': '/internal', 

440 'query_string': b'', 

441 'headers': Headers({}).raw, 

442 'client': ('127.0.0.1', 12345), 

443 'server': ('127.0.0.1', 80), 

444 'scheme': 'http', 

445 'app': app, 

446 } 

447 ) 

448 

449 log.info('Initializing tool servers...') 

450 try: 

451 await set_tool_servers(mock_request) 

452 log.info('Initialized %s tool server(s)', len(app.state.TOOL_SERVERS)) 

453 except Exception as e: 

454 log.warning(f'Failed to initialize tool servers at startup: {e}') 

455 

456 try: 

457 await set_terminal_servers(mock_request) 

458 log.info('Initialized %s terminal server(s)', len(app.state.TERMINAL_SERVERS)) 

459 except Exception as e: 

460 log.warning(f'Failed to initialize terminal servers at startup: {e}') 

461 

462 # Mark application as ready to accept traffic from a startup perspective. 

463 if license_task: 463 ↛ 464line 463 didn't jump to line 464 because the condition on line 463 was never true

464 try: 

465 await asyncio.wait_for(asyncio.shield(license_task), timeout=2) 

466 except asyncio.TimeoutError: 

467 log.warning('License data retrieval is still pending; continuing startup without it') 

468 except Exception as e: 

469 log.warning(f'License data retrieval failed during startup: {e}') 

470 

471 app.state.startup_complete = True 

472 await publish_event(app, EVENTS.SYSTEM_STARTUP_COMPLETED, source='system') 

473 

474 yield 

475 

476 await publish_event(app, EVENTS.SYSTEM_SHUTDOWN_STARTED, source='system') 

477 

478 # Shutdown: clean up shared resources 

479 from open_webui.utils.session_pool import close_session 

480 

481 await close_session() 

482 

483 if hasattr(app.state, 'redis_task_command_listener'): 483 ↛ 484line 483 didn't jump to line 484 because the condition on line 483 was never true

484 app.state.redis_task_command_listener.cancel() 

485 

486 if hasattr(app.state, 'redis_task_heartbeat'): 486 ↛ 487line 486 didn't jump to line 487 because the condition on line 486 was never true

487 app.state.redis_task_heartbeat.cancel() 

488 

489 if hasattr(app.state, 'redis_event_listener'): 489 ↛ 490line 489 didn't jump to line 490 because the condition on line 489 was never true

490 app.state.redis_event_listener.cancel() 

491 

492 app.state.periodic_usage_pool_cleanup.cancel() 

493 app.state.periodic_session_pool_cleanup.cancel() 

494 app.state.scheduler_worker_loop.cancel() 

495 

496 await publish_event(app, EVENTS.SYSTEM_SHUTDOWN_COMPLETED, source='system') 

497 

498 

499# Opt-in (ENABLE_ORJSON): orjson for request-body parsing and JSONResponse bodies; 

500# response_model routes keep FastAPI's Pydantic fast path either way. 

501apply_orjson_http_json() 

502 

503# LICENSE covers this Open WebUI API metadata identifier. 

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

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

506app = FastAPI( 

507 title='Open WebUI', 

508 docs_url='/docs' if ENV == 'dev' else None, 

509 openapi_url='/openapi.json' if ENV == 'dev' else None, 

510 redoc_url=None, 

511 lifespan=lifespan, 

512) 

513 

514 

515@app.exception_handler(RecurrenceEvaluationTimeout) 

516async def recurrence_timeout_handler(request: Request, exc: RecurrenceEvaluationTimeout): 

517 return JSONResponse(status_code=400, content={'detail': str(exc)}) 

518 

519 

520# Used by readiness checks to gate traffic until startup work is done. 

521app.state.startup_complete = False 

522 

523# For Open WebUI OIDC/OAuth2 

524oauth_manager = OAuthManager(app) 

525app.state.oauth_manager = oauth_manager 

526 

527# For Integrations 

528oauth_client_manager = OAuthClientManager(app) 

529app.state.oauth_client_manager = oauth_client_manager 

530 

531app.state.instance_id = None 

532app.state.redis = None 

533 

534# LICENSE covers this Open WebUI branding surface, including name, logo, 

535# visual, textual, symbolic identifiers, metadata, and surrounding UI. 

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

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

538app.state.WEBUI_NAME = WEBUI_NAME 

539app.state.LICENSE_METADATA = None 

540app.state.USER_COUNT = None 

541app.state.EXTERNAL_PWA_MANIFEST_URL = EXTERNAL_PWA_MANIFEST_URL 

542 

543 

544######################################## 

545# 

546# OPENTELEMETRY 

547# 

548######################################## 

549 

550if ENABLE_OTEL: 550 ↛ 551line 550 didn't jump to line 551 because the condition on line 550 was never true

551 from open_webui.utils.telemetry.setup import setup as setup_opentelemetry 

552 

553 setup_opentelemetry(app=app, db_engine=engine) 

554 

555 

556######################################## 

557# 

558# OLLAMA 

559# 

560######################################## 

561 

562 

563app.state.OLLAMA_MODELS = {} 

564 

565######################################## 

566# 

567# OPENAI 

568# 

569######################################## 

570 

571 

572app.state.OPENAI_MODELS = {} 

573 

574######################################## 

575# 

576# TOOL SERVERS 

577# 

578######################################## 

579 

580app.state.TOOL_SERVERS = [] 

581 

582######################################## 

583# 

584# TERMINAL SERVER 

585# 

586######################################## 

587 

588app.state.TERMINAL_SERVERS = [] 

589 

590######################################## 

591# 

592# DIRECT CONNECTIONS 

593# 

594######################################## 

595 

596 

597######################################## 

598# 

599# SCIM 

600# 

601######################################## 

602 

603app.state.ENABLE_SCIM = ENABLE_SCIM 

604app.state.SCIM_TOKEN = SCIM_TOKEN 

605 

606######################################## 

607# 

608# MODELS 

609# 

610######################################## 

611 

612app.state.BASE_MODELS = [] 

613 

614######################################## 

615# 

616# WEBUI 

617# 

618######################################## 

619 

620 

621async def initialize_runtime_config(app: FastAPI): 

622 # Migrate legacy access_control → access_grants on boot. 

623 from open_webui.utils.access_control import migrate_access_control 

624 

625 connections = await Config.get('tool_server.connections', []) or [] 

626 if any('access_control' in c.get('config', {}) for c in connections): 626 ↛ 627line 626 didn't jump to line 627 because the condition on line 626 was never true

627 for connection in connections: 

628 migrate_access_control(connection.get('config', {})) 

629 await Config.upsert({'tool_server.connections': connections}) 

630 

631 for tool_server_connection in connections: 631 ↛ 632line 631 didn't jump to line 632 because the loop on line 631 never started

632 if tool_server_connection.get('type', 'openapi') == 'mcp': 

633 server_id = (tool_server_connection.get('info') or {}).get('id') 

634 auth_type = tool_server_connection.get('auth_type', 'none') 

635 

636 if server_id and auth_type in ('oauth_2.1', 'oauth_2.1_static'): 

637 try: 

638 oauth_client_info = resolve_oauth_client_info(tool_server_connection) 

639 oauth_client_info = await recover_static_oauth_client_metadata( 

640 tool_server_connection, oauth_client_info 

641 ) 

642 oauth_client_info = apply_connection_oauth_options(tool_server_connection, oauth_client_info) 

643 app.state.oauth_client_manager.add_client( 

644 f'mcp:{server_id}', 

645 OAuthClientInformationFull(**oauth_client_info), 

646 ) 

647 except InvalidToken: 

648 log.error( 

649 'Error adding OAuth client for MCP tool server %s: InvalidToken. ' 

650 'Stored OAuth client data is invalid; reconnect this tool server.', 

651 server_id, 

652 ) 

653 except Exception as e: 

654 log.error( 

655 'Error adding OAuth client for MCP tool server %s: %s', 

656 server_id, 

657 f'{type(e).__name__}: {e}' if str(e) else type(e).__name__, 

658 ) 

659 

660 arena_models = await Config.get('evaluation.arena.models', []) or [] 

661 if any('access_control' in m.get('meta', {}) for m in arena_models): 661 ↛ 662line 661 didn't jump to line 662 because the condition on line 661 was never true

662 for model in arena_models: 

663 migrate_access_control(model.get('meta', {})) 

664 await Config.upsert({'evaluation.arena.models': arena_models}) 

665 

666 app.state.EMBEDDING_FUNCTION = None 

667 app.state.RERANKING_FUNCTION = None 

668 app.state.ef = None 

669 app.state.rf = None 

670 app.state.YOUTUBE_LOADER_TRANSLATION = None 

671 

672 try: 

673 rag_config = await Config.get_many( 

674 'rag.embedding_engine', 

675 'rag.embedding_model', 

676 'rag.enable_hybrid_search', 

677 'rag.bypass_embedding_and_retrieval', 

678 'rag.reranking_engine', 

679 'rag.reranking_model', 

680 'rag.external_reranker_url', 

681 'rag.external_reranker_api_key', 

682 'rag.external_reranker_timeout', 

683 ) 

684 app.state.ef = get_ef(rag_config.get('rag.embedding_engine'), rag_config.get('rag.embedding_model')) 

685 if rag_config.get('rag.enable_hybrid_search') and not rag_config.get('rag.bypass_embedding_and_retrieval'): 685 ↛ 686line 685 didn't jump to line 686 because the condition on line 685 was never true

686 app.state.rf = get_rf( 

687 rag_config.get('rag.reranking_engine'), 

688 rag_config.get('rag.reranking_model'), 

689 rag_config.get('rag.external_reranker_url'), 

690 rag_config.get('rag.external_reranker_api_key'), 

691 rag_config.get('rag.external_reranker_timeout'), 

692 ) 

693 else: 

694 app.state.rf = None 

695 except Exception as e: 

696 log.error(f'Error updating models: {e}') 

697 app.state.rf = None 

698 

699 rag_config = await Config.get_many( 

700 'rag.embedding_engine', 

701 'rag.embedding_model', 

702 'rag.openai.api_base_url', 

703 'rag.ollama.base_url', 

704 'rag.azure_openai.base_url', 

705 'rag.openai.api_key', 

706 'rag.ollama.api_key', 

707 'rag.azure_openai.api_key', 

708 'rag.embedding_batch_size', 

709 'rag.azure_openai.api_version', 

710 'rag.enable_async_embedding', 

711 'rag.embedding_concurrent_requests', 

712 'rag.reranking_engine', 

713 'rag.reranking_model', 

714 'rag.reranking_batch_size', 

715 ) 

716 embedding_engine = rag_config.get('rag.embedding_engine') 

717 app.state.EMBEDDING_FUNCTION = get_embedding_function( 

718 embedding_engine, 

719 rag_config.get('rag.embedding_model'), 

720 embedding_function=app.state.ef, 

721 url=( 

722 rag_config.get('rag.openai.api_base_url') 

723 if embedding_engine == 'openai' 

724 else ( 

725 rag_config.get('rag.ollama.base_url') 

726 if embedding_engine == 'ollama' 

727 else rag_config.get('rag.azure_openai.base_url') 

728 ) 

729 ), 

730 key=( 

731 rag_config.get('rag.openai.api_key') 

732 if embedding_engine == 'openai' 

733 else ( 

734 rag_config.get('rag.ollama.api_key') 

735 if embedding_engine == 'ollama' 

736 else rag_config.get('rag.azure_openai.api_key') 

737 ) 

738 ), 

739 embedding_batch_size=rag_config.get('rag.embedding_batch_size'), 

740 azure_api_version=( 

741 rag_config.get('rag.azure_openai.api_version') if embedding_engine == 'azure_openai' else None 

742 ), 

743 enable_async=rag_config.get('rag.enable_async_embedding'), 

744 concurrent_requests=rag_config.get('rag.embedding_concurrent_requests'), 

745 ) 

746 

747 app.state.RERANKING_FUNCTION = get_reranking_function( 

748 rag_config.get('rag.reranking_engine'), 

749 rag_config.get('rag.reranking_model'), 

750 reranking_function=app.state.rf, 

751 reranking_batch_size=rag_config.get('rag.reranking_batch_size'), 

752 ) 

753 

754 

755######################################## 

756# 

757# CODE EXECUTION 

758# 

759######################################## 

760 

761 

762######################################## 

763# 

764# IMAGES 

765# 

766######################################## 

767 

768 

769######################################## 

770# 

771# AUDIO 

772# 

773######################################## 

774 

775 

776app.state.faster_whisper_model = None 

777app.state.speech_synthesiser = None 

778app.state.speech_speaker_embeddings_dataset = None 

779 

780 

781######################################## 

782# 

783# TASKS 

784# 

785######################################## 

786 

787 

788######################################## 

789# 

790# WEBUI 

791# 

792######################################## 

793 

794app.state.MODELS = MODELS 

795 

796# Add the middleware to the app 

797try: 

798 audit_level = AuditLevel(AUDIT_LOG_LEVEL) 

799except ValueError as e: 

800 logger.error(f'Invalid audit level: {AUDIT_LOG_LEVEL}. Error: {e}') 

801 audit_level = AuditLevel.NONE 

802 

803# Added before CompressMiddleware so audit sits inside compression and 

804# captures response bodies before they are compressed (last added runs 

805# outermost). 

806if audit_level != AuditLevel.NONE: 806 ↛ 807line 806 didn't jump to line 807 because the condition on line 806 was never true

807 app.add_middleware( 

808 AuditLoggingMiddleware, 

809 audit_level=audit_level, 

810 excluded_paths=AUDIT_EXCLUDED_PATHS, 

811 included_paths=AUDIT_INCLUDED_PATHS, 

812 audit_get_requests=ENABLE_AUDIT_GET_REQUESTS, 

813 max_body_size=MAX_BODY_LOG_SIZE, 

814 ) 

815 

816if ENABLE_COMPRESSION_MIDDLEWARE: 816 ↛ 828line 816 didn't jump to line 828 because the condition on line 816 was always true

817 app.add_middleware(CompressMiddleware) 

818 

819 

820# All HTTP middlewares below are pure-ASGI implementations. The previous 

821# `BaseHTTPMiddleware` / `@app.middleware('http')` versions wrapped the 

822# downstream app in an anyio task group whose cancel scope cancelled 

823# in-flight DB calls (and any other awaits) on client disconnect / 

824# response completion — which surfaced as noisy SQLAlchemy 

825# `terminate_force_close` tracebacks under aiosqlite and as random 

826# CancelledError storms across the request path. See 

827# `open_webui.utils.asgi_middleware` for the rationale. 

828app.add_middleware(AppHTTPMiddleware) 

829 

830 

831app.add_middleware( 

832 CORSMiddleware, 

833 allow_origins=CORS_ALLOW_ORIGIN, 

834 allow_credentials=True, 

835 allow_methods=['*'], 

836 allow_headers=['*'], 

837) 

838 

839 

840app.mount('/ws', socket_app) 

841 

842 

843app.include_router(ollama.router, prefix='/ollama', tags=['ollama']) 

844app.include_router(openai.router, prefix='/openai', tags=['openai']) 

845 

846 

847app.include_router(pipelines.router, prefix='/api/v1/pipelines', tags=['pipelines']) 

848app.include_router(tasks.router, prefix='/api/v1/tasks', tags=['tasks']) 

849app.include_router(images.router, prefix='/api/v1/images', tags=['images']) 

850 

851app.include_router(audio.router, prefix='/api/v1/audio', tags=['audio']) 

852app.include_router(retrieval.router, prefix='/api/v1/retrieval', tags=['retrieval']) 

853 

854app.include_router(configs.router, prefix='/api/v1/configs', tags=['configs']) 

855 

856app.include_router(auths.router, prefix='/api/v1/auths', tags=['auths']) 

857app.include_router(users.router, prefix='/api/v1/users', tags=['users']) 

858 

859 

860app.include_router(channels.router, prefix='/api/v1/channels', tags=['channels']) 

861app.include_router(chats.router, prefix='/api/v1/chats', tags=['chats']) 

862app.include_router(notes.router, prefix='/api/v1/notes', tags=['notes']) 

863 

864 

865app.include_router(models.router, prefix='/api/v1/models', tags=['models']) 

866app.include_router(notifications.router, prefix='/api/v1/notifications', tags=['notifications']) 

867app.include_router(knowledge.router, prefix='/api/v1/knowledge', tags=['knowledge']) 

868app.include_router(prompts.router, prefix='/api/v1/prompts', tags=['prompts']) 

869app.include_router(tools.router, prefix='/api/v1/tools', tags=['tools']) 

870app.include_router(skills.router, prefix='/api/v1/skills', tags=['skills']) 

871 

872app.include_router(memories.router, prefix='/api/v1/memories', tags=['memories']) 

873app.include_router(folders.router, prefix='/api/v1/folders', tags=['folders']) 

874app.include_router(groups.router, prefix='/api/v1/groups', tags=['groups']) 

875app.include_router(files.router, prefix='/api/v1/files', tags=['files']) 

876app.include_router(functions.router, prefix='/api/v1/functions', tags=['functions']) 

877app.include_router(evaluations.router, prefix='/api/v1/evaluations', tags=['evaluations']) 

878if ENABLE_ADMIN_ANALYTICS: 878 ↛ 880line 878 didn't jump to line 880 because the condition on line 878 was always true

879 app.include_router(analytics.router, prefix='/api/v1/analytics', tags=['analytics']) 

880app.include_router(utils.router, prefix='/api/v1/utils', tags=['utils']) 

881app.include_router(terminals.router, prefix='/api/v1/terminals', tags=['terminals']) 

882app.include_router(automations.router, prefix='/api/v1/automations', tags=['automations']) 

883app.include_router(calendar.router, prefix='/api/v1/calendars', tags=['calendars']) 

884 

885# SCIM 2.0 API for identity management 

886if ENABLE_SCIM: 886 ↛ 887line 886 didn't jump to line 887 because the condition on line 886 was never true

887 app.include_router(scim.router, prefix='/api/v1/scim/v2', tags=['scim']) 

888 

889 

890################################## 

891# 

892# Chat Endpoints 

893# 

894################################## 

895 

896 

897@app.get('/api/models') 

898@app.get('/api/v1/models') # Experimental: Compatibility with OpenAI API 

899async def get_models(request: Request, refresh: bool = False, user=Depends(get_verified_user)): 

900 all_models = await get_all_models(request, refresh=refresh, user=user) 

901 

902 # Filter out filter pipelines 

903 models = [ 

904 model for model in all_models if not ('pipeline' in model and model['pipeline'].get('type', None) == 'filter') 

905 ] 

906 

907 # Chat requests resolve models by ID from request.app.state.MODELS, where 

908 # duplicate IDs collapse to the last model. Return the same effective list. 

909 models = list({model['id']: model for model in models}.values()) 

910 

911 # Access-filter first so the per-model payload work below only runs for 

912 # models the caller can actually see. 

913 models = await get_filtered_models(models, user) 

914 

915 for model in models: 915 ↛ 916line 915 didn't jump to line 916 because the loop on line 915 never started

916 info = model.get('info') if isinstance(model.get('info'), dict) else {} 

917 meta = info.get('meta') if isinstance(info.get('meta'), dict) else {} 

918 

919 # Remove profile image URL to reduce payload size 

920 meta.pop('profile_image_url', None) 

921 

922 if 'tags' in meta: 

923 meta['tags'] = normalize_model_tags(meta['tags']) 

924 

925 tags = normalize_model_tags(meta.get('tags')) + normalize_model_tags(model.get('tags')) 

926 model['tags'] = list({tag['name']: tag for tag in tags}.values()) 

927 

928 model_order_list = await Config.get('ui.model_order_list') 

929 if model_order_list: 929 ↛ 930line 929 didn't jump to line 930 because the condition on line 929 was never true

930 model_order_dict = {model_id: i for i, model_id in enumerate(model_order_list)} 

931 # Sort models by order list priority, with fallback for those not in the list 

932 models.sort( 

933 key=lambda model: ( 

934 model_order_dict.get(model.get('id', ''), float('inf')), 

935 (model.get('name', '') or ''), 

936 ) 

937 ) 

938 

939 if log.isEnabledFor(logging.DEBUG): 939 ↛ 940line 939 didn't jump to line 940 because the condition on line 939 was never true

940 log.debug( 

941 f'/api/models returned filtered models accessible to the user: {JSONCodec.dumps([model.get("id") for model in models])}' 

942 ) 

943 return {'data': models} 

944 

945 

946@app.get('/api/models/base') 

947async def get_base_models(request: Request, user=Depends(get_admin_user)): 

948 models = await get_all_base_models(request, user=user) 

949 return {'data': models} 

950 

951 

952class ModelUnloadForm(BaseModel): 

953 model: str 

954 

955 

956@app.post('/api/models/unload') 

957async def unload_model(request: Request, form_data: ModelUnloadForm, user=Depends(get_admin_user)): 

958 """ 

959 Unified model unload endpoint. 

960 Resolves the provider that owns the model and calls its native unload mechanism. 

961 Supports: Ollama (keep_alive=0) and llama.cpp (/models/unload). 

962 """ 

963 model_id = form_data.model 

964 

965 ollama_models = getattr(request.app.state, 'OLLAMA_MODELS', None) or {} 

966 openai_models = getattr(request.app.state, 'OPENAI_MODELS', None) or {} 

967 

968 seen = set() 

969 while model_id not in ollama_models and model_id not in openai_models and model_id not in seen: 969 ↛ 977line 969 didn't jump to line 977 because the condition on line 969 was always true

970 seen.add(model_id) 

971 model_info = await Models.get_model_by_id(model_id) 

972 if not model_info or not model_info.base_model_id: 972 ↛ 974line 972 didn't jump to line 974 because the condition on line 972 was always true

973 break 

974 model_id = model_info.base_model_id 

975 

976 # --- Ollama provider --- 

977 if model_id in ollama_models: 977 ↛ 978line 977 didn't jump to line 978 because the condition on line 977 was never true

978 ollama_config = await Config.get_many('ollama.base_urls', 'ollama.api_configs') 

979 ollama_base_urls = ollama_config.get('ollama.base_urls') or [] 

980 ollama_api_configs = ollama_config.get('ollama.api_configs') or {} 

981 url_indices = ollama_models[model_id].get('urls', []) 

982 errors = [] 

983 for idx in url_indices: 

984 url = ollama_base_urls[idx] 

985 api_config = ollama_api_configs.get( 

986 str(idx), 

987 ollama_api_configs.get(url, {}), 

988 ) 

989 key = api_config.get('key', None) 

990 

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

992 actual_model = strip_provider_model_prefix(model_id, prefix_id) 

993 

994 payload = JSONCodec.dumps({'model': actual_model, 'keep_alive': 0, 'prompt': ''}) 

995 

996 try: 

997 timeout = aiohttp.ClientTimeout(total=30) 

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

999 headers = { 

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

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

1002 } 

1003 async with session.post( 

1004 f'{url}/api/generate', 

1005 data=payload, 

1006 headers=headers, 

1007 ) as r: 

1008 if not r.ok: 

1009 errors.append({'url_idx': idx, 'error': await r.text()}) 

1010 except Exception as e: 

1011 log.exception(f'Failed to unload model on Ollama node {idx}: {e}') 

1012 errors.append({'url_idx': idx, 'error': str(e)}) 

1013 

1014 if errors: 

1015 raise HTTPException( 

1016 status_code=500, 

1017 detail=f'Failed to unload model on {len(errors)} node(s): {errors}', 

1018 ) 

1019 return {'status': True} 

1020 

1021 # --- OpenAI-compatible providers --- 

1022 if model_id in openai_models: 1022 ↛ 1023line 1022 didn't jump to line 1023 because the condition on line 1022 was never true

1023 openai_config = await Config.get_many('openai.api_configs', 'openai.api_base_urls', 'openai.api_keys') 

1024 openai_api_configs = openai_config.get('openai.api_configs') or {} 

1025 openai_base_urls = openai_config.get('openai.api_base_urls') or [] 

1026 openai_api_keys = openai_config.get('openai.api_keys') or [] 

1027 model_info = openai_models[model_id] 

1028 idx = model_info.get('urlIdx') 

1029 api_config = openai_api_configs.get(str(idx), {}) 

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

1031 base_url = openai_base_urls[idx] 

1032 key = openai_api_keys[idx] if idx < len(openai_api_keys) else '' 

1033 

1034 if provider == 'llama.cpp': 1034 ↛ 1059line 1034 didn't jump to line 1059 because the condition on line 1034 was always true

1035 root_url = base_url.rstrip('/').removesuffix('/v1') 

1036 actual_model = strip_provider_model_prefix(model_id, api_config.get('prefix_id')) 

1037 try: 

1038 timeout = aiohttp.ClientTimeout(total=30) 

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

1040 headers = { 

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

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

1043 } 

1044 async with session.post( 

1045 f'{root_url}/models/unload', 

1046 json={'model': actual_model}, 

1047 headers=headers, 

1048 ) as r: 

1049 if not r.ok: 

1050 detail = await r.text() 

1051 raise HTTPException(status_code=r.status, detail=detail) 

1052 return await r.json() 

1053 except HTTPException: 

1054 raise 

1055 except Exception as e: 

1056 log.exception(f'Failed to unload model via llama.cpp: {e}') 

1057 raise HTTPException(status_code=500, detail=str(e)) 

1058 else: 

1059 raise HTTPException( 

1060 status_code=400, 

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

1062 ) 

1063 

1064 raise HTTPException(status_code=404, detail=f'Model "{model_id}" not found') 

1065 

1066 

1067################################## 

1068# Embeddings 

1069################################## 

1070 

1071 

1072@app.post('/api/embeddings') 

1073@app.post('/api/v1/embeddings') # Experimental: Compatibility with OpenAI API 

1074async def embeddings(request: Request, form_data: dict, user=Depends(get_verified_user)): 

1075 """ 

1076 OpenAI-compatible embeddings endpoint. 

1077 

1078 This handler: 

1079 - Performs user/model checks and dispatches to the correct backend. 

1080 - Supports OpenAI, Ollama, arena models, pipelines, and any compatible provider. 

1081 

1082 Args: 

1083 request (Request): Request context. 

1084 form_data (dict): OpenAI-like payload (e.g., {"model": "...", "input": [...]}) 

1085 user (UserModel): Authenticated user. 

1086 

1087 Returns: 

1088 dict: OpenAI-compatible embeddings response. 

1089 """ 

1090 # Make sure models are loaded in app state 

1091 if not request.app.state.MODELS: 1091 ↛ 1094line 1091 didn't jump to line 1094 because the condition on line 1091 was always true

1092 await get_all_models(request, user=user) 

1093 # Use generic dispatcher in utils.embeddings 

1094 return await generate_embeddings(request, form_data, user) 

1095 

1096 

1097async def _set_direct_model(request: Request, model_item: dict, user) -> None: 

1098 model_meta = (model_item.get('info') or {}).get('meta') or {} 

1099 knowledge_items = model_meta.get('knowledge') 

1100 if knowledge_items: 

1101 from open_webui.utils.access_control.files import get_accessible_folder_files 

1102 

1103 model_meta['knowledge'] = await get_accessible_folder_files(knowledge_items, user) 

1104 request.state.direct = True 

1105 request.state.model = model_item 

1106 

1107 

1108@app.post('/api/chat/completions') 

1109@app.post('/api/v1/chat/completions') # Experimental: Compatibility with OpenAI API 

1110async def chat_completion( 

1111 request: Request, 

1112 form_data: dict, 

1113 user=Depends(get_verified_user), 

1114): 

1115 if not request.app.state.MODELS: 1115 ↛ 1118line 1115 didn't jump to line 1118 because the condition on line 1115 was always true

1116 await get_all_models(request, user=user) 

1117 

1118 model_id = form_data.get('model', None) 

1119 model_item = form_data.pop('model_item', {}) 

1120 tasks = form_data.pop('background_tasks', None) 

1121 

1122 metadata = {} 

1123 try: 

1124 model_info = None 

1125 fallback_model = None 

1126 missing_base_model = False 

1127 if not model_item.get('direct', False): 1127 ↛ 1163line 1127 didn't jump to line 1163 because the condition on line 1127 was always true

1128 if model_id not in request.app.state.MODELS: 1128 ↛ 1131line 1128 didn't jump to line 1131 because the condition on line 1128 was always true

1129 raise Exception('Model not found') 

1130 

1131 model = request.app.state.MODELS[model_id] 

1132 model_info = await Models.get_model_by_id(model_id) 

1133 missing_base_model = bool( 

1134 model_info and model_info.base_model_id and model_info.base_model_id not in request.app.state.MODELS 

1135 ) 

1136 

1137 if missing_base_model and ENABLE_CUSTOM_MODEL_FALLBACK: 

1138 fallback_model_id = next( 

1139 ( 

1140 model_id.strip() 

1141 for model_id in ((await Config.get('ui.default_models')) or '').split(',') 

1142 if model_id.strip() 

1143 ), 

1144 None, 

1145 ) 

1146 if fallback_model_id: 

1147 fallback_model = request.app.state.MODELS.get(fallback_model_id) 

1148 

1149 # Check if user has access to the model 

1150 if not BYPASS_MODEL_ACCESS_CONTROL and (user.role != 'admin' or not BYPASS_ADMIN_ACCESS_CONTROL): 

1151 try: 

1152 access_model_info = ( 

1153 model_info.model_copy(update={'base_model_id': None}) 

1154 if fallback_model is not None 

1155 else model_info 

1156 ) 

1157 await check_model_access(user, model, model_info=access_model_info) 

1158 if fallback_model is not None: 

1159 await check_model_access(user, fallback_model) 

1160 except Exception as e: 

1161 raise e 

1162 else: 

1163 model = model_item 

1164 await _set_direct_model(request, model, user) 

1165 

1166 # Read before the fallback below can rebind model to a different one. 

1167 model_capabilities = ((model.get('info') or {}).get('meta') or {}).get('capabilities') or {} 

1168 

1169 # Model params: global defaults as base, per-model overrides win 

1170 default_model_params = copy.deepcopy(await Config.get('models.default_params', {}) or {}) 

1171 model_info_params = merge_model_params( 

1172 default_model_params, 

1173 model_info.params.model_dump() if model_info and model_info.params else {}, 

1174 ) 

1175 request_params = {key: value for key, value in (form_data.get('params') or {}).items() if value is not None} 

1176 if model_info_params or request_params: 

1177 form_data['params'] = merge_model_params(model_info_params, request_params) 

1178 

1179 # Check base model existence for custom models 

1180 if missing_base_model: 

1181 if fallback_model is None: 

1182 raise Exception('Model not found') 

1183 # Update model and form_data so routing uses the fallback model's type 

1184 model = fallback_model 

1185 form_data['model'] = fallback_model['id'] 

1186 

1187 # Chat Params 

1188 stream_delta_chunk_size = form_data.get('params', {}).get('stream_delta_chunk_size') 

1189 reasoning_tags = form_data.get('params', {}).get('reasoning_tags') 

1190 compact_token_threshold = form_data.get('params', {}).get('compact_token_threshold') 

1191 

1192 # Model Params 

1193 if model_info_params.get('stream_response') is not None: 

1194 form_data['stream'] = model_info_params.get('stream_response') 

1195 

1196 # Providers only report token counts when asked, so ask on every caller's behalf. 

1197 if form_data.get('stream') and model_capabilities.get('usage'): 

1198 form_data['stream_options'] = {**(form_data.get('stream_options') or {}), 'include_usage': True} 

1199 

1200 if model_info_params.get('stream_delta_chunk_size'): 

1201 stream_delta_chunk_size = model_info_params.get('stream_delta_chunk_size') 

1202 

1203 if model_info_params.get('reasoning_tags') is not None: 

1204 reasoning_tags = model_info_params.get('reasoning_tags') 

1205 

1206 if model_info_params.get('compact_token_threshold') is not None: 

1207 compact_token_threshold = model_info_params.get('compact_token_threshold') 

1208 

1209 # parent_id signals intent: 

1210 # null → new chat (root message, no parent) 

1211 # value → follow-up (user message's parentId = prev assistant) 

1212 # absent → legacy caller, no chat management 

1213 is_new_chat = 'parent_id' in form_data and form_data['parent_id'] is None and not form_data.get('chat_id') 

1214 parent_id = form_data.pop('parent_id', None) 

1215 form_data.pop('new_chat', None) # Legacy field 

1216 

1217 # Multi-model message_ids: list of {model_id, message_id} entries. 

1218 # Supports both the new array format and legacy dict format for backward compat. 

1219 message_ids = form_data.pop('message_ids', None) 

1220 if isinstance(message_ids, list): 

1221 # New format: [{"model_id": ..., "message_id": ...}, ...] 

1222 form_data.pop('id', None) 

1223 elif isinstance(message_ids, dict): 

1224 # Legacy dict format: {model_id: message_id} — convert to list 

1225 message_ids = [{'model_id': k, 'message_id': v} for k, v in message_ids.items()] 

1226 form_data.pop('id', None) 

1227 else: 

1228 # Single-model fallback 

1229 message_ids = [{'model_id': model_id, 'message_id': form_data.pop('id', None)}] 

1230 

1231 user_message = form_data.pop('user_message', None) or form_data.pop('parent_message', None) 

1232 chat_id = form_data.pop('chat_id', None) or '' 

1233 chat_variables = form_data.pop('chat_variables', None) 

1234 if chat_variables is None: 

1235 existing_chat = await Chats.get_chat_by_id(chat_id) if is_saved_chat_id(chat_id) else None 

1236 chat_variables = existing_chat.variables if existing_chat else {} 

1237 

1238 chat_variables = normalize_chat_variables(chat_variables) 

1239 

1240 # Drop tool_servers if caller lacks features.direct_tool_servers — 

1241 # mirrors the storage-side strip in user/settings/update. 

1242 tool_servers = form_data.pop('tool_servers', None) 

1243 if ( 

1244 tool_servers 

1245 and user.role != 'admin' 

1246 and not await has_permission( 

1247 user.id, 

1248 'features.direct_tool_servers', 

1249 await Config.get('user.permissions'), 

1250 ) 

1251 ): 

1252 tool_servers = None 

1253 

1254 automation_id = form_data.pop('automation_id', None) 

1255 tool_approval_mode = ( 

1256 'full' 

1257 if automation_id or chat_id.startswith('channel:') 

1258 else ( 

1259 form_data.get('params', {}).get('tool_approval_mode') 

1260 if await Config.get('chat.tool_permissions.enable', False) 

1261 else 'full' 

1262 ) 

1263 or 'full' 

1264 ) 

1265 

1266 metadata = { 

1267 'user_id': user.id, 

1268 'user_agent': request.headers.get('user-agent', '') or '', 

1269 'internal': getattr(request.state, 'internal', False) is True, 

1270 'chat_id': chat_id, 

1271 'user_message': user_message, 

1272 'user_message_id': user_message.get('id') if user_message else None, 

1273 'assistant_message_id': form_data.pop('assistant_message_id', None), 

1274 'session_id': form_data.pop('session_id', None), 

1275 'automation_id': automation_id, 

1276 'folder_id': form_data.pop('folder_id', None), 

1277 'filter_ids': form_data.pop('filter_ids', []), 

1278 'tool_ids': form_data.get('tool_ids', None), 

1279 'tool_servers': tool_servers, 

1280 'files': form_data.get('files', None), 

1281 'features': form_data.get('features', {}), 

1282 'variables': form_data.get('variables', {}), 

1283 'chat_variables': chat_variables, 

1284 'model': model, 

1285 'direct': model_item.get('direct', False), 

1286 'params': { 

1287 'stream_delta_chunk_size': stream_delta_chunk_size, 

1288 'reasoning_tags': reasoning_tags, 

1289 'compact_token_threshold': compact_token_threshold, 

1290 'function_calling': ( 

1291 form_data.get('params', {}).get('function_calling') 

1292 or model_info_params.get('function_calling') 

1293 or 'native' 

1294 ), 

1295 'tool_approval_mode': tool_approval_mode, 

1296 }, 

1297 } 

1298 

1299 if is_new_chat: 

1300 metadata['chat_id'] = str(uuid4()) 

1301 

1302 initial_title_generation = None 

1303 if is_new_chat and tasks and TASKS.TITLE_GENERATION in tasks: 

1304 initial_title_generation = tasks.pop(TASKS.TITLE_GENERATION) 

1305 

1306 if metadata.get('chat_id') and user: 

1307 chat_id = metadata['chat_id'] 

1308 

1309 # Gate channel: branch — caller needs write access on the channel, and the 

1310 # supplied message_id must belong to that channel and be the caller's own. 

1311 if chat_id.startswith('channel:'): 

1312 channel_id = chat_id.removeprefix('channel:') 

1313 channel = await Channels.get_channel_by_id(channel_id) 

1314 if not channel: 

1315 raise HTTPException( 

1316 status_code=status.HTTP_404_NOT_FOUND, 

1317 detail=ERROR_MESSAGES.NOT_FOUND, 

1318 ) 

1319 if user.role != 'admin': 

1320 if channel.type in ['group', 'dm']: 

1321 if not await Channels.is_user_channel_member(channel.id, user.id): 

1322 raise HTTPException( 

1323 status_code=status.HTTP_403_FORBIDDEN, 

1324 detail=ERROR_MESSAGES.DEFAULT(), 

1325 ) 

1326 else: 

1327 if not await AccessGrants.has_access( 

1328 user_id=user.id, 

1329 resource_type='channel', 

1330 resource_id=channel.id, 

1331 permission='write', 

1332 ): 

1333 raise HTTPException( 

1334 status_code=status.HTTP_403_FORBIDDEN, 

1335 detail=ERROR_MESSAGES.DEFAULT(), 

1336 ) 

1337 for entry in message_ids: 

1338 target_message_id = entry.get('message_id') 

1339 if not target_message_id: 

1340 continue 

1341 target_message = await Messages.get_message_by_id(target_message_id) 

1342 if target_message and ( 

1343 target_message.channel_id != channel.id 

1344 # Write access is not authorship — block cross-member edits. 

1345 or (user.role != 'admin' and target_message.user_id != user.id) 

1346 ): 

1347 raise HTTPException( 

1348 status_code=status.HTTP_403_FORBIDDEN, 

1349 detail=ERROR_MESSAGES.DEFAULT(), 

1350 ) 

1351 

1352 if is_saved_chat_id(chat_id): 

1353 if is_new_chat: 

1354 # The chat created below is persisted with this folder_id. 

1355 folder_id = metadata['folder_id'] 

1356 if folder_id is not None and not await has_folder_write_access(user.id, folder_id): 

1357 raise HTTPException( 

1358 status_code=status.HTTP_404_NOT_FOUND, 

1359 detail=ERROR_MESSAGES.NOT_FOUND, 

1360 ) 

1361 

1362 # Build the full history upfront with ALL assistant placeholders 

1363 user_message = metadata.get('user_message') or {} 

1364 user_message_id = user_message.get('id') if user_message else None 

1365 

1366 history_messages = {} 

1367 all_assistant_ids = [entry['message_id'] for entry in message_ids if entry.get('message_id')] 

1368 

1369 if user_message_id and user_message: 

1370 user_message['childrenIds'] = all_assistant_ids 

1371 history_messages[user_message_id] = user_message 

1372 

1373 for entry in message_ids: 

1374 target_model_id = entry['model_id'] 

1375 assistant_message_id = entry['message_id'] 

1376 if assistant_message_id: 

1377 assistant_message = { 

1378 'id': assistant_message_id, 

1379 'parentId': user_message_id, 

1380 'childrenIds': [], 

1381 'role': 'assistant', 

1382 'content': '', 

1383 'done': False, 

1384 'model': target_model_id, 

1385 'timestamp': int(time.time()), 

1386 } 

1387 # Preserve the side-by-side column index so duplicate 

1388 # models don't collapse into one another on reload. 

1389 if entry.get('modelIdx') is not None: 

1390 assistant_message['modelIdx'] = entry['modelIdx'] 

1391 history_messages[assistant_message_id] = assistant_message 

1392 

1393 await Chats.insert_new_chat( 

1394 chat_id, 

1395 user.id, 

1396 ChatForm( 

1397 chat={ 

1398 'id': chat_id, 

1399 'title': 'New Chat', 

1400 'models': [entry['model_id'] for entry in message_ids], 

1401 'history': { 

1402 'currentId': all_assistant_ids[0] if all_assistant_ids else user_message_id, 

1403 'messages': history_messages, 

1404 }, 

1405 'messages': [ 

1406 {'role': 'user', 'content': user_message.get('content', '')}, 

1407 ] 

1408 if user_message_id 

1409 else [], 

1410 'files': metadata.get('files') or [], 

1411 'tags': [], 

1412 'timestamp': int(time.time() * 1000), 

1413 }, 

1414 variables=chat_variables, 

1415 folder_id=metadata.get('folder_id'), 

1416 ), 

1417 ) 

1418 await publish_event( 

1419 request, 

1420 EVENTS.CHAT_CREATED, 

1421 actor=user, 

1422 subject_id=chat_id, 

1423 data={'title': 'New Chat'}, 

1424 ) 

1425 await emit_chat_list_event(metadata, chat_id) 

1426 if user_message_id: 1426 ↛ 1438line 1426 didn't jump to line 1438 because the condition on line 1426 was always true

1427 await publish_event( 

1428 request, 

1429 EVENTS.MESSAGE_CREATED, 

1430 actor=user, 

1431 subject_id=user_message_id, 

1432 data={ 

1433 'chat_id': chat_id, 

1434 'role': 'user', 

1435 'content_preview': user_message.get('content', '')[:300], 

1436 }, 

1437 ) 

1438 for entry in message_ids: 

1439 assistant_message_id = entry.get('message_id') 

1440 if assistant_message_id: 

1441 await publish_event( 

1442 request, 

1443 EVENTS.MESSAGE_CREATED, 

1444 actor=user, 

1445 subject_id=assistant_message_id, 

1446 data={ 

1447 'chat_id': chat_id, 

1448 'role': 'assistant', 

1449 'model': entry.get('model_id'), 

1450 }, 

1451 ) 

1452 

1453 # Insert chat files from user message if any 

1454 user_message_files = user_message.get('files', []) 

1455 if user_message_files: 

1456 try: 

1457 await Chats.insert_chat_files( 

1458 chat_id, 

1459 user_message_id, 

1460 [ 

1461 file_item.get('id') 

1462 for file_item in user_message_files 

1463 if file_item.get('type') == 'file' 

1464 ], 

1465 user.id, 

1466 ) 

1467 except Exception as e: 

1468 log.debug('Error inserting chat files: %s', e) 

1469 pass 

1470 

1471 if initial_title_generation is not None and all_assistant_ids: 

1472 title_metadata = { 

1473 **metadata, 

1474 'message_id': all_assistant_ids[0], 

1475 } 

1476 event_emitter = await get_event_emitter(title_metadata, update_db=False) 

1477 title_ctx = { 

1478 'request': request, 

1479 'form_data': form_data, 

1480 'user': user, 

1481 'metadata': title_metadata, 

1482 'tasks': {TASKS.TITLE_GENERATION: initial_title_generation}, 

1483 'event_emitter': event_emitter, 

1484 } 

1485 

1486 async def run_initial_title_generation(): 

1487 try: 

1488 await background_tasks_handler(title_ctx) 

1489 except Exception: 

1490 log.exception('Error generating initial chat title') 

1491 

1492 asyncio.create_task(run_initial_title_generation()) 

1493 else: 

1494 # Existing chat — verify ownership 

1495 if not await Chats.is_chat_owner(chat_id, user.id) and user.role != 'admin': 

1496 raise HTTPException( 

1497 status_code=status.HTTP_404_NOT_FOUND, 

1498 detail=ERROR_MESSAGES.DEFAULT(), 

1499 ) 

1500 

1501 user_message = metadata.get('user_message') or {} 

1502 selected_chat_models = user_message.get('models') if isinstance(user_message, dict) else None 

1503 if not isinstance(selected_chat_models, list) or not selected_chat_models: 

1504 selected_chat_models = [entry.get('model_id') for entry in message_ids if entry.get('model_id')] 

1505 

1506 # Persist chat-level fields the frontend used to save on every message. 

1507 # The old frontend saveChatHandler did this on every message; 

1508 # now the backend owns persistence. 

1509 chat_files = metadata.get('files') 

1510 chat_fields = {} 

1511 if chat_files is not None: 

1512 chat_fields['files'] = chat_files 

1513 if selected_chat_models: 

1514 chat_fields['models'] = selected_chat_models 

1515 if chat_fields: 

1516 await Chats.update_chat_by_id(chat_id, chat_fields, touch=False) 

1517 

1518 await Chats.update_chat_variables_by_id(chat_id, chat_variables) 

1519 

1520 # Save user message to DB 

1521 if user_message and user_message.get('id'): 

1522 await Chats.upsert_message_to_chat_by_id_and_message_id( 

1523 chat_id, 

1524 user_message['id'], 

1525 user_message, 

1526 ) 

1527 await emit_chat_list_event({**metadata, 'message_id': user_message['id']}, chat_id) 

1528 await publish_event( 

1529 request, 

1530 EVENTS.MESSAGE_CREATED, 

1531 actor=user, 

1532 subject_id=user_message['id'], 

1533 data={ 

1534 'chat_id': chat_id, 

1535 'role': user_message.get('role', 'user'), 

1536 'content_preview': user_message.get('content', '')[:300], 

1537 }, 

1538 ) 

1539 if not getattr(request.state, 'internal', False) and not (user_message.get('meta') or {}).get( 

1540 'internal' 

1541 ): 

1542 try: 

1543 from open_webui.utils.timers import cancel_timers_for_chat 

1544 

1545 await cancel_timers_for_chat(chat_id, 'chat.user_message', user.id) 

1546 except Exception: 

1547 log.exception('Failed to cancel chat.user_message timers for chat %s', chat_id) 

1548 

1549 # Link grandparent → user message (childrenIds) 

1550 grandparent_id = user_message.get('parentId') 

1551 if grandparent_id: 

1552 grandparent = await Chats.get_message_by_id_and_message_id(chat_id, grandparent_id) 

1553 if grandparent: 

1554 child_ids = grandparent.get('childrenIds', []) 

1555 if user_message['id'] not in child_ids: 

1556 child_ids.append(user_message['id']) 

1557 await Chats.upsert_message_to_chat_by_id_and_message_id( 

1558 chat_id, grandparent_id, {'childrenIds': child_ids} 

1559 ) 

1560 

1561 # Insert chat files from user message if any 

1562 user_message_files = user_message.get('files', []) 

1563 if user_message_files: 

1564 try: 

1565 await Chats.insert_chat_files( 

1566 chat_id, 

1567 user_message.get('id'), 

1568 [ 

1569 file_item.get('id') 

1570 for file_item in user_message_files 

1571 if file_item.get('type') == 'file' 

1572 ], 

1573 user.id, 

1574 ) 

1575 except Exception as e: 

1576 log.debug('Error inserting chat files: %s', e) 

1577 pass 

1578 

1579 # Save ALL assistant placeholders 

1580 user_message_id = metadata.get('user_message_id') 

1581 all_assistant_ids = [entry['message_id'] for entry in message_ids if entry.get('message_id')] 

1582 

1583 # Link user message → all assistant messages (childrenIds) 

1584 if user_message_id and all_assistant_ids: 

1585 existing_user_message = await Chats.get_message_by_id_and_message_id(chat_id, user_message_id) 

1586 if existing_user_message: 

1587 child_ids = existing_user_message.get('childrenIds', []) 

1588 for assistant_id in all_assistant_ids: 

1589 if assistant_id not in child_ids: 

1590 child_ids.append(assistant_id) 

1591 await Chats.upsert_message_to_chat_by_id_and_message_id( 

1592 chat_id, 

1593 user_message_id, 

1594 {'childrenIds': child_ids}, 

1595 ) 

1596 

1597 # Save each assistant placeholder 

1598 for entry in message_ids: 

1599 target_model_id = entry['model_id'] 

1600 assistant_message_id = entry['message_id'] 

1601 if assistant_message_id and assistant_message_id == metadata.get('assistant_message_id'): 

1602 continue 

1603 if assistant_message_id: 

1604 assistant_message = { 

1605 'id': assistant_message_id, 

1606 'parentId': user_message_id, 

1607 'childrenIds': [], 

1608 'role': 'assistant', 

1609 'content': '', 

1610 'done': False, 

1611 'model': target_model_id, 

1612 'timestamp': int(time.time()), 

1613 } 

1614 # Preserve the side-by-side column index so duplicate 

1615 # models don't collapse into one another on reload. 

1616 if entry.get('modelIdx') is not None: 

1617 assistant_message['modelIdx'] = entry['modelIdx'] 

1618 await Chats.upsert_message_to_chat_by_id_and_message_id( 

1619 chat_id, 

1620 assistant_message_id, 

1621 assistant_message, 

1622 ) 

1623 await publish_event( 

1624 request, 

1625 EVENTS.MESSAGE_CREATED, 

1626 actor=user, 

1627 subject_id=assistant_message_id, 

1628 data={ 

1629 'chat_id': chat_id, 

1630 'role': 'assistant', 

1631 'model': target_model_id, 

1632 }, 

1633 ) 

1634 

1635 request.state.metadata = metadata 

1636 form_data['metadata'] = metadata 

1637 

1638 except HTTPException: 

1639 raise 

1640 except Exception as e: 

1641 log.warning(f'Error processing chat metadata: {e}') 

1642 raise HTTPException( 

1643 status_code=status.HTTP_400_BAD_REQUEST, 

1644 detail=str(e), 

1645 ) 

1646 

1647 async def process_chat(request, form_data, user, metadata, model, tasks=None): 

1648 try: 

1649 ctx = None 

1650 if metadata.get('assistant_message_id'): 

1651 ctx = await build_chat_response_context(request, form_data, user, model, metadata, tasks, []) 

1652 form_data, metadata, events = await process_chat_payload(request, form_data, user, metadata, model) 

1653 

1654 if await drain_approved_tool_calls(request, form_data, user, model, metadata): 

1655 return {'status': True, 'chat_id': metadata.get('chat_id'), 'paused': True} 

1656 

1657 response = await chat_completion_handler(request, form_data, user) 

1658 

1659 # When the upstream provider returns an error (e.g. HTTP 400 

1660 # content-filter, quota exceeded), generate_chat_completion 

1661 # returns a JSONResponse instead of raising. Detect this and 

1662 # raise so the except-block below emits a terminal 

1663 # chat:message:error, unblocking the frontend. 

1664 if isinstance(response, JSONResponse) and response.status_code >= 400: 1664 ↛ anywhereline 1664 didn't jump anywhere: it always raised an exception.

1665 raise Exception(get_response_error_detail(response)) 

1666 

1667 if ctx is None: 

1668 ctx = await build_chat_response_context(request, form_data, user, model, metadata, tasks, events) 

1669 else: 

1670 ctx.update(form_data=form_data, metadata=metadata, events=events) 

1671 

1672 return await process_chat_response(response, ctx) 

1673 except asyncio.CancelledError: 

1674 log.info('Chat processing was cancelled') 

1675 try: 

1676 

1677 async def emit_cancel_event(): 

1678 event_emitter = await get_event_emitter(metadata) 

1679 if event_emitter: 

1680 await event_emitter({'type': 'chat:tasks:cancel'}) 

1681 

1682 await asyncio.shield(emit_cancel_event()) 

1683 except Exception: 

1684 pass 

1685 raise # re-raise to ensure proper task cancellation handling 

1686 except Exception as e: 

1687 error_detail = e.detail if isinstance(e, HTTPException) else str(e) 

1688 log.error('Error processing chat payload: %s', error_detail) 

1689 if metadata.get('chat_id') and metadata.get('message_id'): 

1690 # Update the chat message with the error 

1691 try: 

1692 if is_saved_chat_id(metadata.get('chat_id')): 

1693 await Chats.upsert_message_to_chat_by_id_and_message_id( 

1694 metadata['chat_id'], 

1695 metadata['message_id'], 

1696 { 

1697 'parentId': metadata.get('user_message_id', None), 

1698 'error': {'content': error_detail}, 

1699 'done': True, 

1700 }, 

1701 ) 

1702 

1703 event_emitter = await get_event_emitter(metadata) 

1704 if event_emitter: 

1705 await event_emitter( 

1706 { 

1707 'type': 'chat:message:error', 

1708 'data': {'error': {'content': error_detail}, 'done': True}, 

1709 } 

1710 ) 

1711 

1712 except Exception: 

1713 pass 

1714 else: 

1715 # No chat_id/message_id → legacy/direct API path with no 

1716 # WebSocket error channel. We must surface the error as 

1717 # a proper HTTP response; without this the function would 

1718 # return None which FastAPI serializes as null. #23924 

1719 raise HTTPException( 

1720 status_code=status.HTTP_400_BAD_REQUEST, 

1721 detail=error_detail, 

1722 ) 

1723 finally: 

1724 # Clean up MCP clients. Each client is isolated so one 

1725 # failure doesn't skip the rest. 

1726 # 

1727 # NOTE: asyncio.wait_for() / asyncio.shield() must NOT be used 

1728 # here — they create new asyncio Tasks, which violate anyio 

1729 # cancel-scope task-ownership rules when the MCPClient's 

1730 # exit_stack contains anyio transport resources (streamable_http). 

1731 # Exiting those cancel scopes from the wrong task raises 

1732 # "Attempted to exit a cancel scope that isn't the current 

1733 # task's current cancel scope", which propagates as a 

1734 # BaseException through the finally block, discards the response 

1735 # return value, and surfaces as a 500 "No response returned." 

1736 # MCPClient.disconnect() suppresses known transport teardown errors 

1737 # while still propagating real task cancellation. 

1738 try: 

1739 if mcp_clients := metadata.get('mcp_clients'): 

1740 for client in reversed(list(mcp_clients.values())): 

1741 try: 

1742 await client.disconnect() 

1743 except BaseException as e: 

1744 log.debug('Error disconnecting MCP client: %s', e) 

1745 except BaseException as e: 

1746 log.debug('Error cleaning up MCP clients: %s', e) 

1747 

1748 # Deregister this task, then emit chat:active=false if no others remain 

1749 try: 

1750 chat_id = metadata.get('chat_id') 

1751 task_id = metadata.get('task_id') 

1752 if chat_id and task_id: 

1753 await cleanup_task(request.app.state.redis, task_id, chat_id) 

1754 if not await has_active_tasks(request.app.state.redis, chat_id): 

1755 event_emitter = await get_event_emitter(metadata, update_db=False) 

1756 if event_emitter: 

1757 try: 

1758 folder_id = metadata.get('folder_id') or await Chats.get_chat_folder_id( 

1759 chat_id, user.id 

1760 ) 

1761 await asyncio.shield( 

1762 event_emitter( 

1763 { 

1764 'type': 'chat:active', 

1765 'data': {'active': False, 'folder_id': folder_id}, 

1766 } 

1767 ) 

1768 ) 

1769 except asyncio.CancelledError: 

1770 pass 

1771 except Exception: 

1772 pass 

1773 

1774 try: 

1775 chat_id = metadata.get('chat_id') 

1776 if ( 

1777 chat_id 

1778 and getattr(request.state, 'internal', False) is not True 

1779 and not await has_active_tasks(request.app.state.redis, chat_id) 

1780 ): 

1781 from open_webui.utils.subagents import process_pending_internal_messages 

1782 

1783 await process_pending_internal_messages( 

1784 request, 

1785 chat_id, 

1786 user.id, 

1787 { 

1788 'model_id': metadata.get('model_id') or form_data.get('model'), 

1789 'session_id': metadata.get('session_id'), 

1790 'tool_ids': metadata.get('tool_ids') or [], 

1791 'skill_ids': metadata.get('skill_ids') or [], 

1792 'system_prompt': metadata.get('system_prompt'), 

1793 'filter_ids': metadata.get('filter_ids') or [], 

1794 'terminal_id': metadata.get('terminal_id'), 

1795 'features': metadata.get('features') or {}, 

1796 'variables': metadata.get('variables') or {}, 

1797 }, 

1798 ) 

1799 except Exception: 

1800 log.exception('Failed to process pending internal messages for chat %s', metadata.get('chat_id')) 

1801 

1802 # Fan out: one task per model 

1803 if metadata.get('session_id') and metadata.get('chat_id'): 

1804 task_ids = [] 

1805 subagent_results = [] 

1806 is_internal = getattr(request.state, 'internal', False) is True 

1807 chat_id = metadata['chat_id'] 

1808 

1809 for idx, entry in enumerate(message_ids): 

1810 target_model_id = entry['model_id'] 

1811 assistant_message_id = entry['message_id'] 

1812 if not assistant_message_id: 

1813 continue 

1814 

1815 # Per-model metadata: own message_id + model 

1816 per_model_metadata = { 

1817 **metadata, 

1818 'message_id': assistant_message_id, 

1819 'task_id': str(uuid4()), 

1820 } 

1821 

1822 # Per-model form_data: own model 

1823 model_form_data = { 

1824 **form_data, 

1825 'model': target_model_id, 

1826 'metadata': per_model_metadata, 

1827 } 

1828 

1829 # Resolve the model object for this specific model 

1830 resolved_model = request.app.state.MODELS.get(target_model_id, model) 

1831 

1832 # Only the first model runs chat-level background tasks; 

1833 # subsequent models only run follow-ups. 

1834 process = process_chat( 

1835 request, 

1836 model_form_data, 

1837 user, 

1838 per_model_metadata, 

1839 resolved_model, 

1840 tasks 

1841 if idx == 0 

1842 else { 

1843 k: v for k, v in (tasks or {}).items() if k not in (TASKS.TITLE_GENERATION, TASKS.TAGS_GENERATION) 

1844 } 

1845 or None, 

1846 ) 

1847 if is_internal: 

1848 subagent_results.append(await process) 

1849 continue 

1850 

1851 task_id, _ = await create_task( 

1852 request.app.state.redis, 

1853 process, 

1854 id=chat_id, 

1855 task_id=per_model_metadata['task_id'], 

1856 ) 

1857 task_ids.append(task_id) 

1858 

1859 if is_internal: 

1860 return { 

1861 'status': True, 

1862 'task_ids': [], 

1863 'chat_id': chat_id, 

1864 'results': subagent_results, 

1865 } 

1866 

1867 # Emit chat:active=true 

1868 if task_ids: 

1869 event_emitter = await get_event_emitter( 

1870 {**metadata, 'message_id': message_ids[0]['message_id']}, 

1871 update_db=False, 

1872 ) 

1873 if event_emitter: 

1874 folder_id = metadata.get('folder_id') or await Chats.get_chat_folder_id(chat_id, user.id) 

1875 await event_emitter({'type': 'chat:active', 'data': {'active': True, 'folder_id': folder_id}}) 

1876 

1877 return { 

1878 'status': True, 

1879 'task_ids': task_ids, 

1880 'chat_id': chat_id, 

1881 } 

1882 else: 

1883 # Legacy/direct: single model, synchronous 

1884 metadata['message_id'] = message_ids[0]['message_id'] 

1885 return await process_chat(request, form_data, user, metadata, model, tasks) 

1886 

1887 

1888# Alias for chat_completion (Legacy) 

1889generate_chat_completions = chat_completion 

1890generate_chat_completion = chat_completion 

1891 

1892 

1893@app.post('/api/v1/chats/{id}/messages/{message_id}/resolve') 

1894async def resolve_chat_message_tool_call( 

1895 request: Request, 

1896 id: str, 

1897 message_id: str, 

1898 form_data: ResolveToolCallForm, 

1899 user=Depends(get_verified_user), 

1900 db: AsyncSession = Depends(get_async_session), 

1901): 

1902 resolution = await resolve_tool_call_output(id, message_id, form_data, user, db=db) 

1903 payload = await build_tool_approval_resume_payload(id, message_id, chat=resolution['chat']) 

1904 result = await chat_completion(request, payload, user) 

1905 return { 

1906 'status': True, 

1907 'chat_id': id, 

1908 'message_id': message_id, 

1909 **(result if isinstance(result, dict) else {}), 

1910 } 

1911 

1912 

1913# Expose as app.state so internal callers (e.g. automations) can 

1914# use the full pipeline without importing from main.py (avoids circular deps). 

1915app.state.CHAT_COMPLETION_HANDLER = chat_completion 

1916 

1917 

1918################################## 

1919# 

1920# Anthropic Messages API Compatible Endpoint 

1921# 

1922################################## 

1923 

1924 

1925from open_webui.utils.anthropic import ( 

1926 convert_anthropic_to_openai_payload, 

1927 convert_openai_to_anthropic_response, 

1928 is_anthropic_messages_passthrough, 

1929 openai_stream_to_anthropic_stream, 

1930) 

1931 

1932 

1933@app.post('/api/message/count_tokens') 

1934@app.post('/api/v1/messages/count_tokens') # Anthropic Messages token-count endpoint 

1935async def count_message_tokens( 

1936 request: Request, 

1937 form_data: dict, 

1938 user=Depends(get_verified_user), 

1939): 

1940 return {'input_tokens': await openai.count_anthropic_tokens(request, form_data, user)} 

1941 

1942 

1943async def passthrough_anthropic_messages(request: Request, form_data: dict, user) -> Response | dict: 

1944 requested_model, payload, url, key, headers, cookies = await openai.get_anthropic_request_target( 

1945 request, form_data, user 

1946 ) 

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

1948 response = None 

1949 streaming = False 

1950 

1951 try: 

1952 session = await get_session() 

1953 response = await session.request( 

1954 method='POST', 

1955 url=request_url, 

1956 data=JSONCodec.dumps(payload), 

1957 headers=headers, 

1958 cookies=cookies, 

1959 ssl=AIOHTTP_CLIENT_SESSION_SSL, 

1960 timeout=get_client_timeout(stream=bool(payload.get('stream'))), 

1961 ) 

1962 

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

1964 streaming = True 

1965 return StreamingResponse( 

1966 stream_wrapper(response), 

1967 status_code=response.status, 

1968 headers=openai._clean_proxy_headers(response.headers), 

1969 ) 

1970 

1971 try: 

1972 response_data = await response.json() 

1973 except Exception: 

1974 response_data = await response.text() 

1975 

1976 if response.status >= 400: 

1977 await openai.publish_model_provider_request_failed( 

1978 request, 

1979 actor=user, 

1980 provider='openai-compatible', 

1981 base_url=url, 

1982 api_key=key, 

1983 status=response.status, 

1984 requested_model=requested_model, 

1985 upstream_error=response_data, 

1986 ) 

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

1988 return JSONResponse(status_code=response.status, content=response_data) 

1989 return Response(status_code=response.status, content=response_data) 

1990 

1991 return response_data 

1992 except HTTPException: 

1993 raise 

1994 except Exception: 

1995 log.exception('Failed to passthrough Anthropic Messages request for model %s', requested_model) 

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

1997 finally: 

1998 if not streaming: 

1999 await cleanup_response(response) 

2000 

2001 

2002@app.post('/api/message') 

2003@app.post('/api/v1/messages') # Anthropic Messages API compatible endpoint 

2004async def generate_messages( 

2005 request: Request, 

2006 form_data: dict, 

2007 user=Depends(get_verified_user), 

2008): 

2009 """ 

2010 Anthropic Messages API compatible endpoint. 

2011 

2012 Accepts the Anthropic Messages API format, converts internally to OpenAI 

2013 Chat Completions format, routes through the existing chat completion 

2014 pipeline, then converts the response back to Anthropic Messages format. 

2015 

2016 Supports both streaming and non-streaming requests. 

2017 All models configured in Open WebUI are accessible via this endpoint. 

2018 

2019 Authentication: Supports both standard Authorization header and 

2020 Anthropic's x-api-key header (via middleware translation). 

2021 """ 

2022 requested_model = form_data.get('model', '') 

2023 input_tokens = None 

2024 try: 

2025 input_tokens = await openai.count_anthropic_tokens(request, form_data, user) 

2026 except Exception: 

2027 # Counting must not turn a compatible generation request into an outage. 

2028 log.warning('Unable to count Anthropic input tokens for model %s', requested_model, exc_info=True) 

2029 

2030 model_id = requested_model 

2031 model_info = await Models.get_model_by_id(model_id) 

2032 if model_info and model_info.base_model_id: 2032 ↛ 2033line 2032 didn't jump to line 2033 because the condition on line 2032 was never true

2033 model_id = model_info.base_model_id 

2034 

2035 passthrough_params = [] 

2036 models = request.app.state.OPENAI_MODELS 

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

2038 await openai.get_all_models(request, user=user) 

2039 models = request.app.state.OPENAI_MODELS 

2040 model = models.get(model_id) 

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

2042 url, _, api_config = await openai.get_openai_connection(model['urlIdx']) 

2043 if is_anthropic_messages_passthrough(url, api_config): 

2044 return await passthrough_anthropic_messages(request, form_data, user) 

2045 passthrough_params = api_config.get('passthrough_params') or [] 

2046 

2047 # Convert Anthropic payload to OpenAI format 

2048 openai_payload = convert_anthropic_to_openai_payload(form_data, passthrough_params) 

2049 

2050 # Route through the existing chat_completion handler 

2051 response = await chat_completion(request, openai_payload, user) 

2052 

2053 # Convert response back to Anthropic format 

2054 if isinstance(response, StreamingResponse): 

2055 # Streaming response: wrap the generator to convert SSE format 

2056 return StreamingResponse( 

2057 openai_stream_to_anthropic_stream(response.body_iterator, model=requested_model, input_tokens=input_tokens), 

2058 media_type='text/event-stream', 

2059 headers={ 

2060 'Cache-Control': 'no-cache', 

2061 'Connection': 'keep-alive', 

2062 }, 

2063 ) 

2064 elif isinstance(response, dict): 

2065 return convert_openai_to_anthropic_response(response, model=requested_model, input_tokens=input_tokens) 

2066 else: 

2067 # Passthrough for error responses (JSONResponse, PlainTextResponse, etc.) 

2068 return response 

2069 

2070 

2071async def verify_chat_ownership(chat_id: str | None, user) -> None: 

2072 """Temporary chats are per-socket and unsaved, so they have no owner to check.""" 

2073 if not chat_id or is_temporary_chat_id(chat_id): 2073 ↛ 2077line 2073 didn't jump to line 2077 because the condition on line 2073 was always true

2074 return 

2075 

2076 # Channel messages need the membership and write-access gate that only /api/chat/completions has. 

2077 if chat_id.startswith('channel:'): 

2078 raise HTTPException( 

2079 status_code=status.HTTP_400_BAD_REQUEST, 

2080 detail='Channel chats are not supported on this endpoint', 

2081 ) 

2082 

2083 if user.role != 'admin' and not await Chats.is_chat_owner(chat_id, user.id): 

2084 raise HTTPException( 

2085 status_code=status.HTTP_404_NOT_FOUND, 

2086 detail=ERROR_MESSAGES.DEFAULT(), 

2087 ) 

2088 

2089 

2090@app.post('/api/chat/completed') 

2091async def chat_completed(request: Request, form_data: dict, user=Depends(get_verified_user)): 

2092 """Deprecated: outlet filters now run inline during chat completion. 

2093 Kept for backward compatibility with external integrations.""" 

2094 await verify_chat_ownership(form_data.get('chat_id'), user) 

2095 

2096 try: 

2097 model_item = form_data.pop('model_item', {}) 

2098 

2099 if model_item.get('direct', False): 2099 ↛ 2100line 2099 didn't jump to line 2100 because the condition on line 2099 was never true

2100 await _set_direct_model(request, model_item, user) 

2101 

2102 return await chat_completed_handler(request, form_data, user) 

2103 except Exception as e: 

2104 raise HTTPException( 

2105 status_code=status.HTTP_400_BAD_REQUEST, 

2106 detail=str(e), 

2107 ) 

2108 

2109 

2110@app.post('/api/chat/actions/{action_id}') 

2111async def chat_action(request: Request, action_id: str, form_data: dict, user=Depends(get_verified_user)): 

2112 await verify_chat_ownership(form_data.get('chat_id'), user) 

2113 

2114 try: 

2115 model_item = form_data.pop('model_item', {}) 

2116 

2117 if model_item.get('direct', False): 2117 ↛ 2118line 2117 didn't jump to line 2118 because the condition on line 2117 was never true

2118 await _set_direct_model(request, model_item, user) 

2119 

2120 return await chat_action_handler(request, action_id, form_data, user) 

2121 except Exception as e: 

2122 raise HTTPException( 

2123 status_code=status.HTTP_400_BAD_REQUEST, 

2124 detail=str(e), 

2125 ) 

2126 

2127 

2128@app.post('/api/tasks/stop/{task_id}') 

2129async def stop_task_endpoint(request: Request, task_id: str, user=Depends(get_admin_user)): 

2130 try: 

2131 result = await stop_task(request.app.state.redis, task_id) 

2132 return result 

2133 except ValueError as e: 

2134 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=str(e)) 

2135 

2136 

2137@app.get('/api/tasks') 

2138async def list_tasks_endpoint(request: Request, user=Depends(get_admin_user)): 

2139 return {'tasks': await list_tasks(request.app.state.redis)} 

2140 

2141 

2142@app.get('/api/tasks/chat/{chat_id:path}') 

2143async def list_tasks_by_chat_id_endpoint(request: Request, chat_id: str, user=Depends(get_verified_user)): 

2144 socket_id = get_temporary_chat_session_id(chat_id) 

2145 if socket_id: 2145 ↛ 2146line 2145 didn't jump to line 2146 because the condition on line 2145 was never true

2146 owner_id = get_user_id_from_session_pool(socket_id) 

2147 if owner_id != user.id and user.role != 'admin': 

2148 return {'task_ids': []} 

2149 else: 

2150 chat = await Chats.get_chat_by_id(chat_id) 

2151 if chat is None or (chat.user_id != user.id and user.role != 'admin'): 2151 ↛ 2154line 2151 didn't jump to line 2154 because the condition on line 2151 was always true

2152 return {'task_ids': []} 

2153 

2154 task_ids = await list_task_ids_by_item_id(request.app.state.redis, chat_id) 

2155 

2156 log.debug('Task IDs for chat %s: %s', chat_id, task_ids) 

2157 return {'task_ids': task_ids} 

2158 

2159 

2160@app.post('/api/tasks/chat/{chat_id:path}/stop') 

2161async def stop_tasks_by_chat_id_endpoint(request: Request, chat_id: str, user=Depends(get_verified_user)): 

2162 socket_id = get_temporary_chat_session_id(chat_id) 

2163 chat = None 

2164 if socket_id: 2164 ↛ 2165line 2164 didn't jump to line 2165 because the condition on line 2164 was never true

2165 owner_id = get_user_id_from_session_pool(socket_id) 

2166 if owner_id != user.id and user.role != 'admin': 

2167 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND) 

2168 else: 

2169 chat = await Chats.get_chat_by_id(chat_id) 

2170 if chat is None or (chat.user_id != user.id and user.role != 'admin'): 

2171 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND) 

2172 result = await stop_item_tasks(request.app.state.redis, chat_id) 

2173 

2174 if not socket_id and str(result.get('message', '')).startswith('No tasks found'): 2174 ↛ 2214line 2174 didn't jump to line 2214 because the condition on line 2174 was always true

2175 messages_map = await Chats.get_messages_map_by_chat_id(chat_id) or {} 

2176 for message_id, message in messages_map.items(): 

2177 if message.get('role') != 'assistant' or message.get('done') is not False: 2177 ↛ 2180line 2177 didn't jump to line 2180 because the condition on line 2177 was always true

2178 continue 

2179 

2180 output = message.get('output') 

2181 if isinstance(output, list): 

2182 for item in output: 

2183 if item.get('type') == 'function_call' and item.get('status') in { 

2184 'pending', 

2185 'queued', 

2186 'requires_approval', 

2187 }: 

2188 item['status'] = 'rejected' 

2189 item.pop('approved', None) 

2190 

2191 await Chats.upsert_message_to_chat_by_id_and_message_id( 

2192 chat_id, 

2193 message_id, 

2194 {'done': True, **({'output': output} if isinstance(output, list) else {})}, 

2195 touch=False, 

2196 ) 

2197 result = { 

2198 'status': True, 

2199 'message': 'Finalized pending approval message.', 

2200 } 

2201 

2202 event_emitter = await get_event_emitter( 

2203 { 

2204 'user_id': chat.user_id, 

2205 'chat_id': chat_id, 

2206 'message_id': message_id, 

2207 }, 

2208 update_db=False, 

2209 ) 

2210 if event_emitter: 

2211 await event_emitter({'type': 'chat:completion', 'data': {'done': True, 'output': output}}) 

2212 await event_emitter({'type': 'chat:tasks:cancel'}) 

2213 

2214 return result 

2215 

2216 

2217################################## 

2218# 

2219# Config Endpoints 

2220# 

2221################################## 

2222 

2223 

2224@app.get('/api/config') 

2225async def get_app_config(request: Request): 

2226 user = None 

2227 token = None 

2228 

2229 auth_header = request.headers.get('Authorization') 

2230 if auth_header: 2230 ↛ 2231line 2230 didn't jump to line 2231 because the condition on line 2230 was never true

2231 cred = get_http_authorization_cred(auth_header) 

2232 if cred: 

2233 token = cred.credentials 

2234 

2235 if not token and 'token' in request.cookies: 2235 ↛ 2236line 2235 didn't jump to line 2236 because the condition on line 2235 was never true

2236 token = request.cookies.get('token') 

2237 

2238 if token: 2238 ↛ 2239line 2238 didn't jump to line 2239 because the condition on line 2238 was never true

2239 try: 

2240 data = decode_token(token) 

2241 except Exception as e: 

2242 log.debug(e) 

2243 raise HTTPException( 

2244 status_code=status.HTTP_401_UNAUTHORIZED, 

2245 detail='Invalid token', 

2246 ) 

2247 if data is not None and 'id' in data: 

2248 user = await Users.get_user_by_id(data['id']) 

2249 

2250 onboarding = False 

2251 if user is None: 2251 ↛ 2254line 2251 didn't jump to line 2254 because the condition on line 2251 was always true

2252 onboarding = not await Users.has_users() 

2253 

2254 license_metadata = getattr(app.state, 'LICENSE_METADATA', None) 

2255 user_count = await Users.get_num_users() if license_metadata else None 

2256 config = await Config.get_many( 

2257 'oauth.enable', 

2258 'oauth.auto_redirect', 

2259 'ldap.enable', 

2260 'ui.enable_signup', 

2261 'ui.enable_login_form', 

2262 'auth.enable_api_keys', 

2263 'ui.enable_password_change_form', 

2264 'direct.enable', 

2265 'direct.integrations.enable', 

2266 'folders.enable', 

2267 'folders.max_file_count', 

2268 'channels.enable', 

2269 'calendar.enable', 

2270 'automations.enable', 

2271 'notes.enable', 

2272 'chat.context_compaction.enable', 

2273 'chat.tool_permissions.enable', 

2274 'web.search.enable', 

2275 'web.search.confirmation.enable', 

2276 'web.search.confirmation.content', 

2277 'code_execution.enable', 

2278 'code_interpreter.enable', 

2279 'image_generation.enable', 

2280 'task.autocomplete.enable', 

2281 'ui.enable_community_sharing', 

2282 'ui.enable_message_rating', 

2283 'ui.enable_user_webhooks', 

2284 'users.enable_status', 

2285 'google_drive.enable', 

2286 'onedrive.enable', 

2287 'memories.enable', 

2288 'ui.default_models', 

2289 'ui.default_pinned_models', 

2290 'ui.default_interface_settings', 

2291 'ui.i18n', 

2292 'ui.prompt_suggestions', 

2293 'ui.prompt_suggestions_i18n', 

2294 'code_execution.engine', 

2295 'code_interpreter.engine', 

2296 'audio.tts.engine', 

2297 'audio.tts.voice', 

2298 'audio.tts.split_on', 

2299 'audio.stt.engine', 

2300 'rag.file.max_size', 

2301 'rag.file.max_count', 

2302 'file.image_compression_width', 

2303 'file.image_compression_height', 

2304 'user.permissions', 

2305 'ui.pending_user_overlay_title', 

2306 'ui.pending_user_overlay_content', 

2307 'ui.watermark', 

2308 ) 

2309 

2310 return { 

2311 **({'onboarding': True} if onboarding else {}), 

2312 'status': True, 

2313 'name': app.state.WEBUI_NAME, 

2314 'version': VERSION, 

2315 'default_locale': str(DEFAULT_LOCALE), 

2316 'i18n': config.get('ui.i18n') or {}, 

2317 'oauth': { 

2318 # Hide providers (and thus the login buttons / auto-redirect) when OAuth 

2319 # is disabled, without clearing the admin's provider configuration. 

2320 'providers': ( 

2321 {name: provider.get('name', name) for name, provider in OAUTH_PROVIDERS.items()} 

2322 if config.get('oauth.enable', True) 

2323 else {} 

2324 ), 

2325 'auto_redirect': config.get('oauth.auto_redirect'), 

2326 }, 

2327 'features': { 

2328 'slim': USE_SLIM, 

2329 # --- Public: required by login/signup page pre-auth --- 

2330 'auth': WEBUI_AUTH, 

2331 'auth_trusted_header': bool(WEBUI_AUTH_TRUSTED_EMAIL_HEADER), 

2332 'enable_signup_password_confirmation': ENABLE_SIGNUP_PASSWORD_CONFIRMATION, 

2333 'enable_ldap': config.get('ldap.enable'), 

2334 'enable_signup': config.get('ui.enable_signup'), 

2335 'enable_login_form': config.get('ui.enable_login_form'), 

2336 'enable_websocket': ENABLE_WEBSOCKET_SUPPORT, 

2337 **( 

2338 {'websocket_heartbeat_interval': WEBSOCKET_HEARTBEAT_INTERVAL} 

2339 if WEBSOCKET_HEARTBEAT_INTERVAL is not None 

2340 else {} 

2341 ), 

2342 # --- Authenticated: only consumed by logged-in frontend --- 

2343 **( 

2344 { 

2345 'enable_api_keys': config.get('auth.enable_api_keys'), 

2346 'enable_password_change_form': config.get('ui.enable_password_change_form'), 

2347 'enable_version_update_check': ENABLE_VERSION_UPDATE_CHECK, 

2348 'enable_pyodide_file_persistence': ENABLE_PYODIDE_FILE_PERSISTENCE, 

2349 'enable_public_active_users_count': ENABLE_PUBLIC_ACTIVE_USERS_COUNT, 

2350 'enable_easter_eggs': ENABLE_EASTER_EGGS, 

2351 'enable_direct_connections': config.get('direct.enable'), 

2352 'enable_direct_integrations': config.get('direct.integrations.enable', False), 

2353 'enable_plugins': ENABLE_PLUGINS, 

2354 'enable_folders': config.get('folders.enable'), 

2355 'folder_max_file_count': config.get('folders.max_file_count'), 

2356 'enable_channels': config.get('channels.enable'), 

2357 'enable_calendar': config.get('calendar.enable'), 

2358 'enable_automations': config.get('automations.enable'), 

2359 'enable_notes': config.get('notes.enable'), 

2360 'enable_context_compaction': config.get('chat.context_compaction.enable'), 

2361 'enable_tool_permissions': config.get('chat.tool_permissions.enable'), 

2362 'enable_web_search': config.get('web.search.enable'), 

2363 'enable_web_search_confirmation': config.get('web.search.confirmation.enable'), 

2364 'web_search_confirmation_content': config.get('web.search.confirmation.content'), 

2365 'enable_code_execution': config.get('code_execution.enable'), 

2366 'enable_code_interpreter': config.get('code_interpreter.enable'), 

2367 'enable_image_generation': config.get('image_generation.enable'), 

2368 'enable_autocomplete_generation': config.get('task.autocomplete.enable'), 

2369 'enable_community_sharing': config.get('ui.enable_community_sharing'), 

2370 'enable_message_rating': config.get('ui.enable_message_rating'), 

2371 'enable_user_webhooks': config.get('ui.enable_user_webhooks'), 

2372 'enable_user_status': config.get('users.enable_status'), 

2373 'enable_admin_export': ENABLE_ADMIN_EXPORT, 

2374 'enable_admin_chat_access': ENABLE_ADMIN_CHAT_ACCESS, 

2375 'enable_admin_analytics': ENABLE_ADMIN_ANALYTICS, 

2376 'enable_google_drive_integration': config.get('google_drive.enable'), 

2377 'enable_onedrive_integration': config.get('onedrive.enable'), 

2378 'enable_memories': config.get('memories.enable'), 

2379 **( 

2380 { 

2381 'enable_onedrive_personal': ENABLE_ONEDRIVE_PERSONAL, 

2382 'enable_onedrive_business': ENABLE_ONEDRIVE_BUSINESS, 

2383 } 

2384 if config.get('onedrive.enable') 

2385 else {} 

2386 ), 

2387 } 

2388 if user is not None 

2389 else {} 

2390 ), 

2391 }, 

2392 **( 

2393 { 

2394 'default_models': config.get('ui.default_models'), 

2395 'default_pinned_models': config.get('ui.default_pinned_models'), 

2396 'default_prompt_suggestions': config.get('ui.prompt_suggestions'), 

2397 'default_prompt_suggestions_i18n': config.get('ui.prompt_suggestions_i18n'), 

2398 **({'user_count': user_count} if user_count is not None else {}), 

2399 'code': { 

2400 'engine': config.get('code_execution.engine'), 

2401 'interpreter_engine': config.get('code_interpreter.engine'), 

2402 }, 

2403 'audio': { 

2404 'tts': { 

2405 'engine': config.get('audio.tts.engine'), 

2406 'voice': config.get('audio.tts.voice'), 

2407 'split_on': config.get('audio.tts.split_on'), 

2408 }, 

2409 'stt': { 

2410 'engine': config.get('audio.stt.engine'), 

2411 }, 

2412 }, 

2413 'file': { 

2414 'max_size': config.get('rag.file.max_size'), 

2415 'max_count': config.get('rag.file.max_count'), 

2416 'image_compression': { 

2417 'width': config.get('file.image_compression_width'), 

2418 'height': config.get('file.image_compression_height'), 

2419 }, 

2420 }, 

2421 'permissions': {**(config.get('user.permissions') or {})}, 

2422 'google_drive': { 

2423 'client_id': GOOGLE_DRIVE_CLIENT_ID, 

2424 'api_key': GOOGLE_DRIVE_API_KEY, 

2425 }, 

2426 'onedrive': { 

2427 'client_id_personal': ONEDRIVE_CLIENT_ID_PERSONAL, 

2428 'client_id_business': ONEDRIVE_CLIENT_ID_BUSINESS, 

2429 'sharepoint_url': ONEDRIVE_SHAREPOINT_URL, 

2430 'sharepoint_tenant_id': ONEDRIVE_SHAREPOINT_TENANT_ID, 

2431 }, 

2432 'ui': { 

2433 'default_interface_settings': config.get('ui.default_interface_settings'), 

2434 'pending_user_overlay_title': config.get('ui.pending_user_overlay_title'), 

2435 'pending_user_overlay_content': config.get('ui.pending_user_overlay_content'), 

2436 'response_watermark': config.get('ui.watermark'), 

2437 'iframe_csp': IFRAME_CSP, 

2438 }, 

2439 'license_metadata': license_metadata, 

2440 **( 

2441 { 

2442 'active_entries': user_count, 

2443 } 

2444 if user.role == 'admin' and user_count is not None 

2445 else {} 

2446 ), 

2447 } 

2448 if user is not None and (user.role in ['admin', 'user']) 

2449 else { 

2450 **( 

2451 { 

2452 'ui': { 

2453 'pending_user_overlay_title': config.get('ui.pending_user_overlay_title'), 

2454 'pending_user_overlay_content': config.get('ui.pending_user_overlay_content'), 

2455 } 

2456 } 

2457 if user and user.role == 'pending' 

2458 else {} 

2459 ), 

2460 **( 

2461 { 

2462 'metadata': { 

2463 'login_footer': license_metadata.get('login_footer', ''), 

2464 'auth_logo_position': license_metadata.get('auth_logo_position', ''), 

2465 } 

2466 } 

2467 if license_metadata 

2468 else {} 

2469 ), 

2470 } 

2471 ), 

2472 } 

2473 

2474 

2475class EventWebhookForm(BaseModel): 

2476 name: str | None = None 

2477 url: str 

2478 enabled: bool = True 

2479 events: list[str] | None = None 

2480 targets: list[dict[str, str]] | None = None 

2481 

2482 

2483class EventWebhookUpdateForm(BaseModel): 

2484 name: str | None = None 

2485 url: str | None = None 

2486 enabled: bool | None = None 

2487 events: list[str] | None = None 

2488 targets: list[dict[str, str]] | None = None 

2489 

2490 

2491@app.get('/api/events') 

2492async def get_event_catalog(user=Depends(get_admin_user)): 

2493 return { 

2494 'schema': VERSION, 

2495 'events': get_event_catalog_items(), 

2496 } 

2497 

2498 

2499@app.get('/api/events/webhooks') 

2500async def get_event_webhooks_api(user=Depends(get_admin_user)): 

2501 return await get_event_webhooks() 

2502 

2503 

2504@app.post('/api/events/webhooks') 

2505async def create_event_webhook(form_data: EventWebhookForm, user=Depends(get_admin_user)): 

2506 try: 

2507 webhook = await upsert_event_webhook( 

2508 { 

2509 'name': form_data.name, 

2510 'url': form_data.url, 

2511 'enabled': form_data.enabled, 

2512 'events': form_data.events, 

2513 'targets': form_data.targets, 

2514 } 

2515 ) 

2516 except ValueError as e: 

2517 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=str(e)) 

2518 

2519 await publish_event( 

2520 app, 

2521 EVENTS.CONFIG_WEBHOOK_UPDATED, 

2522 actor=user, 

2523 subject_id=webhook['id'], 

2524 subject_type='config', 

2525 data={ 

2526 'action': 'created', 

2527 'enabled': webhook.get('enabled'), 

2528 'events': webhook.get('events'), 

2529 'targets': webhook.get('targets'), 

2530 }, 

2531 ) 

2532 return webhook 

2533 

2534 

2535@app.put('/api/events/webhooks/{webhook_id}') 

2536async def update_event_webhook(webhook_id: str, form_data: EventWebhookUpdateForm, user=Depends(get_admin_user)): 

2537 webhooks = await get_event_webhooks() 

2538 existing = next((webhook for webhook in webhooks if webhook.get('id') == webhook_id), None) 

2539 if not existing: 2539 ↛ 2542line 2539 didn't jump to line 2542 because the condition on line 2539 was always true

2540 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail='Webhook not found') 

2541 

2542 try: 

2543 webhook = await upsert_event_webhook( 

2544 { 

2545 **existing, 

2546 **form_data.model_dump(exclude_unset=True), 

2547 'id': webhook_id, 

2548 } 

2549 ) 

2550 except ValueError as e: 

2551 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=str(e)) 

2552 

2553 await publish_event( 

2554 app, 

2555 EVENTS.CONFIG_WEBHOOK_UPDATED, 

2556 actor=user, 

2557 subject_id=webhook_id, 

2558 subject_type='config', 

2559 data={ 

2560 'action': 'updated', 

2561 'enabled': webhook.get('enabled'), 

2562 'events': webhook.get('events'), 

2563 'targets': webhook.get('targets'), 

2564 }, 

2565 ) 

2566 return webhook 

2567 

2568 

2569@app.delete('/api/events/webhooks/{webhook_id}') 

2570async def delete_event_webhook_api(webhook_id: str, user=Depends(get_admin_user)): 

2571 deleted = await delete_event_webhook(webhook_id) 

2572 if not deleted: 2572 ↛ 2575line 2572 didn't jump to line 2575 because the condition on line 2572 was always true

2573 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail='Webhook not found') 

2574 

2575 await publish_event( 

2576 app, 

2577 EVENTS.CONFIG_WEBHOOK_UPDATED, 

2578 actor=user, 

2579 subject_id=webhook_id, 

2580 subject_type='config', 

2581 data={'action': 'deleted'}, 

2582 ) 

2583 return {'status': True} 

2584 

2585 

2586@app.get('/api/version') 

2587async def get_app_version(): 

2588 return { 

2589 'version': VERSION, 

2590 'deployment_id': DEPLOYMENT_ID, 

2591 } 

2592 

2593 

2594@app.get('/api/version/updates') 

2595async def get_app_latest_release_version(user=Depends(get_verified_user)): 

2596 if not ENABLE_VERSION_UPDATE_CHECK: 2596 ↛ 2599line 2596 didn't jump to line 2599 because the condition on line 2596 was always true

2597 log.debug(f'Version update check is disabled, returning current version as latest version') 

2598 return {'current': VERSION, 'latest': VERSION} 

2599 try: 

2600 timeout = aiohttp.ClientTimeout(total=1) 

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

2602 async with session.get( 

2603 'https://api.github.com/repos/open-webui/open-webui/releases/latest', 

2604 ssl=AIOHTTP_CLIENT_SESSION_SSL, 

2605 ) as response: 

2606 response.raise_for_status() 

2607 data = await response.json() 

2608 latest_version = data['tag_name'] 

2609 

2610 return {'current': VERSION, 'latest': latest_version[1:]} 

2611 except Exception as e: 

2612 log.warning(f'Version update check failed: {e}') 

2613 return {'current': VERSION, 'latest': None} 

2614 

2615 

2616@app.get('/api/changelog') 

2617async def get_app_changelog(): 

2618 return {key: CHANGELOG[key] for idx, key in enumerate(CHANGELOG) if idx < 5} 

2619 

2620 

2621@app.get('/api/usage') 

2622async def get_current_usage(user=Depends(get_verified_user)): 

2623 """ 

2624 Get current usage statistics for Open WebUI. 

2625 This is an experimental endpoint and subject to change. 

2626 """ 

2627 try: 

2628 # If public visibility is disabled, only allow admins to access this endpoint 

2629 if not ENABLE_PUBLIC_ACTIVE_USERS_COUNT and user.role != 'admin': 2629 ↛ 2630line 2629 didn't jump to line 2630 because the condition on line 2629 was never true

2630 raise HTTPException( 

2631 status_code=status.HTTP_403_FORBIDDEN, 

2632 detail='Access denied. Only administrators can view usage statistics.', 

2633 ) 

2634 

2635 return { 

2636 'model_ids': get_models_in_use(), 

2637 'user_count': await Users.get_active_user_count(), 

2638 } 

2639 except HTTPException: 

2640 raise 

2641 except Exception as e: 

2642 log.error(f'Error getting usage statistics: {e}') 

2643 raise HTTPException(status_code=500, detail='Internal Server Error') 

2644 

2645 

2646# --- OAuth Login & Callback --- 

2647 

2648 

2649try: 

2650 if ENABLE_STAR_SESSIONS_MIDDLEWARE: 2650 ↛ 2651line 2650 didn't jump to line 2651 because the condition on line 2650 was never true

2651 redis_session_store = RedisStore( 

2652 url=REDIS_URL, 

2653 prefix=(f'{REDIS_KEY_PREFIX}:session:' if REDIS_KEY_PREFIX else 'session:'), 

2654 ) 

2655 

2656 app.add_middleware(SessionAutoloadMiddleware) 

2657 app.add_middleware( 

2658 StarSessionsMiddleware, 

2659 store=redis_session_store, 

2660 cookie_name='owui-session', 

2661 cookie_same_site=WEBUI_SESSION_COOKIE_SAME_SITE, 

2662 cookie_https_only=WEBUI_SESSION_COOKIE_SECURE, 

2663 ) 

2664 log.info('Using Redis for session') 

2665 else: 

2666 raise ValueError('No Redis URL provided') 

2667except Exception as e: 

2668 app.add_middleware( 

2669 SessionMiddleware, 

2670 secret_key=WEBUI_SECRET_KEY, 

2671 session_cookie='owui-session', 

2672 same_site=WEBUI_SESSION_COOKIE_SAME_SITE, 

2673 https_only=WEBUI_SESSION_COOKIE_SECURE, 

2674 ) 

2675 

2676 

2677async def register_client(request, client_id: str) -> bool: 

2678 server_type, server_id = client_id.split(':', 1) 

2679 

2680 connection = None 

2681 connection_idx = None 

2682 

2683 tool_server_connections = await Config.get('tool_server.connections', []) or [] 

2684 for idx, conn in enumerate(tool_server_connections): 

2685 if conn.get('type', 'openapi') == server_type: 

2686 info = conn.get('info') or {} 

2687 if info.get('id') == server_id: 

2688 connection = conn 

2689 connection_idx = idx 

2690 break 

2691 

2692 if connection is None or connection_idx is None: 

2693 log.warning(f'Unable to locate MCP tool server configuration for client {client_id} during re-registration') 

2694 return False 

2695 

2696 server_url = connection.get('url') 

2697 auth_type = connection.get('auth_type', 'none') 

2698 oauth_scope = (connection.get('info') or {}).get('oauth_scope') or (connection.get('config') or {}).get( 

2699 'oauth_scope' 

2700 ) 

2701 oauth_server_key = (connection.get('config') or {}).get('oauth_server_key') 

2702 

2703 try: 

2704 if auth_type == 'oauth_2.1_static': 

2705 # Static credentials: rebuild from admin-provided credentials + fresh metadata 

2706 info = connection.get('info') or {} 

2707 oauth_client_id = info.get('oauth_client_id') or '' 

2708 oauth_client_secret = info.get('oauth_client_secret') or '' 

2709 if not oauth_client_id or not oauth_client_secret: 

2710 # Fall back to blob for backward compatibility 

2711 existing_client_info = info.get('oauth_client_info', '') 

2712 if not existing_client_info: 

2713 log.error(f'No stored OAuth client info for static client {client_id}') 

2714 return False 

2715 existing_data = decrypt_data(existing_client_info) 

2716 oauth_client_id = oauth_client_id or existing_data.get('client_id', '') 

2717 oauth_client_secret = oauth_client_secret or existing_data.get('client_secret', '') 

2718 oauth_client_info = await get_oauth_client_info_with_static_credentials( 

2719 request, 

2720 client_id, 

2721 server_url, 

2722 oauth_client_id=oauth_client_id, 

2723 oauth_client_secret=oauth_client_secret, 

2724 oauth_scope=oauth_scope, 

2725 ) 

2726 else: 

2727 oauth_client_info = await get_oauth_client_info_with_dynamic_client_registration( 

2728 request, 

2729 client_id, 

2730 server_url, 

2731 oauth_server_key, 

2732 oauth_scope=oauth_scope, 

2733 ) 

2734 except InvalidToken: 

2735 log.error( 

2736 'OAuth client re-registration failed for %s: InvalidToken. ' 

2737 'Stored OAuth client data is invalid; reconnect this tool server.', 

2738 client_id, 

2739 ) 

2740 return False 

2741 except Exception as e: 

2742 log.error( 

2743 'OAuth client re-registration failed for %s: %s', 

2744 client_id, 

2745 f'{type(e).__name__}: {e}' if str(e) else type(e).__name__, 

2746 ) 

2747 return False 

2748 

2749 try: 

2750 connections = await Config.get('tool_server.connections', []) or [] 

2751 connections[connection_idx] = { 

2752 **connection, 

2753 'info': { 

2754 **(connection.get('info') or {}), 

2755 'oauth_client_info': encrypt_data(oauth_client_info.model_dump(mode='json')), 

2756 }, 

2757 } 

2758 await Config.upsert({'tool_server.connections': connections}) 

2759 except Exception as e: 

2760 log.error(f'Failed to persist updated OAuth client info for tool server {client_id}: {e}') 

2761 return False 

2762 

2763 oauth_client_manager.remove_client(client_id) 

2764 oauth_client_info = OAuthClientInformationFull( 

2765 **apply_connection_oauth_options(connection, oauth_client_info.model_dump(mode='json')) 

2766 ) 

2767 oauth_client_manager.add_client(client_id, oauth_client_info) 

2768 log.info('Re-registered OAuth client %s for tool server', client_id) 

2769 return True 

2770 

2771 

2772@app.get('/oauth/clients/{client_id}/authorize') 

2773async def oauth_client_authorize( 

2774 client_id: str, 

2775 request: Request, 

2776 response: Response, 

2777 user=Depends(get_verified_user), 

2778): 

2779 # ensure_valid_client_registration 

2780 client = await oauth_client_manager.get_client(client_id) 

2781 client_info = await oauth_client_manager.get_client_info(client_id) 

2782 if client is None or client_info is None: 2782 ↛ 2785line 2782 didn't jump to line 2785 because the condition on line 2782 was always true

2783 raise HTTPException(status.HTTP_404_NOT_FOUND) 

2784 

2785 if not await oauth_client_manager._preflight_authorization_url(client, client_info): 

2786 log.info( 

2787 'Detected invalid OAuth client %s; attempting re-registration', 

2788 client_id, 

2789 ) 

2790 

2791 registered = await register_client(request, client_id) 

2792 if not registered: 

2793 raise HTTPException( 

2794 status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, 

2795 detail='Failed to re-register OAuth client', 

2796 ) 

2797 

2798 client = await oauth_client_manager.get_client(client_id) 

2799 client_info = await oauth_client_manager.get_client_info(client_id) 

2800 if client is None or client_info is None: 

2801 raise HTTPException( 

2802 status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, 

2803 detail='OAuth client unavailable after re-registration', 

2804 ) 

2805 

2806 if not await oauth_client_manager._preflight_authorization_url(client, client_info): 

2807 raise HTTPException( 

2808 status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, 

2809 detail='OAuth client registration is still invalid after re-registration', 

2810 ) 

2811 

2812 return await oauth_client_manager.handle_authorize(request, client_id=client_id, user_id=user.id) 

2813 

2814 

2815@app.get('/oauth/clients/{client_id}/callback') 

2816async def oauth_client_callback( 

2817 client_id: str, 

2818 request: Request, 

2819 response: Response, 

2820): 

2821 return await oauth_client_manager.handle_callback( 

2822 request, 

2823 client_id=client_id, 

2824 response=response, 

2825 ) 

2826 

2827 

2828@app.get('/oauth/{provider}/login') 

2829async def oauth_login(provider: str, request: Request): 

2830 return await oauth_manager.handle_login(request, provider) 

2831 

2832 

2833@app.get('/oauth/{provider}/login/callback') 

2834@app.get('/oauth/{provider}/callback') # Legacy endpoint 

2835async def oauth_login_callback( 

2836 provider: str, 

2837 request: Request, 

2838 response: Response, 

2839 db: AsyncSession = Depends(get_async_session), 

2840): 

2841 """Handle the OAuth provider callback. 

2842 

2843 Resolution order: 

2844 1. Match by subject ID bound to the provider. 

2845 2. If ``OAUTH_MERGE_ACCOUNTS_BY_EMAIL`` is enabled, match by email 

2846 (note: some providers do not verify email addresses). 

2847 3. If no match and ``ENABLE_OAUTH_SIGNUP`` is enabled, create a new user 

2848 (fails if the email is already registered). 

2849 """ 

2850 return await oauth_manager.handle_callback(request, provider, response, db=db) 

2851 

2852 

2853############################ 

2854# OIDC Back-Channel Logout 

2855############################ 

2856 

2857 

2858@app.post('/oauth/backchannel-logout') 

2859async def oauth_backchannel_logout( 

2860 request: Request, 

2861 db: AsyncSession = Depends(get_async_session), 

2862): 

2863 if not ENABLE_OAUTH_BACKCHANNEL_LOGOUT: 2863 ↛ 2865line 2863 didn't jump to line 2865 because the condition on line 2863 was always true

2864 raise HTTPException(status_code=404) 

2865 return await oauth_manager.handle_backchannel_logout(request, db=db) 

2866 

2867 

2868@app.get('/manifest.json') 

2869async def get_manifest_json(): 

2870 external_pwa_manifest_url = getattr(app.state, 'EXTERNAL_PWA_MANIFEST_URL', None) 

2871 if external_pwa_manifest_url: 2871 ↛ 2876line 2871 didn't jump to line 2876 because the condition on line 2871 was never true

2872 # LICENSE covers this install-time Open WebUI branding surface, including 

2873 # names, logos, manifests, metadata, and surrounding UI. 

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

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

2876 session = await get_session() 

2877 async with session.get( 

2878 external_pwa_manifest_url, 

2879 ssl=AIOHTTP_CLIENT_SESSION_SSL, 

2880 ) as r: 

2881 r.raise_for_status() 

2882 return await r.json() 

2883 else: 

2884 # LICENSE covers this generated Open WebUI install branding surface, 

2885 # including names, logos, manifests, metadata, and surrounding UI. 

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

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

2888 return { 

2889 'name': app.state.WEBUI_NAME, 

2890 'short_name': app.state.WEBUI_NAME, 

2891 'description': f'{app.state.WEBUI_NAME} is an open, extensible, user-friendly interface for AI that adapts to your workflow.', 

2892 'start_url': '/', 

2893 'display': 'standalone', 

2894 'background_color': '#343541', 

2895 'icons': [ 

2896 # LICENSE covers this Open WebUI install icon. 

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

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

2899 { 

2900 'src': '/static/logo.png', 

2901 'type': 'image/png', 

2902 'sizes': '500x500', 

2903 'purpose': 'any', 

2904 }, 

2905 { 

2906 'src': '/static/logo.png', 

2907 'type': 'image/png', 

2908 'sizes': '500x500', 

2909 'purpose': 'maskable', 

2910 }, 

2911 ], 

2912 'share_target': { 

2913 'action': '/', 

2914 'method': 'GET', 

2915 'params': {'text': 'shared'}, 

2916 }, 

2917 } 

2918 

2919 

2920@app.get('/opensearch.xml') 

2921async def get_opensearch_xml(): 

2922 webui_url = await Config.get('webui.url') 

2923 # LICENSE covers this Open WebUI search identifier. 

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

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

2926 xml_content = rf""" 

2927 <OpenSearchDescription xmlns="http://a9.com/-/spec/opensearch/1.1/" xmlns:moz="http://www.mozilla.org/2006/browser/search/"> 

2928 <ShortName>{app.state.WEBUI_NAME}</ShortName> 

2929 <Description>Search {app.state.WEBUI_NAME}</Description> 

2930 <InputEncoding>UTF-8</InputEncoding> 

2931 <Image width="16" height="16" type="image/x-icon">{webui_url}/static/favicon.png</Image> 

2932 <Url type="text/html" method="get" template="{webui_url}/?q={'{searchTerms}'}"/> 

2933 <moz:SearchForm>{webui_url}</moz:SearchForm> 

2934 </OpenSearchDescription> 

2935 """ 

2936 return Response(content=xml_content, media_type='application/xml') 

2937 

2938 

2939def _sync_db_ping() -> None: 

2940 """Verify the database is reachable with a simple SELECT 1. 

2941 

2942 Uses a raw connection from the engine pool instead of the thread-local 

2943 ScopedSession. This is necessary because AppHTTPMiddleware 

2944 deliberately skips healthcheck paths (/health, /ready, /health/db), 

2945 so any ScopedSession opened on a healthcheck worker thread is never 

2946 rolled back or removed. If the session ever enters an invalid state 

2947 (e.g. after a transient connection error), it stays broken on that 

2948 thread permanently, causing PendingRollbackError on every subsequent 

2949 probe — exactly the failure reported in #24605. 

2950 

2951 A raw ``engine.connect()`` context manager obtains a fresh connection 

2952 from the pool, executes the ping, and deterministically returns the 

2953 connection regardless of success or failure. 

2954 """ 

2955 with engine.connect() as conn: 

2956 conn.execute(text('SELECT 1')) 

2957 

2958 

2959async def async_db_ping() -> None: 

2960 await asyncio.to_thread(_sync_db_ping) 

2961 

2962 

2963@app.get('/health') 

2964async def healthcheck(): 

2965 return {'status': True} 

2966 

2967 

2968@app.get('/ready') 

2969async def readiness_check(): 

2970 """ 

2971 Returns 200 only when the application is ready to accept traffic. 

2972 """ 

2973 

2974 # Ensure application startup work has completed 

2975 if not getattr(app.state, 'startup_complete', False): 2975 ↛ 2976line 2975 didn't jump to line 2976 because the condition on line 2975 was never true

2976 log.info('Readiness check failed: startup not complete') 

2977 raise HTTPException( 

2978 status_code=status.HTTP_503_SERVICE_UNAVAILABLE, 

2979 detail='Startup not complete', 

2980 ) 

2981 

2982 # Check database connectivity 

2983 try: 

2984 await async_db_ping() 

2985 except Exception as e: 

2986 log.warning(f'Readiness check DB ping failed: {e!r}') 

2987 raise HTTPException( 

2988 status_code=status.HTTP_503_SERVICE_UNAVAILABLE, 

2989 detail='Database not ready', 

2990 ) 

2991 

2992 # Check Redis connectivity if configured 

2993 redis = app.state.redis 

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

2995 try: 

2996 pong = await redis.ping() 

2997 if pong is False: 

2998 raise Exception('Redis PING returned False') 

2999 except Exception as e: 

3000 log.warning(f'Readiness check Redis ping failed: {e!r}') 

3001 raise HTTPException( 

3002 status_code=status.HTTP_503_SERVICE_UNAVAILABLE, 

3003 detail='Redis not ready', 

3004 ) 

3005 

3006 return {'status': True} 

3007 

3008 

3009@app.get('/health/db') 

3010async def check_db_health(): 

3011 """Verify database connectivity by issuing a lightweight ping.""" 

3012 await async_db_ping() 

3013 return {'status': True} 

3014 

3015 

3016# --- static assets & files --- 

3017# Windows registry entries can override these with text/plain, which breaks module and wasm loading 

3018mimetypes.add_type('text/javascript', '.js') 

3019mimetypes.add_type('text/javascript', '.mjs') 

3020mimetypes.add_type('application/wasm', '.wasm') 

3021 

3022# Serve build-time static assets (CSS, JS, images, favicon, etc.) 

3023app.mount('/static', StaticFiles(directory=STATIC_DIR), name='static') 

3024 

3025 

3026@app.get('/cache/{path:path}') 

3027async def serve_cache_file( 

3028 path: str, 

3029 user=Depends(get_verified_user), 

3030): 

3031 """Serve cached files (e.g. tool outputs) with path-traversal protection. 

3032 

3033 Only ``image/*``, ``audio/*``, and ``video/*`` MIME types are served inline; 

3034 everything else gets a ``Content-Disposition: attachment`` header to prevent 

3035 XSS from user-generated HTML stored in the cache directory. 

3036 """ 

3037 file_path = os.path.abspath(os.path.join(CACHE_DIR, path)) 

3038 # trailing os.sep is required: without it, a path resolving to a sibling 

3039 # whose name starts with the cache-dir basename (e.g. cache_backup) passes 

3040 cache_root = os.path.abspath(CACHE_DIR) + os.sep 

3041 if not file_path.startswith(cache_root): 3041 ↛ 3042line 3041 didn't jump to line 3042 because the condition on line 3041 was never true

3042 raise HTTPException(status_code=404, detail='File not found') 

3043 if not os.path.isfile(file_path): 3043 ↛ 3046line 3043 didn't jump to line 3046 because the condition on line 3043 was always true

3044 raise HTTPException(status_code=404, detail='File not found') 

3045 

3046 mime, _ = mimetypes.guess_type(file_path) 

3047 inline_safe = mime and mime.split('/', 1)[0] in {'image', 'audio', 'video'} 

3048 headers = {'X-Content-Type-Options': 'nosniff'} 

3049 if not inline_safe: 

3050 headers['Content-Disposition'] = f'attachment; filename="{os.path.basename(file_path)}"' 

3051 return FileResponse(file_path, headers=headers) 

3052 

3053 

3054def swagger_ui_html(*args, **kwargs): 

3055 return get_swagger_ui_html( 

3056 *args, 

3057 **kwargs, 

3058 swagger_js_url='/static/swagger-ui/swagger-ui-bundle.js', 

3059 swagger_css_url='/static/swagger-ui/swagger-ui.css', 

3060 swagger_favicon_url='/static/swagger-ui/favicon.png', 

3061 ) 

3062 

3063 

3064applications.get_swagger_ui_html = swagger_ui_html 

3065 

3066if os.path.exists(FRONTEND_BUILD_DIR): 3066 ↛ 3077line 3066 didn't jump to line 3077 because the condition on line 3066 was always true

3067 pyodide_dir = FRONTEND_BUILD_DIR / 'pyodide' 

3068 if os.path.exists(pyodide_dir): 3068 ↛ 3071line 3068 didn't jump to line 3071 because the condition on line 3068 was always true

3069 app.mount('/pyodide', CORSStaticFiles(directory=pyodide_dir), name='pyodide') 

3070 

3071 app.mount( 

3072 '/', 

3073 SPAStaticFiles(directory=FRONTEND_BUILD_DIR, html=True), 

3074 name='spa-static-files', 

3075 ) 

3076else: 

3077 log.warning(f"Frontend build directory not found at '{FRONTEND_BUILD_DIR}'. Serving API only.")