diff --git a/backend/open_webui/models/chats.py b/backend/open_webui/models/chats.py index 258d5b7cc6..d25400dbae 100644 --- a/backend/open_webui/models/chats.py +++ b/backend/open_webui/models/chats.py @@ -6,6 +6,7 @@ import logging import re import time import uuid +from typing import Literal # local imports from open_webui.internal.db import Base, JSONField, get_async_db_context @@ -1047,6 +1048,23 @@ class ChatTable: return chat.chat.get('history', {}).get('messages', {}).get(message_id, {}) + async def get_message_list_field( + self, id: str, message_id: str, field: Literal['files', 'sources', 'embeds'] + ) -> list[dict]: + # Read the column, not the row: the write path commits before validating, so rows the model rejects exist. + async with get_async_db_context() as db: + result = await db.execute(select(getattr(ChatMessage, field)).where(ChatMessage.id == f'{id}-{message_id}')) + row = result.first() + if row is not None: + return row[0] or [] + + # Legacy chats have no chat_message rows; fall back to the embedded history. + chat = await self.get_chat_by_id(id) + if chat is None: + return [] + stored = chat.chat.get('history', {}).get('messages', {}).get(message_id, {}) + return stored.get(field) or [] + async def upsert_message_to_chat_by_id_and_message_id( self, id: str, message_id: str, message: dict, *, touch: bool = True ) -> ChatModel | None: diff --git a/backend/open_webui/socket/main.py b/backend/open_webui/socket/main.py index 9224501767..b6f8db0394 100644 --- a/backend/open_webui/socket/main.py +++ b/backend/open_webui/socket/main.py @@ -1073,11 +1073,11 @@ async def get_event_emitter(request_info, update_db=True): embeds = event_payload.get('embeds', []) if not event_payload.get('replace', False): - message = await Chats.get_message_by_id_and_message_id( - request_info['chat_id'], - request_info['message_id'], + embeds.extend( + await Chats.get_message_list_field( + request_info['chat_id'], request_info['message_id'], 'embeds' + ) ) - embeds.extend(message.get('embeds', [])) await Chats.upsert_message_to_chat_by_id_and_message_id( request_info['chat_id'], @@ -1089,13 +1089,10 @@ async def get_event_emitter(request_info, update_db=True): ) elif event_type == 'files': - message = await Chats.get_message_by_id_and_message_id( - request_info['chat_id'], - request_info['message_id'], - ) - files = event_data.get('data', {}).get('files', []) - files.extend(message.get('files', [])) + files.extend( + await Chats.get_message_list_field(request_info['chat_id'], request_info['message_id'], 'files') + ) await Chats.upsert_message_to_chat_by_id_and_message_id( request_info['chat_id'], @@ -1109,12 +1106,9 @@ async def get_event_emitter(request_info, update_db=True): elif event_type in ('source', 'citation'): data = event_data.get('data', {}) if data.get('type') is None: - message = await Chats.get_message_by_id_and_message_id( - request_info['chat_id'], - request_info['message_id'], + sources = await Chats.get_message_list_field( + request_info['chat_id'], request_info['message_id'], 'sources' ) - - sources = message.get('sources', []) sources.append(data) await Chats.upsert_message_to_chat_by_id_and_message_id(