diff --git a/backend/open_webui/internal/db.py b/backend/open_webui/internal/db.py index a7326aeafc..e0d9c9ede6 100644 --- a/backend/open_webui/internal/db.py +++ b/backend/open_webui/internal/db.py @@ -232,6 +232,27 @@ def _make_async_url(url: str) -> str: return url +def _json_codec_kwargs(kwargs: dict) -> dict: + """Default an engine to JSONCodec for native ``JSON`` columns. + + Unlike ``JSONField``, those serialize through the engine, which otherwise uses + stdlib ``json``. With ``ENABLE_ORJSON`` off JSONCodec is stdlib ``json`` anyway. + """ + kwargs.setdefault('json_serializer', JSONCodec.dumps) + kwargs.setdefault('json_deserializer', JSONCodec.loads) + return kwargs + + +def _create_engine(*args, **kwargs): + """``create_engine`` with the app JSON codec wired in.""" + return create_engine(*args, **_json_codec_kwargs(kwargs)) + + +def _create_async_engine(*args, **kwargs): + """``create_async_engine`` with the app JSON codec wired in.""" + return create_async_engine(*args, **_json_codec_kwargs(kwargs)) + + # ============================================================ # SYNC ENGINE (used only for: startup migrations, config loading, # Alembic, peewee migration, health checks) @@ -260,7 +281,7 @@ if SQLALCHEMY_DATABASE_URL.startswith('sqlite+sqlcipher://'): # in the native sqlcipher3 C library. Use NullPool by default for safety, # or QueuePool if DATABASE_POOL_SIZE is explicitly configured. if isinstance(DATABASE_POOL_SIZE, int) and DATABASE_POOL_SIZE > 0: - engine = create_engine( + engine = _create_engine( 'sqlite://', creator=create_sqlcipher_connection, pool_size=DATABASE_POOL_SIZE, @@ -272,7 +293,7 @@ if SQLALCHEMY_DATABASE_URL.startswith('sqlite+sqlcipher://'): echo=False, ) else: - engine = create_engine( + engine = _create_engine( 'sqlite://', creator=create_sqlcipher_connection, poolclass=NullPool, @@ -282,7 +303,7 @@ if SQLALCHEMY_DATABASE_URL.startswith('sqlite+sqlcipher://'): log.info('Connected to encrypted SQLite database using SQLCipher') elif 'sqlite' in SQLALCHEMY_DATABASE_URL: - engine = create_engine(SQLALCHEMY_DATABASE_URL, connect_args={'check_same_thread': False}) + engine = _create_engine(SQLALCHEMY_DATABASE_URL, connect_args={'check_same_thread': False}) def _apply_sqlite_pragmas(dbapi_connection): """Apply all configured SQLite PRAGMAs to a raw DBAPI connection.""" @@ -314,7 +335,7 @@ elif 'sqlite' in SQLALCHEMY_DATABASE_URL: else: if isinstance(DATABASE_POOL_SIZE, int): if DATABASE_POOL_SIZE > 0: - engine = create_engine( + engine = _create_engine( SQLALCHEMY_DATABASE_URL, pool_size=DATABASE_POOL_SIZE, max_overflow=DATABASE_POOL_MAX_OVERFLOW, @@ -324,9 +345,9 @@ else: poolclass=QueuePool, ) else: - engine = create_engine(SQLALCHEMY_DATABASE_URL, pool_pre_ping=True, poolclass=NullPool) + engine = _create_engine(SQLALCHEMY_DATABASE_URL, pool_pre_ping=True, poolclass=NullPool) else: - engine = create_engine(SQLALCHEMY_DATABASE_URL, pool_pre_ping=True) + engine = _create_engine(SQLALCHEMY_DATABASE_URL, pool_pre_ping=True) enable_iam_token_auth(engine) @@ -373,7 +394,7 @@ if 'sqlite' in ASYNC_SQLALCHEMY_DATABASE_URL: # No pool_pre_ping: a local SQLite file cannot drop connections, and the # ping costs a worker-thread hop plus a SELECT 1 on every checkout. _sqlite_pool_size = DATABASE_POOL_SIZE if isinstance(DATABASE_POOL_SIZE, int) and DATABASE_POOL_SIZE > 0 else 512 - async_engine = create_async_engine( + async_engine = _create_async_engine( ASYNC_SQLALCHEMY_DATABASE_URL, connect_args={'check_same_thread': False}, pool_size=_sqlite_pool_size, @@ -387,7 +408,7 @@ if 'sqlite' in ASYNC_SQLALCHEMY_DATABASE_URL: else: if isinstance(DATABASE_POOL_SIZE, int): if DATABASE_POOL_SIZE > 0: - async_engine = create_async_engine( + async_engine = _create_async_engine( ASYNC_SQLALCHEMY_DATABASE_URL, pool_size=DATABASE_POOL_SIZE, max_overflow=DATABASE_POOL_MAX_OVERFLOW, @@ -396,13 +417,13 @@ else: pool_pre_ping=True, ) else: - async_engine = create_async_engine( + async_engine = _create_async_engine( ASYNC_SQLALCHEMY_DATABASE_URL, pool_pre_ping=True, poolclass=NullPool, ) else: - async_engine = create_async_engine( + async_engine = _create_async_engine( ASYNC_SQLALCHEMY_DATABASE_URL, pool_pre_ping=True, )