From 3b61562c82448cf83710d8b6ed29b797991aa83a Mon Sep 17 00:00:00 2001 From: Timothy Jaeryang Baek Date: Fri, 13 Feb 2026 17:26:54 -0600 Subject: [PATCH] refac --- backend/open_webui/utils/middleware.py | 277 +++++++++++++++---------- backend/open_webui/utils/misc.py | 32 ++- 2 files changed, 185 insertions(+), 124 deletions(-) diff --git a/backend/open_webui/utils/middleware.py b/backend/open_webui/utils/middleware.py index e39787f71..1fa562e0d 100644 --- a/backend/open_webui/utils/middleware.py +++ b/backend/open_webui/utils/middleware.py @@ -411,6 +411,34 @@ def serialize_output(output: list) -> str: if content and not content.endswith("\n"): content += "\n" + # Render the code_interpreter item as a
block + # so the frontend Collapsible renders "Analyzing..."/"Analyzed". + code = item.get("code", "").strip() + lang = item.get("lang", "python") + status = item.get("status", "in_progress") + duration = item.get("duration") + is_last_item = idx == len(output) - 1 + + # Build inner content: code block + display = "" + if code: + display = f"```{lang}\n{code}\n```" + + # Build output attribute as HTML-escaped JSON for CodeBlock.svelte + ci_output = item.get("output") + output_attr = "" + if ci_output: + if isinstance(ci_output, dict): + output_json = json.dumps(ci_output, ensure_ascii=False) + else: + output_json = json.dumps({"result": str(ci_output)}, ensure_ascii=False) + output_attr = f' output="{html.escape(output_json)}"' + + if status == "completed" or duration is not None or not is_last_item: + content += f'
\nAnalyzed\n{display}\n
\n' + else: + content += f'
\nAnalyzing…\n{display}\n
\n' + return content.strip() @@ -2930,11 +2958,14 @@ async def streaming_chat_response_handler(response, ctx): # Handle as a background task async def response_handler(response, events): - def tag_output_handler(content_type, tags, content, output): + def tag_output_handler(content_type, tags, output): """ Detect special tags (reasoning, solution, code_interpreter) in streaming content and create corresponding OR-aligned output items directly. Operates on output items instead of content_blocks. + + Uses the text from the output items themselves for tag detection, + eliminating state divergence between accumulated content and items. """ end_flag = False @@ -2974,6 +3005,8 @@ async def streaming_chat_response_handler(response, ctx): last_type = output[-1].get("type", "") if output else "" if last_type == "message": + # Use the output item's own text for tag detection + item_text = get_last_text(output) for start_tag, end_tag in tags: start_tag_pattern = rf"{re.escape(start_tag)}" @@ -2982,7 +3015,7 @@ async def streaming_chat_response_handler(response, ctx): rf"<{re.escape(start_tag[1:-1])}(\s.*?)?>" ) - match = re.search(start_tag_pattern, content) + match = re.search(start_tag_pattern, item_text) if match: try: attr_content = match.group(1) if match.group(1) else "" @@ -2991,20 +3024,13 @@ async def streaming_chat_response_handler(response, ctx): attributes = extract_attributes(attr_content) - before_tag = content[: match.start()] - after_tag = content[match.end() :] + before_tag = item_text[: match.start()] + after_tag = item_text[match.end() :] - # Remove the start tag and everything after from last message - current_text = get_last_text(output) - set_last_text( - output, - current_text.replace(match.group(0) + after_tag, ""), - ) + # Keep only text before the tag in the message + set_last_text(output, before_tag) - if before_tag: - set_last_text(output, before_tag) - - if not get_last_text(output).strip(): + if not before_tag.strip(): # Remove empty message item if output and output[-1].get("type") == "message": output.pop() @@ -3069,9 +3095,11 @@ async def streaming_chat_response_handler(response, ctx): else: set_last_text(output, after_tag) - tag_output_handler( - content_type, tags, after_tag, output + _, recursive_end = tag_output_handler( + content_type, tags, output ) + if recursive_end: + end_flag = True break @@ -3090,24 +3118,23 @@ async def streaming_chat_response_handler(response, ctx): start_tag = item.get("start_tag", "") end_tag = item.get("end_tag", "") - if end_tag.startswith("<") and end_tag.endswith(">"): - end_tag_pattern = rf"{re.escape(end_tag)}" - else: - end_tag_pattern = rf"{re.escape(end_tag)}" + end_tag_pattern = rf"{re.escape(end_tag)}" - if re.search(end_tag_pattern, content): + # Get the block content from the item itself + if last_type == "reasoning": + parts = item.get("content", []) + block_content = "" + if parts and parts[-1].get("type") == "output_text": + block_content = parts[-1].get("text", "") + elif last_type == "open_webui:code_interpreter": + block_content = item.get("code", "") + else: + block_content = get_last_text(output) + + if re.search(end_tag_pattern, block_content): end_flag = True - # Get the block content - if last_type == "reasoning": - parts = item.get("content", []) - block_content = "" - if parts and parts[-1].get("type") == "output_text": - block_content = parts[-1].get("text", "") - elif last_type == "open_webui:code_interpreter": - block_content = item.get("code", "") - else: - block_content = get_last_text(output) + # Strip start and end tags from content start_tag_pattern = rf"{re.escape(start_tag)}" @@ -3151,36 +3178,20 @@ async def streaming_chat_response_handler(response, ctx): item["ended_at"] = time.time() # Reset by appending a new message item for leftover - if content_type != "code_interpreter": - output.append( - { - "type": "message", - "id": output_id("msg"), - "status": "in_progress", - "role": "assistant", - "content": [ - { - "type": "output_text", - "text": leftover_content, - } - ], - } - ) - else: - output.append( - { - "type": "message", - "id": output_id("msg"), - "status": "in_progress", - "role": "assistant", - "content": [ - { - "type": "output_text", - "text": leftover_content, - } - ], - } - ) + output.append( + { + "type": "message", + "id": output_id("msg"), + "status": "in_progress", + "role": "assistant", + "content": [ + { + "type": "output_text", + "text": leftover_content, + } + ], + } + ) else: # Remove the block if content is empty output.pop() @@ -3199,19 +3210,7 @@ async def streaming_chat_response_handler(response, ctx): } ) - # Clean processed content - start_tag_clean = rf"{re.escape(start_tag)}" - if start_tag.startswith("<") and start_tag.endswith(">"): - start_tag_clean = rf"<{re.escape(start_tag[1:-1])}(\s.*?)?>" - - content = re.sub( - rf"{start_tag_clean}(.|\n)*?{re.escape(end_tag)}", - "", - content, - flags=re.DOTALL, - ) - - return content, output, end_flag + return output, end_flag message = Chats.get_message_by_id_and_message_id( metadata["chat_id"], metadata["message_id"] @@ -3674,58 +3673,122 @@ async def streaming_chat_response_handler(response, ctx): ) content = f"{content}{value}" - if ( - not output - or output[-1].get("type") != "message" - ): - output.append( - { - "type": "message", - "id": output_id("msg"), - "status": "in_progress", - "role": "assistant", - "content": [ + + # Check if we're inside a tag-based block + # (reasoning, code_interpreter, or solution). + # If so, append to the existing in-progress + # item instead of creating a new message — + # otherwise tag_output_handler re-detects the + # start tag on every chunk and fragments the + # output. + last_item = output[-1] if output else None + last_item_type = ( + last_item.get("type", "") if last_item else "" + ) + inside_tag_block = ( + last_item is not None + and last_item.get("status") == "in_progress" + and last_item.get("attributes", {}).get("type") + != "reasoning_content" + and ( + last_item_type == "reasoning" + or last_item_type + == "open_webui:code_interpreter" + or ( + last_item_type == "message" + and last_item.get("_tag_type") + is not None + ) + ) + ) + + if inside_tag_block: + # Append to the existing tag-based item + if last_item_type == "open_webui:code_interpreter": + last_item["code"] = ( + last_item.get("code", "") + value + ) + elif last_item_type == "reasoning": + parts = last_item.get("content", []) + if ( + parts + and parts[-1].get("type") + == "output_text" + ): + parts[-1]["text"] += value + else: + last_item["content"] = [ { "type": "output_text", - "text": "", + "text": value, } - ], - } - ) - - # Append value to last message item's text - msg_parts = output[-1].get("content", []) - if ( - msg_parts - and msg_parts[-1].get("type") - == "output_text" - ): - msg_parts[-1]["text"] += value + ] + else: + # solution or other _tag_type message + msg_parts = last_item.get("content", []) + if ( + msg_parts + and msg_parts[-1].get("type") + == "output_text" + ): + msg_parts[-1]["text"] += value + else: + last_item["content"] = [ + { + "type": "output_text", + "text": value, + } + ] else: - output[-1]["content"] = [ - {"type": "output_text", "text": value} - ] + if ( + not output + or output[-1].get("type") != "message" + ): + output.append( + { + "type": "message", + "id": output_id("msg"), + "status": "in_progress", + "role": "assistant", + "content": [ + { + "type": "output_text", + "text": "", + } + ], + } + ) + + # Append value to last message item's text + msg_parts = output[-1].get("content", []) + if ( + msg_parts + and msg_parts[-1].get("type") + == "output_text" + ): + msg_parts[-1]["text"] += value + else: + output[-1]["content"] = [ + {"type": "output_text", "text": value} + ] if DETECT_REASONING_TAGS: - content, output, _ = tag_output_handler( + output, _ = tag_output_handler( "reasoning", reasoning_tags, - content, output, ) - content, output, _ = tag_output_handler( + output, _ = tag_output_handler( "solution", DEFAULT_SOLUTION_TAGS, - content, output, ) if DETECT_CODE_INTERPRETER: - content, output, end = tag_output_handler( + output, end = tag_output_handler( "code_interpreter", DEFAULT_CODE_INTERPRETER_TAGS, - content, output, ) diff --git a/backend/open_webui/utils/misc.py b/backend/open_webui/utils/misc.py index a192f7b66..9d9cfa1d0 100644 --- a/backend/open_webui/utils/misc.py +++ b/backend/open_webui/utils/misc.py @@ -230,25 +230,23 @@ def convert_output_to_messages(output: list, raw: bool = False) -> list[dict]: # else: skip reasoning blocks for normal LLM messages elif item_type == "open_webui:code_interpreter": - if raw: - # Include code interpreter content for LLM re-processing - code = item.get("code", "") - code_output = item.get("output", "") + # Always include code interpreter content so the LLM knows + # the code was already executed and doesn't retry. + code = item.get("code", "") + code_output = item.get("output", "") - if code: - lang = item.get("lang", "python") - pending_content.append(f"```{lang}\n{code}\n```") + if code: + pending_content.append(f"\n{code}\n") - if code_output: - if isinstance(code_output, dict): - stdout = code_output.get("stdout", "") - result = code_output.get("result", "") - output_text = stdout or result - else: - output_text = str(code_output) - if output_text: - pending_content.append(f"Output:\n{output_text}") - # else: skip extension types + if code_output: + if isinstance(code_output, dict): + stdout = code_output.get("stdout", "") + result = code_output.get("result", "") + output_text = stdout or result + else: + output_text = str(code_output) + if output_text: + pending_content.append(f"\n{output_text}\n") elif item_type.startswith("open_webui:"): # Skip other extension types