diff --git a/backend/open_webui/models/chats.py b/backend/open_webui/models/chats.py index 957492d817..571a3d8378 100644 --- a/backend/open_webui/models/chats.py +++ b/backend/open_webui/models/chats.py @@ -569,7 +569,18 @@ class ChatTable: **message, } else: - history['messages'][message_id] = message + now = int(time.time()) + # This upsert is also used for partial streaming/final updates. + # If a concurrent whole-chat write dropped the assistant placeholder, + # never persist the partial payload as a malformed history node. + history['messages'][message_id] = { + 'id': message_id, + 'parentId': message.get('parentId'), + 'childrenIds': message.get('childrenIds', []), + 'role': message.get('role', 'assistant'), + 'timestamp': message.get('timestamp', now), + **message, + } history['currentId'] = message_id diff --git a/backend/open_webui/utils/middleware.py b/backend/open_webui/utils/middleware.py index 56226fc226..617f6983ec 100644 --- a/backend/open_webui/utils/middleware.py +++ b/backend/open_webui/utils/middleware.py @@ -3582,6 +3582,17 @@ async def streaming_chat_response_handler(response, ctx): task_id = str(uuid4()) # Create a unique task ID. model_id = form_data.get('model', '') + def build_assistant_message_update(**fields): + return { + 'id': metadata.get('message_id'), + 'parentId': metadata.get('user_message_id'), + 'childrenIds': [], + 'role': 'assistant', + 'model': model_id, + 'timestamp': int(time.time()), + **fields, + } + # Handle as a background task async def response_handler(response, events): def tag_output_handler(content_type, tags, output): @@ -5057,24 +5068,24 @@ async def streaming_chat_response_handler(response, ctx): await Chats.upsert_message_to_chat_by_id_and_message_id( metadata['chat_id'], metadata['message_id'], - { - 'done': True, - 'content': serialize_output(output), - 'output': output, + build_assistant_message_update( + done=True, + content=serialize_output(output), + output=output, **({'usage': usage} if usage else {}), - }, + ), ) elif usage: await Chats.upsert_message_to_chat_by_id_and_message_id( metadata['chat_id'], metadata['message_id'], - {'done': True, 'usage': usage}, + build_assistant_message_update(done=True, usage=usage), ) else: await Chats.upsert_message_to_chat_by_id_and_message_id( metadata['chat_id'], metadata['message_id'], - {'done': True}, + build_assistant_message_update(done=True), ) # Send a webhook notification if the user is not active @@ -5127,17 +5138,17 @@ async def streaming_chat_response_handler(response, ctx): await Chats.upsert_message_to_chat_by_id_and_message_id( metadata['chat_id'], metadata['message_id'], - { - 'done': True, - 'content': serialize_output(output), - 'output': output, - }, + build_assistant_message_update( + done=True, + content=serialize_output(output), + output=output, + ), ) else: await Chats.upsert_message_to_chat_by_id_and_message_id( metadata['chat_id'], metadata['message_id'], - {'done': True}, + build_assistant_message_update(done=True), ) try: