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
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 05:07 +0000
1from __future__ import annotations
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
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
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
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
284logging.basicConfig(stream=sys.stdout, level=GLOBAL_LOG_LEVEL)
285log = logging.getLogger(__name__)
288async def emit_chat_list_event(metadata: dict, chat_id: str):
289 if not is_saved_chat_id(chat_id):
290 return
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}})
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
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
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 ╚═════╝ ╚═╝ ╚══════╝╚═╝ ╚═══╝ ╚══╝╚══╝ ╚══════╝╚═════╝ ╚═════╝ ╚═╝
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')
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()
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 )
360 app.state.instance_id = INSTANCE_ID
361 start_logger()
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()
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')
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))
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})
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()
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()
390 app.state.redis = get_redis_client(async_mode=True)
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))
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())
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())
403 from open_webui.utils.automations import scheduler_worker_loop
405 app.state.scheduler_worker_loop = asyncio.create_task(scheduler_worker_loop(app))
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}')
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 )
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}')
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}')
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}')
471 app.state.startup_complete = True
472 await publish_event(app, EVENTS.SYSTEM_STARTUP_COMPLETED, source='system')
474 yield
476 await publish_event(app, EVENTS.SYSTEM_SHUTDOWN_STARTED, source='system')
478 # Shutdown: clean up shared resources
479 from open_webui.utils.session_pool import close_session
481 await close_session()
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()
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()
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()
492 app.state.periodic_usage_pool_cleanup.cancel()
493 app.state.periodic_session_pool_cleanup.cancel()
494 app.state.scheduler_worker_loop.cancel()
496 await publish_event(app, EVENTS.SYSTEM_SHUTDOWN_COMPLETED, source='system')
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()
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)
515@app.exception_handler(RecurrenceEvaluationTimeout)
516async def recurrence_timeout_handler(request: Request, exc: RecurrenceEvaluationTimeout):
517 return JSONResponse(status_code=400, content={'detail': str(exc)})
520# Used by readiness checks to gate traffic until startup work is done.
521app.state.startup_complete = False
523# For Open WebUI OIDC/OAuth2
524oauth_manager = OAuthManager(app)
525app.state.oauth_manager = oauth_manager
527# For Integrations
528oauth_client_manager = OAuthClientManager(app)
529app.state.oauth_client_manager = oauth_client_manager
531app.state.instance_id = None
532app.state.redis = None
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
544########################################
545#
546# OPENTELEMETRY
547#
548########################################
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
553 setup_opentelemetry(app=app, db_engine=engine)
556########################################
557#
558# OLLAMA
559#
560########################################
563app.state.OLLAMA_MODELS = {}
565########################################
566#
567# OPENAI
568#
569########################################
572app.state.OPENAI_MODELS = {}
574########################################
575#
576# TOOL SERVERS
577#
578########################################
580app.state.TOOL_SERVERS = []
582########################################
583#
584# TERMINAL SERVER
585#
586########################################
588app.state.TERMINAL_SERVERS = []
590########################################
591#
592# DIRECT CONNECTIONS
593#
594########################################
597########################################
598#
599# SCIM
600#
601########################################
603app.state.ENABLE_SCIM = ENABLE_SCIM
604app.state.SCIM_TOKEN = SCIM_TOKEN
606########################################
607#
608# MODELS
609#
610########################################
612app.state.BASE_MODELS = []
614########################################
615#
616# WEBUI
617#
618########################################
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
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})
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')
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 )
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})
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
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
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 )
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 )
755########################################
756#
757# CODE EXECUTION
758#
759########################################
762########################################
763#
764# IMAGES
765#
766########################################
769########################################
770#
771# AUDIO
772#
773########################################
776app.state.faster_whisper_model = None
777app.state.speech_synthesiser = None
778app.state.speech_speaker_embeddings_dataset = None
781########################################
782#
783# TASKS
784#
785########################################
788########################################
789#
790# WEBUI
791#
792########################################
794app.state.MODELS = MODELS
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
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 )
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)
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)
831app.add_middleware(
832 CORSMiddleware,
833 allow_origins=CORS_ALLOW_ORIGIN,
834 allow_credentials=True,
835 allow_methods=['*'],
836 allow_headers=['*'],
837)
840app.mount('/ws', socket_app)
843app.include_router(ollama.router, prefix='/ollama', tags=['ollama'])
844app.include_router(openai.router, prefix='/openai', tags=['openai'])
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'])
851app.include_router(audio.router, prefix='/api/v1/audio', tags=['audio'])
852app.include_router(retrieval.router, prefix='/api/v1/retrieval', tags=['retrieval'])
854app.include_router(configs.router, prefix='/api/v1/configs', tags=['configs'])
856app.include_router(auths.router, prefix='/api/v1/auths', tags=['auths'])
857app.include_router(users.router, prefix='/api/v1/users', tags=['users'])
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'])
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'])
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'])
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'])
890##################################
891#
892# Chat Endpoints
893#
894##################################
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)
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 ]
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())
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)
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 {}
919 # Remove profile image URL to reduce payload size
920 meta.pop('profile_image_url', None)
922 if 'tags' in meta:
923 meta['tags'] = normalize_model_tags(meta['tags'])
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())
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 )
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}
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}
952class ModelUnloadForm(BaseModel):
953 model: str
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
965 ollama_models = getattr(request.app.state, 'OLLAMA_MODELS', None) or {}
966 openai_models = getattr(request.app.state, 'OPENAI_MODELS', None) or {}
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
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)
991 prefix_id = api_config.get('prefix_id', None)
992 actual_model = strip_provider_model_prefix(model_id, prefix_id)
994 payload = JSONCodec.dumps({'model': actual_model, 'keep_alive': 0, 'prompt': ''})
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)})
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}
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 ''
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 )
1064 raise HTTPException(status_code=404, detail=f'Model "{model_id}" not found')
1067##################################
1068# Embeddings
1069##################################
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.
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.
1082 Args:
1083 request (Request): Request context.
1084 form_data (dict): OpenAI-like payload (e.g., {"model": "...", "input": [...]})
1085 user (UserModel): Authenticated user.
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)
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
1103 model_meta['knowledge'] = await get_accessible_folder_files(knowledge_items, user)
1104 request.state.direct = True
1105 request.state.model = model_item
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)
1118 model_id = form_data.get('model', None)
1119 model_item = form_data.pop('model_item', {})
1120 tasks = form_data.pop('background_tasks', None)
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')
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 )
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)
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)
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 {}
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)
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']
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')
1192 # Model Params
1193 if model_info_params.get('stream_response') is not None:
1194 form_data['stream'] = model_info_params.get('stream_response')
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}
1200 if model_info_params.get('stream_delta_chunk_size'):
1201 stream_delta_chunk_size = model_info_params.get('stream_delta_chunk_size')
1203 if model_info_params.get('reasoning_tags') is not None:
1204 reasoning_tags = model_info_params.get('reasoning_tags')
1206 if model_info_params.get('compact_token_threshold') is not None:
1207 compact_token_threshold = model_info_params.get('compact_token_threshold')
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
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)}]
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 {}
1238 chat_variables = normalize_chat_variables(chat_variables)
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
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 )
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 }
1299 if is_new_chat:
1300 metadata['chat_id'] = str(uuid4())
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)
1306 if metadata.get('chat_id') and user:
1307 chat_id = metadata['chat_id']
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 )
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 )
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
1366 history_messages = {}
1367 all_assistant_ids = [entry['message_id'] for entry in message_ids if entry.get('message_id')]
1369 if user_message_id and user_message:
1370 user_message['childrenIds'] = all_assistant_ids
1371 history_messages[user_message_id] = user_message
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
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 )
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
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 }
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')
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 )
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')]
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)
1518 await Chats.update_chat_variables_by_id(chat_id, chat_variables)
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
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)
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 )
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
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')]
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 )
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 )
1635 request.state.metadata = metadata
1636 form_data['metadata'] = metadata
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 )
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)
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}
1657 response = await chat_completion_handler(request, form_data, user)
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))
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)
1672 return await process_chat_response(response, ctx)
1673 except asyncio.CancelledError:
1674 log.info('Chat processing was cancelled')
1675 try:
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'})
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 )
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 )
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)
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
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
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'))
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']
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
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 }
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 }
1829 # Resolve the model object for this specific model
1830 resolved_model = request.app.state.MODELS.get(target_model_id, model)
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
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)
1859 if is_internal:
1860 return {
1861 'status': True,
1862 'task_ids': [],
1863 'chat_id': chat_id,
1864 'results': subagent_results,
1865 }
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}})
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)
1888# Alias for chat_completion (Legacy)
1889generate_chat_completions = chat_completion
1890generate_chat_completion = chat_completion
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 }
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
1918##################################
1919#
1920# Anthropic Messages API Compatible Endpoint
1921#
1922##################################
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)
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)}
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
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 )
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 )
1971 try:
1972 response_data = await response.json()
1973 except Exception:
1974 response_data = await response.text()
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)
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)
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.
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.
2016 Supports both streaming and non-streaming requests.
2017 All models configured in Open WebUI are accessible via this endpoint.
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)
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
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 []
2047 # Convert Anthropic payload to OpenAI format
2048 openai_payload = convert_anthropic_to_openai_payload(form_data, passthrough_params)
2050 # Route through the existing chat_completion handler
2051 response = await chat_completion(request, openai_payload, user)
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
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
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 )
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 )
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)
2096 try:
2097 model_item = form_data.pop('model_item', {})
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)
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 )
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)
2114 try:
2115 model_item = form_data.pop('model_item', {})
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)
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 )
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))
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)}
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': []}
2154 task_ids = await list_task_ids_by_item_id(request.app.state.redis, chat_id)
2156 log.debug('Task IDs for chat %s: %s', chat_id, task_ids)
2157 return {'task_ids': task_ids}
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)
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
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)
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 }
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'})
2214 return result
2217##################################
2218#
2219# Config Endpoints
2220#
2221##################################
2224@app.get('/api/config')
2225async def get_app_config(request: Request):
2226 user = None
2227 token = None
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
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')
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'])
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()
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 )
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 }
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
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
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 }
2499@app.get('/api/events/webhooks')
2500async def get_event_webhooks_api(user=Depends(get_admin_user)):
2501 return await get_event_webhooks()
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))
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
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')
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))
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
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')
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}
2586@app.get('/api/version')
2587async def get_app_version():
2588 return {
2589 'version': VERSION,
2590 'deployment_id': DEPLOYMENT_ID,
2591 }
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']
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}
2616@app.get('/api/changelog')
2617async def get_app_changelog():
2618 return {key: CHANGELOG[key] for idx, key in enumerate(CHANGELOG) if idx < 5}
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 )
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')
2646# --- OAuth Login & Callback ---
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 )
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 )
2677async def register_client(request, client_id: str) -> bool:
2678 server_type, server_id = client_id.split(':', 1)
2680 connection = None
2681 connection_idx = None
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
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
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')
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
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
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
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)
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 )
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 )
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 )
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 )
2812 return await oauth_client_manager.handle_authorize(request, client_id=client_id, user_id=user.id)
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 )
2828@app.get('/oauth/{provider}/login')
2829async def oauth_login(provider: str, request: Request):
2830 return await oauth_manager.handle_login(request, provider)
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.
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)
2853############################
2854# OIDC Back-Channel Logout
2855############################
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)
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 }
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')
2939def _sync_db_ping() -> None:
2940 """Verify the database is reachable with a simple SELECT 1.
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.
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'))
2959async def async_db_ping() -> None:
2960 await asyncio.to_thread(_sync_db_ping)
2963@app.get('/health')
2964async def healthcheck():
2965 return {'status': True}
2968@app.get('/ready')
2969async def readiness_check():
2970 """
2971 Returns 200 only when the application is ready to accept traffic.
2972 """
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 )
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 )
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 )
3006 return {'status': True}
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}
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')
3022# Serve build-time static assets (CSS, JS, images, favicon, etc.)
3023app.mount('/static', StaticFiles(directory=STATIC_DIR), name='static')
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.
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')
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)
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 )
3064applications.get_swagger_ui_html = swagger_ui_html
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')
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.")