refac: async db

This commit is contained in:
Timothy Jaeryang Baek
2026-04-12 14:22:11 -05:00
parent b618d84065
commit 27169124f2
74 changed files with 4831 additions and 4479 deletions
+180 -163
View File
@@ -4,8 +4,9 @@ import time
from typing import Optional
import uuid
from sqlalchemy.orm import Session
from open_webui.internal.db import Base, JSONField, get_db, get_db_context
from sqlalchemy import select, delete, update, or_, func
from sqlalchemy.ext.asyncio import AsyncSession
from open_webui.internal.db import Base, JSONField, get_async_db_context
from open_webui.models.files import (
File,
@@ -27,7 +28,6 @@ from sqlalchemy import (
Text,
JSON,
UniqueConstraint,
or_,
)
log = logging.getLogger(__name__)
@@ -134,25 +134,25 @@ class KnowledgeFileListResponse(BaseModel):
class KnowledgeTable:
def _get_access_grants(self, knowledge_id: str, db: Optional[Session] = None) -> list[AccessGrantModel]:
return AccessGrants.get_grants_by_resource('knowledge', knowledge_id, db=db)
async def _get_access_grants(self, knowledge_id: str, db: Optional[AsyncSession] = None) -> list[AccessGrantModel]:
return await AccessGrants.get_grants_by_resource('knowledge', knowledge_id, db=db)
def _to_knowledge_model(
async def _to_knowledge_model(
self,
knowledge: Knowledge,
access_grants: Optional[list[AccessGrantModel]] = None,
db: Optional[Session] = None,
db: Optional[AsyncSession] = None,
) -> KnowledgeModel:
knowledge_data = KnowledgeModel.model_validate(knowledge).model_dump(exclude={'access_grants'})
knowledge_data['access_grants'] = (
access_grants if access_grants is not None else self._get_access_grants(knowledge_data['id'], db=db)
access_grants if access_grants is not None else await self._get_access_grants(knowledge_data['id'], db=db)
)
return KnowledgeModel.model_validate(knowledge_data)
def insert_new_knowledge(
self, user_id: str, form_data: KnowledgeForm, db: Optional[Session] = None
async def insert_new_knowledge(
self, user_id: str, form_data: KnowledgeForm, db: Optional[AsyncSession] = None
) -> Optional[KnowledgeModel]:
with get_db_context(db) as db:
async with get_async_db_context(db) as db:
knowledge = KnowledgeModel(
**{
**form_data.model_dump(exclude={'access_grants'}),
@@ -167,27 +167,28 @@ class KnowledgeTable:
try:
result = Knowledge(**knowledge.model_dump(exclude={'access_grants'}))
db.add(result)
db.commit()
db.refresh(result)
AccessGrants.set_access_grants('knowledge', result.id, form_data.access_grants, db=db)
await db.commit()
await db.refresh(result)
await AccessGrants.set_access_grants('knowledge', result.id, form_data.access_grants, db=db)
if result:
return self._to_knowledge_model(result, db=db)
return await self._to_knowledge_model(result, db=db)
else:
return None
except Exception:
return None
def get_knowledge_bases(
self, skip: int = 0, limit: int = 30, db: Optional[Session] = None
async def get_knowledge_bases(
self, skip: int = 0, limit: int = 30, db: Optional[AsyncSession] = None
) -> list[KnowledgeUserModel]:
with get_db_context(db) as db:
all_knowledge = db.query(Knowledge).order_by(Knowledge.updated_at.desc()).all()
async with get_async_db_context(db) as db:
result = await db.execute(select(Knowledge).order_by(Knowledge.updated_at.desc()))
all_knowledge = result.scalars().all()
user_ids = list(set(knowledge.user_id for knowledge in all_knowledge))
knowledge_ids = [knowledge.id for knowledge in all_knowledge]
users = Users.get_users_by_user_ids(user_ids, db=db) if user_ids else []
users = await Users.get_users_by_user_ids(user_ids, db=db) if user_ids else []
users_dict = {user.id: user for user in users}
grants_map = AccessGrants.get_grants_by_resources('knowledge', knowledge_ids, db=db)
grants_map = await AccessGrants.get_grants_by_resources('knowledge', knowledge_ids, db=db)
knowledge_bases = []
for knowledge in all_knowledge:
@@ -195,33 +196,33 @@ class KnowledgeTable:
knowledge_bases.append(
KnowledgeUserModel.model_validate(
{
**self._to_knowledge_model(
**(await self._to_knowledge_model(
knowledge,
access_grants=grants_map.get(knowledge.id, []),
db=db,
).model_dump(),
)).model_dump(),
'user': user.model_dump() if user else None,
}
)
)
return knowledge_bases
def search_knowledge_bases(
async def search_knowledge_bases(
self,
user_id: str,
filter: dict,
skip: int = 0,
limit: int = 30,
db: Optional[Session] = None,
db: Optional[AsyncSession] = None,
) -> KnowledgeListResponse:
try:
with get_db_context(db) as db:
query = db.query(Knowledge, User).outerjoin(User, User.id == Knowledge.user_id)
async with get_async_db_context(db) as db:
stmt = select(Knowledge, User).outerjoin(User, User.id == Knowledge.user_id)
if filter:
query_key = filter.get('query')
if query_key:
query = query.filter(
stmt = stmt.filter(
or_(
Knowledge.name.ilike(f'%{query_key}%'),
Knowledge.description.ilike(f'%{query_key}%'),
@@ -233,42 +234,46 @@ class KnowledgeTable:
view_option = filter.get('view_option')
if view_option == 'created':
query = query.filter(Knowledge.user_id == user_id)
stmt = stmt.filter(Knowledge.user_id == user_id)
elif view_option == 'shared':
query = query.filter(Knowledge.user_id != user_id)
stmt = stmt.filter(Knowledge.user_id != user_id)
query = AccessGrants.has_permission_filter(
stmt = AccessGrants.has_permission_filter(
db=db,
query=query,
query=stmt,
DocumentModel=Knowledge,
filter=filter,
resource_type='knowledge',
permission='read',
)
query = query.order_by(Knowledge.updated_at.desc(), Knowledge.id.asc())
stmt = stmt.order_by(Knowledge.updated_at.desc(), Knowledge.id.asc())
total = query.count()
count_result = await db.execute(
select(func.count()).select_from(stmt.subquery())
)
total = count_result.scalar()
if skip:
query = query.offset(skip)
stmt = stmt.offset(skip)
if limit:
query = query.limit(limit)
stmt = stmt.limit(limit)
items = query.all()
result = await db.execute(stmt)
items = result.all()
knowledge_ids = [kb.id for kb, _ in items]
grants_map = AccessGrants.get_grants_by_resources('knowledge', knowledge_ids, db=db)
grants_map = await AccessGrants.get_grants_by_resources('knowledge', knowledge_ids, db=db)
knowledge_bases = []
for knowledge_base, user in items:
knowledge_bases.append(
KnowledgeUserModel.model_validate(
{
**self._to_knowledge_model(
**(await self._to_knowledge_model(
knowledge_base,
access_grants=grants_map.get(knowledge_base.id, []),
db=db,
).model_dump(),
)).model_dump(),
'user': (UserModel.model_validate(user).model_dump() if user else None),
}
)
@@ -279,28 +284,27 @@ class KnowledgeTable:
print(e)
return KnowledgeListResponse(items=[], total=0)
def search_knowledge_files(
self, filter: dict, skip: int = 0, limit: int = 30, db: Optional[Session] = None
async def search_knowledge_files(
self, filter: dict, skip: int = 0, limit: int = 30, db: Optional[AsyncSession] = None
) -> KnowledgeFileListResponse:
"""
Scalable version: search files across all knowledge bases the user has
READ access to, without loading all KBs or using large IN() lists.
"""
try:
with get_db_context(db) as db:
async with get_async_db_context(db) as db:
# Base query: join Knowledge → KnowledgeFile → File
query = (
db.query(File, User, Knowledge)
stmt = (
select(File, User, Knowledge)
.join(KnowledgeFile, File.id == KnowledgeFile.file_id)
.join(Knowledge, KnowledgeFile.knowledge_id == Knowledge.id)
.outerjoin(User, User.id == KnowledgeFile.user_id)
)
# Apply access-control directly to the joined query
# This makes the database handle filtering, even with 10k+ KBs
query = AccessGrants.has_permission_filter(
stmt = AccessGrants.has_permission_filter(
db=db,
query=query,
query=stmt,
DocumentModel=Knowledge,
filter=filter,
resource_type='knowledge',
@@ -311,20 +315,24 @@ class KnowledgeTable:
if filter:
q = filter.get('query')
if q:
query = query.filter(File.filename.ilike(f'%{q}%'))
stmt = stmt.filter(File.filename.ilike(f'%{q}%'))
# Order by file changes
query = query.order_by(File.updated_at.desc(), File.id.asc())
stmt = stmt.order_by(File.updated_at.desc(), File.id.asc())
# Count before pagination
total = query.count()
count_result = await db.execute(
select(func.count()).select_from(stmt.subquery())
)
total = count_result.scalar()
if skip:
query = query.offset(skip)
stmt = stmt.offset(skip)
if limit:
query = query.limit(limit)
stmt = stmt.limit(limit)
rows = query.all()
result = await db.execute(stmt)
rows = result.all()
items = []
for file, user, knowledge in rows:
@@ -332,7 +340,7 @@ class KnowledgeTable:
FileUserResponse(
**FileModel.model_validate(file).model_dump(),
user=(UserResponse(**UserModel.model_validate(user).model_dump()) if user else None),
collection=self._to_knowledge_model(knowledge, db=db).model_dump(),
collection=(await self._to_knowledge_model(knowledge, db=db)).model_dump(),
)
)
@@ -342,14 +350,15 @@ class KnowledgeTable:
print('search_knowledge_files error:', e)
return KnowledgeFileListResponse(items=[], total=0)
def check_access_by_user_id(self, id, user_id, permission='write', db: Optional[Session] = None) -> bool:
knowledge = self.get_knowledge_by_id(id, db=db)
async def check_access_by_user_id(self, id, user_id, permission='write', db: Optional[AsyncSession] = None) -> bool:
knowledge = await self.get_knowledge_by_id(id, db=db)
if not knowledge:
return False
if knowledge.user_id == user_id:
return True
user_group_ids = {group.id for group in Groups.get_groups_by_member_id(user_id, db=db)}
return AccessGrants.has_access(
user_groups = await Groups.get_groups_by_member_id(user_id, db=db)
user_group_ids = {group.id for group in user_groups}
return await AccessGrants.has_access(
user_id=user_id,
resource_type='knowledge',
resource_id=knowledge.id,
@@ -358,45 +367,50 @@ class KnowledgeTable:
db=db,
)
def get_knowledge_bases_by_user_id(
self, user_id: str, permission: str = 'write', db: Optional[Session] = None
async def get_knowledge_bases_by_user_id(
self, user_id: str, permission: str = 'write', db: Optional[AsyncSession] = None
) -> list[KnowledgeUserModel]:
knowledge_bases = self.get_knowledge_bases(db=db)
user_group_ids = {group.id for group in Groups.get_groups_by_member_id(user_id, db=db)}
return [
knowledge_base
for knowledge_base in knowledge_bases
if knowledge_base.user_id == user_id
or AccessGrants.has_access(
knowledge_bases = await self.get_knowledge_bases(db=db)
user_groups = await Groups.get_groups_by_member_id(user_id, db=db)
user_group_ids = {group.id for group in user_groups}
result = []
for knowledge_base in knowledge_bases:
if knowledge_base.user_id == user_id:
result.append(knowledge_base)
elif await AccessGrants.has_access(
user_id=user_id,
resource_type='knowledge',
resource_id=knowledge_base.id,
permission=permission,
user_group_ids=user_group_ids,
db=db,
)
]
):
result.append(knowledge_base)
return result
def get_knowledge_by_id(self, id: str, db: Optional[Session] = None) -> Optional[KnowledgeModel]:
async def get_knowledge_by_id(self, id: str, db: Optional[AsyncSession] = None) -> Optional[KnowledgeModel]:
try:
with get_db_context(db) as db:
knowledge = db.query(Knowledge).filter_by(id=id).first()
return self._to_knowledge_model(knowledge, db=db) if knowledge else None
async with get_async_db_context(db) as db:
result = await db.execute(select(Knowledge).filter_by(id=id))
knowledge = result.scalars().first()
return await self._to_knowledge_model(knowledge, db=db) if knowledge else None
except Exception:
return None
def get_knowledge_by_id_and_user_id(
self, id: str, user_id: str, db: Optional[Session] = None
async def get_knowledge_by_id_and_user_id(
self, id: str, user_id: str, db: Optional[AsyncSession] = None
) -> Optional[KnowledgeModel]:
knowledge = self.get_knowledge_by_id(id, db=db)
knowledge = await self.get_knowledge_by_id(id, db=db)
if not knowledge:
return None
if knowledge.user_id == user_id:
return knowledge
user_group_ids = {group.id for group in Groups.get_groups_by_member_id(user_id, db=db)}
if AccessGrants.has_access(
user_groups = await Groups.get_groups_by_member_id(user_id, db=db)
user_group_ids = {group.id for group in user_groups}
if await AccessGrants.has_access(
user_id=user_id,
resource_type='knowledge',
resource_id=knowledge.id,
@@ -407,19 +421,19 @@ class KnowledgeTable:
return knowledge
return None
def get_knowledges_by_file_id(self, file_id: str, db: Optional[Session] = None) -> list[KnowledgeModel]:
async def get_knowledges_by_file_id(self, file_id: str, db: Optional[AsyncSession] = None) -> list[KnowledgeModel]:
try:
with get_db_context(db) as db:
knowledges = (
db.query(Knowledge)
async with get_async_db_context(db) as db:
result = await db.execute(
select(Knowledge)
.join(KnowledgeFile, Knowledge.id == KnowledgeFile.knowledge_id)
.filter(KnowledgeFile.file_id == file_id)
.all()
)
knowledges = result.scalars().all()
knowledge_ids = [k.id for k in knowledges]
grants_map = AccessGrants.get_grants_by_resources('knowledge', knowledge_ids, db=db)
grants_map = await AccessGrants.get_grants_by_resources('knowledge', knowledge_ids, db=db)
return [
self._to_knowledge_model(
await self._to_knowledge_model(
knowledge,
access_grants=grants_map.get(knowledge.id, []),
db=db,
@@ -429,19 +443,19 @@ class KnowledgeTable:
except Exception:
return []
def search_files_by_id(
async def search_files_by_id(
self,
knowledge_id: str,
user_id: str,
filter: dict,
skip: int = 0,
limit: int = 30,
db: Optional[Session] = None,
db: Optional[AsyncSession] = None,
) -> KnowledgeFileListResponse:
try:
with get_db_context(db) as db:
query = (
db.query(File, User)
async with get_async_db_context(db) as db:
stmt = (
select(File, User)
.join(KnowledgeFile, File.id == KnowledgeFile.file_id)
.outerjoin(User, User.id == KnowledgeFile.user_id)
.filter(KnowledgeFile.knowledge_id == knowledge_id)
@@ -453,13 +467,13 @@ class KnowledgeTable:
if filter:
query_key = filter.get('query')
if query_key:
query = query.filter(or_(File.filename.ilike(f'%{query_key}%')))
stmt = stmt.filter(or_(File.filename.ilike(f'%{query_key}%')))
view_option = filter.get('view_option')
if view_option == 'created':
query = query.filter(KnowledgeFile.user_id == user_id)
stmt = stmt.filter(KnowledgeFile.user_id == user_id)
elif view_option == 'shared':
query = query.filter(KnowledgeFile.user_id != user_id)
stmt = stmt.filter(KnowledgeFile.user_id != user_id)
order_by = filter.get('order_by')
direction = filter.get('direction')
@@ -473,17 +487,21 @@ class KnowledgeTable:
primary_sort = File.updated_at.asc() if is_asc else File.updated_at.desc()
# Apply sort with secondary key for deterministic pagination
query = query.order_by(primary_sort, File.id.asc())
stmt = stmt.order_by(primary_sort, File.id.asc())
# Count BEFORE pagination
total = query.count()
count_result = await db.execute(
select(func.count()).select_from(stmt.subquery())
)
total = count_result.scalar()
if skip:
query = query.offset(skip)
stmt = stmt.offset(skip)
if limit:
query = query.limit(limit)
stmt = stmt.limit(limit)
items = query.all()
result = await db.execute(stmt)
items = result.all()
files = []
for file, user in items:
@@ -499,35 +517,34 @@ class KnowledgeTable:
print(e)
return KnowledgeFileListResponse(items=[], total=0)
def get_files_by_id(self, knowledge_id: str, db: Optional[Session] = None) -> list[FileModel]:
async def get_files_by_id(self, knowledge_id: str, db: Optional[AsyncSession] = None) -> list[FileModel]:
try:
with get_db_context(db) as db:
files = (
db.query(File)
async with get_async_db_context(db) as db:
result = await db.execute(
select(File)
.join(KnowledgeFile, File.id == KnowledgeFile.file_id)
.filter(KnowledgeFile.knowledge_id == knowledge_id)
.all()
)
files = result.scalars().all()
return [FileModel.model_validate(file) for file in files]
except Exception:
return []
def get_file_metadatas_by_id(self, knowledge_id: str, db: Optional[Session] = None) -> list[FileMetadataResponse]:
async def get_file_metadatas_by_id(self, knowledge_id: str, db: Optional[AsyncSession] = None) -> list[FileMetadataResponse]:
try:
with get_db_context(db) as db:
files = self.get_files_by_id(knowledge_id, db=db)
return [FileMetadataResponse(**file.model_dump()) for file in files]
files = await self.get_files_by_id(knowledge_id, db=db)
return [FileMetadataResponse(**file.model_dump()) for file in files]
except Exception:
return []
def add_file_to_knowledge_by_id(
async def add_file_to_knowledge_by_id(
self,
knowledge_id: str,
file_id: str,
user_id: str,
db: Optional[Session] = None,
db: Optional[AsyncSession] = None,
) -> Optional[KnowledgeFileModel]:
with get_db_context(db) as db:
async with get_async_db_context(db) as db:
knowledge_file = KnowledgeFileModel(
**{
'id': str(uuid.uuid4()),
@@ -542,8 +559,8 @@ class KnowledgeTable:
try:
result = KnowledgeFile(**knowledge_file.model_dump())
db.add(result)
db.commit()
db.refresh(result)
await db.commit()
await db.refresh(result)
if result:
return KnowledgeFileModel.model_validate(result)
else:
@@ -551,103 +568,103 @@ class KnowledgeTable:
except Exception:
return None
def has_file(self, knowledge_id: str, file_id: str, db: Optional[Session] = None) -> bool:
async def has_file(self, knowledge_id: str, file_id: str, db: Optional[AsyncSession] = None) -> bool:
"""Check whether a file belongs to a knowledge base."""
try:
with get_db_context(db) as db:
return db.query(KnowledgeFile).filter_by(knowledge_id=knowledge_id, file_id=file_id).first() is not None
async with get_async_db_context(db) as db:
result = await db.execute(
select(KnowledgeFile).filter_by(knowledge_id=knowledge_id, file_id=file_id).limit(1)
)
return result.scalars().first() is not None
except Exception:
return False
def remove_file_from_knowledge_by_id(self, knowledge_id: str, file_id: str, db: Optional[Session] = None) -> bool:
async def remove_file_from_knowledge_by_id(self, knowledge_id: str, file_id: str, db: Optional[AsyncSession] = None) -> bool:
try:
with get_db_context(db) as db:
db.query(KnowledgeFile).filter_by(knowledge_id=knowledge_id, file_id=file_id).delete()
db.commit()
async with get_async_db_context(db) as db:
await db.execute(delete(KnowledgeFile).filter_by(knowledge_id=knowledge_id, file_id=file_id))
await db.commit()
return True
except Exception:
return False
def reset_knowledge_by_id(self, id: str, db: Optional[Session] = None) -> Optional[KnowledgeModel]:
async def reset_knowledge_by_id(self, id: str, db: Optional[AsyncSession] = None) -> Optional[KnowledgeModel]:
try:
with get_db_context(db) as db:
async with get_async_db_context(db) as db:
# Delete all knowledge_file entries for this knowledge_id
db.query(KnowledgeFile).filter_by(knowledge_id=id).delete()
db.commit()
await db.execute(delete(KnowledgeFile).filter_by(knowledge_id=id))
await db.commit()
# Update the knowledge entry's updated_at timestamp
db.query(Knowledge).filter_by(id=id).update(
{
'updated_at': int(time.time()),
}
await db.execute(
update(Knowledge).filter_by(id=id).values(updated_at=int(time.time()))
)
db.commit()
await db.commit()
return self.get_knowledge_by_id(id=id, db=db)
return await self.get_knowledge_by_id(id=id, db=db)
except Exception as e:
log.exception(e)
return None
def update_knowledge_by_id(
async def update_knowledge_by_id(
self,
id: str,
form_data: KnowledgeForm,
overwrite: bool = False,
db: Optional[Session] = None,
db: Optional[AsyncSession] = None,
) -> Optional[KnowledgeModel]:
try:
with get_db_context(db) as db:
knowledge = self.get_knowledge_by_id(id=id, db=db)
db.query(Knowledge).filter_by(id=id).update(
{
async with get_async_db_context(db) as db:
await db.execute(
update(Knowledge).filter_by(id=id).values(
**form_data.model_dump(exclude={'access_grants'}),
'updated_at': int(time.time()),
}
updated_at=int(time.time()),
)
)
db.commit()
await db.commit()
if form_data.access_grants is not None:
AccessGrants.set_access_grants('knowledge', id, form_data.access_grants, db=db)
return self.get_knowledge_by_id(id=id, db=db)
await AccessGrants.set_access_grants('knowledge', id, form_data.access_grants, db=db)
return await self.get_knowledge_by_id(id=id, db=db)
except Exception as e:
log.exception(e)
return None
def update_knowledge_data_by_id(
self, id: str, data: dict, db: Optional[Session] = None
async def update_knowledge_data_by_id(
self, id: str, data: dict, db: Optional[AsyncSession] = None
) -> Optional[KnowledgeModel]:
try:
with get_db_context(db) as db:
knowledge = self.get_knowledge_by_id(id=id, db=db)
db.query(Knowledge).filter_by(id=id).update(
{
'data': data,
'updated_at': int(time.time()),
}
async with get_async_db_context(db) as db:
await db.execute(
update(Knowledge).filter_by(id=id).values(
data=data,
updated_at=int(time.time()),
)
)
db.commit()
return self.get_knowledge_by_id(id=id, db=db)
await db.commit()
return await self.get_knowledge_by_id(id=id, db=db)
except Exception as e:
log.exception(e)
return None
def delete_knowledge_by_id(self, id: str, db: Optional[Session] = None) -> bool:
async def delete_knowledge_by_id(self, id: str, db: Optional[AsyncSession] = None) -> bool:
try:
with get_db_context(db) as db:
AccessGrants.revoke_all_access('knowledge', id, db=db)
db.query(Knowledge).filter_by(id=id).delete()
db.commit()
async with get_async_db_context(db) as db:
await AccessGrants.revoke_all_access('knowledge', id, db=db)
await db.execute(delete(Knowledge).filter_by(id=id))
await db.commit()
return True
except Exception:
return False
def delete_all_knowledge(self, db: Optional[Session] = None) -> bool:
with get_db_context(db) as db:
async def delete_all_knowledge(self, db: Optional[AsyncSession] = None) -> bool:
async with get_async_db_context(db) as db:
try:
knowledge_ids = [row[0] for row in db.query(Knowledge.id).all()]
result = await db.execute(select(Knowledge.id))
knowledge_ids = [row[0] for row in result.all()]
for knowledge_id in knowledge_ids:
AccessGrants.revoke_all_access('knowledge', knowledge_id, db=db)
db.query(Knowledge).delete()
db.commit()
await AccessGrants.revoke_all_access('knowledge', knowledge_id, db=db)
await db.execute(delete(Knowledge))
await db.commit()
return True
except Exception: